From 89e5619b01128a4b548bb3cbf835a9f1087b213c Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 11:45:56 -0400 Subject: [PATCH 1/6] Give NativeCaptureStorageConfig the publisher lease knobs The storage service never passed a lease TTL or publish timeout to the native writer, so every capture process ran on the writer's built-in 30 s lease: a crashed job held its catalog for up to 30 s, and nothing let a deployment choose otherwise. The config now carries lease_ttl_s (default 15), publish_timeout_s (5), clock_skew_s (0) and start_lease_wait_s (None, meaning lease_ttl_s + publish_timeout_s), and NativeCaptureStorage hands them to the service in nanoseconds. The reader takes no lease and is not given them. Validation is the native writer's rule, checked in nanoseconds so that anything accepted here is accepted at start(): the TTL must exceed publish_timeout_s + clock_skew_s by at least 0.1 s (equality is allowed, as catalog_writer.cpp allows it), and publish_timeout_s must be whole seconds because it is sent as max_execution_time. clock_skew_s is not in the milestone's field list; it is added because it is a term of that rule and the binding already reads clock_skew_ns. The binding already reads lease_ttl_ns, publish_timeout_ns and clock_skew_ns. start_lease_wait_ns is read from the next commit, which makes start() wait for an expiring predecessor. Evidence: tests/test_native_capture_storage_wiring.py, 30 new cases failing (TypeError on the unknown fields, KeyError on the service config) before the change; 67 passed after. --- src/dmi/storage/native_capture.py | 79 +++++++++++++- tests/test_native_capture_storage_wiring.py | 113 ++++++++++++++++++++ 2 files changed, 190 insertions(+), 2 deletions(-) diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index db2c90b6e..73b9a0c6a 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -19,7 +19,8 @@ Deployment shape: the catalog has ONE publisher lease per (``database``, ``table_prefix``), so run one capture process per catalog. A second engine -on the same catalog fails at ``create_record_runtime`` with the lease held. +on the same catalog waits ``start_lease_wait_s`` for the lease, then fails +at ``create_record_runtime`` with the lease held, naming the holder. """ from __future__ import annotations @@ -55,6 +56,22 @@ def _load_native_store_extension() -> Any: ) from exc +def _finite(name: str, value: Any) -> None: + if type(value) not in (int, float): + raise TypeError(f"{name} must be float") + if not math.isfinite(value): + raise ValueError(f"{name} must be finite") + + +def _ns(seconds: float) -> int: + return int(round(seconds * 1_000_000_000)) + + +# The native writer's MINIMUM_FENCE_MARGIN_NS: what must remain of the lease +# once the publish statement cap and the skew bound are spent. +_FENCE_MARGIN_NS = 100_000_000 + + def _positive(name: str, value: Any, kind: type) -> None: if type(value) is not kind and not (kind is float and type(value) is int): raise TypeError(f"{name} must be {kind.__name__}") @@ -109,6 +126,21 @@ class NativeCaptureStorageConfig: # open pack, then getting every staged pack into the catalog. close_flush_timeout_s: float = 60.0 + # The publisher lease. A crashed process holds the catalog for up to + # lease_ttl_s; a ClickHouse error of unknown outcome sets the lease aside + # for as long, after which the service takes a fresh one. The TTL must + # exceed publish_timeout_s + clock_skew_s by at least 0.1 s: that margin + # is what keeps a publish statement from outliving its lease. + lease_ttl_s: float = 15.0 + # The server-side cap on each fenced publish statement, whole seconds. + publish_timeout_s: float = 5.0 + # The bound on host clock disagreement across a replicated catalog. + clock_skew_s: float = 0.0 + # How long start() waits for a predecessor's lease to expire before + # failing with it held. None waits lease_ttl_s + publish_timeout_s, + # enough to outlast a crashed predecessor; 0 fails at once. + start_lease_wait_s: Optional[float] = None + def __post_init__(self) -> None: for name in ("s3_endpoint", "s3_bucket", "s3_access_key", "s3_secret_key", "store_id", "database", "table_prefix", @@ -148,6 +180,44 @@ def __post_init__(self) -> None: raise ValueError("reconcile_interval_s must be finite") if self.reconcile_interval_s < 0: raise ValueError("reconcile_interval_s must be non-negative") + self._validate_lease() + + def _validate_lease(self) -> None: + for name in ("lease_ttl_s", "publish_timeout_s", "clock_skew_s"): + _finite(name, getattr(self, name)) + if self.start_lease_wait_s is not None: + _finite("start_lease_wait_s", self.start_lease_wait_s) + if self.start_lease_wait_s < 0: + raise ValueError("start_lease_wait_s must be non-negative") + _positive("lease_ttl_s", self.lease_ttl_s, float) + _positive("publish_timeout_s", self.publish_timeout_s, float) + if self.clock_skew_s < 0: + raise ValueError("clock_skew_s must be non-negative") + if float(self.publish_timeout_s) != int(self.publish_timeout_s): + raise ValueError( + "publish_timeout_s must be a whole number of seconds: the " + "catalog sends it as max_execution_time in seconds, where a " + "fraction truncates and 0 means no limit") + # In nanoseconds, as the native writer checks it, so a pairing + # accepted here is never refused at start(). + margin = (_ns(self.lease_ttl_s) - _ns(self.publish_timeout_s) + - _ns(self.clock_skew_s)) + if margin < _FENCE_MARGIN_NS: + raise ValueError( + "lease_ttl_s must exceed publish_timeout_s + clock_skew_s by " + "at least 0.1 s, or a publish statement can still be running " + "when its lease becomes takeable") + + def _lease_native(self) -> dict[str, int]: + wait = self.start_lease_wait_s + if wait is None: + wait = self.lease_ttl_s + self.publish_timeout_s + return { + "lease_ttl_ns": _ns(self.lease_ttl_s), + "publish_timeout_ns": _ns(self.publish_timeout_s), + "clock_skew_ns": _ns(self.clock_skew_s), + "start_lease_wait_ns": _ns(wait), + } def _native_dict(self) -> dict[str, Any]: return { @@ -193,12 +263,17 @@ def __init__( reconcile_prefix=config.reconcile_prefix, reconcile_interval_ns=int(config.reconcile_interval_s * 1e9), sweep_spool_on_start=sweep_spool, + **config._lease_native(), ) self._config = config self._service = module.StorageService(native) def start(self) -> None: - """Sweep the spool, ensure the catalog schema, take the lease.""" + """Ensure the catalog schema, take the lease, sweep the spool. + + Waits up to ``start_lease_wait_s`` for another holder's lease to + expire, then raises naming the holder. + """ self._service.start() def flush(self, timeout_s: float) -> None: diff --git a/tests/test_native_capture_storage_wiring.py b/tests/test_native_capture_storage_wiring.py index ce59315b4..4c247dc80 100644 --- a/tests/test_native_capture_storage_wiring.py +++ b/tests/test_native_capture_storage_wiring.py @@ -431,3 +431,116 @@ def test_replacing_a_record_ring_drains_and_stops_the_service( engine.create_record_runtime(_record_format()) assert len(services) == 2 assert engine._capture_storage is not None + + +# --- the publisher lease knobs ------------------------------------------------- + + +def test_lease_knobs_default_to_a_fifteen_second_lease(): + config = _storage_config() + assert config.lease_ttl_s == 15.0 + assert config.publish_timeout_s == 5.0 + assert config.clock_skew_s == 0.0 + # None waits out a predecessor for the TTL plus the publish timeout. + assert config.start_lease_wait_s is None + + +@pytest.mark.parametrize("ttl, publish, skew", [ + (5.0, 5.0, 0.0), # no margin at all + (5.09, 5.0, 0.0), # under the 0.1 s the renewed lease needs + (6.0, 5.0, 0.95), # the skew bound spends the rest + (3.0, 4.0, 0.0), # the statement cap outlives the lease +]) +def test_the_lease_must_outlast_the_publish_fence(ttl, publish, skew): + # The native writer's rule, checked here so a bad pairing fails at + # construction rather than at start(): the TTL must exceed the publish + # timeout plus the skew bound by at least 0.1 s. + with pytest.raises(ValueError, match="lease_ttl_s"): + _storage_config(lease_ttl_s=ttl, publish_timeout_s=publish, + clock_skew_s=skew) + + +@pytest.mark.parametrize("ttl, publish, skew", [ + (5.1, 5.0, 0.0), (3.0, 1.0, 0.0), (7.0, 5.0, 1.5), (2, 1, 0)]) +def test_a_lease_that_outlasts_the_fence_is_accepted(ttl, publish, skew): + config = _storage_config(lease_ttl_s=ttl, publish_timeout_s=publish, + clock_skew_s=skew) + assert config.lease_ttl_s == ttl + + +@pytest.mark.parametrize("publish", [1.5, 0.5, 4.999]) +def test_the_publish_timeout_is_whole_seconds(publish): + # The native writer sends it as max_execution_time in whole seconds; a + # fraction would truncate, and 0 is no limit at all. + with pytest.raises(ValueError, match="publish_timeout_s must be a whole"): + _storage_config(lease_ttl_s=30.0, publish_timeout_s=publish) + + +@pytest.mark.parametrize("name, value", [ + ("lease_ttl_s", 0.0), ("lease_ttl_s", -3.0), ("publish_timeout_s", 0), + ("publish_timeout_s", -1.0), ("clock_skew_s", -0.5), + ("start_lease_wait_s", -1.0)]) +def test_lease_knobs_refuse_out_of_range_values(name, value): + with pytest.raises(ValueError, match=name): + _storage_config(**{name: value}) + + +@pytest.mark.parametrize("name", [ + "lease_ttl_s", "publish_timeout_s", "clock_skew_s", "start_lease_wait_s"]) +@pytest.mark.parametrize("value", [float("inf"), float("nan")]) +def test_lease_knobs_must_be_finite(name, value): + with pytest.raises(ValueError, match=f"{name} must be finite"): + _storage_config(**{name: value}) + + +def test_a_zero_start_wait_is_accepted(): + assert _storage_config(start_lease_wait_s=0.0).start_lease_wait_s == 0.0 + + +def test_the_service_gets_the_lease_knobs_in_nanoseconds(monkeypatch, tmp_path): + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + + engine.create_record_runtime(_record_format()) + + config = services[0].config + assert config["lease_ttl_ns"] == 15_000_000_000 + assert config["publish_timeout_ns"] == 5_000_000_000 + assert config["clock_skew_ns"] == 0 + assert config["start_lease_wait_ns"] == 20_000_000_000 + + +def test_explicit_lease_knobs_reach_the_service(monkeypatch, tmp_path): + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + engine._capture_storage_config = _storage_config( + lease_ttl_s=3, publish_timeout_s=1.0, clock_skew_s=0.25, + start_lease_wait_s=0.5) + + engine.create_record_runtime(_record_format()) + + config = services[0].config + assert config["lease_ttl_ns"] == 3_000_000_000 + assert config["publish_timeout_ns"] == 1_000_000_000 + assert config["clock_skew_ns"] == 250_000_000 + assert config["start_lease_wait_ns"] == 500_000_000 + for name in ("lease_ttl_ns", "publish_timeout_ns", "clock_skew_ns", + "start_lease_wait_ns"): + assert type(config[name]) is int, name + + +def test_the_reader_is_not_handed_the_lease_knobs(monkeypatch): + # A reader takes no lease; the knobs are the service's alone. + from dmi.storage import native_capture + + seen = {} + + class _Reader: + def __init__(self, config): + seen.update(config) + + monkeypatch.setattr( + native_capture, "_load_native_store_extension", + lambda: SimpleNamespace(CaptureReader=_Reader, SEARCH_ITEM_COLUMNS=())) + native_capture.NativeCaptureReader(_storage_config(lease_ttl_s=3.0, + publish_timeout_s=1)) + assert not {"lease_ttl_ns", "publish_timeout_ns", "clock_skew_ns", + "start_lease_wait_ns"} & set(seen) From afb87fa85e05dc5a4c2f70d56f57cea6f5d3ec91 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 11:59:45 -0400 Subject: [PATCH 2/6] Wait for an expiring predecessor's lease at start, before sweeping A SIGKILLed capture process leaves its publisher lease live for up to a TTL, and CaptureStorageService::start() refused the lease at once, so a restart inside that window failed instead of waiting a few seconds. start() now retries a lease refusal (kHeld, or a contested claim) every ttl/10, clamped to 50-500 ms, for up to start_lease_wait_ns. A live publisher keeps renewing, so a genuine second publisher still fails, only later, with the refusal naming the holder and the time waited. The binding reads start_lease_wait_ns; NativeCaptureStorageConfig already sends it (lease_ttl_s + publish_timeout_s unless set). The C++ default stays 0, so a caller of the raw binding keeps the fail-at-once behaviour. start() also swept the spool before taking the lease. Recover() deletes every .open file its Spool does not own, so a second process pointed at a live process's spool deleted that sink's in-progress packs and only then learned the catalog was held. The order is now schema, lease, sweep, reconcile, as the plan's start order has it, and a failed sweep releases the lease it just took. Evidence, local ClickHouse 25.12, tests/test_native_capture_storage_live.py: - test_a_restart_within_the_ttl_of_a_killed_predecessor_succeeds (a subprocess holds the lease at TTL 5 s and is SIGKILLed): before, start() raised "held by 'crashed-publisher' ... for another 4784239704 ns"; after, it returns after waiting, within 1 s < elapsed < 8 s. - test_a_start_refused_the_lease_leaves_the_spool_unswept: before, the refused start deleted the .open file; after, it survives the refusal and is swept by the next start that gets the lease. - The whole file: 15 passed. test_one_publisher_per_catalog now takes ~20 s, the default wait for a lease its holder keeps renewing. --- native/csrc/catalog/bindings_store.cpp | 2 + native/csrc/catalog/storage_service.cpp | 59 ++++++++++++-- native/csrc/catalog/storage_service.h | 22 ++++-- tests/test_native_capture_storage_live.py | 95 +++++++++++++++++++++++ 4 files changed, 164 insertions(+), 14 deletions(-) diff --git a/native/csrc/catalog/bindings_store.cpp b/native/csrc/catalog/bindings_store.cpp index 5866d9f79..792269b39 100644 --- a/native/csrc/catalog/bindings_store.cpp +++ b/native/csrc/catalog/bindings_store.cpp @@ -72,6 +72,8 @@ dc::StorageServiceConfig service_config(const py::dict& d) { c.writer.publish_timeout_ns = get(d, "publish_timeout_ns", c.writer.publish_timeout_ns); c.writer.clock_skew_ns = get(d, "clock_skew_ns", c.writer.clock_skew_ns); + c.start_lease_wait_ns = + get(d, "start_lease_wait_ns", c.start_lease_wait_ns); c.indexer.max_packs = get(d, "indexer_max_packs", c.indexer.max_packs); c.indexer.max_estimated_bytes = get(d, "indexer_max_estimated_bytes", c.indexer.max_estimated_bytes); diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index efc1f2a99..ad3d2c36b 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -5,6 +5,7 @@ #include #include #include +#include #include #include "catalog/schema.h" @@ -52,6 +53,11 @@ bool is_pack_id(const std::string& value) { return dmi_pack::ParseUuid(value, &bytes, &canonical) && canonical == value; } +bool is_lease_refusal(const CatalogError& exc) { + return exc.kind() == CatalogError::Kind::kHeld || + exc.kind() == CatalogError::Kind::kLease; +} + } // namespace CaptureStorageService::CaptureStorageService(StorageServiceConfig config) @@ -91,10 +97,25 @@ void CaptureStorageService::start() { std::lock_guard cycle(cycle_mutex_); if (started_) throw std::logic_error("storage service: already started"); + CatalogSchema(clickhouse_, config_.writer.database, config_.writer.table_prefix) + .ensure(&writer_.leases(), config_.schema_retry_sleep_ns); + { + std::lock_guard lease(lease_mutex_); + acquire_lease_at_start(); // throws kHeld if another publisher keeps it + } + + // After the lease, never before: a second process pointed at this spool + // would otherwise delete a live sink's .open files and only then learn + // that the catalog is held. if (config_.sweep_spool_on_start) { std::vector recovered; std::string error; if (spool_.Recover(&recovered, &error) != dmi_store::SpoolStatus::kOk) { + try { + std::lock_guard lease(lease_mutex_); + if (writer_.held_lease() != nullptr) writer_.release_lease(); + } catch (...) { + } throw std::runtime_error("storage service: spool recovery failed: " + error); } @@ -102,14 +123,6 @@ void CaptureStorageService::start() { state_.swept_on_start = recovered.size(); } - CatalogSchema(clickhouse_, config_.writer.database, config_.writer.table_prefix) - .ensure(&writer_.leases(), config_.schema_retry_sleep_ns); - { - std::lock_guard lease(lease_mutex_); - writer_.acquire_lease(config_.holder); // throws kHeld if another publisher holds it - last_renew_ns_ = steady_ns(); - } - // A failed pass is not fatal -- the bucket is still there next time -- but a // lost lease is, and start() must not return holding a lease stop() will // never release. @@ -579,6 +592,36 @@ void CaptureStorageService::renew_lease_if_due() { ++state_.lease_renewals; } +void CaptureStorageService::acquire_lease_at_start() { + // A crashed predecessor's lease stays live for up to its TTL. Waiting it + // out here turns a restart inside that window into a short delay instead + // of a failed start; a live publisher keeps renewing, so the wait ends in + // the same refusal as before, naming the holder. + const uint64_t deadline = steady_ns() + config_.start_lease_wait_ns; + const uint64_t poll = std::clamp( + config_.writer.lease_ttl_ns / 10, 50'000'000ull, 500'000'000ull); + while (true) { + try { + writer_.acquire_lease(config_.holder); + last_renew_ns_ = steady_ns(); + return; + } catch (const CatalogError& exc) { + if (!is_lease_refusal(exc) || config_.start_lease_wait_ns == 0) throw; + const uint64_t now = steady_ns(); + if (now >= deadline) { + throw CatalogError( + exc.kind(), + "storage service: waited " + + std::to_string(config_.start_lease_wait_ns / 1'000'000) + + " ms for the publisher lease and it is still held: " + + exc.what()); + } + std::this_thread::sleep_for( + std::chrono::nanoseconds(std::min(poll, deadline - now))); + } + } +} + void CaptureStorageService::record_error(const std::string& message) { std::lock_guard lock(state_mutex_); state_.last_error = message.substr(0, 512); diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 10372b022..e05bac2c2 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -21,9 +21,10 @@ // periodically when reconcile_interval_ns is non-zero. // // Deployment shape: the service holds the catalog's single publisher lease, so -// run ONE service per (database, table_prefix). A second one fails to acquire -// the lease at start(). That is the in-process mode; a standalone daemon can -// reuse this class unchanged. +// 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. +// That is the in-process mode; a standalone daemon can reuse this class +// unchanged. #pragma once #include @@ -82,10 +83,16 @@ struct StorageServiceConfig { int max_index_attempts = 5; uint64_t schema_retry_sleep_ns = 500'000'000ull; + // How long start() waits for another holder's lease to expire before it + // fails with the lease held. A crashed predecessor's lease stays live for + // up to its TTL; 0 fails at once. + uint64_t start_lease_wait_ns = 0; + // Sweep a crashed sink's stale .open files before anything writes to the // spool. Recover() deletes every .open file this object does not own, so it // is only safe while no writer is live: start() must run before the sink - // opens the spool. + // opens the spool. It runs after the lease is taken, so a start refused + // the catalog never touches the spool. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; @@ -118,9 +125,10 @@ class CaptureStorageService { CaptureStorageService(const CaptureStorageService&) = delete; CaptureStorageService& operator=(const CaptureStorageService&) = delete; - // Sweep the spool, ensure the catalog schema, take the publisher lease, + // Ensure the catalog schema, take the publisher lease (waiting up to + // start_lease_wait_ns for another holder's to expire), sweep the spool, // reconcile once, then start the background cycle. Throws if the lease is - // held by another publisher. + // still held by another publisher when the wait ends. void start(); // Run cycles until one finds the spool empty with every uploaded pack @@ -156,6 +164,8 @@ class CaptureStorageService { void reconcile(); void keep_lease(); // the lease thread's body void renew_lease_if_due(); // requires lease_mutex_ + // Takes the lease at start(), waiting for an expiring predecessor. + void acquire_lease_at_start(); // requires lease_mutex_ // Sets a pack aside for good; flush() reports it. Requires cycle_mutex_. void reject(const PackRefData& ref, const std::string& reason); void record_error(const std::string& message); diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index 390af9945..3a8569a2a 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -768,3 +768,98 @@ def _native(spool, holder): assert snapshot["upload_failures"] >= 1, snapshot assert snapshot["lease_renewals"] >= 3, snapshot + + +# --- the publisher lease through ClickHouse errors and restarts --------------- + + +_HOLD_LEASE = """ +import json, sys, time +from dmi.storage.native_capture import ( + NativeCaptureStorage, NativeCaptureStorageConfig) + +service = NativeCaptureStorage( + NativeCaptureStorageConfig(**json.loads(sys.argv[1])), + spool_root=sys.argv[2], spool_max_bytes=1 << 40, sweep_spool=True) +service.start() +print("started", flush=True) +time.sleep(600) +""" + + +def test_a_restart_within_the_ttl_of_a_killed_predecessor_succeeds( + fake_s3, tmp_path): + """A SIGKILLed publisher leaves its lease live for up to a TTL, and a + restart in that window failed at once with the lease held. start() now + waits, within start_lease_wait_s, for the lease to expire.""" + import os + import signal + import sys + + with _catalog() as (_client, catalog): + knobs = dict( + s3_endpoint=fake_s3, s3_bucket=BUCKET, s3_region=REGION, + s3_access_key=ACCESS, s3_secret_key=SECRET, + s3_allow_insecure_http=True, clickhouse_host=CLICKHOUSE_HOST, + clickhouse_port=CLICKHOUSE_HTTP_PORT, database=DATABASE, + table_prefix=catalog.table_prefix, reconcile_on_start=False, + lease_ttl_s=5.0, publish_timeout_s=1) + env = dict(os.environ) + env["PYTHONPATH"] = os.pathsep.join( + [str(REPO / "src")] + ([env["PYTHONPATH"]] + if env.get("PYTHONPATH") else [])) + predecessor = subprocess.Popen( + [sys.executable, "-c", _HOLD_LEASE, + json.dumps(dict(knobs, holder="crashed-publisher")), + str(tmp_path / "predecessor")], + stdout=subprocess.PIPE, text=True, env=env) + try: + assert predecessor.stdout.readline().strip() == "started" + finally: + predecessor.send_signal(signal.SIGKILL) + predecessor.wait(timeout=30) + assert predecessor.returncode == -signal.SIGKILL + + from dmi.storage.native_capture import NativeCaptureStorageConfig + + successor = _service(NativeCaptureStorageConfig( + **knobs, holder="restarted-publisher", start_lease_wait_s=8.0), + tmp_path / "successor") + started = time.monotonic() + successor.start() + elapsed = time.monotonic() - started + try: + snapshot = successor.snapshot() + finally: + successor.stop() + + assert snapshot["running"] is True, snapshot + # It really waited on the dead holder's lease, and no longer than it. + assert 1.0 < elapsed < 8.0, elapsed + + +def test_a_start_refused_the_lease_leaves_the_spool_unswept(fake_s3, tmp_path): + """start() swept the spool before taking the lease, so a second process + pointed at a live process's spool deleted its in-progress .open files + and only then found the catalog held. The lease now comes first.""" + with _catalog() as (_client, catalog): + first = _service(_storage_config( + fake_s3, catalog.table_prefix, reconcile_on_start=False), + tmp_path / "first") + spool_root = tmp_path / "second" + spool_root.mkdir() + in_progress = spool_root / f".{uuid.uuid4()}.1234.open" + in_progress.write_bytes(b"a pack being written") + second = _service(_storage_config( + fake_s3, catalog.table_prefix, reconcile_on_start=False, + start_lease_wait_s=0.0), spool_root) + first.start() + try: + with pytest.raises(RuntimeError, match="held"): + second.start() + assert in_progress.exists() + finally: + first.stop() + second.start() # the lease is free now, and the sweep runs + second.stop() + assert not in_progress.exists() From c13a828718999c355af675150393163de25cc8ad Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 12:04:35 -0400 Subject: [PATCH 3/6] Recover the publisher lease after a ClickHouse error instead of latching Any catalog statement whose outcome is unknown -- a renewal or a publish that hits a transport error or timeout -- quarantines the CatalogWriter, which drops its lease without a tombstone and refuses to publish for one TTL. When the window passed, the next renewal or index found no local lease and threw kLease, and the lease thread or run_cycle latched "publisher lease lost": the loop exited for good, although ClickHouse was back and nobody else held the catalog. A ClickHouse blip longer than the renewal interval stopped indexing permanently. The loss is now recoverable, and only a real takeover is fatal: - ensure_publisher_lease() runs at the start of every cycle and on every lease-thread tick. While the writer is quarantined it claims nothing; the cycle skips the catalog phase and keeps pending_index_, and uploads continue only as far as the existing "nothing new while a pack is owed" rule allows. Once the window has passed it acquires a fresh lease_id, the recovery the Python oracle's writer documents (clickhouse_catalog.py publish_snapshot: wait out the window, acquire a fresh lease, re-index). - A claim or renewal refused by another holder (kHeld, or a contested claim) is retried every ttl/6. It latches only once the refusals have lasted 2 x TTL, since our own dropped row is dead within one TTL and a handover ends sooner. A refusal inside run_cycle is recorded, and the next ensure_publisher_lease() makes the call. - On a latch: running=false, one line on stderr ("dmi capture storage: indexing stopped: ..." naming the holder), and flush()/rethrow_if_failed raise the refusal. The snapshot gains failed, lease_state (none, held, quarantined, reacquiring, failed, released), quarantined_until (time.monotonic() seconds, 0.0 when not quarantined) and lease_reacquisitions. - A re-acquired lease kicks the loop out of its backoff, so what is owed is indexed at once. This replaces the plan's "cap backoff at ttl/3", which predates the separate lease thread: renewal no longer depends on the cycle's wait, and only the post-recovery index latency did. The writer, the lease coordinator and every statement are unchanged: tests/test_native_catalog_lease_live.py, the byte-identity gate, passes unmodified (54 passed). Evidence, local ClickHouse 25.12, tests/test_native_capture_storage_live.py, TTL 3 s and publish_timeout 1 s with the TCP switch in front of ClickHouse: - A 2.5 s cut spanning a renewal and an index attempt: before, a reproduction left last_error "publisher lease lost: the catalog writer holds no publisher lease" and flush raised; the new test failed on the missing snapshot fields. After: lease_state goes held -> quarantined -> held, lease_reacquisitions >= 1, failed False, the background loop indexes the owed pack without a flush, flush succeeds, and all 6 captures read back byte-identical. - A rival that stops within 2 x TTL: before, flush raised "no publisher lease is held" once the rival had gone; after, the lease is re-acquired and nothing latches (no stderr line). - A rival that takes over during the cut and keeps renewing: the service latches at least 2 x TTL after the rival started, running False, lease_state "failed", flush raises naming "rival-publisher", exactly one stderr line. (It latched before this change too, on the kLease path; the test pins that recovery does not paper over a takeover.) - The cut and rival tests passed 3 runs out of 3; whole file 18 passed. --- native/csrc/catalog/bindings_store.cpp | 6 + native/csrc/catalog/storage_service.cpp | 207 ++++++++++++++++++---- native/csrc/catalog/storage_service.h | 40 ++++- tests/test_native_capture_storage_live.py | 165 +++++++++++++++++ 4 files changed, 384 insertions(+), 34 deletions(-) diff --git a/native/csrc/catalog/bindings_store.cpp b/native/csrc/catalog/bindings_store.cpp index 792269b39..38f20ca2a 100644 --- a/native/csrc/catalog/bindings_store.cpp +++ b/native/csrc/catalog/bindings_store.cpp @@ -107,6 +107,12 @@ 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["failed"] = s.failed; + out["lease_state"] = s.lease_state; + // Seconds on the monotonic clock, comparable with time.monotonic() (both + // CLOCK_MONOTONIC on Linux); 0.0 when not quarantined. + out["quarantined_until"] = static_cast(s.quarantined_until_ns) / 1e9; + out["lease_reacquisitions"] = s.lease_reacquisitions; out["last_error"] = s.last_error; return out; } diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index ad3d2c36b..f8add13ed 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -58,6 +59,13 @@ bool is_lease_refusal(const CatalogError& exc) { exc.kind() == CatalogError::Kind::kLease; } +// The lease thread's tick, which is also the retry interval for a claim +// another holder refused: a sixth of the TTL, so a renewal due at ttl/3 is +// never more than a tick late. +uint64_t lease_tick_ns(uint64_t ttl_ns) { + return std::max(ttl_ns / 6, 10'000'000ull); +} + } // namespace CaptureStorageService::CaptureStorageService(StorageServiceConfig config) @@ -130,8 +138,7 @@ void CaptureStorageService::start() { try { reconcile(); } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { + if (is_lease_refusal(exc)) { try { std::lock_guard lease(lease_mutex_); if (writer_.held_lease() != nullptr) writer_.release_lease(); @@ -149,6 +156,11 @@ void CaptureStorageService::start() { { std::lock_guard lock(wake_mutex_); stop_requested_ = false; + kick_ = false; + } + { + std::lock_guard lease(lease_mutex_); + publish_lease_state(); } { std::lock_guard lock(state_mutex_); @@ -170,12 +182,22 @@ void CaptureStorageService::stop() { std::lock_guard cycle(cycle_mutex_); if (started_) { started_ = false; + // A quarantined writer holds no lease, so it writes no tombstone: the + // outcome-unknown statement may still be running, and its row must stay + // live until the TTL keeps a successor out of that window. + bool released = false; try { std::lock_guard lease(lease_mutex_); - if (writer_.held_lease() != nullptr) writer_.release_lease(); + if (writer_.held_lease() != nullptr) { + writer_.release_lease(); + released = true; + } } catch (const std::exception& exc) { record_error(std::string("lease release failed: ") + exc.what()); } + std::lock_guard lock(state_mutex_); + if (!state_.failed) state_.lease_state = released ? "released" : "none"; + state_.quarantined_until_ns = 0; } std::lock_guard lock(state_mutex_); state_.running = false; @@ -229,12 +251,13 @@ void CaptureStorageService::loop() { { std::unique_lock lock(wake_mutex_); wake_.wait_for(lock, std::chrono::nanoseconds(wait_ns), - [this] { return stop_requested_; }); + [this] { return stop_requested_ || kick_; }); if (stop_requested_) return; + kick_ = false; } { std::lock_guard lock(state_mutex_); - if (failure_) return; // the lease is gone; nothing more can publish + if (failure_) return; // another publisher holds the catalog } std::lock_guard cycle(cycle_mutex_); run_cycle(); @@ -251,9 +274,21 @@ void CaptureStorageService::loop() { CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { CycleOutcome outcome; + // The catalog phase needs the lease. Without one -- quarantined after an + // unknown outcome, or refused by another holder -- the cycle still + // uploads as far as the owed-pack rule below allows, and owes the rest. + bool catalog = false; + { + std::lock_guard lease(lease_mutex_); + 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. - const auto index_or_owe = [this](std::vector refs) { + const auto index_or_owe = [this, catalog](std::vector refs) { + if (!catalog) { + pending_index_.insert(pending_index_.end(), refs.begin(), refs.end()); + return; + } std::vector unindexed; try { if (!refs.empty()) index_bounded(std::move(refs), &unindexed); @@ -270,7 +305,7 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { // of it is still owed, the catalog is down or refusing: upload // nothing new, so new packs stay in the durable spool rather than // joining a list that only this process remembers. - if (!pending_index_.empty()) { + if (catalog && !pending_index_.empty()) { std::vector owed; owed.swap(pending_index_); index_or_owe(std::move(owed)); @@ -310,7 +345,7 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { index_or_owe(std::move(to_index)); // 4. Reconcile on its interval. The lease thread keeps the lease alive. - if (config_.reconcile_interval_ns > 0 && + if (catalog && config_.reconcile_interval_ns > 0 && steady_ns() - last_reconcile_ns_ >= config_.reconcile_interval_ns) { reconcile(); last_reconcile_ns_ = steady_ns(); @@ -324,7 +359,10 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { // A cycle that already failed is not drained whatever the spool holds, // so it skips the listing: an empty spool lists for free, but a backlog // would be re-hashed on every cycle of an outage. - outcome.failed = upload_failures != 0 || !pending_index_.empty(); + // Without the lease nothing can be confirmed in the catalog, so the + // cycle is not drained, and it counts towards the backoff. + outcome.failed = + !catalog || upload_failures != 0 || !pending_index_.empty(); bool nothing_pending = !batch.refs.empty(); if (batch.refs.empty() && !outcome.failed) { std::vector pending; @@ -341,13 +379,12 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { std::lock_guard lock(state_mutex_); ++state_.cycles; } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { - latch_failure(std::current_exception(), - std::string("publisher lease lost: ") + exc.what()); - } else { - record_error(exc.what()); - } + // A lease refusal here is not fatal either: the writer has dropped or + // been fenced out of its lease, and the next ensure_publisher_lease() + // decides between a fresh lease and a latch. + record_error(is_lease_refusal(exc) + ? std::string("publisher lease lost: ") + exc.what() + : std::string(exc.what())); } catch (const std::exception& exc) { record_error(exc.what()); } @@ -411,8 +448,7 @@ void CaptureStorageService::index_bounded(std::vector refs, reject(batch.front(), exc.what()); continue; } - const bool lease_lost = exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease; + const bool lease_lost = is_lease_refusal(exc); give_up(batch, std::string("index failed: ") + exc.what()); if (lease_lost) throw; return; @@ -550,10 +586,10 @@ void CaptureStorageService::keep_lease() { // The lease renews only inside a publish, and the cycle loop backs off up // to max_backoff_ns while the object store is down -- past the lease TTL. // Renewing from the cycle let the lease lapse during an outage and a rival - // take the catalog. This thread renews on its own schedule instead. - const uint64_t ttl = config_.writer.lease_ttl_ns; - const auto tick = std::chrono::nanoseconds( - std::max(ttl / 6, 10'000'000ull)); + // take the catalog. This thread renews on its own schedule instead, and + // takes a fresh lease once a lost one can be replaced. + const auto tick = + std::chrono::nanoseconds(lease_tick_ns(config_.writer.lease_ttl_ns)); while (true) { { std::unique_lock lock(wake_mutex_); @@ -564,20 +600,27 @@ void CaptureStorageService::keep_lease() { std::lock_guard lock(state_mutex_); if (failure_) return; } + std::lock_guard lease(lease_mutex_); + if (writer_.held_lease() == nullptr) { + ensure_publisher_lease(); + continue; + } try { - std::lock_guard lease(lease_mutex_); renew_lease_if_due(); } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { - latch_failure(std::current_exception(), - std::string("publisher lease lost: ") + exc.what()); - return; + if (is_lease_refusal(exc)) { + // The coordinator dropped the lease: a live foreign head refused + // the renewal claim, or the claim was contested. + lease_held_elsewhere(exc); + } else { + record_error(std::string("lease renewal failed: ") + exc.what()); } - record_error(std::string("lease renewal failed: ") + exc.what()); } catch (const std::exception& exc) { + // An unknown outcome: the writer quarantined itself and dropped the + // lease. ensure_publisher_lease() replaces it after the window. record_error(std::string("lease renewal failed: ") + exc.what()); } + publish_lease_state(); } } @@ -588,6 +631,7 @@ void CaptureStorageService::renew_lease_if_due() { if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); last_renew_ns_ = steady_ns(); + held_elsewhere_since_ns_ = 0; std::lock_guard lock(state_mutex_); ++state_.lease_renewals; } @@ -622,6 +666,91 @@ void CaptureStorageService::acquire_lease_at_start() { } } +bool CaptureStorageService::ensure_publisher_lease() { + if (writer_.held_lease() != nullptr) return true; + { + std::lock_guard lock(state_mutex_); + if (failure_) return false; + } + if (writer_.quarantined()) { + // The unknown-outcome statement may still be running under the lease + // it dropped; no claim, not even a fresh lease_id, until the window + // (one TTL) has passed and that lease's row has expired with it. + publish_lease_state(); + return false; + } + const uint64_t now = steady_ns(); + if (now < next_claim_ns_) return false; + try { + // A fresh lease_id: the coordinator mints one because the quarantine + // or refusal already dropped the old lease. + writer_.acquire_lease(config_.holder); + } catch (const CatalogError& exc) { + if (is_lease_refusal(exc)) { + lease_held_elsewhere(exc); + } else { + record_error(std::string("publisher lease acquisition failed: ") + + exc.what()); + } + publish_lease_state(); + return false; + } catch (const std::exception& exc) { + // Unknown outcome again: the writer quarantined itself for a TTL. + record_error(std::string("publisher lease acquisition failed: ") + + exc.what()); + publish_lease_state(); + return false; + } + last_renew_ns_ = steady_ns(); + held_elsewhere_since_ns_ = 0; + next_claim_ns_ = 0; + { + std::lock_guard lock(state_mutex_); + ++state_.lease_reacquisitions; + } + publish_lease_state(); + { + std::lock_guard lock(wake_mutex_); + kick_ = true; + } + wake_.notify_all(); + return true; +} + +void CaptureStorageService::lease_held_elsewhere(const CatalogError& refusal) { + const uint64_t now = steady_ns(); + const uint64_t ttl = config_.writer.lease_ttl_ns; + if (held_elsewhere_since_ns_ == 0) held_elsewhere_since_ns_ = now; + next_claim_ns_ = now + lease_tick_ns(ttl); + // Our own dropped row is dead within one TTL of the loss, and a handover + // (a rival that stops, releasing with a tombstone) ends sooner still. A + // refusal that has lasted 2 x TTL is a publisher that means to stay. + if (now - held_elsewhere_since_ns_ >= 2 * ttl) { + latch_failure(std::make_exception_ptr(refusal), + std::string("publisher lease held by another publisher " + "for over 2 x TTL: ") + refusal.what()); + } else { + record_error(std::string("publisher lease held elsewhere: ") + + refusal.what()); + } +} + +void CaptureStorageService::publish_lease_state() { + uint64_t until = 0; + const char* lease_state = "reacquiring"; + if (writer_.held_lease() != nullptr) { + lease_state = "held"; + } else if (writer_.quarantined(&until)) { + lease_state = "quarantined"; + } else { + until = 0; + } + std::lock_guard lock(state_mutex_); + if (state_.failed) return; + state_.lease_state = lease_state; + state_.quarantined_until_ns = until; +} + void CaptureStorageService::record_error(const std::string& message) { std::lock_guard lock(state_mutex_); state_.last_error = message.substr(0, 512); @@ -629,9 +758,23 @@ void CaptureStorageService::record_error(const std::string& message) { void CaptureStorageService::latch_failure(std::exception_ptr failure, const std::string& message) { - std::lock_guard lock(state_mutex_); - if (!failure_) failure_ = std::move(failure); - state_.last_error = message.substr(0, 512); + std::string line = message.substr(0, 512); + std::replace(line.begin(), line.end(), '\n', ' '); + { + std::lock_guard lock(state_mutex_); + if (failure_) return; + failure_ = std::move(failure); + state_.last_error = line; + state_.failed = true; + state_.running = false; + state_.lease_state = "failed"; + state_.quarantined_until_ns = 0; + } + // Once, and on one line: the loop and the lease thread both stop here, + // and nothing else will say why indexing did. + std::fprintf(stderr, "dmi capture storage: indexing stopped: %s\n", + line.c_str()); + std::fflush(stderr); } } // namespace dmi_catalog diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index e05bac2c2..d909bfd31 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -25,6 +25,17 @@ // start_lease_wait_ns for the lease at start(), then fails naming the holder. // That is the in-process mode; a standalone daemon can reuse this class // unchanged. +// +// The lease through ClickHouse errors. A catalog statement whose outcome is +// unknown (a transport error or timeout on a renewal or a publish) +// quarantines the writer: it drops its lease without a tombstone and refuses +// to publish for one TTL (catalog_writer.cpp). That is recoverable, not +// fatal: while the writer is quarantined or holds no lease the cycle skips +// the catalog phase and keeps pending_index_, and once the window passes the +// service acquires a FRESH lease_id, as the Python oracle's writer documents +// (clickhouse_catalog.py, publish_snapshot). Only a foreign lease that stays +// live for 2 x TTL is fatal: the service stops (snapshot().failed), writes +// one line to stderr, and flush() rethrows the refusal naming the holder. #pragma once #include @@ -115,6 +126,15 @@ struct StorageServiceSnapshot { uint64_t swept_on_start = 0; // ready packs Recover() found at start uint64_t pending_index = 0; // uploaded packs awaiting a retried index uint64_t rejected_packs = 0; // set aside: cannot be indexed (see flush) + // A foreign lease outlived 2 x TTL: the service stopped for good. + bool failed = false; + // "none" before start, "held", "quarantined" (an unknown outcome set the + // lease aside), "reacquiring" (no lease, trying for a fresh one), + // "failed", or "released" (stop() wrote the tombstone). + std::string lease_state = "none"; + // steady_clock ns at which the quarantine ends; 0 when not quarantined. + uint64_t quarantined_until_ns = 0; + uint64_t lease_reacquisitions = 0; // fresh leases taken after a loss std::string last_error; }; @@ -145,8 +165,8 @@ class CaptureStorageService { StorageServiceSnapshot snapshot() const; - // Rethrows a fatal failure latched by the background cycle: losing the - // publisher lease to another holder. + // Rethrows a fatal failure latched by the background cycle: another + // publisher holding the lease for longer than 2 x TTL. void rethrow_if_failed() const; private: @@ -166,6 +186,13 @@ class CaptureStorageService { void renew_lease_if_due(); // requires lease_mutex_ // Takes the lease at start(), waiting for an expiring predecessor. void acquire_lease_at_start(); // requires lease_mutex_ + // Whether the writer holds a lease, taking a fresh one when it has none + // and is no longer quarantined. Never throws. Requires lease_mutex_. + bool ensure_publisher_lease(); + // Another holder refused a claim or renewal; latches once that has lasted + // 2 x TTL. Call from the catch block. Requires lease_mutex_. + void lease_held_elsewhere(const CatalogError& refusal); + void publish_lease_state(); // requires lease_mutex_ // Sets a pack aside for good; flush() reports it. Requires cycle_mutex_. void reject(const PackRefData& ref, const std::string& reason); void record_error(const std::string& message); @@ -187,6 +214,12 @@ class CaptureStorageService { // while cycles publish. Taken inside cycle_mutex_, never the other way. std::mutex lease_mutex_; uint64_t last_renew_ns_ = 0; // guarded by lease_mutex_ + // When a claim or renewal was first refused by another holder since the + // lease was last held; 0 while none has been. Guarded by lease_mutex_. + uint64_t held_elsewhere_since_ns_ = 0; + // Earliest next claim after a refusal, so flush()'s fast cycles do not + // hammer the lease table. Guarded by lease_mutex_. + uint64_t next_claim_ns_ = 0; uint64_t last_reconcile_ns_ = 0; int failure_streak_ = 0; // consecutive failed cycles, for the backoff // Uploaded, so gone from the spool, but not yet in the catalog. @@ -201,6 +234,9 @@ class CaptureStorageService { std::mutex wake_mutex_; std::condition_variable wake_; bool stop_requested_ = false; + // Set when a lease is re-acquired, so a loop in a long backoff indexes + // what is owed now rather than after its wait. Guarded by wake_mutex_. + bool kick_ = false; bool started_ = false; mutable std::mutex state_mutex_; diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index 3a8569a2a..a7c643e15 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -773,6 +773,170 @@ def _native(spool, holder): # --- the publisher lease through ClickHouse errors and restarts --------------- +def test_a_catalog_cut_spanning_a_renewal_and_a_publish_recovers( + fake_s3, tmp_path): + """A ClickHouse error of unknown outcome -- a renewal or a publish that + cannot reach the server -- quarantines the catalog writer and drops its + lease without a tombstone. When the quarantine ended the service found + no lease, latched "publisher lease lost", and indexing stopped for good, + although ClickHouse was back and nobody else held the catalog. It now + waits the quarantine out and takes a fresh lease.""" + spool_root = tmp_path / "spool" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + reconcile_on_start=False, lease_ttl_s=3.0, publish_timeout_s=1) + service = _service(config, spool_root) + service.start() + try: + tensors = _stage(spool_root, range(2)) + service.flush(30.0) + snapshot = service.snapshot() + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_reacquisitions"] == 0, snapshot + assert snapshot["quarantined_until"] == 0.0, snapshot + + switch.cut() + cut_at = time.monotonic() + # Staged during the cut: uploaded, then owed, since the index + # (and its publish) cannot reach the catalog. + tensors.update(_stage(spool_root, range(2, 4))) + # The renewal due every ttl/3 fails and quarantines the lease. + _wait_for(lambda: service.snapshot()["lease_state"] == "quarantined", + timeout_s=5.0) + snapshot = service.snapshot() + assert snapshot["quarantined_until"] > time.monotonic(), snapshot + assert snapshot["running"] is True, snapshot + assert snapshot["index_failures"] >= 1, snapshot # it tried + assert snapshot["pending_index"] == 1, snapshot + with pytest.raises(TimeoutError, match="1 uploaded but unindexed"): + service.flush(0.5) + # Longer than ttl/3, so the cut spans a renewal and a publish. + time.sleep(max(0.0, cut_at + 2.5 - time.monotonic())) + switch.restore() + + # The background loop recovers on its own, without a flush. + _wait_for(lambda: service.snapshot()["indexed_packs"] == 2, + timeout_s=15.0) + snapshot = service.snapshot() + assert snapshot["failed"] is False, snapshot + assert snapshot["running"] is True, snapshot + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["lease_reacquisitions"] >= 1, snapshot + assert snapshot["quarantined_until"] == 0.0, snapshot + + tensors.update(_stage(spool_root, range(4, 6))) + service.flush(30.0) + service.rethrow_if_failed() + snapshot = service.snapshot() + finally: + service.stop() + switch.close() + + assert snapshot["indexed_packs"] == 3, snapshot + assert snapshot["pending_index"] == 0, snapshot + direct = _storage_config(fake_s3, catalog.table_prefix) + captures = _read_all(direct) + assert sorted(captures) == sorted(tensors) + for capture_id, tensor in tensors.items(): + assert captures[capture_id].payload == tensor.numpy().tobytes() + + +def _latch_lines(err: str) -> list[str]: + return [line for line in err.splitlines() if "indexing stopped" in line] + + +def test_a_rival_that_takes_over_during_a_cut_still_latches( + fake_s3, tmp_path, capfd): + """Recovery must not paper over a real takeover: when another publisher + holds the catalog for longer than two lease TTLs, the service stops, + says so once on stderr, and flush raises naming the holder.""" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + knobs = dict(reconcile_on_start=False, lease_ttl_s=3.0, + publish_timeout_s=1) + first = _service(_storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + holder="first-publisher", **knobs), tmp_path / "first") + rival = _service(_storage_config( + fake_s3, catalog.table_prefix, holder="rival-publisher", + start_lease_wait_s=10.0, **knobs), tmp_path / "rival") + first.start() + try: + switch.cut() + rival.start() # waits out the lease the cut first cannot renew + rival_started = time.monotonic() + switch.restore() + + _wait_for(lambda: first.snapshot()["failed"], timeout_s=20.0) + latched_at = time.monotonic() + snapshot = first.snapshot() + assert snapshot["running"] is False, snapshot + assert snapshot["lease_state"] == "failed", snapshot + assert "rival-publisher" in snapshot["last_error"], snapshot + # Only a foreign lease that outlives two TTLs latches. + assert latched_at - rival_started >= 2 * 3.0 - 0.5 + with pytest.raises(RuntimeError, match="rival-publisher"): + first.flush(1.0) + with pytest.raises(RuntimeError, match="rival-publisher"): + first.rethrow_if_failed() + + rival.flush(10.0) + rival_snapshot = rival.snapshot() + assert rival_snapshot["lease_state"] == "held", rival_snapshot + assert rival_snapshot["failed"] is False, rival_snapshot + finally: + first.stop() + rival.stop() + switch.close() + + lines = _latch_lines(capfd.readouterr().err) + assert len(lines) == 1, lines + assert "rival-publisher" in lines[0] + + +def test_a_rival_that_stops_within_two_ttls_does_not_latch( + fake_s3, tmp_path, capfd): + """A foreign lease that is gone again before two TTLs pass is a + handover, not a takeover: the service takes the lease back.""" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + spool_root = tmp_path / "first" + with _catalog() as (_client, catalog): + knobs = dict(reconcile_on_start=False, lease_ttl_s=3.0, + publish_timeout_s=1) + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + holder="first-publisher", **knobs) + first = _service(config, spool_root) + rival = _service(_storage_config( + fake_s3, catalog.table_prefix, holder="rival-publisher", + start_lease_wait_s=10.0, **knobs), tmp_path / "rival") + first.start() + try: + switch.cut() + rival.start() + switch.restore() + time.sleep(2.0) + rival.stop() # releases with a tombstone + + tensors = _stage(spool_root, range(2)) + first.flush(20.0) + snapshot = first.snapshot() + finally: + first.stop() + rival.stop() + switch.close() + + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["lease_reacquisitions"] >= 1, snapshot + assert sorted(_read_all(_storage_config( + fake_s3, catalog.table_prefix))) == sorted(tensors) + assert _latch_lines(capfd.readouterr().err) == [] + + _HOLD_LEASE = """ import json, sys, time from dmi.storage.native_capture import ( @@ -834,6 +998,7 @@ def test_a_restart_within_the_ttl_of_a_killed_predecessor_succeeds( successor.stop() assert snapshot["running"] is True, snapshot + assert snapshot["lease_state"] == "held", snapshot # It really waited on the dead holder's lease, and no longer than it. assert 1.0 < elapsed < 8.0, elapsed From 204a8d2f4627d4298fa7c136db24b4c9bd2734a0 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 12:34:13 -0400 Subject: [PATCH 4/6] Count clock skew in the default start wait, and say what sweep-after-lease covers A replica whose clock lags by clock_skew_s sees a crashed predecessor's row as live for that long after lease_ttl_s, so the default start wait is now lease_ttl_s + publish_timeout_s + clock_skew_s. The config comment also notes that a predecessor with a longer TTL (the native 30 s default) can outlast the default wait. The storage service comment no longer claims sweeping after the lease protects a live sink's spool: a quarantined holder lets its row lapse, and a second process can take the lease in that gap. The spool owner lock is the real fix. test_one_publisher_per_catalog passes start_lease_wait_s=0.0, since it tests the refusal and the restart tests cover the wait (20.7 s -> 2.1 s). --- native/csrc/catalog/storage_service.cpp | 9 ++++++--- src/dmi/storage/native_capture.py | 11 ++++++++--- tests/test_native_capture_storage_live.py | 5 ++++- tests/test_native_capture_storage_wiring.py | 12 ++++++++++++ 4 files changed, 30 insertions(+), 7 deletions(-) diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index f8add13ed..1c81f90cb 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -112,9 +112,12 @@ void CaptureStorageService::start() { acquire_lease_at_start(); // throws kHeld if another publisher keeps it } - // After the lease, never before: a second process pointed at this spool - // would otherwise delete a live sink's .open files and only then learn - // that the catalog is held. + // 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. if (config_.sweep_spool_on_start) { std::vector recovered; std::string error; diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index 73b9a0c6a..e3d3fe01b 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -137,8 +137,12 @@ class NativeCaptureStorageConfig: # The bound on host clock disagreement across a replicated catalog. clock_skew_s: float = 0.0 # How long start() waits for a predecessor's lease to expire before - # failing with it held. None waits lease_ttl_s + publish_timeout_s, - # enough to outlast a crashed predecessor; 0 fails at once. + # failing with it held. None waits lease_ttl_s + publish_timeout_s + + # clock_skew_s, enough to outlast a crashed predecessor that ran with the + # same knobs; 0 fails at once. A predecessor with a longer TTL (the + # native default is 30 s, which processes predating these knobs used) + # can outlast it: set this explicitly for the first restart after such + # a process. start_lease_wait_s: Optional[float] = None def __post_init__(self) -> None: @@ -211,7 +215,8 @@ def _validate_lease(self) -> None: def _lease_native(self) -> dict[str, int]: wait = self.start_lease_wait_s if wait is None: - wait = self.lease_ttl_s + self.publish_timeout_s + wait = (self.lease_ttl_s + self.publish_timeout_s + + self.clock_skew_s) return { "lease_ttl_ns": _ns(self.lease_ttl_s), "publish_timeout_ns": _ns(self.publish_timeout_s), diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index a7c643e15..a6917d128 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -514,7 +514,10 @@ def test_a_failed_head_is_an_error_not_a_foreign_object(fake_s3, tmp_path): def test_one_publisher_per_catalog(fake_s3, tmp_path): with _catalog() as (_client, catalog): - config = _storage_config(fake_s3, catalog.table_prefix) + # No start wait: the refusal is the point here, and waiting out the + # holder is covered by the restart tests. + config = _storage_config(fake_s3, catalog.table_prefix, + start_lease_wait_s=0.0) first = _service(config, tmp_path / "first") second = _service(config, tmp_path / "second") first.start() diff --git a/tests/test_native_capture_storage_wiring.py b/tests/test_native_capture_storage_wiring.py index 4c247dc80..4e233250c 100644 --- a/tests/test_native_capture_storage_wiring.py +++ b/tests/test_native_capture_storage_wiring.py @@ -509,6 +509,18 @@ def test_the_service_gets_the_lease_knobs_in_nanoseconds(monkeypatch, tmp_path): assert config["start_lease_wait_ns"] == 20_000_000_000 +def test_the_default_start_wait_outlasts_a_skewed_predecessor(monkeypatch, + tmp_path): + # A replica whose clock lags by clock_skew_s still sees the crashed + # predecessor's row as live for that long after lease_ttl_s. + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + engine._capture_storage_config = _storage_config(clock_skew_s=2.0) + + engine.create_record_runtime(_record_format()) + + assert services[0].config["start_lease_wait_ns"] == 22_000_000_000 + + def test_explicit_lease_knobs_reach_the_service(monkeypatch, tmp_path): engine, _events, services = _capture_engine(monkeypatch, tmp_path) engine._capture_storage_config = _storage_config( From b03be3605e6b7bb548fa56eea8ce5d45c07eab9f Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 15:06:14 -0400 Subject: [PATCH 5/6] Leave new packs in the spool while the service cannot index Without the publisher lease -- quarantined after a ClickHouse statement of unknown outcome, or refused by another holder -- a cycle with nothing owed still uploaded the whole spool, then owed every pack in pending_index_, which only this process remembers. A crash inside the quarantine window then left those packs in the bucket and never in the catalog whenever reconcile_on_start is off (the documented shared-bucket setting). The cycle now uploads only when it holds the lease and owes nothing, so packs staged during a quarantine stay in the durable spool and go up once the lease is back. This replaces the plan's "uploads continue" during a quarantine, at the author's decision. Test first in tests/test_native_capture_storage_live.py: a quarantined service with nothing owed and one pack staged uploaded it (uploaded_packs 1, expected 0); now it stays in the spool through the quarantine and is uploaded and indexed after the cut is restored. Live storage and lease suites: 73 passed. --- native/csrc/catalog/storage_service.cpp | 14 +++++-- tests/test_native_capture_storage_live.py | 49 +++++++++++++++++++++++ 2 files changed, 59 insertions(+), 4 deletions(-) diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 1c81f90cb..63be51074 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -278,8 +278,8 @@ void CaptureStorageService::loop() { CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { CycleOutcome outcome; // The catalog phase needs the lease. Without one -- quarantined after an - // unknown outcome, or refused by another holder -- the cycle still - // uploads as far as the owed-pack rule below allows, and owes the rest. + // unknown outcome, or refused by another holder -- the cycle uploads + // nothing either: whatever it uploaded it could only owe, in memory. bool catalog = false; { std::lock_guard lease(lease_mutex_); @@ -314,9 +314,15 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { index_or_owe(std::move(owed)); } - // 2. Upload everything the sink has staged. + // 2. Upload everything the sink has staged -- but only with the lease + // and nothing owed. Without the lease an uploaded pack could only be + // owed, and pending_index_ dies with the process: with + // reconcile_on_start off, a crash would leave it in the bucket and + // never in the catalog. Left in the spool it survives the crash. dmi_store::UploadBatchResult batch; - if (pending_index_.empty()) batch = uploader_->UploadPending(-1); + if (catalog && pending_index_.empty()) { + batch = uploader_->UploadPending(-1); + } std::vector to_index; uint64_t uploaded_packs = 0; uint64_t uploaded_bytes = 0; diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index a6917d128..81b8a9bce 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -847,6 +847,55 @@ def test_a_catalog_cut_spanning_a_renewal_and_a_publish_recovers( assert captures[capture_id].payload == tensor.numpy().tobytes() +def test_a_quarantined_service_leaves_new_packs_in_the_spool( + fake_s3, tmp_path): + """While quarantined the service cannot index, so it must not upload + either: an uploaded but unindexed pack is remembered only in memory, and + a crash then orphans it in the bucket when reconcile_on_start is off. + Before, a quarantined cycle with nothing owed uploaded the whole spool + into that in-memory list; now new packs stay in the durable spool until + the lease is back.""" + spool_root = tmp_path / "spool" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + reconcile_on_start=False, lease_ttl_s=3.0, publish_timeout_s=1) + service = _service(config, spool_root) + service.start() + try: + switch.cut() + # Nothing is owed: the renewal alone fails and quarantines. + _wait_for(lambda: service.snapshot()["lease_state"] == "quarantined", + timeout_s=5.0) + assert service.snapshot()["pending_index"] == 0 + + tensors = _stage(spool_root, range(2)) + assert len(_ready(spool_root)) == 1 + # flush() drives cycles; none of them may upload. + with pytest.raises(TimeoutError): + service.flush(0.5) + snapshot = service.snapshot() + assert snapshot["lease_state"] == "quarantined", snapshot + assert snapshot["uploaded_packs"] == 0, snapshot + assert snapshot["pending_index"] == 0, snapshot + assert len(_ready(spool_root)) == 1 + + switch.restore() + service.flush(30.0) + service.rethrow_if_failed() + snapshot = service.snapshot() + finally: + service.stop() + switch.close() + + assert snapshot["uploaded_packs"] == 1, snapshot + assert snapshot["indexed_packs"] == 1, snapshot + assert _ready(spool_root) == [] + captures = _read_all(_storage_config(fake_s3, catalog.table_prefix)) + assert sorted(captures) == sorted(tensors) + + def _latch_lines(err: str) -> list[str]: return [line for line in err.splitlines() if "indexing stopped" in line] From c0361d766ce6c127bce78d83043b46c50dae8e34 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 19:18:02 -0400 Subject: [PATCH 6/6] Say in the service header that a quarantined cycle uploads nothing The header still described the quarantined cycle as skipping the catalog phase and keeping pending_index_; since the previous commit it also leaves new packs in the spool. --- native/csrc/catalog/storage_service.h | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index d909bfd31..965387c91 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -31,8 +31,10 @@ // quarantines the writer: it drops its lease without a tombstone and refuses // to publish for one TTL (catalog_writer.cpp). That is recoverable, not // fatal: while the writer is quarantined or holds no lease the cycle skips -// the catalog phase and keeps pending_index_, and once the window passes the -// service acquires a FRESH lease_id, as the Python oracle's writer documents +// the catalog phase, keeps pending_index_ and uploads nothing, so new packs +// stay in the durable spool rather than in a list only this process +// remembers. Once the window passes the service acquires a FRESH lease_id, +// as the Python oracle's writer documents // (clickhouse_catalog.py, publish_snapshot). Only a foreign lease that stays // live for 2 x TTL is fatal: the service stops (snapshot().failed), writes // one line to stderr, and flush() rethrows the refusal naming the holder.