Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
11 changes: 6 additions & 5 deletions native/csrc/catalog/storage_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<std::mutex> state(state_mutex_);
state_.live_siblings = live;
Expand All @@ -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<std::mutex> state(state_mutex_);
++state_.live_siblings;
return;
Expand Down
18 changes: 9 additions & 9 deletions native/csrc/catalog/storage_service.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<std::string> adoption_queue_;
std::unique_ptr<Adoption> adopting_;
std::set<std::string> 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.
Expand Down
3 changes: 2 additions & 1 deletion native/csrc/store/spool.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
};

Expand Down
52 changes: 52 additions & 0 deletions tests/test_native_spool_adoption_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading