Recover the publisher lease after ClickHouse errors and wait for a crashed predecessor at start - #150
Merged
Merged
Conversation
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.
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.
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.
…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).
zaoxing
requested review from
Samfisheryu
and
a lite review from Copilot
and removed request for
Copilot
September 24, 2026 16:34
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.
zaoxing
force-pushed
the
fix/lease-recovery
branch
from
September 24, 2026 22:59
352d84f to
b03be36
Compare
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.
zaoxing
added a commit
that referenced
this pull request
Sep 25, 2026
…k.sh An independent review found the specs README claiming more than the models show. This corrects the claims, adds the configs that back or bound them, and adds a script that re-runs the checks and compares every verdict. - Ring: only payload_compute_spans's span arithmetic is checked, not a ring protocol. The CBMC harness gains P5, that no span byte lands in the unconsumed region [tail, head); it holds with the precondition and fails without it. P5 takes head's offset from off1 rather than recomputing head % cap, which does not finish. - LeaseLifecycle settles every ClickHouse call in zero time, while keep_lease() holds lease_mutex_ across up to three requests bounded only by request_s (60 s against a 15 s TTL). O1_slowreq3 (MaxLate 3, no cycle) holds and O1_slowreq (MaxLate 4) is refuted, so the O1 HOLDS verdicts are qualified: they need each lease request to finish in about half the TTL. A follow-up PR will bound lease request time. - LIMITS 5 cited catalog_writer.cpp:519's publish-timeout cap to dismiss a late-landing unknown outcome, but that cap covers only publish_snapshot; the lease INSERT has no max_execution_time. The model lands such a row at once, so O2_quar and O3_false hold only under that assumption, and of the two suggested self-latch fixes only resetting held_elsewhere_since_ns_ is robust. - NoSelfRefusal flagged ticks on which the service was quarantined and took no claim. It now counts only a claim actually taken and refused by the service's own row; O2_selfref (Skew 1) is still refuted, and it holds at Skew 0. The history variable changed, so the LeaseLifecycle state counts are re-recorded; no verdict changed. - PublisherLease.cfg's AllSafety holds vacuously for NoOverlappingAdmit at the base constants. base5_ovr (refuted) proves base5 non-vacuous, and noovr1 / ovr1 add the one-chunk pair beside noovr0 / ovr0. - Limitations now also say: Linearizable needs insert_quorum as well as select_sequential_consistency on a replicated catalog; skew runs one way; Stop is never enabled and its tombstone differs from the code's; each allocator allocates once, with no cross-call monotonicity check. - Line references move to #150's head, c0361d7 (storage_service.cpp +6 from run_cycle on, storage_service.h sweep comment at :104-108), and commit references from the orphaned 2b74d14 to 204a8d2, the same tree. - Nits: seven Z3 checks, not six; O5_cosweep's config allows one cut (the trace takes none, and MaxCuts 0 is refuted too); the .tla no longer points at an LLruns/ directory or a vac_cosweep config; apt's cbmc on Ubuntu 20.04 is too old. specs/check.sh runs the fast set (every config but PublisherLease.cfg, believers, holderssafe, overrun, nonlin_fence, stalepid and noovr1, which --all adds), z3/clock_skew.py and both CBMC builds, compares each result with a table of expected verdicts, and exits non-zero on a mismatch or on a .cfg the table does not list. Tools come from TLA2TOOLS_JAR, CBMC and PYTHON. TLC runs in a scratch copy, so nothing lands in the tree. It is not wired into CI.
# Conflicts: # tests/test_native_capture_storage_live.py
zaoxing
added a commit
that referenced
this pull request
Sep 26, 2026
Conflicts were side-by-side additions: helpers in native_capture.py and new tests appended at the ends of the storage wiring and live test files; both sides kept. One semantic fix: this branch required clickhouse_request_timeout_s to be at least twice a fixed 5 s _PUBLISH_TIMEOUT_S, written when the config did not expose the publish timeout. #150 made publish_timeout_s a config field, so the rule is now twice the configured publish_timeout_s, checked after the lease fields are validated; the integration doc says so, and a wiring test pins it (a 7 s publish cap needs a 14 s request timeout).
zaoxing
added a commit
that referenced
this pull request
Sep 26, 2026
One conflict, in src/dmi/storage/native_capture.py: this branch added NativeSinkConfig where main added the connection and lease validation helpers; both kept, helpers first. On the merged tree: pytest -m cpu 2520 passed; the capture storage, catalog lease, capture chain and reader parity live suites 126 passed. The ring, sink and engine code is identical to 52b9627, which passed the GPU suites (test_ring_engine 174/174, test_record_failure_policy_gpu 8/8).
zaoxing
added a commit
that referenced
this pull request
Sep 28, 2026
…k.sh An independent review found the specs README claiming more than the models show. This corrects the claims, adds the configs that back or bound them, and adds a script that re-runs the checks and compares every verdict. - Ring: only payload_compute_spans's span arithmetic is checked, not a ring protocol. The CBMC harness gains P5, that no span byte lands in the unconsumed region [tail, head); it holds with the precondition and fails without it. P5 takes head's offset from off1 rather than recomputing head % cap, which does not finish. - LeaseLifecycle settles every ClickHouse call in zero time, while keep_lease() holds lease_mutex_ across up to three requests bounded only by request_s (60 s against a 15 s TTL). O1_slowreq3 (MaxLate 3, no cycle) holds and O1_slowreq (MaxLate 4) is refuted, so the O1 HOLDS verdicts are qualified: they need each lease request to finish in about half the TTL. A follow-up PR will bound lease request time. - LIMITS 5 cited catalog_writer.cpp:519's publish-timeout cap to dismiss a late-landing unknown outcome, but that cap covers only publish_snapshot; the lease INSERT has no max_execution_time. The model lands such a row at once, so O2_quar and O3_false hold only under that assumption, and of the two suggested self-latch fixes only resetting held_elsewhere_since_ns_ is robust. - NoSelfRefusal flagged ticks on which the service was quarantined and took no claim. It now counts only a claim actually taken and refused by the service's own row; O2_selfref (Skew 1) is still refuted, and it holds at Skew 0. The history variable changed, so the LeaseLifecycle state counts are re-recorded; no verdict changed. - PublisherLease.cfg's AllSafety holds vacuously for NoOverlappingAdmit at the base constants. base5_ovr (refuted) proves base5 non-vacuous, and noovr1 / ovr1 add the one-chunk pair beside noovr0 / ovr0. - Limitations now also say: Linearizable needs insert_quorum as well as select_sequential_consistency on a replicated catalog; skew runs one way; Stop is never enabled and its tombstone differs from the code's; each allocator allocates once, with no cross-call monotonicity check. - Line references move to #150's head, c0361d7 (storage_service.cpp +6 from run_cycle on, storage_service.h sweep comment at :104-108), and commit references from the orphaned 2b74d14 to 204a8d2, the same tree. - Nits: seven Z3 checks, not six; O5_cosweep's config allows one cut (the trace takes none, and MaxCuts 0 is refuted too); the .tla no longer points at an LLruns/ directory or a vac_cosweep config; apt's cbmc on Ubuntu 20.04 is too old. specs/check.sh runs the fast set (every config but PublisherLease.cfg, believers, holderssafe, overrun, nonlin_fence, stalepid and noovr1, which --all adds), z3/clock_skew.py and both CBMC builds, compares each result with a table of expected verdicts, and exits non-zero on a mismatch or on a .cfg the table does not list. Tools come from TLA2TOOLS_JAR, CBMC and PYTHON. TLC runs in a scratch copy, so nothing lands in the tree. It is not wired into CI.
zaoxing
added a commit
that referenced
this pull request
Sep 28, 2026
#150 was squash-merged as 7419fd0, and main then took #151, #152 and #139. The specs cited #150's head c0361d7; they now cite main at 71be2af. - storage_service.cpp moved up two lines (#151 builds the ClickHouse client from one ClickHouseConnection), native_capture.py moved with #149/#151/#152, and deciding_read() is now clickhouse_client.cpp:374. - #151 also made execute() retry a read after a transient failure, up to max_attempts (3 by default), and never a write that may have reached the server. The README's O1 caveat, its Limitations entry and LIMITS 3 in LeaseLifecycle.tla now say a lease request's reads can take up to three request timeouts, and RenewIfDue says the quarantining exception is the first to outlast those retries. No modelled outcome changes. - Refs that missed the code they describe, in files main did not change: the O1a quote is storage_service.h:233-234, not storage_service.cpp; the renewal in publish_snapshot is catalog_writer.cpp:490 and :579; publish_snapshot ends at :669; the config check with the quorum rule is :148-169; the chunk loop is :531; the watermark read-back is :608-628 (:611-627 for its refusal); the version allocator's statement lines; reject_live's comparison is lease_coordinator.cpp:222; and the start wait's knob checks are native_capture.py:351-354. - The README says what 204a8d2 is now that #150's branch is squashed. Comment and prose changes only; every verdict is unchanged.
zaoxing
added a commit
that referenced
this pull request
Sep 28, 2026
…ithmetic (#157) * Add TLA+, Z3 and CBMC specs for the lease, allocator and ring protocols * Say in the specs README that a failed renewal is never retried * Expect O1_tries5 to be refuted, and say which O1 caveats are model results * Say what the specs leave out, and check every verdict with specs/check.sh * Point the specs' line references at main after #150's squash * Correct the specs' retry wording and three stale citations * Point the specs' line references past #154's comments * Check the four-wake claim at TTL 6, and fix four stale spec citations
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Milestone B3 of the native capture production plan. It fixes three ways the publisher lease on the capture catalog could stop indexing, or fail a restart, when nothing was really wrong.
What changes
1. Recovery after a ClickHouse error (the main fix).
flush()raises.failed,lease_state(none / held / quarantined / reacquiring / failed / released),quarantined_untilandlease_reacquisitions.2. Restart after a crash. A SIGKILLed process leaves its lease live for up to one TTL.
start()used to refuse at once. It now waits up tostart_lease_wait_s, and still sweeps the spool only after it holds the lease.3. The lease settings are now configurable from Python.
NativeCaptureStorageConfiggains four fields:lease_ttl_s(15 s);publish_timeout_s(5 s);clock_skew_s(0 s);start_lease_wait_s. The default islease_ttl_s + publish_timeout_s + clock_skew_s, and 0 means fail at once.Validation requires TTL > publish_timeout + skew + 0.1 s.
Unchanged: the lease coordinator, the catalog writer and every lease statement.
tests/test_native_catalog_lease_live.py, which checks that the lease statements are byte-identical, passes without changes. Nothing undersrc/dmi/storage/capture/changed.Evidence
The live tests ran against a local ClickHouse, with a TCP switch in front of it to cut connections, a 3 s TTL and a 1 s publish timeout:
flushraisedflushraised "no publisher lease is held"test_native_capture_storage_live.pyandtest_native_catalog_lease_live.pygave 72 passed.test_native_capture_storage_wiring.pygave 68 passed.No uploads while the service cannot index (commit after review)
The plan said uploads continue during a quarantine. But an uploaded, unindexed pack is remembered only in the in-memory
pending_index_. A crash inside the window then leaves it in the bucket and never in the catalog whenreconcile_on_start=False. By the author's decision, a cycle now uploads only when it holds the lease and owes nothing, so packs staged during a quarantine stay in the spool, which survives a crash. The new live testtest_a_quarantined_service_leaves_new_packs_in_the_spoolfailed before the change (uploaded_packs1, expected 0) and passes now. The live storage and lease suites give 73 passed.Independent review
Verdict: ship, with no majors.
The reviewer found 5 minors. This PR fixes three (commit 2b74d14):
clock_skew_s. It also documents that a predecessor with a longer TTL (the old native 30 s default) can outlast the wait.test_one_publisher_per_catalogpassesstart_lease_wait_s=0(20.7 s → 2.1 s in CI).Left open: