From b6db813852829497ddc27d979dff8656f55b210f Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Tue, 6 Oct 2026 15:58:30 -0400 Subject: [PATCH] Look at the sibling spools again whether or not the last look found a live one adopt_step re-scanned the siblings on adoption_recheck_interval_ns only while its last scan had found a live one. A sibling claimed after that scan -- a rank or a restart that starts later than this service -- and then left with packs by a crash was never adopted while the service ran, and Spool's charge_dead_siblings kept charging it against the live sink's budget until the next restart on the node, contrary to spool.h's "the capacity comes back as adoption drains them". Rescan on the interval regardless. A pass lists the parent directory and probes each live sibling's lock without blocking, so the cost stays one directory listing per interval. live_siblings_ had no other reader and goes; the snapshot's live_siblings count is unchanged. --- native/csrc/catalog/storage_service.cpp | 11 ++--- native/csrc/catalog/storage_service.h | 18 ++++---- native/csrc/store/spool.h | 3 +- tests/test_native_spool_adoption_live.py | 52 ++++++++++++++++++++++++ 4 files changed, 69 insertions(+), 15 deletions(-) diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 16c747e88..fdc1e21e0 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -823,10 +823,13 @@ bool CaptureStorageService::adopt_step(uint64_t deadline_ns, size_t* deferred, *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. + // Look at the siblings when that is owed (from start() on), and again + // on the recheck interval -- whatever the last look found: a sibling + // claimed after it, by a rank or a restart that then dies, is dead and + // charged against this spool's budget without ever having been seen + // alive. const bool recheck_due = - live_siblings_ && config_.adoption_recheck_interval_ns > 0 && + config_.adoption_recheck_interval_ns > 0 && started - last_adoption_scan_ns_ >= config_.adoption_recheck_interval_ns; if (!adoption_scan_owed_ && !recheck_due) return true; @@ -951,7 +954,6 @@ bool CaptureStorageService::scan_siblings() { } } adoption_scan_owed_ = false; - live_siblings_ = live != 0; last_adoption_scan_ns_ = steady_ns(); std::lock_guard state(state_mutex_); state_.live_siblings = live; @@ -966,7 +968,6 @@ void CaptureStorageService::begin_adoption(const std::string& directory) { 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; diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index c1b95a28d..17e970a98 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -174,12 +174,13 @@ struct StorageServiceConfig { // are owed like the service's own. Live siblings -- another rank or job on this // node, a predecessor still closing -- are left alone, and are not owed. bool adopt_sibling_spools = false; - // While a pass found a live sibling, the loop passes over the siblings - // again this often, so one whose owner dies later -- a predecessor that - // was still inside close() when this service started, a rank that - // crashes while this one runs -- is adopted then, not at the next - // restart on the node. A live sibling costs one non-blocking lock probe - // per pass. 0 never looks again after the first pass. + // The loop passes over the siblings again this often, so one whose + // owner dies later -- a predecessor that was still inside close() when + // this service started, a rank that crashes while this one runs, or one + // that claimed its directory after the last pass and died -- is adopted + // then, not at the next restart on the node. A pass lists the parent + // directory and costs one non-blocking lock probe per live sibling. + // 0 never looks again after the first pass. uint64_t adoption_recheck_interval_ns = 30'000'000'000ull; // How long one cycle may spend adopting before it lets go of the cycle -- // to the service's own uploads, stop() -- and carries on in the next. @@ -579,13 +580,12 @@ class CaptureStorageService { bool reconcile_owed_ = false; // Adoption's state, guarded by cycle_mutex_: whether a look at the // siblings is due (from start() on), the dead ones the last look found, - // the one being adopted, whether the last look found a live one and when - // it ran (the loop looks again every adoption_recheck_interval_ns). + // the one being adopted, and when the last look ran (the loop looks + // again every adoption_recheck_interval_ns). bool adoption_scan_owed_ = false; std::deque adoption_queue_; std::unique_ptr adopting_; std::set blocked_siblings_; - bool live_siblings_ = false; uint64_t last_adoption_scan_ns_ = 0; // flush() calls in progress. A cycle adopting lets go of the cycle at // its next step while one is, so a flush never waits out a slice. diff --git a/native/csrc/store/spool.h b/native/csrc/store/spool.h index 31b4a10db..2331f629d 100644 --- a/native/csrc/store/spool.h +++ b/native/csrc/store/spool.h @@ -104,7 +104,8 @@ struct SpoolConfig { // the node's spool; before the layout every restart reused one directory // and one budget. The charge is refreshed wherever the committed account // is (Open, and before a stage is refused), so the capacity comes back as - // adoption drains them. + // adoption drains them; the service finds a sibling that died after its + // last look at its next one (adoption_recheck_interval_ns). bool charge_dead_siblings = false; }; diff --git a/tests/test_native_spool_adoption_live.py b/tests/test_native_spool_adoption_live.py index fceccff69..c539c5e50 100644 --- a/tests/test_native_spool_adoption_live.py +++ b/tests/test_native_spool_adoption_live.py @@ -633,6 +633,58 @@ def test_a_sibling_whose_owner_dies_after_start_is_adopted_by_a_recheck( lock.release_and_remove_if_empty() +def test_a_sibling_that_appears_after_start_and_dies_is_adopted( + fake_s3, tmp_path): + """The siblings are looked at again on the recheck interval whether or + not the last look found a live one. A rank or a restart that claims its + directory after this service's first look, stages, and dies is adopted + while the service runs -- not left, charged against the live sink's + budget, until the next restart on the node. Under per-rank services + that is any rank that starts later than this one and crashes.""" + base = tmp_path / "spool" + with _catalog() as prefix: + config = _storage_config(fake_s3, prefix) + lock = _claim(base, config) + native = config._native_dict() + native.update( + spool_root=lock.directory, spool_max_bytes=1 << 30, + holder="late-sibling-test", poll_interval_ns=50_000_000, + sweep_spool_on_start=True, spool_owner_lock="held_by_caller", + adopt_sibling_spools=True, + adoption_recheck_interval_ns=200_000_000, + **config._lease_native()) + service = _store().StorageService(native) + service.start() + try: + # The first look finds no sibling at all. + _wait_for(lambda: not service.snapshot()["adoption_owed"] + and service.snapshot()["cycles"] >= 1) + snapshot = service.snapshot() + assert snapshot["live_siblings"] == 0, snapshot + assert snapshot["adopted_spools"] == 0, snapshot + + sibling = _claim(base, config) # a rank that starts later + sibling_directory = Path(sibling.directory) + _stage_into(sibling.directory, STAGED_BY_THE_DEAD) + staged = sorted(sibling_directory.rglob("*.dmi-pack.ready")) + assert len(staged) == len(STAGED_BY_THE_DEAD) // RECORDS_PER_PACK + sibling.release() # and dies with its packs staged + + _wait_for(_adopted(service, 1)) + snapshot = service.snapshot() + assert snapshot["adopted_packs"] == len(staged), snapshot + assert not sibling_directory.exists() + assert service.flush(60.0) + expected = { + capture_id: tensor.contiguous().view(-1).numpy().tobytes() + for capture_id, tensor in _envelope( + STAGED_BY_THE_DEAD).expected.items()} + assert _read_all(config) == expected + finally: + service.stop() + lock.release_and_remove_if_empty() + + def test_flush_adopts_nothing_while_the_loop_is_idle(fake_s3, tmp_path): """flush() covers this process's records, and its cycles never adopt: a dead backlog is the loop's. With the object store down, a flush that