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