From 7e73115f352f6a62d518c2ad76670d0fdf9da9f6 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Tue, 6 Oct 2026 15:46:49 -0400 Subject: [PATCH] End a refusal run at a claim that was sent, not only at our own late row A claim sends its INSERT only once its head read found no live holder, so whoever refused the service before had left by then. ensure_publisher_lease reset the refusal clock for a refusal by the service's own late claim row, but not for a claim whose INSERT timed out: a rival that left, such a claim, and a second rival refusing the claims after the quarantine added up to 2 x TTL and latched the service over a catalog nobody held in between. The unknown-outcome branch now restarts the clock whenever the writer quarantined, which in that branch is exactly when the INSERT may have been sent. --- native/csrc/catalog/storage_service.cpp | 11 +++- tests/test_native_capture_storage_live.py | 72 +++++++++++++++++++++++ 2 files changed, 82 insertions(+), 1 deletion(-) diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index 16c747e88..f7333bd9e 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -1654,7 +1654,16 @@ bool CaptureStorageService::ensure_publisher_lease() { record_error(std::string("publisher lease acquisition failed: ") + exc.what()); note_lease_failure(exc); - if (!writer_.quarantined()) { + if (writer_.quarantined()) { + // A claim sends its INSERT only once its head read found no live + // holder, so whoever refused this service before had left by then: + // the run of refusals ends here, as it does at a refusal by our own + // late row (lease_held_elsewhere). Otherwise a rival that left, this + // claim timing out and a second rival refusing the claims after the + // quarantine add up to 2 x TTL, and latch the service over a catalog + // nobody held in between. + held_elsewhere_since_ns_ = 0; + } else { next_claim_ns_ = steady_ns() + lease_tick_ns(config_.writer.lease_ttl_ns); } publish_lease_state(); diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index 9457919e8..cae1cae83 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -3410,6 +3410,78 @@ def test_a_late_landing_claim_of_its_own_does_not_latch_the_service( assert _latch_lines(capfd.readouterr().err) == [] +def test_a_claim_sent_while_nobody_held_the_lease_restarts_the_refusal_clock( + fake_s3, tmp_path, capfd): + """A claim only sends its INSERT once its head read found no live + holder, so a rival's run of refusals ended there, even when the INSERT + then timed out. The refusal clock ran on through such a claim: a rival + that left, a claim that timed out, and a second rival refusing the + claim after the quarantine added up to 2 x TTL, and the service latched + "held by another publisher" over a catalog nobody had held between the + two.""" + 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") + second = _service(_storage_config( + fake_s3, catalog.table_prefix, holder="second-publisher", + **knobs), tmp_path / "second") + first.start() + try: + # A cut quarantines the first service; the rival takes over. + switch.cut() + rival.start() + switch.restore() + # The first service's quarantine ends and the rival refuses it: + # the refusal clock starts. + _wait_for(lambda: first.snapshot()["lease_state"] == "reacquiring", + timeout_s=10.0) + refused_at = time.monotonic() + # The rival leaves (a handover), and the first service's next + # claim finds nobody holding the lease, sends its INSERT, and + # hears nothing back: the 1 s bound on a claim made without a + # lease (lease_ttl_s / 3) gives up, and it quarantines. The + # INSERT never reaches the server, so no row of its own is left + # to refuse it later. + switch.stall_requests(_lease_insert) + rival.stop() + _wait_for(lambda: first.snapshot()["lease_state"] == "quarantined", + timeout_s=5.0) + assert switch.stalled, "the claim INSERT was never sent" + switch.restore() + # A second publisher takes the free lease, and refuses the first + # service's claims once its quarantine ends -- more than 2 x TTL + # after the first rival's first refusal, but a fresh run. + second.start() + time.sleep(max(0.0, refused_at + 2 * 3.0 + 2.0 - time.monotonic())) + snapshot = first.snapshot() + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_state"] == "reacquiring", snapshot + assert "second-publisher" in snapshot["last_error"], snapshot + # A handover: the second publisher leaves, and the first service + # takes the lease back. + second.stop() + _wait_for(lambda: first.snapshot()["lease_state"] == "held" + or first.snapshot()["failed"], timeout_s=10.0) + snapshot = first.snapshot() + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_state"] == "held", snapshot + first.rethrow_if_failed() + finally: + first.stop() + rival.stop() + second.stop() + switch.close() + + assert _latch_lines(capfd.readouterr().err) == [] + + _HOLD_LEASE = """ import json, sys, time from dmi.storage.native_capture import (