Skip to content

close() delivers the tail pack; flush() and stop() return on time while the store or catalog stalls - #162

Merged
zaoxing merged 33 commits into
mainfrom
feat/close-flush-backstop
Sep 29, 2026
Merged

zaoxing merged 33 commits into
mainfrom
feat/close-flush-backstop

Conversation

@zaoxing

@zaoxing zaoxing commented Sep 29, 2026

Copy link
Copy Markdown
Collaborator

Milestone B5 of the native capture production plan: close() delivers the tail pack, and flush(timeout) and stop() return on time even when the object store or the catalog stops answering.

What changes

The tail reaches the catalog on close (648f6ce, 7a60f35).

  • The problem: NativePackSink::on_engine_release now stages the sink's open pack with a bounded, noexcept Flush. Before, a sink with a long linger left its last pack unsealed, so the audit probe "ready packs right after stop" reported 0.
  • The fix: the new flush is bounded by release_flush_timeout_s (0 disables it; nan or negative values are refused). It fails soft, with one line on stderr, and waits only its own timeout if another flush is in flight.

Uploads can be cut short (fa250a1, b61db70, 4d90b5a, 41398b8, 8bc3fb7, 119cea8).

  • stop() and a flush past its deadline trip a Cancellation.
  • CURLOPT_XFERINFOFUNCTION aborts the transfer, and a cut multipart upload is aborted. The abort itself isn't cancellable and is bounded at 5 s; that overrun is documented.
  • Both the S3 and uploader backoffs wait on an interruptible condition variable.
  • A pack is booked as cancelled only when the cancel actually cut its request. A real failure on the last attempt stays a failure and keeps its last_error.

Flush honours its deadline, and stop() returns promptly (5f882cd, b13c7ec, 0121538, efe9279, 5fecd88, f563341, d677f07, 9acff4d, e49739f).

  • run_cycle(deadline, allow_reconcile): a flush cycle never reconciles; only the loop does.
  • Index reads go through their own Cancellation. stop() cancels it for good, and a flush arms it one catalog request timeout past its deadline.
  • A read the store never answers ends the index pass without counting against the pack. So a read outage no longer sets sound packs aside, and the new s3_read_timeout_s and s3_max_attempts settings bound those reads.
  • Each cycle lists the spool once, then uploads and indexes a chunk of indexer.max_packs at a time. A cycle cut at its deadline therefore owes at most one chunk, instead of everything it uploaded. Past a flush's deadline, only a full first batch is indexed.
  • A spool listing stops between packs at stop() or at a flush's deadline, instead of hashing the whole backlog first.
  • No loop cycle runs between a flush and the stop() after it.
  • Catalog statements are never aborted mid-flight. Bound how long the publisher lease can go unrenewed: a deadline for every request under it #159's per-request deadlines bound them.

Docs (af4f4f2, 90d00b0, 58453c8, bb774c4, 8bc3fb7, 7525aac, 1303d85). The close_flush_timeout_s comment, engine.close() and flush_and_wait in engine.py, config.py and docs/integration-api-v1.md now state what the budget actually bounds:

  • the drain can outlast it by the catalog work in flight at the deadline;
  • stop()'s lease release, or a lease renewal in flight;
  • the ≤ 5 s multipart abort.

Deliberately stricter than the Python reference: chunked upload and index, cancellable index reads, and cancel booking in pack_index. The oracle has none of these. The parity suites are unchanged and pass.

Evidence

  • CPU tier: pytest -m cpu, 2596 passed. The 1 skip is the non-CPU case with no CUDA.
  • Live suites: capture storage 62, catalog lease 54 (file unmodified), capture chain 4, reader parity 52. The changed CPU files (wiring, pack read cancel, pack sink timeout, S3 client, sink release, uploader) gave 200 passed in each of 3 runs.
  • Repeats: the 22 new or changed live tests passed 22/22 in each of 3 runs.
  • GPU, on one RTX 4090 through an idle-GPU gate: tests/test_native_capture_storage_gpu_e2e.py passed in each of 3 runs. That covers close() without flush_and_wait leaving every capture queryable, and stopping the ring staging the open pack without a service.
  • Red first:
    • flush(1.0) against stalled index reads took 40.3 s and is now about 5 s;
    • stop() during stalled index reads went from 26.8 s to under 3 s;
    • a flush against a slow catalog went from 59.3 s to 7.8 s;
    • a stopped-then-restarted service used to cancel every upload.
  • Mutations, each turning tests red:
    • dropping the release Flush;
    • an XFERINFO callback that never aborts;
    • removing read_cancel_ from stop();
    • disabling the no-cycle-before-stop skip;
    • stopping the lease thread before the loop's join.

Review

  • First round: three independent lenses (shutdown correctness, cancellation and concurrency, tests and compatibility) produced 14 findings, including two majors. Every one survived adversarial verification, and all are fixed.
  • Second round: a re-review of the whole branch found 13 more, including one major, now fixed in d677f07. When close() ran out of budget during uploads, most of the packs it had just uploaded never reached the catalog. All 13 are fixed. The final check passed.

Known, unchanged: tests/test_native_catalog_lease_live.py's _catalog_drop_only never drops its tables. That's already on main, and the file is the lease SQL byte-identity gate.

