diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md
index 8fdc5c4ed..f63fc389b 100644
--- a/docs/capture-storage-design.md
+++ b/docs/capture-storage-design.md
@@ -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 `
/.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².
diff --git a/docs/integration-api-v1.md b/docs/integration-api-v1.md
index b34179302..b952dae58 100644
--- a/docs/integration-api-v1.md
+++ b/docs/integration-api-v1.md
@@ -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`,
+`//r-/` (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: ` 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 `/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
+(`//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
diff --git a/native/csrc/catalog/bindings_store.cpp b/native/csrc/catalog/bindings_store.cpp
index dcaf4cdc6..10a8bb1c6 100644
--- a/native/csrc/catalog/bindings_store.cpp
+++ b/native/csrc/catalog/bindings_store.cpp
@@ -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
@@ -10,12 +14,14 @@
#include
#include
+#include
#include
#include
#include "catalog/hydration.h"
#include "catalog/reader.h"
#include "catalog/storage_service.h"
+#include "store/spool.h"
namespace py = pybind11;
namespace dc = dmi_catalog;
@@ -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();
+ out.clickhouse_port = d["clickhouse_port"].cast();
+ out.database = d["database"].cast();
+ out.table_prefix = d["table_prefix"].cast();
+ out.s3_endpoint = d["s3_endpoint"].cast();
+ out.s3_bucket = d["s3_bucket"].cast();
+ out.store_id = d["store_id"].cast();
+ return out;
+}
+
dc::StorageServiceConfig service_config(const py::dict& d) {
dc::StorageServiceConfig c;
c.spool_root = get(d, "spool_root", "");
c.spool_max_bytes = get(d, "spool_max_bytes", c.spool_max_bytes);
+ c.spool_owner_lock =
+ owner_lock(get(d, "spool_owner_lock", "take"));
+ c.spool_allow_shared_filesystem = get(
+ d, "spool_allow_shared_filesystem", c.spool_allow_shared_filesystem);
+ c.adopt_sibling_spools =
+ get(d, "adopt_sibling_spools", c.adopt_sibling_spools);
+ c.adoption_recheck_interval_ns = get(
+ d, "adoption_recheck_interval_ns", c.adoption_recheck_interval_ns);
+ c.adoption_slice_ns =
+ get(d, "adoption_slice_ns", c.adoption_slice_ns);
c.s3 = s3_config(d);
c.uploader.store_id = get(d, "store_id", c.uploader.store_id);
c.uploader.max_workers = get(d, "uploader_max_workers", c.uploader.max_workers);
@@ -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
@@ -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) {
@@ -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 spool_owned_error(
+ m, "SpoolOwnedError", PyExc_RuntimeError);
+ g_spool_owned_error = spool_owned_error.ptr();
+
+ py::class_(m, "SpoolOwnerLock")
+ .def(py::init([](const std::string& directory,
+ bool allow_shared_filesystem) {
+ auto lock = std::make_unique();
+ 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 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(),
+ "//r-, the section "
+ "2.3 spool layout; a fresh incarnation when none is given.");
+ m.def("_set_spool_filesystem_type_for_testing",
+ [](std::optional 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_(m, "CaptureReader")
.def(py::init(), py::arg("config"))
.def("search", &CaptureReader::search, py::arg("filters"))
diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp
index c1490eec4..16c747e88 100644
--- a/native/csrc/catalog/storage_service.cpp
+++ b/native/csrc/catalog/storage_service.cpp
@@ -4,6 +4,9 @@
#include
#include
#include
+#include
+#include
+#include
#include
#include
#include
@@ -152,6 +155,24 @@ class CaptureStorageService::LeaseScope {
uint64_t counted_on_entry_ = 0;
};
+// The dead sibling a cycle is adopting, kept across cycles: its owner lock,
+// a Spool on it (held_by_caller, the lock being this service's), its
+// recovery -- the packs listed and those validated so far, a pack a step
+// (Spool::BeginRecovery) -- and, once every pack is validated, those not
+// uploaded yet. So a backlog is hashed once, however many cycles its
+// listing and its upload take.
+struct CaptureStorageService::Adoption {
+ std::string directory;
+ dmi_store::SpoolOwnerLock lock;
+ dmi_store::Spool spool;
+ dmi_store::SpoolRecovery recovery;
+ bool validated = false; // every listed pack; `remaining` holds the valid
+ std::deque remaining;
+ // The first pack this service can never upload (UploadFailure::
+ // retryable false); the directory is left, with it, once the rest are up.
+ std::string blocked;
+};
+
CaptureStorageService::CaptureStorageService(StorageServiceConfig config)
: config_(std::move(config)),
s3_(config_.s3),
@@ -190,10 +211,41 @@ CaptureStorageService::CaptureStorageService(StorageServiceConfig config)
" ms, or raise lease_ttl_ns");
}
std::string error;
- if (dmi_store::Spool::Open({config_.spool_root, config_.spool_max_bytes},
- &spool_, &error) != dmi_store::SpoolStatus::kOk) {
+ dmi_store::SpoolConfig spool_config{config_.spool_root,
+ config_.spool_max_bytes};
+ spool_config.owner_lock = config_.spool_owner_lock;
+ spool_config.allow_shared_filesystem = config_.spool_allow_shared_filesystem;
+ if (dmi_store::Spool::Open(spool_config, &spool_, &error) !=
+ dmi_store::SpoolStatus::kOk) {
throw std::runtime_error("storage service: cannot open spool: " + error);
}
+ if (config_.adopt_sibling_spools) {
+ // Siblings are adopted INTO this catalog, so this directory must sit
+ // under this catalog's key: a directory under another catalog's key
+ // would index that catalog's packs here.
+ const std::filesystem::path own(spool_.root());
+ dmi_store::SpoolDestination destination;
+ destination.clickhouse_host = config_.clickhouse.host;
+ destination.clickhouse_port = config_.clickhouse.port;
+ destination.database = config_.writer.database;
+ destination.table_prefix = config_.writer.table_prefix;
+ destination.s3_endpoint = config_.s3.endpoint;
+ destination.s3_bucket = config_.s3.bucket;
+ destination.store_id = config_.uploader.store_id;
+ const std::string key = dmi_store::SpoolCatalogKey(destination);
+ uint64_t rank = 0;
+ std::string incarnation;
+ if (!dmi_store::ParseSpoolRankDirectoryName(own.filename().string(),
+ &rank, &incarnation) ||
+ own.parent_path().filename().string() != key) {
+ throw std::invalid_argument(
+ "storage service: adopt_sibling_spools needs spool_root to be a "
+ "rank directory /" + key + "/r- (this "
+ "catalog's key for its ClickHouse host and port, database, "
+ "table_prefix, S3 endpoint, bucket and store_id), got " +
+ spool_.root());
+ }
+ }
s3_.set_cancellation(&read_cancel_);
upload_s3_.set_cancellation(&upload_cancel_);
uploader_ = std::make_unique(&spool_, &upload_s3_,
@@ -250,8 +302,10 @@ void CaptureStorageService::start() {
last_reconcile_ns_ = steady_ns();
{
+ // With siblings to look at, the loop's first cycle runs at once: it,
+ // not start(), adopts them.
std::lock_guard lock(wake_mutex_);
- kick_ = false;
+ kick_ = adoption_scan_owed_;
}
{
LeaseScope lease(this); // publishes the lease state
@@ -265,12 +319,10 @@ void CaptureStorageService::start() {
}
void CaptureStorageService::sweep_and_reconcile_at_start() {
- // After the lease, never before, so a second process pointed at this spool
- // usually learns that the catalog is held before it can delete a live
- // sink's .open files. Only usually: a holder that is quarantined has let
- // its row lapse, and a second process can take the lease in that gap. The
- // spool itself is not locked; one process per spool is the caller's job
- // until the spool gets an owner lock.
+ // After the lease: a start refused the catalog never touches the spool.
+ // The spool's owner lock (taken at construction, by this service or its
+ // caller) is what keeps another process's writer out of the directory
+ // this deletes .open files in.
if (config_.sweep_spool_on_start) {
std::vector recovered;
std::string error;
@@ -282,6 +334,16 @@ void CaptureStorageService::sweep_and_reconcile_at_start() {
state_.swept_on_start = recovered.size();
}
+ // Dead siblings are the loop's, from its first cycle on: after the lease
+ // and this directory's own sweep, as the plan orders them, but not here,
+ // where a backlog the object store refuses would hold start() -- and
+ // create_record_runtime -- through every pack's retry chain.
+ if (config_.adopt_sibling_spools) {
+ adoption_scan_owed_ = true;
+ std::lock_guard lock(state_mutex_);
+ state_.adoption_owed = true;
+ }
+
// A failed pass is not fatal -- the bucket is still there next time. Nor
// is a lease lost while it runs, to a quarantine or to another holder:
// that is the running service's case, and the loop handles it as it does
@@ -326,6 +388,7 @@ void CaptureStorageService::stop() {
// lease has to keep renewing until that is done.
stop_lease_thread();
std::lock_guard cycle(cycle_mutex_);
+ let_go_of_adoption();
if (started_) {
started_ = false;
// A quarantined writer holds no lease, so it writes no tombstone: the
@@ -359,6 +422,15 @@ void CaptureStorageService::stop_lease_thread() {
}
bool CaptureStorageService::flush(double timeout_s) {
+ // The loop's adoption gives way to this call at its next step, so the
+ // cycle lock is not held for the rest of an adoption slice.
+ struct InProgress {
+ explicit InProgress(std::atomic* count) : count_(count) {
+ count_->fetch_add(1, std::memory_order_acq_rel);
+ }
+ ~InProgress() { count_->fetch_sub(1, std::memory_order_acq_rel); }
+ std::atomic* count_;
+ } in_progress(&flushes_in_progress_);
const auto deadline = std::chrono::steady_clock::now() +
std::chrono::duration_cast(
std::chrono::duration(timeout_s));
@@ -383,8 +455,9 @@ bool CaptureStorageService::flush(double timeout_s) {
if (upload_cancel_.cancelled_for_good()) return false;
// Its own cycle, bounded by the deadline, and without the reconcile:
// the reconcile lists the whole bucket and asks the catalog about
- // every page, which no deadline bounds. The loop runs it.
- const bool drained = run_cycle(deadline_ns, false).drained;
+ // every page, which no deadline bounds. The loop runs it. Nor does it
+ // adopt: a dead backlog is not this process's records.
+ const bool drained = run_cycle(deadline_ns, false, false).drained;
if (!rejected_unreported_.empty()) {
std::string message = "storage service: " +
std::to_string(rejected_unreported_.size()) +
@@ -427,11 +500,22 @@ void CaptureStorageService::loop() {
kicked = kick_;
kick_ = false;
}
+ bool failed = false;
{
std::lock_guard lock(state_mutex_);
- if (failure_) return; // another publisher holds the catalog
+ failed = failure_ != nullptr;
}
std::lock_guard cycle(cycle_mutex_);
+ if (failed) {
+ // Another publisher holds the catalog, and no cycle runs again. A
+ // sibling half adopted would stay locked by this process, with
+ // nobody working on it, until stop() -- the engine's close(), maybe
+ // hours of capture later -- and every other process on the node
+ // would read it as live meanwhile. Nor does it look again.
+ adoption_scan_owed_ = false;
+ let_go_of_adoption();
+ return;
+ }
{
std::lock_guard lock(wake_mutex_);
if (stop_requested_) return;
@@ -442,8 +526,9 @@ void CaptureStorageService::loop() {
// flush -- close()'s order -- by that cycle's catalog work, which stop()
// cannot cut. So wait out the rest of the interval first, unless a
// fresh lease asked for a cycle now -- once: flushes that keep coming
- // do every cycle's work but the reconcile, which only this loop runs,
- // so the next wake runs a cycle whatever they did.
+ // do every cycle's work but the reconcile and the adoption, which only
+ // this loop runs, so the next wake runs a cycle whatever they did.
+ // start() kicks the first cycle when there are siblings to look at.
const uint64_t since_ns = steady_ns() - last_cycle_end_ns_;
if (!kicked && !waited_again && last_cycle_end_ns_ != 0 &&
since_ns < wait_ns) {
@@ -452,7 +537,7 @@ void CaptureStorageService::loop() {
continue;
}
waited_again = false;
- run_cycle(0, true);
+ run_cycle(0, true, true);
// poll_interval * 2^streak, capped: flush() shares the streak, so an
// outage it saw also slows the loop, and a success from either resets it.
wait_ns = config_.poll_interval_ns;
@@ -465,7 +550,7 @@ void CaptureStorageService::loop() {
}
CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
- uint64_t deadline_ns, bool allow_reconcile) {
+ uint64_t deadline_ns, bool allow_reconcile, bool adopt) {
CycleOutcome outcome;
// stop() has begun: it cancelled the uploads and the reads for good, so a
// cycle could only make catalog requests -- a lease claim, a replay guard
@@ -494,30 +579,10 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
LeaseScope lease(this);
catalog = ensure_publisher_lease();
}
- // Indexes refs, keeping whatever does not index owed: it is already gone
- // from the spool, so pending_index_ is the only record of it in-process.
- // Those a cancel or the deadline left owed are counted in `deferred`
- // (index_bounded).
+ // What a cycle uploads and cannot index stays owed in pending_index_
+ // (index_or_owe). Those a cancel or the deadline left owed are counted
+ // here: they make the cycle cut short, not failed.
size_t deferred = 0;
- const auto index_or_owe = [this, catalog, deadline_ns,
- &deferred](std::vector refs) {
- if (!catalog) {
- pending_index_.insert(pending_index_.end(), refs.begin(), refs.end());
- return;
- }
- std::vector unindexed;
- try {
- if (!refs.empty()) {
- deferred += index_bounded(std::move(refs), &unindexed, deadline_ns);
- }
- } catch (...) {
- pending_index_.insert(pending_index_.end(), unindexed.begin(),
- unindexed.end());
- throw;
- }
- pending_index_.insert(pending_index_.end(), unindexed.begin(),
- unindexed.end());
- };
try {
// 1. Retry what earlier cycles uploaded but could not index. While any
// of it is still owed, the catalog is down or refusing: upload
@@ -527,7 +592,7 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
if (catalog && !pending_index_.empty()) {
std::vector owed;
owed.swap(pending_index_);
- index_or_owe(std::move(owed));
+ deferred += index_or_owe(std::move(owed), catalog, deadline_ns);
}
// 2. Upload what the sink has staged -- but only with the lease and
@@ -589,52 +654,43 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
}
}
const size_t end = std::min(staged.size(), next + chunk);
- const dmi_store::UploadBatchResult batch = uploader_->UploadStaged(
+ ChunkOutcome sent = upload_chunk(
+ uploader_.get(),
std::vector(staged.begin() + next,
staged.begin() + end));
next = end;
- std::vector to_index;
- uint64_t uploaded_packs = 0;
- uint64_t uploaded_bytes = 0;
- size_t failed_uploads = 0;
- size_t cancelled_uploads = 0;
- for (size_t i = 0; i < batch.refs.size(); ++i) {
- const dmi_store::PackRef& ref = batch.refs[i];
- if (!ref.pack_id.empty()) {
- to_index.push_back({ref.pack_id, ref.store_id, ref.object_key,
- ref.object_bytes, ref.checksum,
- ref.record_count});
- ++uploaded_packs;
- uploaded_bytes += ref.object_bytes;
- } else if (i < batch.failures.size() &&
- batch.failures[i].cancelled) {
- ++cancelled_uploads; // still staged; not the pack's fault
- } else {
- // A failed upload stays in the spool, so a later cycle
- // retries it.
- ++failed_uploads;
- if (i < batch.failures.size()) {
- record_error("upload failed for " +
- batch.failures[i].object_key + ": " +
- batch.failures[i].error);
- }
+ for (size_t i = 0; i < sent.batch.failures.size(); ++i) {
+ const dmi_store::UploadFailure& failure = sent.batch.failures[i];
+ // A failed upload stays in the spool, so a later cycle retries
+ // it; a cancelled one is still staged, and not the pack's fault.
+ if (!sent.batch.refs[i].pack_id.empty() || failure.cancelled) {
+ continue;
}
+ record_error("upload failed for " + failure.object_key + ": " +
+ failure.error);
}
- if (cancelled_uploads != 0) outcome.cut_short = true;
- upload_failures += failed_uploads;
- {
- std::lock_guard lock(state_mutex_);
- state_.uploaded_packs += uploaded_packs;
- state_.uploaded_bytes += uploaded_bytes;
- state_.upload_failures += failed_uploads;
- state_.cancelled_uploads += cancelled_uploads;
- }
- index_or_owe(std::move(to_index));
+ if (sent.cancelled != 0) outcome.cut_short = true;
+ upload_failures += sent.failed;
+ deferred += index_or_owe(std::move(sent.to_index), true, deadline_ns);
}
uploaded_all = next == staged.size() && upload_failures == 0;
}
}
+ // 3a. The loop's cycles adopt dead siblings, a slice at a time, under
+ // the same rule as the uploads above: only with the lease, every
+ // staged pack of this process's up, nothing owed, and no cancel.
+ // flush()'s cycles do not: a dead backlog is not this process's
+ // records. Adoption uploads and indexes through the same chunks,
+ // and the same Cancellations, as the service's own spool.
+ bool adoption_failed = false;
+ if (adopt && config_.adopt_sibling_spools && uploaded_all &&
+ pending_index_.empty() && !upload_cancel_.cancelled()) {
+ bool adoption_cut = false;
+ adoption_failed = !adopt_step(deadline_ns, &deferred, &adoption_cut);
+ if (adoption_cut) outcome.cut_short = true;
+ }
+
// 4. Reconcile on its interval, or when the pass at start() lost the
// lease before it finished -- in the loop's cycles only, and not
// once a cancel came. The lease thread keeps the lease alive.
@@ -652,10 +708,14 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
// catalog, and nothing pending. Without the lease nothing can be
// confirmed in the catalog, so the cycle is not drained, and it counts
// towards the backoff. Packs a cancel or the deadline left owed make it
- // cut short, not failed.
+ // cut short, not failed. A dead sibling still to adopt is not a
+ // failure, nor undrained: only an adoption upload that failed backs the
+ // loop off, and adopted packs uploaded and not indexed are owed like
+ // the service's own.
if (deferred != 0) outcome.cut_short = true;
outcome.failed = !catalog || listing_failed || lost_lease ||
- upload_failures != 0 || pending_index_.size() > deferred;
+ upload_failures != 0 || adoption_failed ||
+ pending_index_.size() > deferred;
outcome.drained = uploaded_all && !outcome.failed && !outcome.cut_short;
std::lock_guard lock(state_mutex_);
++state_.cycles;
@@ -672,6 +732,7 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
{
std::lock_guard lock(state_mutex_);
state_.pending_index = pending_index_.size();
+ state_.adoption_owed = adoption_owed();
}
// A cycle a cancel cut short proved nothing about the store or the
// catalog: it resets the backoff only by succeeding, and grows it only by
@@ -685,6 +746,384 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
return outcome;
}
+size_t CaptureStorageService::index_or_owe(std::vector refs,
+ bool catalog,
+ uint64_t deadline_ns) {
+ if (!catalog) {
+ pending_index_.insert(pending_index_.end(), refs.begin(), refs.end());
+ return 0;
+ }
+ std::vector unindexed;
+ size_t deferred = 0;
+ try {
+ if (!refs.empty()) {
+ deferred = index_bounded(std::move(refs), &unindexed, deadline_ns);
+ }
+ } catch (...) {
+ pending_index_.insert(pending_index_.end(), unindexed.begin(),
+ unindexed.end());
+ throw;
+ }
+ pending_index_.insert(pending_index_.end(), unindexed.begin(),
+ unindexed.end());
+ return deferred;
+}
+
+CaptureStorageService::ChunkOutcome CaptureStorageService::upload_chunk(
+ dmi_store::SpoolUploader* uploader,
+ std::vector chunk) {
+ ChunkOutcome sent;
+ sent.batch = uploader->UploadStaged(std::move(chunk));
+ uint64_t uploaded_bytes = 0;
+ for (size_t i = 0; i < sent.batch.refs.size(); ++i) {
+ const dmi_store::PackRef& ref = sent.batch.refs[i];
+ if (!ref.pack_id.empty()) {
+ sent.to_index.push_back({ref.pack_id, ref.store_id, ref.object_key,
+ ref.object_bytes, ref.checksum,
+ ref.record_count});
+ uploaded_bytes += ref.object_bytes;
+ } else if (i < sent.batch.failures.size() &&
+ sent.batch.failures[i].cancelled) {
+ ++sent.cancelled; // still staged; not the pack's fault
+ } else {
+ ++sent.failed; // still staged too
+ }
+ }
+ std::lock_guard lock(state_mutex_);
+ state_.uploaded_packs += sent.to_index.size();
+ state_.uploaded_bytes += uploaded_bytes;
+ state_.upload_failures += sent.failed;
+ state_.cancelled_uploads += sent.cancelled;
+ return sent;
+}
+
+void CaptureStorageService::let_go_of_adoption() {
+ // A sibling half adopted keeps what is left of it, for the next process
+ // on the node; what was uploaded from it was indexed, or is owed. Its
+ // directory is not removed: only finish_adoption() removes one, once no
+ // pack of it is left.
+ adopting_.reset();
+ adoption_queue_.clear();
+ std::lock_guard lock(state_mutex_);
+ state_.adoption_owed = adoption_owed();
+}
+
+bool CaptureStorageService::adoption_owed() const {
+ return adoption_scan_owed_ || adopting_ != nullptr ||
+ !adoption_queue_.empty();
+}
+
+bool CaptureStorageService::stop_requested() {
+ std::lock_guard lock(wake_mutex_);
+ return stop_requested_;
+}
+
+bool CaptureStorageService::adopt_step(uint64_t deadline_ns, size_t* deferred,
+ bool* cut_short) {
+ *cut_short = false;
+ const uint64_t started = steady_ns();
+ if (adopting_ == nullptr && adoption_queue_.empty()) {
+ // Look at the siblings when that is owed (from start() on) or, while
+ // the last look found a live one, again on the recheck interval.
+ const bool recheck_due =
+ live_siblings_ && config_.adoption_recheck_interval_ns > 0 &&
+ started - last_adoption_scan_ns_ >=
+ config_.adoption_recheck_interval_ns;
+ if (!adoption_scan_owed_ && !recheck_due) return true;
+ if (!scan_siblings()) return false;
+ }
+ bool ok = true;
+ // Whether this call has taken a step yet: a lock taken and a sweep, one
+ // pack of a listing validated, a round, or a finish. Each call takes one
+ // at least, so flushes that keep coming slow adoption down but never stop
+ // it.
+ bool stepped = false;
+ while (!stop_requested()) {
+ // The service's uploads' Cancellation is adoption's too: stop() cuts
+ // it between steps -- between the packs a listing validates too -- and
+ // in a round (its uploads, and, through the reads' Cancellation, its
+ // index reads).
+ if (upload_cancel_.cancelled()) {
+ *cut_short = true;
+ break;
+ }
+ // A flush() is waiting for the cycle: the rest of the slice is the
+ // next cycle's, which the loop runs once the flush is done.
+ if (stepped &&
+ flushes_in_progress_.load(std::memory_order_acquire) > 0) {
+ break;
+ }
+ stepped = true;
+ if (adopting_ == nullptr) {
+ if (adoption_queue_.empty()) break;
+ const std::string next = adoption_queue_.front();
+ adoption_queue_.pop_front();
+ begin_adoption(next);
+ continue;
+ }
+ if (!adopting_->validated) {
+ // One pack of the listing: validating hashes every byte of it, so
+ // over a dead backlog a whole listing takes as long as the backlog
+ // is big, and neither a flush nor stop() waits for more than the
+ // pack in flight. The slice bounds a listing as it does the rounds.
+ Adoption& adoption = *adopting_;
+ if (adoption.spool.ContinueRecovery(&adoption.recovery)) {
+ // Each pack's identity and object key come from the pack and its
+ // path in the dead directory, exactly as its owner would have
+ // uploaded it.
+ adoption.remaining.assign(
+ std::make_move_iterator(adoption.recovery.valid.begin()),
+ std::make_move_iterator(adoption.recovery.valid.end()));
+ adoption.recovery = dmi_store::SpoolRecovery{};
+ adoption.validated = true;
+ }
+ if (steady_ns() - started >= config_.adoption_slice_ns) break;
+ continue;
+ }
+ if (adopting_->remaining.empty()) {
+ finish_adoption();
+ continue;
+ }
+ // The rules of the cycle's uploads, before every round: none without
+ // the lease, and none while an uploaded pack is still owed to the
+ // catalog -- the dead spool is where the rest are durable.
+ bool catalog = false;
+ {
+ LeaseScope lease(this);
+ catalog = writer_.held_lease() != nullptr;
+ }
+ if (!catalog || !pending_index_.empty()) break;
+ bool cut = false;
+ if (!upload_adopted_round(deadline_ns, deferred, &cut)) {
+ ok = false; // the failed packs stay in the dead spool
+ break;
+ }
+ if (cut) {
+ *cut_short = true; // the cut packs stay in the dead spool
+ break;
+ }
+ if (steady_ns() - started >= config_.adoption_slice_ns) break;
+ }
+ return ok;
+}
+
+bool CaptureStorageService::scan_siblings() {
+ namespace fs = std::filesystem;
+ const fs::path own(spool_.root());
+ std::vector siblings;
+ std::vector claim_staging;
+ std::error_code ec;
+ for (fs::directory_iterator it(own.parent_path(), ec), end;
+ !ec && it != end; it.increment(ec)) {
+ uint64_t rank = 0;
+ std::string incarnation;
+ std::error_code type_ec;
+ if (it->path() == own || it->is_symlink(type_ec) ||
+ !it->is_directory(type_ec)) {
+ continue;
+ }
+ const std::string name = it->path().filename().string();
+ if (dmi_store::IsSpoolClaimStagingName(name)) {
+ claim_staging.push_back(it->path());
+ continue;
+ }
+ if (!dmi_store::ParseSpoolRankDirectoryName(name, &rank, &incarnation) ||
+ blocked_siblings_.count(it->path().string()) != 0) {
+ continue;
+ }
+ siblings.push_back(it->path().string());
+ }
+ if (ec) {
+ record_error("adoption: cannot list " + own.parent_path().string() +
+ ": " + ec.message());
+ return false;
+ }
+ clear_dead_claim_staging(claim_staging);
+ std::sort(siblings.begin(), siblings.end());
+ uint64_t live = 0;
+ for (const std::string& sibling : siblings) {
+ // A live owner answers a non-blocking probe at once; TryAdopt would
+ // retry for a few milliseconds first, on every recheck.
+ if (dmi_store::ReadSpoolOwner(sibling, nullptr)) {
+ ++live;
+ } else {
+ adoption_queue_.push_back(sibling);
+ }
+ }
+ adoption_scan_owed_ = false;
+ live_siblings_ = live != 0;
+ last_adoption_scan_ns_ = steady_ns();
+ std::lock_guard state(state_mutex_);
+ state_.live_siblings = live;
+ return true;
+}
+
+void CaptureStorageService::begin_adoption(const std::string& directory) {
+ auto adoption = std::make_unique();
+ adoption->directory = directory;
+ std::string error;
+ const dmi_store::SpoolStatus locked =
+ dmi_store::SpoolOwnerLock::TryAdopt(directory, &adoption->lock, &error);
+ if (locked == dmi_store::SpoolStatus::kOwned) {
+ // Alive after all (taken since the look): its owner's, and not owed.
+ live_siblings_ = true;
+ std::lock_guard state(state_mutex_);
+ ++state_.live_siblings;
+ return;
+ }
+ if (locked != dmi_store::SpoolStatus::kOk) {
+ // Another adopter drained and removed it meanwhile: nothing is owed.
+ if (!std::filesystem::exists(directory)) return;
+ // Its lock cannot be taken at all -- another user's lock file, say.
+ block_sibling(directory, "cannot lock it: " + error);
+ return;
+ }
+ dmi_store::SpoolConfig config{directory, config_.spool_max_bytes};
+ config.owner_lock = dmi_store::OwnerLock::kHeldByCaller; // adoption->lock
+ config.allow_shared_filesystem = config_.spool_allow_shared_filesystem;
+ // Its .open files swept and its ready packs listed, none hashed yet:
+ // adopt_step validates them, a pack a step.
+ if (dmi_store::Spool::Open(config, &adoption->spool, &error) !=
+ dmi_store::SpoolStatus::kOk ||
+ adoption->spool.BeginRecovery(&adoption->recovery, &error) !=
+ dmi_store::SpoolStatus::kOk) {
+ block_sibling(directory, "cannot open it: " + error, &adoption->lock);
+ return;
+ }
+ adopting_ = std::move(adoption);
+}
+
+bool CaptureStorageService::upload_adopted_round(uint64_t deadline_ns,
+ size_t* deferred,
+ bool* cut) {
+ *cut = false;
+ Adoption& adoption = *adopting_;
+ // A round is a chunk of the service's own upload path (upload_chunk, then
+ // index_or_owe): uploaded, then indexed before the next is uploaded, so
+ // at most one chunk (indexer.max_packs) of it is out of the dead spool
+ // and not yet in the catalog. It is at most uploader.max_workers packs as
+ // well, since the adoption slice is checked between rounds.
+ const size_t round = static_cast(std::min(
+ std::max(1, config_.uploader.max_workers),
+ std::max(1, config_.indexer.max_packs)));
+ std::vector entries;
+ while (!adoption.remaining.empty() && entries.size() < round) {
+ entries.push_back(std::move(adoption.remaining.front()));
+ adoption.remaining.pop_front();
+ }
+ // Through the service's upload client and its Cancellation, so stop()
+ // cuts an adoption's transfers, retries and backoff as it does the
+ // service's own.
+ dmi_store::SpoolUploader uploader(&adoption.spool, &upload_s3_,
+ config_.uploader);
+ uploader.set_cancellation(&upload_cancel_);
+ ChunkOutcome sent = upload_chunk(&uploader, entries);
+ size_t retryable = 0;
+ std::vector cancelled;
+ for (size_t i = 0; i < sent.batch.refs.size(); ++i) {
+ if (!sent.batch.refs[i].pack_id.empty()) continue;
+ // Still in the dead spool either way. One a cancel cut short goes
+ // back to the front, as it was; one a later try could upload is
+ // retried by a later cycle; one no try by this service can is not,
+ // and blocks the directory once the rest are up.
+ const dmi_store::UploadFailure failure =
+ i < sent.batch.failures.size() ? sent.batch.failures[i]
+ : dmi_store::UploadFailure{};
+ const std::string what = failure.object_key + ": " + failure.error;
+ if (failure.cancelled) {
+ cancelled.push_back(entries[i]);
+ } else if (failure.retryable) {
+ ++retryable;
+ adoption.remaining.push_back(entries[i]);
+ record_error("adopting dead spool " + adoption.directory +
+ ": upload failed for " + what);
+ } else if (adoption.blocked.empty()) {
+ adoption.blocked = "it holds a pack this service can never upload, " +
+ what;
+ }
+ }
+ adoption.remaining.insert(adoption.remaining.begin(), cancelled.begin(),
+ cancelled.end());
+ {
+ std::lock_guard state(state_mutex_);
+ state_.adopted_packs += sent.to_index.size();
+ }
+ *cut = !cancelled.empty();
+ // Uploaded, so gone from the dead spool: indexed now, a cancel of the
+ // uploads or not, or owed in pending_index_ like any pack of this
+ // service's own. After the requeue above, since a lost lease throws out
+ // of here and the packs still in the dead spool must stay in remaining.
+ *deferred += index_or_owe(std::move(sent.to_index), true, deadline_ns);
+ return retryable == 0;
+}
+
+void CaptureStorageService::finish_adoption() {
+ std::unique_ptr adoption = std::move(adopting_);
+ const std::string directory = adoption->directory;
+ if (!adoption->blocked.empty()) {
+ // The directory stays, and its lock goes.
+ block_sibling(directory, adoption->blocked, &adoption->lock);
+ return;
+ }
+ {
+ std::lock_guard state(state_mutex_);
+ ++state_.adopted_spools;
+ }
+ std::string error;
+ if (!adoption->lock.ReleaseAndRemoveIfEmpty(&error)) {
+ // Nothing to upload is left, only files that are not its packs (a
+ // quarantined one, a spool directory nested in it): the directory
+ // stays for someone to look at.
+ block_sibling(directory, "it was drained, but still holds files that "
+ "are not its packs (a quarantined pack, or a "
+ "spool directory nested in it)");
+ }
+}
+
+void CaptureStorageService::block_sibling(const std::string& directory,
+ const std::string& reason,
+ dmi_store::SpoolOwnerLock* lock) {
+ if (lock != nullptr) {
+ // Said in its lock file, while this service still holds it, then let
+ // go of: for a person, and for the sinks on the node, which charge a
+ // dead directory against their budget only while an adoption can
+ // drain it (SpoolConfig::charge_dead_siblings). A later take -- by a
+ // process that can adopt it -- rewrites the mark.
+ lock->MarkBlocked(reason);
+ lock->Release();
+ }
+ blocked_siblings_.insert(directory);
+ record_error("dead spool " + directory + " is left in place, not to be "
+ "adopted by this service: " + reason);
+ std::lock_guard state(state_mutex_);
+ state_.blocked_siblings.assign(blocked_siblings_.begin(),
+ blocked_siblings_.end());
+}
+
+void CaptureStorageService::clear_dead_claim_staging(
+ const std::vector& staging) {
+ namespace fs = std::filesystem;
+ // A claim builds its directory's staging copy and renames it into place
+ // within milliseconds, so one this old whose lock nobody holds belongs
+ // to a claim that died before its rename. Younger ones are left alone: a
+ // claim between its mkdir and its flock holds no lock yet, and clearing
+ // its copy would fail it. Nothing is owed either way.
+ constexpr auto kDeadAfter = std::chrono::seconds(60);
+ for (const fs::path& path : staging) {
+ std::error_code ec;
+ const auto written = fs::last_write_time(path, ec);
+ if (ec || fs::file_time_type::clock::now() - written < kDeadAfter) {
+ continue;
+ }
+ dmi_store::SpoolOwnerLock lock;
+ std::string error;
+ if (dmi_store::SpoolOwnerLock::TryAdopt(path.string(), &lock, &error) ==
+ dmi_store::SpoolStatus::kOk) {
+ lock.ReleaseAndRemoveIfEmpty(&error);
+ }
+ }
+}
+
size_t CaptureStorageService::index_bounded(std::vector refs,
std::vector* unindexed,
uint64_t deadline_ns) {
@@ -1312,6 +1751,13 @@ void CaptureStorageService::latch_failure(std::exception_ptr failure,
std::fprintf(stderr, "dmi capture storage: indexing stopped: %s\n",
line.c_str());
std::fflush(stderr);
+ // Woken now, not after its wait (up to max_backoff_ns in an outage), so
+ // the loop lets go of a sibling it was adopting at once.
+ {
+ std::lock_guard lock(wake_mutex_);
+ kick_ = true;
+ }
+ wake_.notify_all();
}
} // namespace dmi_catalog
diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h
index 54c3fa1f0..c1b95a28d 100644
--- a/native/csrc/catalog/storage_service.h
+++ b/native/csrc/catalog/storage_service.h
@@ -23,6 +23,14 @@
// already committed, and indexes the rest. It runs at start(), and
// periodically when reconcile_interval_ns is non-zero.
//
+// One owner per spool directory. The service's spool is opened under the
+// directory's owner lock (store/spool.h), so a second process on it is
+// refused at construction, naming the holder. With adopt_sibling_spools
+// the service's directory is one rank directory of the plan's section 2.3
+// layout, and once started its loop adopts the siblings whose owner has
+// died: a crashed process's spool is recovered by the next process on the
+// node for the same catalog, whatever run it belongs to.
+//
// Deployment shape: the service holds the catalog's single publisher lease, so
// run ONE service per (database, table_prefix). A second one waits up to
// start_lease_wait_ns for the lease at start(), then fails naming the holder.
@@ -85,10 +93,13 @@
#include
#include
#include
+#include
#include
+#include
#include