Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
44 commits
Select commit Hold shift + click to select a range
7b6c9c5
Give each spool directory one owner process: an flock on <dir>/.owner…
zaoxing Sep 28, 2026
dd12fb7
Open a sink and a service on one spool under one lock, and adopt dead…
zaoxing Sep 29, 2026
f5b34b0
Hold the spool owner lock in the engine around its sink and its service
zaoxing Sep 29, 2026
244d33e
Refuse held_by_caller unless this process holds the spool's owner lock
zaoxing Sep 29, 2026
0d1a429
Check a spool's nesting after taking its lock, and only against held …
zaoxing Sep 29, 2026
55ffbd2
Ignore and clear the staging copy a claim killed before its rename le…
zaoxing Sep 29, 2026
6e8724b
Let a child forked without exec drop its copy of every spool owner lock
zaoxing Sep 29, 2026
a56a163
Look again for siblings that were alive when the storage service started
zaoxing Sep 29, 2026
055ce11
Keep the spool directory owned while a sink that did not seal may sti…
zaoxing Sep 29, 2026
553b01b
Key a spool's catalog directory by the servers its packs go to, not o…
zaoxing Sep 29, 2026
db87280
Refuse BeeGFS, CIFS/SMB2 and FUSE spool roots, not only NFS and Lustre
zaoxing Sep 29, 2026
bc7b56f
Give the adoption test's live sibling real packs and a stage in flight
zaoxing Sep 29, 2026
efc42b0
Pin that no dead spool is uploaded while an adopted pack is owed
zaoxing Sep 29, 2026
9932195
Cover the owner-lock paths the review's surviving mutations showed un…
zaoxing Sep 29, 2026
ec57796
Refuse GPFS, 9p, AFS and OrangeFS spool roots as well
zaoxing Sep 29, 2026
a1ded46
Keep spool.cpp portable to macOS: guard statfs and the no-replace rename
zaoxing Sep 29, 2026
1cf61b5
Charge what dead incarnations left beside a spool against its budget
zaoxing Sep 29, 2026
42106b0
Adopt dead spools in the loop, a slice per cycle, not in start() or f…
zaoxing Sep 29, 2026
9cacc53
Leave a dead spool this service can never upload, and say so once
zaoxing Sep 29, 2026
dfff3c3
Warn of packs left in spool_root outside the layout, and say how to h…
zaoxing Sep 29, 2026
01d29a9
Track a spool lock's descriptor for the fork handler from its open()
zaoxing Sep 29, 2026
a285818
Lock the spool directory itself, not only its lock file
zaoxing Sep 29, 2026
2f2a68b
Judge held_by_caller by the owner record where fdinfo lists no flocks
zaoxing Sep 29, 2026
2cb2337
Charge a sink only what an adoption can drain beside it
zaoxing Sep 29, 2026
354cd4e
Say that skipping _refs/ is the C++ spool's alone
zaoxing Sep 29, 2026
72c9929
Pass over nested spool directories, and refuse nesting only while held
zaoxing Sep 29, 2026
e0604af
Let go of a half-adopted sibling once the service has latched
zaoxing Sep 29, 2026
20191d1
Pin that a flush adopts nothing, with the loop idle
zaoxing Sep 29, 2026
633f7e0
Race a scan against new spool directories, not only look at the result
zaoxing Sep 29, 2026
988ee72
Say the fork handler tracks the directory's lock descriptor too
zaoxing Sep 29, 2026
ac868a4
Merge origin/main (B5: close() delivers the tail pack; flush and stop…
zaoxing Sep 29, 2026
10b4fe7
Adopt dead spools through the upload path's chunks, clients and cancels
zaoxing Sep 29, 2026
f8ae2d9
Let a flush wait for the adoption step in flight, not the rest of a s…
zaoxing Sep 29, 2026
87d3a92
Keep the spool lock until the release backstop has sealed the sink
zaoxing Sep 29, 2026
1616682
Say when close() keeps the spool directory, and what a flush waits fo…
zaoxing Sep 29, 2026
9ac51d4
Hold the harness's spool lock for B5's live services as for B6's
zaoxing Sep 29, 2026
74bcac5
Pin on a real ring that close() lets go of the engine's spool directory
zaoxing Sep 29, 2026
9d5364f
Wait for a forked worker to have run before killing its lock's owner
zaoxing Sep 29, 2026
1a10c1e
Forward no request its client did not finish, and store no body short…
zaoxing Sep 29, 2026
92295ca
Give two cut tests the time a loaded runner needs before the cut
zaoxing Sep 29, 2026
45c5c88
Allow the stalled-upload stop tests eight attempts, so an uncut backo…
zaoxing Sep 29, 2026
9081211
Validate an adopted dead spool a pack a step, so a flush waits for on…
zaoxing Sep 29, 2026
ca392bc
Pin that flushes never stop an adoption, and that a cycle's own packs…
zaoxing Sep 30, 2026
5cc69fe
Say adoption lists a dead spool a pack a step, and how big its round is
zaoxing Sep 30, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions docs/capture-storage-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -1827,6 +1827,9 @@ bytes, and retained failure details all have explicit caps.
writer will be native. See *Phase 6 -- Decision: the production writer is
native*.
- One process owns a spool directory; cross-process locking is not implemented.
(The native C++ spool does lock it: an owner lock on `<dir>/.owner.lock`,
B6, documented in `native/csrc/store/spool.h`. The Python reference
deliberately stays unlocked.)
- Durable mode stages synchronously and uploads through a separate explicit
uploader, so remote backpressure is isolated from local commit.
- The pipeline remains opt-in and is not connected to Ring².
Expand Down
65 changes: 65 additions & 0 deletions docs/integration-api-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,71 @@ refuses the persistent backend with `ConfigurationError` rather than generating
with nothing stored. The catalog
takes one publisher per `(database, table_prefix)`, so a second engine on the
same catalog is refused at `create_record_runtime`.
The engine spools into a directory of its own under
`capture_sink_config.spool_root`,
`<spool_root>/<catalog_key>/r<rank>-<incarnation>/` (the key is the first 12
hex digits of a sha256 of where the packs go -- the ClickHouse host and port,
`database`, `table_prefix`, the S3 endpoint and bucket, and `store_id`, as
spelled in the config -- the rank torchrun's
`RANK`, 0 when unset, and the incarnation fresh on every
`create_record_runtime`), and owns it: an flock on its `.owner.lock` and on the
directory itself, taken before the service starts and let go after the sink and
the service are done, when a drained directory is removed. The directory's own
lock keeps it owned should its `.owner.lock` be removed from under it, and
systemd-tmpfiles skips a flocked directory when it ages `/tmp`; other cleaners
may not, so keep `spool_root` out of what they age. If the sink did not seal --
the release backstop's flush, when the stopping ring lets go of it, did not go
through -- it may still be staging, so the directory stays owned by the
process until it exits (a warning names it) and the next process on the node
adopts it; so does the directory of an engine dropped without `close()`.
A second process on a directory is
refused, naming the holder's pid and host. Once started, the service's
background loop adopts the directories under the same catalog key whose owners
have died: their stale `.open` files are swept, their ready packs uploaded and
indexed, a round at a time and a slice of each cycle, and the directory
removed, so a crashed process's packs reach the catalog through the next one on
the node, whatever run it belongs to. Neither `create_record_runtime` nor
`flush_and_wait` waits for that: a flush covers this process's records (an
adopted pack uploaded and not yet indexed is waited for like its own), waits
for no more of an adoption in progress than the step it is in (one round of
its uploads, or one of its packs validated -- a dead directory's packs are
hashed one per step, so a large backlog's listing does not hold a flush up),
and `close()` cuts an adoption as it cuts the service's own uploads, leaving
what it did not upload in the dead directory for the next process. The
storage part of `capture_status()` reports the adoption (`adopted_spools`,
`adopted_packs`, `adoption_owed`, `live_siblings`). A dead directory the
service can never adopt -- one holding a pack it can never upload, such as one
larger than its `uploader_max_in_flight_bytes` or one whose key already holds a
different object, or one it cannot lock -- is left in place once the rest of
its packs are up, listed in `blocked_siblings`, reported once in `last_error`
and in its `.owner.lock` (a `blocked: <why>` line), and neither retried nor
owed. `spool_max_bytes` bounds the directory together with what the dead
directories the service can adopt still hold, so restarts while uploads are
blocked cannot each add a whole budget; the room comes back as they are
adopted. A live process's directory is its own budget, and neither a blocked
directory nor one this process keeps owned itself (an earlier engine's whose
sink did not seal) is charged: no adoption here drains them. The
spool root must be
node-local: NFS, Lustre, BeeGFS, CIFS/SMB2, FUSE, GPFS, 9p, AFS and OrangeFS
are refused (by statfs `f_type`) unless `NativeSinkConfig.spool_allow_shared_filesystem`, which a FUSE
filesystem that is itself local, such as fuse-overlayfs, needs too. Without
`capture_storage_config` the sink owns `spool_root` itself. With an explicit
`record_sink`, the service drains `spool_root` as that sink writes it,
unswept and adopting nothing. Both of those modes pass over the rank
directories under `spool_root` -- what a crashed or undrained default-mode run
left there is the next default-mode start's to adopt -- so switching to them
(the rollback to an explicit `record_sink` included) works after such a run.
They are refused, naming the holder, while a default-mode process on the node
holds a rank directory under that `spool_root`, and a default-mode start is
refused while one of them holds `spool_root`. Upgrading from an engine without this layout: it
spooled into `<spool_root>/v1/...` and its next start uploaded what a crashed
run left there, but nothing adopts packs outside the layout now -- nor those a
sink-only or explicit-`record_sink` run leaves in `spool_root`. Each
`create_record_runtime` logs a warning while any are there, with their count,
one of them, and a directory of the layout nobody owns
(`<spool_root>/<catalog_key>/r0-00000000/`); moved into it with their paths
below `spool_root` kept (`v1/...`), packs bound for this catalog and store are
adopted by the next start on the node.
To reach a secured catalog, set `clickhouse_scheme="https"` (and the server's
TLS HTTP port, usually 8443) on `NativeCaptureStorageConfig`. The client always
verifies the server's certificate and name, against libcurl's built-in CA
Expand Down
153 changes: 153 additions & 0 deletions native/csrc/catalog/bindings_store.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,10 @@
//
// StorageService spool -> object store -> catalog, on a background thread
// CaptureReader search / select / hydrate against the catalog + store
// SpoolOwnerLock the owner lock of one spool directory (store/spool.h),
// which the engine holds around its sink and service
// spool_rank_directory, spool_catalog_key, spool_owner
// the section 2.3 spool layout, and who owns a directory
//
// Pure C++ plus libcurl and libcrypto. It uses pybind11's headers but links
// nothing from torch, and registers no ring types, so it loads beside
Expand All @@ -10,12 +14,14 @@
#include <pybind11/stl.h>

#include <memory>
#include <optional>
#include <string>
#include <vector>

#include "catalog/hydration.h"
#include "catalog/reader.h"
#include "catalog/storage_service.h"
#include "store/spool.h"

namespace py = pybind11;
namespace dc = dmi_catalog;
Expand Down Expand Up @@ -75,10 +81,50 @@ dc::ClickHouseConnection clickhouse_connection(const py::dict& d) {
return c;
}

dmi_store::OwnerLock owner_lock(const std::string& text) {
dmi_store::OwnerLock mode = dmi_store::OwnerLock::kTake;
if (!dmi_store::ParseOwnerLock(text, &mode)) {
throw py::value_error("spool_owner_lock must be 'take' or "
"'held_by_caller', got '" + text + "'");
}
return mode;
}

// Where a spool's packs go (store/spool.h); every field is required.
dmi_store::SpoolDestination spool_destination(const py::dict& d) {
for (const char* key : {"clickhouse_host", "clickhouse_port", "database",
"table_prefix", "s3_endpoint", "s3_bucket",
"store_id"}) {
if (!d.contains(key)) {
throw py::key_error(std::string("spool destination needs '") + key +
"'");
}
}
dmi_store::SpoolDestination out;
out.clickhouse_host = d["clickhouse_host"].cast<std::string>();
out.clickhouse_port = d["clickhouse_port"].cast<uint64_t>();
out.database = d["database"].cast<std::string>();
out.table_prefix = d["table_prefix"].cast<std::string>();
out.s3_endpoint = d["s3_endpoint"].cast<std::string>();
out.s3_bucket = d["s3_bucket"].cast<std::string>();
out.store_id = d["store_id"].cast<std::string>();
return out;
}

dc::StorageServiceConfig service_config(const py::dict& d) {
dc::StorageServiceConfig c;
c.spool_root = get<std::string>(d, "spool_root", "");
c.spool_max_bytes = get<uint64_t>(d, "spool_max_bytes", c.spool_max_bytes);
c.spool_owner_lock =
owner_lock(get<std::string>(d, "spool_owner_lock", "take"));
c.spool_allow_shared_filesystem = get<bool>(
d, "spool_allow_shared_filesystem", c.spool_allow_shared_filesystem);
c.adopt_sibling_spools =
get<bool>(d, "adopt_sibling_spools", c.adopt_sibling_spools);
c.adoption_recheck_interval_ns = get<uint64_t>(
d, "adoption_recheck_interval_ns", c.adoption_recheck_interval_ns);
c.adoption_slice_ns =
get<uint64_t>(d, "adoption_slice_ns", c.adoption_slice_ns);
c.s3 = s3_config(d);
c.uploader.store_id = get<std::string>(d, "store_id", c.uploader.store_id);
c.uploader.max_workers = get<int>(d, "uploader_max_workers", c.uploader.max_workers);
Expand Down Expand Up @@ -129,6 +175,11 @@ py::dict snapshot_dict(const dc::StorageServiceSnapshot& s) {
out["swept_on_start"] = s.swept_on_start;
out["pending_index"] = s.pending_index;
out["rejected_packs"] = s.rejected_packs;
out["adopted_spools"] = s.adopted_spools;
out["adopted_packs"] = s.adopted_packs;
out["adoption_owed"] = s.adoption_owed;
out["live_siblings"] = s.live_siblings;
out["blocked_siblings"] = s.blocked_siblings;
out["failed"] = s.failed;
out["lease_state"] = s.lease_state;
// Seconds on the monotonic clock, comparable with time.monotonic() (both
Expand Down Expand Up @@ -275,6 +326,23 @@ class CaptureReader {
dc::NativeCaptureReader reader_;
};

// A refused SpoolStatus as the Python exception that says what it is: a
// directory another process owns, a refused configuration, or an I/O error.
PyObject* g_spool_owned_error = nullptr; // SpoolOwnedError, set at import

[[noreturn]] void raise_spool_status(dmi_store::SpoolStatus status,
const std::string& error) {
if (status == dmi_store::SpoolStatus::kOwned) {
PyErr_SetString(g_spool_owned_error, error.c_str());
throw py::error_already_set();
}
if (status == dmi_store::SpoolStatus::kBadArgument) {
throw py::value_error(error);
}
PyErr_SetString(PyExc_OSError, error.c_str());
throw py::error_already_set();
}

} // namespace

PYBIND11_MODULE(_dmi_native_store, m) {
Expand All @@ -296,6 +364,91 @@ PYBIND11_MODULE(_dmi_native_store, m) {
[](const dc::CaptureStorageService& s) { return snapshot_dict(s.snapshot()); })
.def("rethrow_if_failed", &dc::CaptureStorageService::rethrow_if_failed);

// Raised when another owner holds a spool directory; a RuntimeError, so
// callers catching the service's refusals keep catching it.
static py::exception<std::runtime_error> spool_owned_error(
m, "SpoolOwnedError", PyExc_RuntimeError);
g_spool_owned_error = spool_owned_error.ptr();

py::class_<dmi_store::SpoolOwnerLock>(m, "SpoolOwnerLock")
.def(py::init([](const std::string& directory,
bool allow_shared_filesystem) {
auto lock = std::make_unique<dmi_store::SpoolOwnerLock>();
std::string error;
const dmi_store::SpoolStatus status =
dmi_store::SpoolOwnerLock::Acquire(
directory, allow_shared_filesystem, lock.get(), &error);
if (status != dmi_store::SpoolStatus::kOk) {
raise_spool_status(status, error);
}
return lock;
}),
py::arg("directory"), py::arg("allow_shared_filesystem") = false)
.def_property_readonly("directory", &dmi_store::SpoolOwnerLock::directory)
.def_property_readonly("held", &dmi_store::SpoolOwnerLock::held)
.def("release", &dmi_store::SpoolOwnerLock::Release)
.def("release_and_remove_if_empty",
[](dmi_store::SpoolOwnerLock& self) {
std::string error;
return self.ReleaseAndRemoveIfEmpty(&error);
})
.def("__enter__",
[](dmi_store::SpoolOwnerLock& self) -> dmi_store::SpoolOwnerLock& {
return self;
}, py::return_value_policy::reference)
.def("__exit__", [](dmi_store::SpoolOwnerLock& self, const py::args&) {
self.Release();
return false;
});

m.def("spool_owner",
[](const std::string& directory) -> py::object {
dmi_store::SpoolOwner owner;
if (!dmi_store::ReadSpoolOwner(directory, &owner)) return py::none();
py::dict out;
out["host"] = owner.host;
out["pid"] = owner.pid;
return out;
},
py::arg("directory"),
"Who holds a spool directory's owner lock, or None when nothing "
"does.");
m.def("spool_catalog_key",
[](const py::dict& destination) {
return dmi_store::SpoolCatalogKey(spool_destination(destination));
},
py::arg("destination"),
"The catalog key of a destination dict (clickhouse_host, "
"clickhouse_port, database, table_prefix, s3_endpoint, s3_bucket, "
"store_id).");
m.def("spool_rank_directory",
[](const std::string& base, const py::dict& destination,
uint64_t producer_rank, std::optional<std::string> incarnation) {
const std::string fresh =
incarnation ? *incarnation : dmi_store::NewSpoolIncarnation();
uint64_t parsed_rank = 0;
std::string parsed;
if (!dmi_store::ParseSpoolRankDirectoryName(
dmi_store::SpoolRankDirectoryName(producer_rank, fresh),
&parsed_rank, &parsed)) {
throw py::value_error("incarnation must be 8 lowercase hex "
"digits");
}
return dmi_store::SpoolRankDirectory(
base, spool_destination(destination), producer_rank, fresh);
},
py::arg("base"), py::arg("destination"), py::arg("producer_rank"),
py::arg("incarnation") = py::none(),
"<base>/<catalog_key>/r<producer_rank>-<incarnation>, the section "
"2.3 spool layout; a fresh incarnation when none is given.");
m.def("_set_spool_filesystem_type_for_testing",
[](std::optional<int64_t> f_type) {
dmi_store::SetFilesystemTypeForTesting(f_type ? *f_type : -1);
},
py::arg("f_type"),
"Test seam: node-local checks in this module read f_type instead "
"of statfs(2); None restores statfs.");

py::class_<CaptureReader>(m, "CaptureReader")
.def(py::init<const py::dict&>(), py::arg("config"))
.def("search", &CaptureReader::search, py::arg("filters"))
Expand Down
Loading
Loading