Sink refusals and stalls never kill the forward: failure policy, per-step stall budget, bounded admission - #152
Merged
Merged
Conversation
Decision 3 wants vLLM serving to keep running when capture storage refuses or fails: capture stops and says why. Until now the record consumer had one behaviour, raise_at_producer: the first sink refusal (a dropped, timed-out or oversized record, or a pipeline failure) latched, and the next descriptor push raised inside the model forward. RecordConsumer now takes a RecordFailurePolicy. raise_at_producer stays the default and is unchanged. Under disable_capture the latch drops the descriptors still queued, and every push and payload after it is discarded and counted instead of raised; the p2p worker keeps draining, so the ring keeps its protocol. The failure still surfaces at wait_until_idle, finish and rethrow_if_failed, and in a new snapshot() (policy, failure text, discarded descriptors and payloads). The policy is threaded through RingEngine and RecordP2PThread; nothing selects it yet. Test-first in tests/native/ring/test_record_consumer.cpp: with the snapshot API in place but no policy behaviour, the three new cases gave 12 FAILs and then an uncaught throw from a push after the latch; with the policy, 48 passed, 0 failed.
When a record reservation does not fit, reserve_record synchronised the
stream and called the unbounded force_flush_and_wait. The drain is held
back by the record worker and the worker by the sink, so a slow or stuck
sink stalled the model forward for as long as the sink took.
RingEnginePy takes RecordRuntimeOptions {failure_policy,
step_stall_budget_ms}. Both reserve slow paths now wait through
force_flush_and_wait_until, bounded by what is left of the step's budget
(begin_record_step starts a step). Past it the policy applies:
raise_at_producer raises from the reservation, with nothing reserved;
disable_capture latches the failure in the consumer, which then
discards instead of submitting, so the remaining wait is at most the one
sink admission already in progress, and the reservation completes.
Under raise_at_producer a latched failure now raises at the start of
reserve_record, before a reservation no producer will use, rather than
from the push after it. A budget of 0 keeps the old unbounded wait and is
the default. record_capture_status() reports the policy, the failure,
the discard counters, budget exhaustions and reserve wait time (total and
worst step).
Test-first in tests/native/ring/test_ring_engine.cu, on GPU 1: a 4 KiB
ring filled twice while a sink sits in a 400 ms admission, 50 ms budget.
Red was a compile failure (no options/status API); a mutation that
ignores the budget fails 13 checks (the raise case waits the full
admission queue and never raises; the disable case stalls past
budget + one admission). Green: 138 passed, 0 failed.
…code
The ring-fed native sink was built with the C++ SinkConfig admission
default, drop_newest, which the Python binding did not let anyone change.
Measured on this branch: one envelope of 64 x 1 MiB rows against the
default 16 MiB queue dropped a record after 18-25 were admitted (three
trials), and the refusal latches the record ring.
bindings_sink.cpp now takes `overload` ("block" | "drop_newest") and
`admission_timeout_s` (None waits without bound), exposes both as
properties, and adds timed_out_records and rejected_closed_records to the
snapshot. The binding's defaults and the C++ SinkConfig stay drop_newest
with no timeout, so PackSink parity with the reference is untouched.
NativeSinkConfig, the ring-fed entry point, moves to
dmi.storage.native_capture (live code) with a re-export at the old path,
gains the two fields with the plan's default of block with 2 s, and
validates them. Its docstring no longer claims a field-by-field match with
the reference pipeline config, which it never had. MonitoringConfig and
the engine now import it from the live module, so the public config does
not load the backup capture package (checked in a fresh interpreter).
Test-first: tests/test_native_sink_admission.py was 19 failed (missing
type, attributes and constructor arguments), then 19 passed; the burst
test persists all 64 with no drop or timeout. The neighbouring cpu suites
(rollback, storage wiring, backend choices, engine runtime API, adapter,
capture storage, v1 API, pack sink, pack sink timeout): 333 passed,
1 skipped (no CUDA device).
Admission screens a record's payload against max_pack_bytes, but a pack also holds a header, the record's footer row and a trailer. A payload of exactly max_pack_bytes clears admission and fits no empty pack. The reference assembler raises OversizedRecordError for it and the pipeline counts it oversized and keeps going. The native worker opened a fresh builder, got kCapacity, and sealed that empty builder, so "seal failed: cannot seal an empty pack" failed the whole sink and every record queued behind it; on a ring that is a latched record runtime. Found while checking the framing reserve for validate_capture_bounds: through the torch binding, one 1 MiB float32 record into a 1 MiB pack failed the sink. Now a kCapacity on a builder with no records drops the record and counts it, as the reference does. Test-first in tests/test_native_pack_sink.py (conformance_sink driver): the new case failed with the next record refused as 'closed'; after the fix the sink suites (pack sink, pack sink timeout, sink admission, capture bounds, rollback) are 164 passed.
…udget Three configurations lose captures only once a forward runs: a record larger than the sink's max_queue_bytes is refused on the record worker (latching the record runtime), one that fits the queue but no empty pack is admitted and then dropped as oversized, and a pack larger than the uploader's in-flight budget is staged and never uploaded. All three are decided by configuration and the largest record an adapter emits. validate_capture_bounds(sink_config, max_record_bytes, storage_config=) in dmi.storage.native_capture refuses each with a ConfigurationError that names the bound to raise; max_pack_bytes must leave PACK_FRAMING_RESERVE_BYTES (64 KiB) of framing room above the record. MonitoringEngine.validate_capture_bounds(max_record_bytes) applies it to the engine's own configs, for adapters to call at attach (no adapter drives the persistent path yet; D2 wires it). With no sink config there is nothing to check. NativeCaptureStorageConfig gains uploader_max_in_flight_bytes (default 256 MiB, the native UploaderConfig default). bindings_store.cpp already read the key, but _native_dict never sent it. Test-first: tests/test_capture_bounds.py was 16 failed (missing field, function and method), then 16 passed; with the storage wiring and sink admission suites, 72 passed.
The native ring can now stop capture instead of raising and bound the forward's wait for a slow sink, but nothing above it could ask for either, and nothing reported a stopped capture short of flush raising. create_record_runtime(failure_policy="raise" | "disable_capture", step_stall_budget_ms=None | int) passes both to RingEngine.create_record. The default stays "raise" with no budget, so existing callers are unchanged; C2's YAML will choose (decision 3: disable_capture for vLLM serving). Bad values are refused before the live ring is touched, and disable_capture is refused with a NativeSinkConfig whose block admission has no timeout, since the post-latch wait is one sink admission. RecordRuntime.begin_step() starts a stall-budget step; integrations call it once per model step (it is explicit because piecewise CUDA graphs replay several plans per step). MonitoringEngine.capture_status() returns the policy, whether capture is still active and why not, discard counters, budget exhaustions, the reserve wait (total and worst step) and the native sink and storage snapshots. close(), and enable_ring_transport replacing a record ring, log a stopped capture at WARNING. The integration API document describes the policies, the stall bound (budget plus one admission timeout), the status fields and validate_capture_bounds. Test-first: tests/test_record_failure_policy.py was 14 failed, then 14 passed; the v1 doc-coverage test, extended with begin_step, capture_status and validate_capture_bounds, failed until they were documented. Seven create_record fakes in the engine and storage wiring suites now accept the options as keywords. Engine suites (runtime API, storage wiring, record runtime, v1 API, backend choices, HF refusal, policy): 162 passed.
B2's live acceptance: each case is either stored or refused at configuration time, never lost in the forward. Through the real NativePackSink, the storage service, a live ClickHouse catalog and the fake S3: - a 17 MiB row is refused by validate_capture_bounds under the default 16 MiB queue, naming max_queue_bytes; with the queue raised it is stored and reads back exactly; - one envelope of 64 x 1 MiB rows, four times the default queue, is stored whole under NativeSinkConfig's defaults (block, 2 s): no drop, no timeout, every row read back. Written after the implementation; the burst's red was measured on the sink directly (the old drop_newest default refused it after 18-25 rows, three trials). tests/test_native_capture_chain_live.py against local ClickHouse 25.12: 4 passed, unique table prefixes, none left behind.
The native suites pin the mechanism; this drives it as an integration does (create_record_runtime, a HookPointV1 on a CUDA tensor, begin_step per step) and checks what serving cares about: whether a hook call raises, and the slowest step's wall time. The stalling sink is a Python target behind the native reference bridge whose admission sleeps 400 ms or refuses; the burst uses the real NativePackSink. - disable_capture, 50 ms budget, 4 KiB ring: no step raises, the worst step is within budget + one admission (measured 0.401 s; 1.606 s with no budget), one record reached the sink, the other 15 were discarded, flush raises the stall-budget failure and close logs it. - raise, 50 ms budget: the stalled step raises at the budget (measured 0.050 s), every later step raises at once, flush raises. - a refusal on the 4th record: disable_capture never raises in a hook and reports "sink refused durable admission"; raise surfaces it from a later step and every step after. - 64 x 1 MiB records through NativePackSink's default config: all 64 persisted, capture active, no drop or timeout. On GPU 1, gated on it being idle: 5 passed, and 5 passed on each of five repeats. The first draft read counters before the drain had delivered anything (3 failures, test-side); the ring now drains each record promptly and the counters are read after the flush.
Three sink-side fixes the record ring's failure policy depends on. One deadline per envelope. The ring's record worker hands NativePackSink one envelope of N rows per submit, and each row's PackSink::Submit started its own admission_timeout_s, so a steadily slow sink held an N-row envelope (and through the drain, the forward) for up to N timeouts. EnvelopeAdmission (record_row.h) takes one deadline when the envelope arrives and admits every row against it, through the new PackSink::SubmitBy; Submit keeps its per-call deadline for its other callers and for parity with the reference. Red first: in tests/native/test_pack_sink_timeout.cpp a stager that takes 0.3 s per pack against a 0.5 s timeout let all 10 rows in over 1.99 s; now the envelope times out at ~0.5 s. Records lost after admission are failures. A record that fits the queue but no empty pack is admitted and then dropped as oversized on the pack worker; flush_and_wait returned True, rethrow_if_failed was silent, and the snapshot said admitted=1 persisted=0. The reference adapter counts oversized (and dropped, timed-out, duplicate, rejected-closed, failures) as losses and raises at flush and rethrow. NativePackSink now keeps its counters at construction and does the same: rethrow_if_failed and a completed flush raise naming the counters that moved, a flush also raises if fewer records were persisted than admitted, and submit checks first, so the ring's record runtime latches on the next envelope instead of storing around the hole. RecordSink::admission_bound(), read by the ring for the stall budget (next commit): zero under drop_newest, the timeout under block, none when block waits forever. The binding exposes it as admission_bound_s on every RecordSink. ClickHouseRecordSink has none, and the reference bridge has none unless its constructor is told what its Python target promises. tests/test_native_sink_admission.py: the bound per policy, the 1 MiB repro (flush, rethrow and the next envelope all raise naming oversized_records), and a stored record still flushing clean. Red: 4 failed; green with the sink and pack-sink suites, 149 passed.
Four review findings on the record ring's stall budget. The budget was lifetime. Nothing in src/ calls begin_step, so the budget spanned the runtime's life: once spent, every later wait gave up at once and max_step_wait reported a lifetime total. The record ring has no reliable step boundary of its own: commit_step is the legacy adapter's and refuses record rings, eager hooks reserve once per output (and output ids repeat within a step), and a CUDA-graph step replays one plan, or several with piecewise graphs. So begin_record_step stays the boundary, and with a budget a reservation before the first one is refused with a logic_error naming RecordRuntime.begin_step(). Budget exhaustion latched capture off for good; the plan skips the rest of that step. RecordConsumer gains a discard window: its start drops the descriptors still queued, pushes until its end are dropped on arrival, and a count of owed payloads discards theirs as the drain delivers them, so descriptor/payload pairing and ring capacity are kept. Exhaustion opens it (counted in skipped_steps and discarded_*), begin_record_step closes it, and capture resumes. Under disable_capture a skipped step is not a failure: capture stays active and the flush succeeds. Under raise the forward no longer raises: the step is skipped the same way, the exhaustion is held and raised at the next begin_record_step (which then latches it) and at flush. Genuine sink failures still latch as before, and begin_record_step now also raises a failure latched during the previous step under raise, so it usually surfaces there rather than in the forward. The post-latch wait was unbounded, and the bound it relied on held only for a persistent-path NativeSinkConfig. The ring now reads the sink's admission_bound() and refuses a budget (either policy) for a sink with none: the ClickHouse host path, the reference bridge, a NativePackSink blocking without a timeout, or no sink. create_record_runtime refuses the same before tearing down the live ring, naming the case. Past the budget the wait is bounded by the sink's bound plus a 2 s drain grace; a sink that holds the ring longer is a ring failure raised from the reservation, rather than a forward that never returns. A sink failure found at the checked flush (NativePackSink's lost-record check) now latches the runtime like a refusal at submit. capture_status() reports skipped_steps, close logs skipped steps, and the integration doc and docstrings state the exact conditions. Test-first. tests/native/ring/test_record_consumer.cpp: a window under both policies (8 failed with a stub, 82 passed). test_ring_engine.cu (GPU, compile-red on the new status field): the skip under both policies, the refused unbounded and sinkless budgets, a reservation before begin_record_step, a sink past its bound, a failing sink flush. tests/test_record_failure_policy.py: 12 failed, then 26 passed. tests/test_record_failure_policy_gpu.py follows the new semantics and adds the begin_step, unbounded-sink and lost-record cases. CPU tier: 2367 passed.
A config pickled from the old class (dmi.storage.capture.native_sink, eight fields) unpickles as the new class through the re-export, but the frozen slots dataclass's generated __setstate__ zips only the fields the pickle has, so .overload raised AttributeError on first use. The new __setstate__ fills the fields a pickle predates with their defaults (block, 2 s) and refuses state with too many values. Test-first in tests/test_native_sink_admission.py with the bytes of an old-class pickle (made from the pre-move module): AttributeError, then passing; a new pickle still round-trips.
The rows of one ring record now share one admission deadline, a record lost after admission latches the runtime like a refusal, and a stall budget is refused for block with no timeout. The docstring said none of it.
A spent budget's discard window drops every descriptor still queued, including records of earlier, finished steps that the sink had not taken yet; the stall bound needs that. The reporting said only "the rest of that step". The raise-mode error, the header comments, capture_status and the v1 contract now say the step's remaining records and every record still queued for the sink, from any step, are discarded. The consumer tags each queued descriptor with its step (begin_step, called from begin_record_step) and counts the distinct steps a window drops records from as steps_with_discards, reported beside skipped_steps (which always equals stall_budget_exhaustions) and logged at close. The drain grace past a spent budget was a fixed 2 s, while the drain still moves the whole ring and staging out before the reservation fits: with the default 4 GiB of each that could exceed it and fail the ring, blaming the sink. ring::record_drain_grace makes it 2 s plus the ring and staging bytes at 1 GB/s (about 10.6 s by default), and the error now says the ring did not drain and names both the sink and the drain as possible causes. The record worker also skips the pageable copy of a payload nobody will store (a skipped step's, or any after a disable_capture latch): discard_next_payload_if_unwanted accounts the discard through the same branch consume_payload uses, and the staging range is freed as before. A test now covers the payloads_to_discard_ term of wait_until_idle: with the queue empty and only owed payloads left, a flush must not go idle. finish() names the owed payloads in its error. The discarded_* comments no longer say kDisableCapture only.
zaoxing
force-pushed
the
feat/sink-failure-policy
branch
from
September 24, 2026 22:59
c74e392 to
b5c5a81
Compare
flush_records_and_wait rethrew the consumer's latched failure before draining, so under disable_capture a record emitted after the latch stayed in the ring: its payload was never delivered to be discarded, discarded_payloads stayed 0 and its ring space was not returned until something else forced a drain. The GPU case 'a sink flush that fails latches the record runtime' failed on exactly that (discarded_payloads == 1), and still failed with a 2 s poll. Under kDisableCapture a flush now waits for the stream prefix and forces a drain (bounded by the flush deadline) before rethrowing, so the discards are counted and the space is back when the caller hears why capture stopped. raise_at_producer is unchanged. test_ring_engine on GPU 1: 174 passed, 0 failed (was 173/1). CPU tier: 2369 passed.
zaoxing
marked this pull request as ready for review
September 24, 2026 23:36
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).
claude Bot
pushed a commit
that referenced
this pull request
Sep 26, 2026
Brings in #149-#152 and #139 (two-phase native search page). One conflict, native/csrc/catalog/reader.cpp search(): #139 split the page into an inner key-only query and an outer argMax query over the same filters, while this branch moved snapshot membership out of `clauses` into the FROM clause (snapshot(), the join that supplies member_version). Resolved by running BOTH phases over snapshot() with the same caller `clauses` (possibly empty), so the inner LIMIT only counts member keys (#139's walk-ends-early guard) and the outer argMax still ranks on (member_version, store_id, pack_id, index_version). The design doc's two-phase paragraph now says both queries read the same snapshot join. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_019oKb8SWwCwPRAtxuWLQSXt
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.
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 B2 of the native capture production plan. When the capture sink refuses records or stalls, the model forward no longer dies or hangs without bound. The failure policy decides what happens instead.
What changes
Failure policy.
create_record_runtime(failure_policy=...)takesraise(the default, which keeps today's behaviour for existing callers) ordisable_capture.disable_capture: a sink failure turns capture off. Later records are discarded and counted, the ring protocol stays intact, and the reason is reported. Inference keeps running.raise: the failure surfaces at the nextRecordRuntime.begin_step()and atflush_and_wait, not inside the forward.Per-step stall budget (
step_stall_budget_ms).skipped_steps), under both policies. Capture resumes at the nextbegin_step(). Underraise, the exhaustion is then raised at the nextbegin_stepand at flush.RecordRuntime.begin_step()once per model step. The runtime has no reliable automatic step boundary. A budget without anybegin_step()is refused at the first reservation, rather than silently turning into a lifetime budget.RecordSinknow reportsadmission_bound_s. A budget is refused, atcreate_record_runtimeand again in the native ring, for these cases:The sink.
NativePackSinkexposesoverload(block | drop_newest) andadmission_timeout_s.NativeSinkConfig, now indmi.storage.native_capturewith a re-export at the old path, defaults to block with 2 s. It still loads pickles made before these fields existed.PackSink::Submitkeeps its per-call deadline, so parity with the reference is unchanged.flush_and_waitandrethrow_if_failed. This covers dropped, timed-out, oversized, duplicate and rejected-closed records, failures, and fewer records persisted than admitted, matching the Python reference adapter. Underdisable_capturesuch a loss turns capture off with the reason.Configuration checks.
validate_capture_bounds(max_record_bytes)refuses, at attach, any configuration that could only lose records once a forward runs: queue too small, pack too small, or pack larger than the uploader's in-flight budget.NativeCaptureStorageConfigexposesuploader_max_in_flight_bytes.Status.
engine.capture_status()reports:skipped_steps(always equal to budget exhaustions) andsteps_with_discards(the distinct steps that lost at least one record to a skip);Nothing Python-only was added under
src/dmi/storage/capture/beyond the re-export.Evidence (CPU and live)
test_record_consumer: 151 passed.test_clickhouse_record_sink: 21 passed.pytest -m cpu: 2368 passed; the 1 skip is the no-CUDA case.tests/test_native_capture_chain_live.py): a 17 MiB row and a 64 × 1 MiB burst through the whole chain, 4 passed.pipeline reported lost records (oversized_records=1).Independent review
The verdict was fix, then ship. All 5 majors are fixed in bc75249, 24efc3e, 4156b37 and c831897:
begin_step. Now it is refused untilbegin_step.raise.Minor #7, loading old pickles, is also fixed.
A second independent review of those fix commits found no majors and 4 minors, all fixed in c74e392:
steps_with_discardscounts the steps affected.That commit also skips the host copy of payloads that are about to be discarded, sharing the discard decision with
consume_payload.The re-review stress-tested the discard window with a producer and a drain thread: 70k records and 1,695 random windows, under both policies. Result: 0 mispairs, and submitted + discarded == total.
test_record_consumernow has 151 checks passing.Plan deviations still in this PR:
commit_step, sobegin_stepstands in for it.begin_steprefuses inside the first forward. It is deterministic, so any integration test catches it.close()logs a pending spent budget underraiseand does not raise.GPU results (GPU 1, RTX 4090)
test_ring_engine: 173 passed and 1 failed on c74e392. The failure, in "a sink flush that fails latches the record runtime", was real. Underdisable_capture, a flush rethrew the latched failure before draining, so a record emitted after the latch stayed in the ring:discarded_payloadsstayed 0 and its space wasn't returned. A 2 s poll didn't help, which shows it wasn't just timing. The fix in 52b9627 drains within the flush deadline, then reports;raiseis unchanged. Now 174 passed, 0 failed.tests/test_record_failure_policy_gpu.py: 8 passed on c74e392 and again on 52b9627.begin_step()every step, the worst step was 0.08 s against a 50 ms budget, and all 40 records were submitted.begin_step(), every reservation is refused with the "call RecordRuntime.begin_step()" error.disable_captureand a 50 ms budget, the worst step was 5.021 s. That is within budget + one admission timeout (5.05 s). One step was skipped, 9 records from 9 steps were discarded, capture stayed active and the flush succeeded.