Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion native/csrc/catalog/storage_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
72 changes: 72 additions & 0 deletions tests/test_native_capture_storage_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Comment on lines +3443 to +3444
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 (
Expand Down
Loading