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
78 changes: 60 additions & 18 deletions native/csrc/catalog/storage_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -344,18 +344,20 @@ void CaptureStorageService::sweep_and_reconcile_at_start() {
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
// there, taking a fresh lease once it can (or latching after 2 x TTL of a
// rival). The pass it cut short is owed, and the loop runs it once it
// holds a lease again.
// A pass that did not finish is not fatal -- the bucket is still there --
// but it is owed, whatever stopped it: a listing or a HEAD the object
// store failed, a catalog error, or a lease lost meanwhile, to a
// quarantine or to another holder. The loop runs it again (reconcile_or_owe)
// once it holds a lease -- taking a fresh one once it can, or latching
// after 2 x TTL of a rival, as it does for the running service. Owing
// only the lease loss left the rest for the next process start, with
// reconcile_interval_ns 0, the default: what a crashed process uploaded
// and never indexed stayed out of the catalog meanwhile.
if (config_.reconcile_on_start) {
try {
reconcile();
reconcile_or_owe();
} catch (const CatalogError& exc) {
if (is_lease_refusal(exc)) {
reconcile_owed_ = true;
record_error(std::string("reconcile at start lost the publisher "
"lease; the loop reconciles once it holds "
"one again: ") +
Expand Down Expand Up @@ -691,17 +693,19 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle(
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.
// 4. Reconcile on its interval, or when a pass that did not finish is
// owed -- start()'s or an earlier cycle's -- once its retry backoff
// is over: in the loop's cycles only, and not once a cancel came.
// The backoff is the pass's own, so an object that never answers a
// HEAD costs a listing of the bucket every max_backoff_ns at most,
// and does not slow the uploads down. The lease thread keeps the
// lease alive.
if (allow_reconcile && catalog && !upload_cancel_.cancelled() &&
steady_ns() >= reconcile_retry_ns_ &&
(reconcile_owed_ ||
(config_.reconcile_interval_ns > 0 &&
steady_ns() - last_reconcile_ns_ >= config_.reconcile_interval_ns))) {
if (reconcile()) {
reconcile_owed_ = false;
last_reconcile_ns_ = steady_ns();
}
reconcile_or_owe();
}

// Drained: every staged pack uploaded, every uploaded pack in the
Expand Down Expand Up @@ -1281,16 +1285,52 @@ void CaptureStorageService::reject(const PackRefData& ref,
++state_.index_failures;
}

bool CaptureStorageService::reconcile_or_owe() {
bool finished = false;
try {
finished = reconcile();
} catch (...) {
owe_reconcile();
throw;
}
if (!finished) {
owe_reconcile();
return false;
}
reconcile_owed_ = false;
reconcile_failures_ = 0;
reconcile_retry_ns_ = 0;
last_reconcile_ns_ = steady_ns();
return true;
}

void CaptureStorageService::owe_reconcile() {
reconcile_owed_ = true;
reconcile_failures_ = std::min(reconcile_failures_ + 1, 32);
// poll_interval * 2^failures, capped, as the loop's own backoff.
uint64_t wait_ns = config_.poll_interval_ns;
for (int i = 0; i < reconcile_failures_ && wait_ns < config_.max_backoff_ns;
++i) {
wait_ns *= 2;
}
wait_ns = std::min(std::max(wait_ns, config_.poll_interval_ns),
std::max(config_.max_backoff_ns, config_.poll_interval_ns));
reconcile_retry_ns_ = steady_ns() + wait_ns;
}

bool CaptureStorageService::reconcile() {
// List every pack under the prefix, ask the catalog which it already
// committed, and HEAD only the rest -- so a steady-state pass over a large
// bucket costs listing pages and one catalog query per page, not a HEAD per
// object. stop() ends it between requests: what it has not reached is
// still in the bucket, for the next pass.
// still in the bucket, for the next pass. So is a pack it could not HEAD
// or index, and then the pass reports it did not finish, so that a next
// pass runs.
std::string token;
uint64_t found = 0;
uint64_t skipped = 0;
uint64_t head_errors = 0;
bool all_indexed = true;
do {
if (upload_cancel_.cancelled()) return false;
dmi_store::ListResult page;
Expand Down Expand Up @@ -1364,17 +1404,19 @@ bool CaptureStorageService::reconcile() {
}
found += missing.size();
// A pack that fails here is still uncommitted in the bucket, so the next
// pass retries it; it is not added to the flush boundary.
// pass retries it; it is not added to the flush boundary. One set aside
// (reject) is not left unindexed: no pass could index it.
std::vector<PackRefData> unindexed;
if (!missing.empty()) index_bounded(std::move(missing), &unindexed);
if (!unindexed.empty()) all_indexed = false;
} while (!token.empty());

std::lock_guard<std::mutex> lock(state_mutex_);
++state_.reconcile_passes;
state_.reconciled_packs += found;
state_.reconcile_skipped_objects += skipped;
state_.reconcile_head_errors += head_errors;
return true;
return head_errors == 0 && all_indexed;
}

void CaptureStorageService::keep_lease() {
Expand Down
20 changes: 17 additions & 3 deletions native/csrc/catalog/storage_service.h
Original file line number Diff line number Diff line change
Expand Up @@ -497,8 +497,17 @@ class CaptureStorageService {
size_t index_bounded(std::vector<PackRefData> refs,
std::vector<PackRefData>* unindexed,
uint64_t deadline_ns = 0);
// False when stop() cut it short, between two of its requests.
// False when the pass did not finish: stop() cut it short, between two
// of its requests, or a pack it listed could not be read (HEAD) or
// indexed. A pass that threw did not finish either.
bool reconcile();
// Runs a pass and books it: one that finished clears what was owed; one
// that did not -- false, or thrown, which it rethrows -- is owed
// (owe_reconcile). Requires cycle_mutex_, or start() before the loop runs.
bool reconcile_or_owe();
// The loop runs an owed pass again no sooner than poll_interval_ns
// doubled per pass that did not finish, capped at max_backoff_ns.
void owe_reconcile();
void keep_lease(); // the lease thread's body
void renew_lease_if_due(); // requires lease_mutex_
// LeaseScope's before_request hook: renew_lease_if_due() before each
Expand Down Expand Up @@ -574,9 +583,14 @@ class CaptureStorageService {
// hammer the lease table. Guarded by lease_mutex_.
uint64_t next_claim_ns_ = 0;
uint64_t last_reconcile_ns_ = 0;
// The reconcile at start() lost the lease before it finished; the loop
// runs one once it holds a lease again. Guarded by cycle_mutex_.
// A reconcile pass -- start()'s or the loop's -- did not finish; the loop
// runs one once it holds a lease, from reconcile_retry_ns_ on. Without
// this, a pass the object store or the catalog failed waited for the
// next process start when reconcile_interval_ns is 0, the default.
// Guarded by cycle_mutex_.
bool reconcile_owed_ = false;
uint64_t reconcile_retry_ns_ = 0;
int reconcile_failures_ = 0; // consecutive passes that did not finish
// 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
Expand Down
86 changes: 86 additions & 0 deletions tests/test_native_capture_storage_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -802,6 +802,92 @@ def test_a_failed_head_is_an_error_not_a_foreign_object(fake_s3, tmp_path):
assert "HEAD" in snapshot["last_error"] and key in snapshot["last_error"]


def test_a_reconcile_that_failed_at_start_is_run_again_by_the_loop(
fake_s3, tmp_path):
"""Only a lease refusal left start()'s pass owed. A listing the object
store failed was reported and dropped, so with reconcile_interval_s 0,
the default, a pack a crashed process uploaded and never indexed waited
for the next process start. Now any pass that did not finish is owed,
and the loop runs it again."""
from dmi.storage.native_capture import _load_native_store_extension

spool_root = tmp_path / "spool"
tensors = _stage(spool_root, range(3), records_per_pack=3)
# The crash window: uploaded and gone from the spool, never indexed.
store = _Driver(STORE_DRIVER)
try:
uploaded = store.call(
op="upload_pending", endpoint=fake_s3, bucket=BUCKET,
region=REGION, access=ACCESS, secret=SECRET, token=None,
insecure=True, connect_timeout=5, read_timeout=15, max_attempts=4,
store_id="s3", root=str(spool_root), spool_max_bytes=1 << 40,
limit=-1, max_workers=4, max_in_flight_bytes=1 << 30)
assert uploaded["ok"], uploaded
finally:
store.close()
assert _ready(spool_root) == []

s3 = _Switch.to_url(fake_s3)
with _catalog() as (_client, catalog):
native = _storage_config(s3.url, catalog.table_prefix)._native_dict()
native.update(
spool_root=str(spool_root), spool_owner_lock=_held(spool_root),
holder="owed-reconcile", poll_interval_ns=20_000_000,
s3_max_attempts=1)
assert native.get("reconcile_interval_ns", 0) == 0 # no periodic pass
service = _load_native_store_extension().StorageService(native)
s3.cut()
service.start() # its reconcile cannot list the bucket
try:
snapshot = service.snapshot()
assert snapshot["reconcile_passes"] == 0, snapshot
assert "listing failed" in snapshot["last_error"], snapshot
s3.restore()
_wait_for(lambda: service.snapshot()["reconciled_packs"] == 1)
snapshot = service.snapshot()
finally:
service.stop()
s3.close()

assert snapshot["reconcile_passes"] == 1, snapshot
direct = _storage_config(fake_s3, catalog.table_prefix)
assert sorted(_read_all(direct)) == sorted(tensors)


def test_a_reconcile_that_missed_a_pack_is_run_again_backing_off(
fake_s3, tmp_path):
"""A pass that could not HEAD a pack finished, and cleared what it owed:
"the next pass retries it" held only with a periodic pass configured.
The loop now runs such a pass again, on a backoff of its own, so an
object that never answers costs a listing every max_backoff_ns at most,
and the loop's uploads do not slow down for it."""
from dmi.storage.native_capture import _load_native_store_extension

key = f"fault/forbidden/{uuid.uuid4()}.dmi-pack"
with STATE.lock:
STATE.objects[key] = {"body": b"x", "meta": {}, "content_type": "",
"etag": '"0"'}
with _catalog() as (_client, catalog):
native = _storage_config(fake_s3, catalog.table_prefix)._native_dict()
native.update(spool_root=str(tmp_path / "spool"),
spool_owner_lock=_held(tmp_path / "spool"),
holder="head-retry-test", reconcile_prefix="fault/",
s3_max_attempts=1, poll_interval_ns=20_000_000,
max_backoff_ns=10_000_000_000)
service = _load_native_store_extension().StorageService(native)
service.start() # reconciles once, and cannot HEAD the object
try:
time.sleep(3.0)
snapshot = service.snapshot()
finally:
service.stop()

# 20 ms doubling reaches ~2.5 s of waits in 6 retries; a retry every
# 20 ms poll would have run ~150.
assert 3 <= snapshot["reconcile_head_errors"] <= 15, snapshot
assert snapshot["cycles"] >= 50, snapshot # the loop kept its own poll


def test_one_publisher_per_catalog(fake_s3, tmp_path):
with _catalog() as (_client, catalog):
# No start wait: the refusal is the point here, and waiting out the
Expand Down
Loading