From e88dc5631803fca566dd51237dc60ed872652bec Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 18:35:08 +0000 Subject: [PATCH 1/4] Say what the lease comments actually guarantee Four comments promised more, or less precisely, than the code does: - renew_lease_if_due() said renewing at ttl/3 leaves two more tries. Four lease-thread wakes do fit before the row expires, but they only cover a renewal that runs late. A failed renewal costs the lease at once: a refusal drops it in the coordinator, and a transport error or timeout quarantines the writer on the first failure. - The sweep_spool_on_start comment in storage_service.h still read as the strong sweep-after-lease guarantee that 2b74d14 weakened in the .cpp. It now says the same: a quarantined holder lets its row lapse, a second process can take the lease and sweep, and the spool is not locked. - start_lease_wait_s said a predecessor with a longer TTL can outlast the default wait. The exact condition is a predecessor TTL above lease_ttl_s + publish_timeout_s (clock_skew_s cancels), 20 s on the defaults, so the native 30 s default falls 10 s short. - The claim read-back takes the expiry from rows[0] of an unordered read, while the rule is the minimum expires_at_ns under (term, lease_id). The comment says why that is safe here, so the next reader need not re-derive it. No code changes. --- native/csrc/catalog/lease_coordinator.cpp | 8 ++++++++ native/csrc/catalog/storage_service.cpp | 9 +++++++-- native/csrc/catalog/storage_service.h | 6 +++++- src/dmi/storage/native_capture.py | 12 +++++++----- 4 files changed, 27 insertions(+), 8 deletions(-) diff --git a/native/csrc/catalog/lease_coordinator.cpp b/native/csrc/catalog/lease_coordinator.cpp index 05633ce4d..8bb8229bf 100644 --- a/native/csrc/catalog/lease_coordinator.cpp +++ b/native/csrc/catalog/lease_coordinator.cpp @@ -119,6 +119,14 @@ PublisherLease LeaseCoordinator::claim_with_rival( std::set owners; for (const Row& row : rows) owners.insert(row[0]); if (owners == std::set{lease_id}) { + // rows[0] is this claim's own row, unordered and ungrouped though the + // read is: every row at this term carries our lease_id, and we wrote one + // (an insert of unknown outcome quarantines rather than retrying). The + // expiry rule is the minimum expires_at_ns under (term, lease_id), which + // differs from any one row only once a release's tombstone shares the + // key, and none can here: a release writes at the holder's own term and + // a claim always goes to head + 1, so no claim lands on a term that + // already holds its own tombstone. lease_ = PublisherLease{ term, lease_id, holder, parse_u64_field(rows[0][1], "lease acquisition"), diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 941a3a038..0c00df0e8 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -632,8 +632,13 @@ void CaptureStorageService::keep_lease() { } void CaptureStorageService::renew_lease_if_due() { - // Renew once a third of the TTL has passed without a publish, which leaves - // two more tries before a rival could claim it. + // Renew once a third of the TTL has passed without a publish. The lease + // thread wakes every ttl/6, so four wakes fall between the renewal coming + // due and the row expiring: room for a renewal that runs late, not for one + // that fails. A failed renewal costs the lease at once whatever the cause. + // A refusal drops it in the coordinator, and a transport error or timeout + // takes renew_for_publish()'s catch, which quarantines the writer on the + // first one. const uint64_t ttl = config_.writer.lease_ttl_ns; if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 13b77dd87..a96f98b80 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -105,7 +105,11 @@ struct StorageServiceConfig { // 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. It runs after the lease is taken, so a start refused - // the catalog never touches the spool. + // the catalog never touches the spool. That keeps a second process off a + // live spool only usually: a holder that is quarantined lets its row lapse, + // and a second process can take the lease and sweep while the first is + // still writing. The spool itself is not locked; one process per spool is + // the caller's job. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index d8e8022e1..599c0ff07 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -269,11 +269,13 @@ class NativeCaptureStorageConfig: 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 + - # 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. + # clock_skew_s; 0 fails at once. That outlasts a crashed predecessor + # exactly when its TTL is at most lease_ttl_s + publish_timeout_s: + # clock_skew_s cancels, since the wait adds it and a lagging replica sees + # the row live that much longer. Set this explicitly for the first + # restart after a predecessor above that threshold -- 20 s on these + # defaults, so one on the native 30 s default (which processes predating + # these knobs used) outlasts the default wait by 10 s. start_lease_wait_s: Optional[float] = None def __post_init__(self) -> None: From 40a4b701d851a9549914d348fb83f03d7cd1e7a0 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 19:37:14 -0400 Subject: [PATCH 2/4] Tighten the lease comments the review found still overclaiming - renew_lease_if_due(): "four wakes fall between the renewal coming due and the row expiring" was nominal at best. last_renew_ns_ is stamped after the round trip, the tick counts from the end of the previous pass, and an index() call can hold lease_mutex_. It now says the renewal fires about a tick after it falls due, leaving roughly half the TTL for it to land. A server error quarantines as well as a transport error or timeout, and "on the first one" now allows for the client's read retries: a write that may have reached the server is never retried. - sweep_spool_on_start: any holder that stops renewing for a TTL, not only a quarantined one, lets a second process take the lease and sweep; one on another (database, table_prefix) never meets the lease; and the second process's Recover() lists the first's sealed packs, so both upload one spool. - start_lease_wait_s: the default wait is guaranteed to outlast a crashed predecessor only when its TTL is at most lease_ttl_s + publish_timeout_s, and clock_skew_s cancels only if it bounds the offset between the replica that stamped the row and the one serving the read. - claim_with_rival(): the client never re-sends an INSERT, and it is CatalogWriter that quarantines after one of unknown outcome. Comments only. --- native/csrc/catalog/lease_coordinator.cpp | 13 +++++++------ native/csrc/catalog/storage_service.cpp | 15 +++++++++------ native/csrc/catalog/storage_service.h | 11 +++++++---- src/dmi/storage/native_capture.py | 16 +++++++++------- 4 files changed, 32 insertions(+), 23 deletions(-) diff --git a/native/csrc/catalog/lease_coordinator.cpp b/native/csrc/catalog/lease_coordinator.cpp index 8bb8229bf..cb3f1be6f 100644 --- a/native/csrc/catalog/lease_coordinator.cpp +++ b/native/csrc/catalog/lease_coordinator.cpp @@ -121,12 +121,13 @@ PublisherLease LeaseCoordinator::claim_with_rival( if (owners == std::set{lease_id}) { // rows[0] is this claim's own row, unordered and ungrouped though the // read is: every row at this term carries our lease_id, and we wrote one - // (an insert of unknown outcome quarantines rather than retrying). The - // expiry rule is the minimum expires_at_ns under (term, lease_id), which - // differs from any one row only once a release's tombstone shares the - // key, and none can here: a release writes at the holder's own term and - // a claim always goes to head + 1, so no claim lands on a term that - // already holds its own tombstone. + // (the client never re-sends an INSERT, and CatalogWriter quarantines + // after one of unknown outcome rather than retrying). The expiry rule is + // the minimum expires_at_ns under (term, lease_id), which differs from + // any one row only once a release's tombstone shares the key, and none + // can here: a release writes at the holder's own term and a claim always + // goes to head + 1, so no claim lands on a term that already holds its + // own tombstone. lease_ = PublisherLease{ term, lease_id, holder, parse_u64_field(rows[0][1], "lease acquisition"), diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 0c00df0e8..11da5044d 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -633,12 +633,15 @@ void CaptureStorageService::keep_lease() { void CaptureStorageService::renew_lease_if_due() { // Renew once a third of the TTL has passed without a publish. The lease - // thread wakes every ttl/6, so four wakes fall between the renewal coming - // due and the row expiring: room for a renewal that runs late, not for one - // that fails. A failed renewal costs the lease at once whatever the cause. - // A refusal drops it in the coordinator, and a transport error or timeout - // takes renew_for_publish()'s catch, which quarantines the writer on the - // first one. + // thread wakes every ttl/6, so the renewal fires about a tick after it + // falls due (later if an index() call holds lease_mutex_), leaving roughly + // half the TTL for it to land before the row expires. That slack covers a + // renewal that runs late, not one that fails. A failed renewal costs the + // lease at once whatever the cause. A refusal drops it in the coordinator, + // and any ClickHouse error (transport, timeout, or a server error) takes + // renew_for_publish()'s catch, which quarantines the writer on the first + // error that survives the client's retries (reads only; a write that may + // have reached the server is never retried). const uint64_t ttl = config_.writer.lease_ttl_ns; if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index a96f98b80..99fe63941 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -106,10 +106,13 @@ struct StorageServiceConfig { // is only safe while no writer is live: start() must run before the sink // opens the spool. It runs after the lease is taken, so a start refused // the catalog never touches the spool. That keeps a second process off a - // live spool only usually: a holder that is quarantined lets its row lapse, - // and a second process can take the lease and sweep while the first is - // still writing. The spool itself is not locked; one process per spool is - // the caller's job. + // live spool only usually: a holder that stops renewing for a TTL + // (quarantined, or stalled) lets its row lapse, and a second process can + // take the lease and sweep while the first is still writing; one on + // another (database, table_prefix) never meets the lease at all. Its + // Recover() also lists the first's sealed packs, so both then upload the + // one spool. The spool itself is not locked; one process per spool is the + // caller's job. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index 599c0ff07..aed635898 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -269,13 +269,15 @@ class NativeCaptureStorageConfig: 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 + - # clock_skew_s; 0 fails at once. That outlasts a crashed predecessor - # exactly when its TTL is at most lease_ttl_s + publish_timeout_s: - # clock_skew_s cancels, since the wait adds it and a lagging replica sees - # the row live that much longer. Set this explicitly for the first - # restart after a predecessor above that threshold -- 20 s on these - # defaults, so one on the native 30 s default (which processes predating - # these knobs used) outlasts the default wait by 10 s. + # clock_skew_s; 0 fails at once. That is guaranteed to outlast a crashed + # predecessor only when its TTL is at most lease_ttl_s + + # publish_timeout_s: clock_skew_s cancels, since the wait adds it and a + # lagging replica sees the row live that much longer -- which assumes + # clock_skew_s bounds the offset between the replica that stamped the + # predecessor's row and the one serving the read. Set this explicitly + # for the first restart after a predecessor above that threshold -- 20 s + # on these defaults, so one on the native 30 s default (which processes + # predating these knobs used) outlasts the default wait by 10 s. start_lease_wait_s: Optional[float] = None def __post_init__(self) -> None: From ad129a1def79d827fce837f1d95cdd9de4c37197 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Sun, 27 Sep 2026 20:27:01 -0400 Subject: [PATCH 3/4] Match the lease comments to the client's retries and a quarantined cycle Rebased onto main, three of the rewritten comments said more than the code now does: - The ClickHouse client retries a write when its connection was never made, since nothing reached the server. The claim read-back no longer says the client never re-sends an INSERT, and the renewal comment no longer calls the retries reads-only. A write that may have reached the server is still never repeated, so the claim still wrote one row. - A quarantined cycle uploads nothing, so the spool-sweep caveat no longer says both processes then upload the one spool. A stalled holder that still holds its lease locally does, until a renewal or publish is refused, and so do two processes on different (database, table_prefix) pairs. Comments only. --- native/csrc/catalog/lease_coordinator.cpp | 5 +++-- native/csrc/catalog/storage_service.cpp | 5 +++-- native/csrc/catalog/storage_service.h | 9 ++++++--- 3 files changed, 12 insertions(+), 7 deletions(-) diff --git a/native/csrc/catalog/lease_coordinator.cpp b/native/csrc/catalog/lease_coordinator.cpp index cb3f1be6f..b54f49f56 100644 --- a/native/csrc/catalog/lease_coordinator.cpp +++ b/native/csrc/catalog/lease_coordinator.cpp @@ -121,8 +121,9 @@ PublisherLease LeaseCoordinator::claim_with_rival( if (owners == std::set{lease_id}) { // rows[0] is this claim's own row, unordered and ungrouped though the // read is: every row at this term carries our lease_id, and we wrote one - // (the client never re-sends an INSERT, and CatalogWriter quarantines - // after one of unknown outcome rather than retrying). The expiry rule is + // (the client repeats an INSERT only when its connection was never made, + // so nothing reached the server, and CatalogWriter quarantines after one + // of unknown outcome rather than retrying). The expiry rule is // the minimum expires_at_ns under (term, lease_id), which differs from // any one row only once a release's tombstone shares the key, and none // can here: a release writes at the holder's own term and a claim always diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 11da5044d..bf4635f3c 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -640,8 +640,9 @@ void CaptureStorageService::renew_lease_if_due() { // lease at once whatever the cause. A refusal drops it in the coordinator, // and any ClickHouse error (transport, timeout, or a server error) takes // renew_for_publish()'s catch, which quarantines the writer on the first - // error that survives the client's retries (reads only; a write that may - // have reached the server is never retried). + // error that survives the client's retries (a write is repeated only when + // its connection was never made; one that may have reached the server + // never is). const uint64_t ttl = config_.writer.lease_ttl_ns; if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 99fe63941..910437c7b 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -110,9 +110,12 @@ struct StorageServiceConfig { // (quarantined, or stalled) lets its row lapse, and a second process can // take the lease and sweep while the first is still writing; one on // another (database, table_prefix) never meets the lease at all. Its - // Recover() also lists the first's sealed packs, so both then upload the - // one spool. The spool itself is not locked; one process per spool is the - // caller's job. + // Recover() also lists the first's sealed packs, which the first may + // upload too: a quarantined holder uploads nothing without the lease, but + // a stalled one still holds it locally and uploads until a renewal or + // publish is refused, and two on different (database, table_prefix) + // pairs each hold a lease and upload freely. The spool itself is not + // locked; one process per spool is the caller's job. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; From 68e6aeee8887b234d5191fc1dc0758ba2e9386d9 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Sun, 27 Sep 2026 20:42:29 -0400 Subject: [PATCH 4/4] Say where a lost lease still uploads and what an index() does to the slack The sweep_spool_on_start comment said a quarantined holder uploads nothing without the lease. run_cycle() checks the lease once, before its upload batch, and UploadPending does not stop when the lease is lost, so a holder quarantined or refused while a batch is in flight finishes that batch (which can outlast the TTL) while a second process may already have taken the lease and swept; only its later cycles upload nothing. The renew_lease_if_due() comment counted the half-TTL slack from the row's last renewal, but the check reads last_renew_ns_, which index_bounded() stamps when index() returns: after the publish's own last renewal, and even when every pack was already committed so nothing was published or renewed. With index() holding lease_mutex_ meanwhile, the slack can be far less than half the TTL, or none. Say so. --- native/csrc/catalog/storage_service.cpp | 15 +++++++++++---- native/csrc/catalog/storage_service.h | 14 +++++++++----- 2 files changed, 20 insertions(+), 9 deletions(-) diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index bf4635f3c..cc1254c96 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -632,10 +632,17 @@ void CaptureStorageService::keep_lease() { } void CaptureStorageService::renew_lease_if_due() { - // Renew once a third of the TTL has passed without a publish. The lease - // thread wakes every ttl/6, so the renewal fires about a tick after it - // falls due (later if an index() call holds lease_mutex_), leaving roughly - // half the TTL for it to land before the row expires. That slack covers a + // Renew once a third of the TTL has passed since last_renew_ns_. The lease + // thread wakes every ttl/6, so with no index() in the way the renewal + // fires within about a tick of falling due, leaving at least roughly half + // the TTL for it to land before the row expires. An index() can leave far + // less, or none. It holds lease_mutex_ throughout, so this thread cannot + // renew until it returns, and index_bounded() then stamps last_renew_ns_ + // whenever it indexed or skipped a pack, as though the row had just been + // renewed. After a publish that stamp trails the publish's own last + // renewal by the watermark INSERT, its read-backs and commit_packs; and + // it is taken even when every pack was already committed, so index() + // published nothing and renewed nothing. Whatever slack is left covers a // renewal that runs late, not one that fails. A failed renewal costs the // lease at once whatever the cause. A refusal drops it in the coordinator, // and any ClickHouse error (transport, timeout, or a server error) takes diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 910437c7b..3c798e04a 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -111,11 +111,15 @@ struct StorageServiceConfig { // take the lease and sweep while the first is still writing; one on // another (database, table_prefix) never meets the lease at all. Its // Recover() also lists the first's sealed packs, which the first may - // upload too: a quarantined holder uploads nothing without the lease, but - // a stalled one still holds it locally and uploads until a renewal or - // publish is refused, and two on different (database, table_prefix) - // pairs each hold a lease and upload freely. The spool itself is not - // locked; one process per spool is the caller's job. + // upload too. A cycle checks the lease once, before its upload batch, and + // UploadPending does not stop when the lease is lost: a holder that is + // quarantined, or refused a renewal or publish, while a batch is in + // flight finishes that batch (which can outlast the TTL), and only its + // later cycles upload nothing while it holds no lease. A stalled one + // still holds the lease locally and starts new batches until a renewal + // or publish is refused. Two on different (database, table_prefix) pairs + // each hold a lease and upload freely. The spool itself is not locked; + // one process per spool is the caller's job. bool sweep_spool_on_start = true; bool reconcile_on_start = true; };