RingEngine::stop drains its record worker into the sink and then releases
the sink's lease without flushing it, and NativePackSink::on_engine_release
was empty. The sink seals a pack only when it fills, when it has lingered
max_linger_ns, or on a flush, so the records of the pack still open at the
stop stayed in memory: the audit's probe counted 0 ready packs right after
the stop (1 once the linger fired), and a process that exited first lost
them. close() flushes the sink before the ring stops only when a storage
service runs (#143); with the sink alone, or after that flush failed,
nothing staged the tail.

on_engine_release now flushes the sink, bounded by release_flush_timeout
(30 s by default; release_flush_timeout_s on the binding, 0 turns it off),
and cannot throw, since it runs in RingEngine::stop and its destructor. A
flush that fails or times out writes one line to stderr, and the failure
stays latched for rethrow_if_failed.

Red first: test_the_open_pack_is_staged_when_the_engine_releases_the_sink,
a 60 s linger and three admitted records, found 0 ready packs after the
release ("assert 0 == 1"); it finds exactly 1, with the sink object still
alive, so its destructor is not what wrote it. The GPU test drives the same
through a real engine's close() with no storage service.
Nothing could interrupt an upload. A PUT to a store that accepted the
connection and never answered held its caller for read_timeout_s (120 s by
default) on each of the S3 client's attempts, with backoff sleeps between,
and the uploader then retried the pack with sleeps of its own. The storage
service's stop() and flush(timeout) both wait on such a cycle, which is
what the next commit fixes; this one gives them the means.

dmi_store::Cancellation (store/cancel.h, header-only): Cancel() for good,
or set_deadline() for one stretch, and a SleepFor that wakes for either.

S3Client::set_cancellation: no request goes out once cancelled; a transfer
in flight is aborted through CURLOPT_XFERINFOFUNCTION, which libcurl calls
at least once a second, connecting included; a retry backoff waits on the
cancellation instead of sleeping. A cancelled request fails with "request
cancelled" and is not retried. A multipart upload cut short is aborted
with a request the cancel does not cut: one attempt, bounded by 5 s, since
whoever cancelled is waiting.

SpoolUploader::set_cancellation: no pack starts once cancelled, and a
pack's retries and backoff end. The pack stays staged and is reported
cancelled -- UploadFailure::cancelled, snapshot cancelled_packs -- not as
a failed pack.

The store driver takes cancel_after_ms, so the client and uploader are
tested at this level. Red first, each against the fake S3:
- a PUT held 5 s by fault/hang completed ok instead of being cut;
- a 10 MiB multipart upload whose parts are held (new fault/hang-parts)
  completed instead of being cut and aborted;
- a HEAD against fault/always-500 with 10 attempts ran all 10 (~26 s);
- upload_one against fault/always-500 ran all 4 upload attempts (~7 s).
Now each ends within 0.5-2 s of the cancel, the multipart upload aborted
and the pack still in the spool.

The Python uploader, the reference, has no cancellation; this is C++ only.
…alls

flush() took the cycle lock by its deadline (#143), but a cycle it ran
itself paid no attention to it. It uploaded everything staged, however
long the store took, and it ran the periodic reconcile when that fell due:
past an index pass that failed on a catalog that accepts connections and
never answers, the reconcile listed the bucket and asked the catalog
again. stop() joined a loop whose cycle could be inside a PUT the store
never answers, for the S3 timeouts of every attempt, with the lease held.

- run_cycle(deadline, allow_reconcile). flush() passes its deadline and
  skips the reconcile, which only the loop runs now.
- The uploads go through an S3 client of their own, sharing one
  Cancellation with the uploader. A flush arms its deadline on it for the
  cycle it runs; stop() cancels it for good before joining the loop. A
  cancelled pack stays in the spool: no upload starts, one in flight is
  aborted, a backoff ends. A flush that finds stop() under way returns
  false instead of indexing until its deadline.
- Nothing else is cut. What a cycle has uploaded it still indexes, since
  until the catalog has it only this process remembers it (the reason
  uploads stop while the service cannot index, #150); the index pass reads
  those packs through the other client. Catalog statements are never cut
  mid-flight: each is bounded by the request timeout, or under the lease
  by the lease deadline (#159), and a pass stops at its first failure. So
  a flush overruns its deadline by the catalog work in flight at it --
  against a catalog that stopped answering, one request timeout.
- A cycle a cancel cut short is neither drained nor failed: it moves the
  backoff neither way. The reconcile stops between requests at stop().
  snapshot() counts cancelled_uploads, apart from upload_failures.
- Past its deadline a flush's cycle still lists the spool, so flush(0) on
  a drained service reports drained.

Red first, against the live catalog and the fake S3 behind a TCP switch:
- flush(1.0) against a black-holed catalog with a reconcile due took
  8.02 s -- the index pass's request timeout (4 s) and the reconcile's --
  over the 1 + 4 + 1 s bound;
- flush(1.0) with the store holding the PUT was still blocked after 15 s;
- stop() with the loop's PUT held was still blocked after 15 s.
Now the first returns within one request timeout of its deadline, the
other two within 3 s; the lease is released, the pack is still staged, and
a successor uploads and indexes it.
The audit found the close() docs promising more than close() does:
close_flush_timeout_s read as "the total budget for ... getting every
staged pack into the catalog", but close() drains best effort -- it logs
what did not drain and stops the service regardless -- and until the
previous commit the drain could outlast the budget by any upload in
flight. close()'s own docstring said only "Tear down backend resources",
and the retire path promised that the next start "uploads or reconciles"
what was left, when only a start with reconcile_on_start finds a pack
that was uploaded but not yet indexed.

Now they say: close() drains within close_flush_timeout_s, best effort;
what misses it stays in the spool (the next start uploads it) or, uploaded
but unindexed, in the bucket (only reconcile_on_start indexes it); the
drain can outlast the budget by the catalog work in flight -- one request
timeout against a catalog that stopped answering -- never by an upload,
and by the sink's release flush (30 s at most) only when the sink itself
is stuck; flush_and_wait is the call that raises. Without a service, the
sink's open pack reaches the spool when the ring releases it. The v1
integration doc says the same in its close() section.

The engine's close order (sink flush, ring stop, service flush, service
stop) was already right (#143) and is pinned by
test_close_flushes_the_sink_before_the_ring_stops; only the words change.
PackSink::Flush took flush_mutex_ with an untimed lock and held it across
its whole barrier wait, so a second flush could not start its own bounded
wait until the first returned. The ring's stop runs the release backstop
(on_engine_release), a flush bounded by release_flush_timeout, 30 s by
default; beside an engine.flush_and_wait on another thread -- both release
the GIL, and nothing forbids close() during it -- the backstop waited out
that flush's timeout, 600 s by default, while the docs promised 30 s at
most. Before the backstop the release was a no-op, so this is new with it.

flush_mutex_ is a std::timed_mutex now, taken with try_lock_until the
flush's own deadline; a flush that cannot get it in time returns false, as
a wait on the barrier would have. A flush with no timeout still waits.

Red first: a third case in tests/native/test_pack_sink_timeout.cpp parks
the stager under a Flush(4.0) on one thread and calls Flush(0.5) on
another: it returned after 3.90 s. It returns within its 0.5 s now, and the
first flush still completes once the stager is released.

The on_engine_release comments said a failed or timed-out release flush
"stays latched for rethrow_if_failed". A released sink refuses that call
as not attached, and a timeout latches nothing. They now say the stderr
line is the only report of a timeout, and that a pipeline failure also
counts in snapshot()["failures"], which rethrow_if_failed reports only
once the sink is attached again.
…ed read

The previous commit cut uploads short but left the index pass reading
through an S3 client no cancel reached, and NativeIndexer::read turned a
read the store never answered into a failure of that one pack and read
the next. So once the store stopped answering GETs, stop() -- and so
engine.close() -- and a flush past its deadline waited out
s3_max_attempts x s3_read_timeout_s plus backoff for every pack of the
pass, one after another: 4 x 120 s + 1.4 s a pack on the defaults, which
NativeCaptureStorageConfig did not let a user lower. The wait gained
nothing: the read failed in the end and the pack was owed either way.
Worse, each unanswered read counted against the pack, so after
max_index_attempts cycles of a store outage the pack was set aside for
good and flush() reported it as never indexable.

- The index reads and the reconcile's requests go through a client with
  a Cancellation of their own (read_cancel_). stop() cancels it for good
  with the uploads'. A flush arms it one catalog request timeout past its
  deadline, where the uploads' is armed at the deadline: a cut upload
  leaves its pack in the durable spool, a cut read leaves an uploaded pack
  owed, which only this process remembers, so what a flush uploaded gets
  the time a catalog statement in flight would. A reconcile request cut
  by stop() ends the pass quietly.
- S3Client::GetRange says whether the store answered for the object
  (a 404, a short body) or not (a transport error, a timeout, a retryable
  status on every attempt, a cancel), and read_pack_descriptor_rows
  throws StoreUnavailableError, a kValue CatalogError, for the latter.
  Under IndexerConfig::end_read_when_store_unavailable, which the service
  sets, NativeIndexer::read lets it propagate, so the pass ends there like
  a batch that threw: the batch and the rest stay owed, and nothing counts
  towards max_index_attempts. A read a cancel cut is not counted as an
  index failure either, and a cycle it cut is cut short, not failed.
  The default keeps the oracle's rule -- CatalogIndexer fails the pack and
  reads the next -- so the conformance driver's index op is unchanged;
  the service is deliberately stricter than the oracle here.
- NativeCaptureStorageConfig exposes s3_read_timeout_s and s3_max_attempts
  (whole seconds and a count; defaults 120 and 4, the native client's),
  which reach the service and the reader.

stop() no longer lets the loop's last cycle finish a read it had begun:
test_the_lease_renews_until_the_loops_last_cycle_is_done (#159) pinned
that a 4 s read past the lease deadline still indexed its pack after
stop(). It keeps its point -- no lease reported held over a dead row
while stop() waits for the loop -- and now asserts that stop() returns
within 2 s, the pack unindexed, and that the next start's reconcile
indexes it from the bucket.

Red first, against the live catalog and the fake S3 behind a TCP switch
that holds every ranged GET, with a 3 s read timeout:
- flush(1.0) over 3 freshly uploaded packs took 40.28 s (3 x 4 attempts);
- stop() with the loop's index read held took 26.79 s;
- with 1 s x 1 attempt and max_index_attempts=2, the one pack was set
  aside after two cycles (rejected_packs 1).
Now the flush returns within 1 s plus the 4 s request timeout with at
most one read's attempts held, the packs owed, none counted or set aside,
and all indexed once the store answers; stop() returns within 3 s and the
next start reconciles the packs it left in the bucket; the outage counts
nothing against the pack and the flush after it drains.
A flush's cycle looked at its deadline only while uploading. Past it, the
cycle still indexed every pack it owed and every pack it had uploaded,
batch after batch, until the work ran out or a request failed. Against a
catalog that answers slowly but inside the request timeout nothing fails,
so the overrun grew with the number of batches instead of staying within
the one request timeout the plan allows, and the docs' "the catalog work
in flight at it" held only for a catalog that had stopped answering.

index_bounded takes the flush's deadline: past it, no batch starts but the
first of the call, and the rest is left owed -- counted as deferred, so
the cycle is cut short rather than failed and the backoff does not move.
The loop's cycles pass no deadline and index everything, and the loop, or
a later flush, indexes what a flush left owed. At most one batch runs
past the deadline per cycle: the one in flight at it, or one of what the
cycle's cancelled uploads left uploaded (a cycle that leaves owed packs
from step 1 uploads nothing). Its reads are cut one request timeout past
the deadline (previous commit) and its catalog statements are never cut,
each bounded by the request timeout. What a flush leaves owed when close()
stops the service right after is in the bucket for the next start's
reconcile, as anything unindexed at stop() is.

Not done: uploading a flush cycle's packs in chunks of indexer.max_packs,
indexing each chunk before the next, which the finding offered to keep
more of what was uploaded indexed. Each UploadPending call lists the spool
again, and listing re-hashes every pending pack, so a backlog of N packs
would be hashed about N / max_packs times over.

Red first, against the live catalog with every INSERT and SELECT held
0.4 s by the TCP switch, eight staged packs, indexer_max_packs=1 and a
60 s request timeout (so the read cut does not end the pass): flush(1.0)
took 59.3 s and indexed all eight. It returns after 7.8 s now, one pack
indexed and seven owed, no index failure counted; the next flush drains
them and every capture reads back.
close() drains with flush(budget) and then stop(). The loop checked for a
stop only before it blocked on the cycle lock, so a loop that woke while
the flush held the lock took it the moment the flush let go -- before
close() could call stop() -- and ran a whole cycle, which stop() then
joined. That cycle uploads nothing once stop() cancels, but its catalog
requests (the owed index's replay guard, a lease claim) are never cut:
against a catalog that accepts connections and never answers, the drain
overran its budget by that request timeout on top of the flush's and the
lease release's.

- The loop waits its interval from the end of the last cycle, anyone's
  (last_cycle_end_ns_, which run_cycle stamps). Having taken the lock
  after a flush's cycle, it waits out the rest of the interval instead
  of running another, unless a fresh lease kicked it; a stop() in that
  wait ends it at once. It gives way so once per wake, so flushes that
  keep coming cannot starve the periodic reconcile, which only the loop
  runs. It also re-checks for a stop once it has the lock.
- A cycle that finds stop() under way (the uploads cancelled for good)
  returns at once, cut short: everything it could do on the object store
  is cancelled, and its catalog requests would only hold stop() up.

The lease release stays: it is one request, bounded by the request
timeout (under the lease, by the lease deadline), and it saves the next
process a TTL's wait. The docs say so in the commit after this one.

Red first, against the live catalog behind the TCP switch, stalled, with
a 4 s request timeout and a 3 s poll interval timed so the loop wakes
during the flush's cycle: flush(1.0) returned after 4.0 s, and the stop()
after it took 8.01 s, the snapshot showing one cycle more than when the
flush returned. The stop() now takes 4.00 s, the lease release alone,
and no cycle runs after the flush's.
close_flush_timeout_s said the drain outlasts its budget only by "the
catalog work in flight when it ends -- against a catalog that stopped
answering, one clickhouse_request_timeout_s -- never by an upload", and
the flush_and_wait comment said the same. That left out the index reads,
which nothing bounded but the S3 timeouts, the index batches that ran
after the deadline against a slow catalog, the loop cycle a stop() after
the flush could wait for, the lease release that stop() makes, and the
uncancellable 5 s abort of a multipart upload the deadline cut. The three
commits before this one bound the first three; the rest is now said.

close_flush_timeout_s: past the budget the drain starts no upload (one in
flight is cut, a multipart one then aborted within 5 s) and at most one
index batch, whose reads are cut one request timeout later and whose
catalog statements are each bounded by the request timeout (or the lease
deadline); stopping the service then releases the lease, one request more.
So close() outlasts the budget by up to about two request timeouts against
a catalog or store that stopped answering, and by one batch of statements
and the release against a slow catalog that still answers; a stuck sink
adds its release flush, 30 s at most (which the first commit of this
series made true beside a concurrent flush_and_wait). The flush_and_wait
comment, NativeCaptureStorage.flush and the v1 integration doc's close()
section say the same in brief.
Past its deadline a flush's cycle starts no upload, but it still called
UploadPending, which lists the spool first, and a listing re-hashes every
staged pack before the cancel turns each one away. flush(0) is what
flush_and_wait and close() pass once the sink's flush has spent the
budget, so after an outage a close() held its caller for as long as
hashing the whole backlog took -- about 0.8 s a GiB here, and the spool
is allowed a TiB by default -- all to report "not drained".

Spool::HasReady answers whether any ready pack is on disk from the file
names alone: a directory walk, nothing hashed, quarantined or
re-accounted. A cycle whose uploads are already past the deadline asks it
instead of listing, and is cut short when a pack is staged; an empty
spool still reports drained, so flush(0) on a drained service returns
True as before (test_a_flush_returns_on_time_while_an_upload_stalls).
A listing that begins before the deadline still hashes to its end.

Red first: one sparse 1 GiB ready pack in the spool, the loop asleep,
flush(0) took 0.79 s. It returns in under a millisecond now, the pack
untouched, and reports it staged.
UploadOne marked a pack cancelled whenever its Cancellation was set once
the retry loop was over -- also when the loop had ended because the
attempts ran out on real failures (a 403, a 5xx, a checksum mismatch), or
at the corrupt-staged-bytes break, and a flush's deadline merely passed
in the meantime. The storage service then counted the pack in
cancelled_uploads rather than upload_failures, never recorded its error,
and left the failure streak alone, so the TimeoutError a flush raised
("last error: ...") named no cause when the flush's cycle was the only
one to see it.

The S3 client's HeadObject, GetRange and PutObject now say, on failure,
whether the Cancellation cut that call short (before an attempt, in its
transfer, in a retry's backoff; a multipart upload's create, parts or
complete). UploadOne books the pack cancelled when the loop left through
one of its own cancel checks or when the last attempt's request was cut;
a failure of the store's own stays a failure. The store driver takes
cancel_uploader_only, which gives the Cancellation to the uploader alone,
so a request the cancel comes in during runs to its own answer -- as one
answered in the gap before libcurl next asks the Cancellation would.

Red first: one upload attempt whose HEAD the fake S3 answers 500 after
1.2 s, the cancel at 0.3 s, the uploader alone holding it: booked
"upload cancelled; the pack stays staged (the attempt before failed:
HeadObject returned HTTP 500)", cancelled true. It is a failure now,
naming the 500; with the client holding the Cancellation too, the cancel
cuts the HEAD and the pack is still booked cancelled.
test_the_lease_renews_until_the_loops_last_cycle_is_done (#159) held the
loop's last index read past the lease deadline, so that a lease thread
stopping before the loop would let the row expire under a lease the
snapshot still called held. Once stop() cut index reads (b13c7ec) that
read ended at once, and the test was changed to expect the pack left for
the next start's reconcile -- but it no longer guarded its property:
with stop() stopping the lease thread before joining the loop, it still
passed.

What stop() leaves running, outside the lease lock, is the abort of a
multipart upload it cut: one uncancelled attempt of up to 5 s. The test
now uploads a 65 MiB pack (a sparse file of zeros, named for its
checksum, over the client's multipart threshold), holds its first part
and the abort, and stops the service with a 3 s lease: the loop's last
cycle then outlives the lease, and the lease thread must renew through
it. It asserts no "held" over a dead row, the lease released, the part
and the abort both held, and the pack cancelled and still staged.

Red first, with stop_lease_thread() moved before the loop's join: the
samples said "held" over an expired row from 2.65 s on. Green as the
code stands, three runs.
_Switch.close() cut the live connections and closed the listener, but
left a stall() in force, and closing a listening socket does not wake a
thread already blocked in accept(): that thread took the next queued
connection and, still stalled, held it open unanswered. So a test that
stalled the catalog and then ran close() before the service's stop()
spent a whole clickhouse_request_timeout_s on stop()'s lease release --
about 4 s of the black-hole tests' 8.4 s -- and its lease release failed
on a timeout instead of at once.

close() now clears every stall and delay before it cuts, so the late
connection is refused. The black-hole flush test takes 4.4 s instead of
8.4 s, and test_flush_returns_on_time_when_the_catalog_stops_answering,
which uses the same pattern, 3.7 s instead of 8.4 s.
No test noticed the S3 client's retry backoff or the uploader's backoff
going back to a plain sleep, main's behaviour. The S3 test cancelled at
0.5 s, inside the 0.2-0.6 s backoff: a plain sleep ended at 0.6 s, where
the check before the next attempt stopped the request anyway, well
inside the test's 3 s. The uploader's test cancelled inside the S3
client's backoff, not its own.

- The S3 test cancels at 1.5 s, early in the 1.6 s backoff after the
  fourth attempt (1.4-3.0 s), and asserts four attempts and a return
  before 2.2 s; slept out, the backoff ends at 3.0 s.
- A new uploader test gives each attempt one transport attempt against
  a 500 and the uploader a 2 s backoff (the store driver's new
  upload_base_backoff_ms), and cancels at 0.3 s, inside the first: one
  attempt, cancelled, back before 1.2 s; slept out, 1.6 s at the least.

Red first, each backoff back to std::this_thread::sleep_for: the S3 test
returned after 3.01 s, the uploader test after 2.23 s. Green as the code
stands.
Only the backstop's good path had a test. test_a_failed_sink_is_released
_without_raising never reached the failure branch: a record dropped as
oversized is a loss the sink counts, not a latched pipeline error, so
its release flush succeeded and wrote nothing to stderr -- its docstring
described a path it did not run. release_flush_timeout_s appeared in no
test at all. With on_engine_release flushing unbounded (Flush(-1)) every
sink test still passed, while a stuck sink would then hold
RingEngine::stop, and engine.close(), for as long as it stayed stuck.

- A stuck pipeline takes a stager that does not move. The binding gets a
  test seam for it, _hold_stages_for_testing(max_hold_s), over the
  spool's existing stage hook: every stage waits, before it writes, until
  the returned function is called or max_hold_s passes -- so a test whose
  code never lets go still gets its sink back.
- Bounded: a 0.5 s release timeout over a held stager returns within it,
  says "timed out after 500 ms" on stderr, and the pack reaches the spool
  once the stager moves.
- Fails soft: a spool too small for the open pack fails the release
  flush; the release still returns at once, says why on stderr, counts
  the failure, and rethrow_if_failed reports it once the sink is attached
  again.
- 0 turns the flush off (the open pack stays the pipeline's, for the next
  flush); nan, inf and -1 are refused.
- The oversized-record case keeps its test, named for what it shows: a
  counted loss does not hold the release.

Red first: with the release flush unbounded, the bounded test's release
took 10.03 s (the seam's own bound); with the stderr line dropped, both
new failure-path tests failed on it. Green as the code stands, three
runs.
Nothing tested that a flush's cycles skip the periodic reconcile
(run_cycle's allow_reconcile). The black-hole flush test, the only one
that flushes with a reconcile due, failed only when that was undone
together with the cancel checks that also stop a reconcile past the
deadline, so a refactor that dropped allow_reconcile alone would let a
flush with time left list the bucket and query the catalog page by page,
and nothing would notice.

A reconcile is due on every cycle here and the loop's first wake is 2 s
off: the flush uploads and indexes the staged pack with no reconcile
pass, and the loop's first cycle then runs one.

Red first, flush passing allow_reconcile=true: reconcile_passes was 1
when the flush returned. Green as the code stands.
Three of the mechanisms the cancellable-upload commits added had no
test, and each could regress with every test green:

- start() resets both Cancellations stop() cancelled for good. Without
  it a service object stopped and started again cancels every upload
  from then on: its flushes return False and the packs pile up in the
  spool. The new test stops, starts, stages and flushes the same object.
- stop() cuts the loop's periodic reconcile. Its listing goes through the
  client the index reads use, so a listing the store never answers held
  stop() for s3_read_timeout_s on every attempt. The new test holds the
  listing and asserts stop() returns within 3 s, the lease released and
  the pass left for later.
- A cancel that cuts CompleteMultipartUpload, every part sent, still
  aborts the upload (fake S3: new fault/hang-complete), since the
  complete may not have taken effect.

Red first: without upload_cancel_.Reset() in start(), the restarted
service's flush(10) returned False with nothing uploaded; with the read
client given no Cancellation, stop() was still blocked after 60 s;
without the abort after a cut complete, no abort reached the store.
Green as the code stands.

Not pinned: a flush returning False once stop() has begun (otherwise it
raises "not started" some 50 ms later -- a test needs stop() held between
the loop's exit and its cycle lock), and the failure streak a cut-short
cycle leaves alone (seen only in the loop's backoff over an outage).
MonitoringConfig's capture_storage_config comment still said close()
drains "within" close_flush_timeout_s, which the storage config's own
comment now bounds as up to about two request timeouts past it. It says
"for" the budget now and points at that comment.

That comment counted stopping the service as one lease release. stop()
first waits for the lease thread, and a renewal in flight then is what
it waits for -- against a catalog that stopped answering, in place of
the release, since a renewal that fails loses the lease and a lost lease
needs no release. It says so, which leaves the two-request bound as it
was.
15db912 kept a flush that is already out of time from listing the spool.
A listing that had begun still ran to its end, though: it hashes every
staged pack before the uploader's workers look at the cancel, so stop()
-- and so engine.close() -- waited for the loop's cycle to hash the whole
backlog, about 0.8 s a GiB here over a spool allowed a TiB, only for the
cancel to turn every pack away after; a flush whose deadline passed
mid-listing likewise returned once the listing was done.

Spool::ListPending takes a Cancellation and stops between packs once it
is cancelled: it reports the listing cut, returns no packs, and leaves the
spool's account alone, since a partial listing must not recount it
(quarantines already made stand). UploadPending passes the uploader's
Cancellation and reports listing_cancelled, trying nothing; the service's
cycle is then cut short, neither drained nor failed. The drained check's
own listing at the end of a cycle is cut the same way. The store driver
reports listing_cancelled for upload_pending.

The Python spool and uploader, the reference, have no cancellation; this
is C++ only.

Red first, 16 sparse 256 MiB packs (zeros named for their checksum):
- the store driver's upload_pending, cancelled at 0.3 s, listed them all
  and returned each "cancelled before it started", after 4.7 s;
- stop() 0.5 s into the loop's listing took 2.58 s;
- flush(0.5) over the same backlog returned after 3.04 s.
Now the listing stops within one pack's hash: the driver reports
listing_cancelled with nothing tried, stop() returns within 1 s and
flush(0.5) within 1.2 s, the packs untouched and none counted cancelled.
Past a flush's deadline index_bounded starts no batch but the first. It
cut its packs into batches of indexer.max_packs from the front and then
took them from the back of the stack, so that first batch was the
remainder chunk, refs.size() % max_packs packs whenever that is not zero.
131 owed packs in batches of 64 indexed 3 and left 128 owed; the batch
meant to guarantee a flush's progress could be a single pack.

The chunks now go on the stack in reverse, so a pass runs them front to
back and starts with a full one. A kBatchTooLarge split still pushes its
halves so the first half runs next. What a cancel, the deadline or a
failure leaves queued is appended to the owed list in the order it would
have run, so a later pass starts from the front again.

Red first: three packs in batches of two, a catalog answering every
statement 0.4 s late and flush(1.0). The one batch past the deadline
indexed 1 pack; it indexes 2 now, and the third follows once the catalog
is fast again.
15db912 kept a flush already out of time from listing the spool through
the uploader, since a listing re-hashed every staged pack: it asked
Spool::HasReady, a walk over the file names, whether anything was staged.
5fecd88 then made the listing stop between packs once cancelled, checked
before each pack's hash, so a listing under a Cancellation already past
its deadline hashes nothing either: it reports itself cut at the first
staged pack, and finds an empty spool empty. The branch gave the same
outcome for the reason its comment no longer stated, and could be deleted
with every test green.

It goes, with Spool::HasReady, which nothing else called: the cycle lists
through the uploader whatever the time, and the step-2 comment says that
the cancelled listing is what keeps flush(0) from hashing. The live test
of flush(0.0) over a sparse 1 GiB pack still returns within 0.25 s with
nothing tried, and a drained spool still reports drained at zero.
…most

A cycle uploaded everything the spool held, then indexed what it had
uploaded. Each upload deletes its pack from the spool once verified, so
until the catalog has it the pack is remembered by this process alone,
in pending_index_, which stop() drops. close() runs flush(budget) and
then stop(): a flush that ran out of budget while uploading had turned
every pack it uploaded into such a debt, indexed one batch of them past
the deadline (0121538), and stop() then discarded the rest. They were in
the bucket only, out of the spool and out of the catalog, for a later
start's reconcile to find -- and with reconcile_on_start off, the
setting for a shared bucket, never. A stop() during the loop's upload of
a backlog did the same, since stop() also cuts the index reads. On main
the same close() indexed everything, only late.

A cycle now lists the spool once and uploads it in chunks of
indexer.max_packs, indexing each chunk before it uploads the next. A
cancel -- a flush's deadline or stop() -- starts no further chunk, so
what was not uploaded by then is still in the spool, where any later
start uploads it whatever its settings; at most the one chunk in flight
is out of the spool and unindexed, and a flush indexes that as its one
batch past the deadline. A chunk left owed stops the uploads, as an owed
pack at the start of a cycle does, and so does a lease lost between
chunks (checked without a request), where a lost lease used to let the
whole batch upload. Listing once keeps chunking linear: repeated
UploadPending(limit) calls would re-list and re-hash the spool per
chunk. SpoolUploader::UploadStaged uploads packs a listing returned;
UploadPending lists and calls it, unchanged for the drivers and the
Python parity suites. Knowing its own listing's outcome, the cycle no
longer lists the spool a second time to decide it is drained.

The Python reference has no storage service; this is C++ only.

Red first:
- twelve packs uploaded one at a time at about 0.3 s each, in chunks of
  four, flush(2.0) then stop(): 6 uploaded, 4 indexed and 2 owed, which
  stop() dropped, so a successor with reconcile_on_start=False left
  those captures unqueryable. Now every uploaded pack is indexed, the
  rest are staged, and the successor gets all 24 captures into the
  catalog.
- eight packs in one-pack batches against a catalog 0.4 s a statement
  (the CC-2 test, now run in close() order): 8 uploaded, 1 indexed, 7
  owed; now 1 uploaded and indexed, 7 staged, nothing owed, and the
  successor, reconcile off, drains all 16 captures.
The full-batch test added with the stack-order fix still passes; with
chunks of indexer.max_packs a pass past a deadline is now handed one
chunk, so that order matters only for owed lists, which no longer
outgrow a chunk.
stop() cancels the index reads for good before it joins the loop, and
run_cycle returns at once when a cycle starts after that. A cycle
already in flight went on, though: the uploads stop() cut returned the
packs uploaded before it, and step 3 began their index pass. Its first
batch took the lease lock and sent the replay guard, a catalog SELECT,
then read, which the client refused at once, and deferred the batch --
throwing the SELECT's answer away. stop() waited for that statement, up
to a request timeout against a slow catalog, before releasing the lease:
outside what its comment says it waits for, and the very request
run_cycle's early return exists to avoid. The flush's read deadline had
the same gap, since only non-first batches checked a deadline at all.

index_bounded now starts no batch, the first included, once read_cancel_
is cancelled -- stop(), or a flush one request timeout past its deadline
-- and leaves everything owed, as the refused read would have. The
reconcile, which checked the cancel before each listing page and each
HEAD, checks it again once a page is listed, before the catalog query.

Red first: four packs staged, one upload worker, the third PUT held, the
replay guard answered 3 s late, then stop(): one guard SELECT went out
after stop() began and stop() took 3.19 s. Now none does, stop() returns
within 2 s with the lease released, the two uploaded packs stay owed for
the successor's reconcile, and every capture reads back. Only the index
half has a test; the reconcile's window, a listing answered between the
cancel and libcurl's next progress poll, cannot be hit on demand.
read_pack_descriptor_rows decided whether a read the store did not
answer was cancelled by asking the S3 client, once the read had failed,
whether its Cancellation was set: the proxy b61db70 removed from the
uploader. A flush's read deadline makes that true from the moment it
passes, however the read ended, so a real store failure -- a 503 on the
last attempt, a connection reset -- answered after the deadline and
before libcurl's next progress poll was booked as a cancel: index_bounded
deferred the batch with nothing recorded, the cycle counted as cut short
rather than failed, and the TimeoutError flush() raised named an older
error. GetRange has said whether the Cancellation cut the call since
b61db70; the trailer and footer reads now pass that out-parameter and
throw StoreUnavailableError with it.

The gap is narrower than libcurl's poll interval, so a test cannot land
a failure in it on demand. S3Client gets a test seam instead,
SetAfterExchangeHookForTesting, run once an exchange has returned and
before the call reads its response; conformance_catalog gets a
session-less read_pack_rows op that reads one pack's rows through a
client holding a Cancellation, cancelled from that seam
(cancel_after_exchange) or armed cancel_after_ms ahead.

Red first: a trailer read answered 500 at once, the cancel coming in as
the answer returned, was reported cancelled; it is the store's failure
now ("HTTP 500"). A read held for 5 s and cut by a deadline 0.2 s in is
still a cancel. The Python reader has no cancellation; C++ only.
b61db70 passes "the Cancellation cut this request" up through the
preflight HEAD, the verifying GET, PutObject -- PutSingle, and
PutMultipart's create, parts and complete -- and the post-upload HEAD,
because on an upload's last attempt nothing after the request sees the
cancel: no further attempt, no backoff. Only the preflight HEAD's report
had a test. Dropping PutObject's out-parameter (&attempt_cut -> nullptr)
passed every uploader, S3 and live test, since those upload with the
default four attempts and the next attempt's check books the cut anyway;
the pack would then count as an upload failure on the last attempt, with
"request cancelled" as the last error.

The new test uploads with one attempt of one transport attempt, its HEAD
answered 404 at once, and holds the PUT -- or, for a pack over a 5 MiB
multipart threshold, its first part -- for 5 s; the cancel at 0.3 s cuts
it, and the upload must come back cancelled, the pack staged and nothing
in the bucket. The fake S3 gains fault/hang-put, which holds a
single-request PUT alone and answers the HEADs.

Mutations checked: PutObject passed nullptr fails both cases; PutMultipart
not reporting a cut part fails the multipart one.
test_a_cancel_stops_the_listing_between_packs asserted the cut listing
returns within 1.0 s. The cut lands 0.3 s in, plus the rest of the
256 MiB hash under way: about 0.5 s on a quiet machine, and hashing slows
with the machine, so on a contended runner (4x oversubscription of two
cores) it failed at 1.00 to 1.11 s, every run.

The bound is what checks the property the test is named for: a listing
that honoured the cancel only once it had hashed all sixteen packs still
reports listing_cancelled with nothing tried, and is told apart only by
taking over 3 s. 2.0 s keeps that apart and leaves room for a loaded
runner; the comment says so.
test_stop_cuts_a_reconcile_whose_listing_stalls says the pass "ends
quietly, to run again on the next start", and asserted only that stop()
returned promptly, the lease was released and no pass completed.
Deleting the check that ends a cut listing quietly left it green: the
listing then threw "reconcile: listing failed: request cancelled", which
run_cycle recorded as the cycle's error, so every stop() during a
reconcile left a spurious last_error. The test now asserts last_error
does not mention the reconcile; with the check deleted it fails on that
message.
…timeout

flush()'s comment, engine.flush_and_wait's and NativeCaptureStorage.flush
gave a flush's overrun as about one request timeout against a catalog or
object store that stopped answering. A multipart upload the deadline cuts
is then aborted, one request of up to kAbortAfterCancelTimeoutS (5 s)
that no cancel cuts, after up to about a second for the stalled transfer
to see the cancel (libcurl's progress poll). The config allows
clickhouse_request_timeout_s down to 2 s (twice publish_timeout_s), and
the default 128 MiB packs are multipart, so the bound was not kept: with
a 2 s request timeout, flush(1.0) returned after 6.4 s where the comments
promised about 3.

The abort stays: an incomplete upload left behind costs storage until a
lifecycle rule reaps it, and 5 s is well under the 60 s default request
timeout. The four comments and the integration guide now name it -- up
to about 6 s past the deadline -- and the storage_service.h paragraph,
which an earlier edit left with one overlong line, is reflowed.

A live test pins it: a 65 MiB sparse pack, its part and its abort both
held, a 2 s request timeout and flush(1.0). The flush returns after
6.4 s, inside the documented 1 + 1 + 5 s, the upload counted cancelled
and the pack still staged.
58453c8 dropped "within" from config.py's comment as an overstatement:
close()'s drain can outlast close_flush_timeout_s by about two request
timeouts, the multipart abort and the sink's 30 s backstop. The
engine.close() docstring and the integration guide, both new on this
branch, still said the drain runs "within" the budget -- the docstring
with stopping the service inside it, though _retire_capture_storage
calls stop() in a finally past the flush, bounded by nothing in the
budget. An integrator sizing a shutdown grace period from either would
kill close() mid-stop and leave the lease row without its tombstone.

Both say "with a budget of" now; the docstring says the drain can
outlast it and points at NativeCaptureStorageConfig.close_flush_timeout_s
for by how much. The guide's paragraph is reflowed.
cancel.h's header still described the state before b13c7ec: the storage
service owning one Cancellation, handed to the client and the uploader
that do its uploads, and a cancelled upload staying in the spool as the
only outcome. b13c7ec added a second one, read_cancel_, which cuts the
index reads and the reconcile through the other S3 client and which a
flush arms one catalog request timeout past its deadline; a read it cuts
leaves an uploaded pack owed, not in the spool. The header names both,
what arms each, and where each leaves the pack it cuts short.
The s3_read_timeout_s / s3_max_attempts comment said an attempt is
retried for "a transport error, a timeout, a 429 or a 5xx, and a backoff
of 0.2 s doubling between them". The native client retries 429, 500,
502, 503 and 504 only -- a 501 or 505 is not -- and nine named curl
failures (timeouts, resolve and connect failures, connections that broke
or answered nothing), not a TLS failure; and its backoff doubles only up
to 5 s. The cap matters now that s3_max_attempts goes up to 1000: read
uncapped, 20 attempts look like 29 hours of backoff where they are 76 s.
The comment says what IsRetryableStatus, IsRetryableCurl and Backoff do.
S3Client's Backoff waits min(5000, 200 << attempt) ms. This branch lets
NativeCaptureStorageConfig.s3_max_attempts go up to 1000, and a retry
backs off with attempt up to max_attempts - 2: from attempt 56 the shift
overflows int64 -- a negative wait, which SleepFor and sleep_for both
end at once, so those retries went back to back against a failing store
-- and from 64 the shift is undefined behaviour. The shift is capped at
5 (6.4 s, past the 5 s cap), as the uploader's own backoff caps its
exponent; every wait below the cap is unchanged.

No test: reaching attempt 56 takes over four minutes of capped backoff
first. The S3 client suite passes unchanged.
storage_service.h's opening paragraph and its lease-check sentence, and
S3Client::Exchange's call, ran past the file's line width. No change
beyond the wrapping.
Copilot AI balanced review requested due to automatic review settings September 29, 2026 17:29

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@zaoxing
zaoxing merged commit 0e1a112 into main Sep 29, 2026
3 checks passed
@zaoxing
zaoxing deleted the feat/close-flush-backstop branch September 29, 2026 17:42
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants