Review the Python oracle: three capture-path fixes, and tests for what was unpinned - #137
Merged
Merged
Conversation
The footer cache is keyed on pack identity -- store, pack id, checksum -- and a hit returned its descriptors straight back. PackIndex.from_store is the only place the tenant/key bind runs on the read path, so a pack PUT at a second key belonging to another tenant was refused cold and served warm; within one hydrate() call, which of the two happened came down to the order of the selection's capture ids, since _plan groups by object key and walks the groups in that order. The cache key is not a content identity either -- for S3 the checksum is the uploader-DECLARED dmi-sha256 -- so a different object at the foreign key could be served as the victim's capture, which neither _require_footer_match nor verify_payload can catch. Adding the object key to the cache key would close this instance while leaving the entry carrying a LOCATION-derived verdict under a content identity, so the next location check added would regress the same way. The bind is recomputed on the hit instead: reject_a_foreign_tenant now takes the tenants rather than the records, so the reader can run it against cached descriptors. One footer read per pack however many keys it legitimately sits at, which the positive-control test pins.
IDENTITY, PREFIX_STRIP and CHUNKED in `RingTransport._record_cpu_tensor` had no test at all: the whole suite called that function exactly twice, both times through SEGMENTED_PACK or SEQ_PREFIX_PACK, so corrupting all three branches at once left the suite byte-identical. Pin them against in-test transcriptions of the kernels they must match (`record_producer_static_kernel`, `record_producer_prefix_kernel` and `record_producer_chunked_kernel` in native/csrc/ring/producer.cu) rather than against more hand-written literal tensors, so the module docstring's claim about matching CUDA producer bytes is actually what is checked. The cases cover the edges that tell the implementations apart: a row count that is negative, zero, exactly `n // row_bytes` and past it; chunk byte counts that are negative, zero and wider than their chunk; and multi-dimensional, multi-dtype IDENTITY tensors. IDENTITY returns the source dtype and shape, not a flat uint8 run, and the test asserts that. No production code changes: the transforms are already byte-for-byte correct against the CUDA kernels.
`object_key_for` divided `captured_at_ns` by 1e9 in float, and the division rounds. For a capture in the 596 ns window before a UTC midnight the rounded second lands on the next day, so the Python sink named a day the capture provably did not happen on. The C++ port (native/csrc/sink/object_key.cpp) truncates and was already right, so the two sinks disagreed on the same input: 1767225599999999404 gave Python 2026-01-01 and native 2025-12-31. `captured_at_ns` is a non-negative nanosecond count, so floor and truncation coincide: `//` is the whole fix, and the C++ side is left alone. Keys written by the old code inside that window are unaffected and stay readable -- nothing parses the `date=` segment; it is a partition hint only -- so this is not an on-disk format change. The parity test now drives the native key builder with the diverging timestamp and both of its neighbours, so the seam cannot be moved instead of removed.
`flush` published its barrier and enqueued it without re-checking `_error`. The worker's failure handler snapshots `_pending_flush` and only then closes the queue, so an interleaving exists where the handler sees no waiter, the publish lands afterwards, and `put_barrier` is accepted onto a still-open queue whose only consumer has already left its loop. Nobody ever completes that barrier. With a finite timeout the caller degraded to a spurious `False`; with `timeout=None` it waited forever while holding `_flush_lock`, so every later flush from any thread returned `False` too. No in-tree caller can pass `timeout=None` (the value comes from bindings.cpp, which rejects non-finite timeouts), but `HostCapturePipeline` is public, so the hang was reachable from outside. Re-read `_error` after a successful `put_barrier`, mirroring the existing CLOSED branch. The publish happens-before that re-read, so observing no error proves the handler's snapshot has not run yet and will therefore see this barrier. The alternative -- closing the queue inside the handler's `_lock` -- would force `_lock` to be held across `put_barrier` and create a lock-ordering constraint against the queue condition, so it is not taken. The test drives the interleaving with events only: a sink that blocks then raises, a `_FlushBarrier` that parks flush between its `_error` read and the publish, and a `_queue.close` that parks the handler between its snapshot and the close. Nothing depends on timing, and the flush timeout is finite so a regression asserts rather than hangs.
The seam test added in 201093c parks flush() before its barrier object exists, so _pending_flush is still None when the worker's failure handler snapshots it. That pins the _error re-read arm and leaves the other one -- the ordinary case, where the sink raises on a record queued ahead of an already-waiting barrier -- covered by nothing: replacing the handler's wake with 'if False: pass' left the whole cpu suite green. Add a sibling that forces the opposite interleaving by gating the worker's failure on flush having entered completed.wait(). Entering the wait is the observable proving the publish, a successful put_barrier and the _error re-read have all already happened with no error latched, so the handler's wake is the only thing left that can complete the barrier. The flush timeout stays finite so a regression asserts rather than hangs. The existing seam test is untouched: the two arms need opposite interleavings, and a control mutation of the re-read still fails it alone.
_CapturePackTarget._submit_capture's dtype and shape checks are the only thing tying a native-produced tensor to the metadata written into the pack footer, and both could be replaced with 'pass' with the entire suite still green. The cause was the e2e stub: _NativeValidationTarget reimplemented the dtype check with its own message, so the only test asserting on it asserted a string the product never emits, and the stub had no shape check at all. Make the stub delegate to a real _CapturePackTarget over a throwaway pipeline and assert the product's own message, then add the two cases with no downstream stand-in, as pure-CPU direct calls so they run in the gate that skips the whole e2e file: an equal-width dtype swap (int32 (3,4) over a float32 (3,4) payload) and a permuted shape (float32 (3,4) over a (4,3) payload). Both are 48 bytes either way, so the pack model's total-bytes compare and verify_payload's length + CRC32 are satisfied by either and the footer would describe bytes it does not match. Parametrize ids are now explicit because the last case carries a full metadata document that pytest would otherwise splice into the test id.
…alse Two dark arms: BoundedRecordQueue.put_barrier's CLOSED return and flush's CLOSED handling. Either regression turns flush()-after-close() into a hang for the timeout=None callers HostCapturePipeline still allows, and record_adapter._flush_capture has no _closed guard to stop it. The pipeline test replaces the barrier's completed event rather than timing the call, so a regression fails an assertion instead of sleeping.
…barrier Two dark arms of flush(). The latched-_error raise is behaviourally redundant (the CLOSED arm re-raises the same error) but its caller state -- flushing after the worker already died -- had no end-to-end test at all. The reuse-path raise is load-bearing: without it a second flush that inherits a barrier which errored, with no records submitted since, falls through to the target comparison and returns True on a failed pipeline. Both tests gate every step on an event and keep flush timeouts finite.
The post-switch rollback arm was dark in cpu and live: an `except BaseException: raise` mutant survived the whole gate. Dropping the Python owners alone does release the lease, so a leak assertion proves nothing; what the arm buys is that the GIL-releasing stop() runs before the destructor path joins the record worker with the GIL held, which is what keeps a Python sink from deadlocking. The test asserts the order first, then the rolled-back engine state.
Samfisheryu
approved these changes
Sep 22, 2026
Samfisheryu
left a comment
Collaborator
There was a problem hiding this comment.
Reviewed at 9aaa871. No blockers found. Independently passed 1,499 CPU tests and 53 live ClickHouse E2E/native-reader parity tests; restoring each of the three original implementations reproduced the corresponding regression. Sink-routing checks also pass. This PR leaves the shared Ring, sink-selection boundary, and existing production ClickHouse path unchanged. The Python reader fixes also apply when reading native-written capture packs. GPU E2E was not rerun for this Python-only production diff.
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.
What
A review pass over
src/dmi/— the Python oracle. Three behaviour fixes andsix tests that pin behaviour which was previously unpinned or pinned by a stub.
The scope was chosen deliberately.
src/dmi/is the reference implementationevery native C++ fix is verified against: it has been the measuring stick for
roughly twenty findings and had never itself been reviewed. A defect here is
worse than one in the port, because the port's conformance tests would
faithfully reproduce it.
Production diff is 52 lines across 3 files; the remaining ~921 lines are
tests.
Based directly on
main(#133), so it is independent of #132, #134 and #135.No production file overlaps any of them — see Overlap below.
The three behaviour fixes
The footer cache returned descriptors without re-checking the tenant bind
reader.pycaches pack footers by content identity. On a hit it returned thecached descriptor directly, skipping
_reject_a_foreign_tenant— so a secondcaller, under a different tenant, could read a descriptor the bind would have
refused.
The fix re-runs the location bind on the hit rather than widening the cache
key. Adding
object_keyto the key was considered and rejected: that leaves acontent-identity cache silently carrying a location-derived result, which is
the same defect one layer down.
_reject_a_foreign_tenantis nowreject_a_foreign_tenant(ref, tenants)so the reader can call it against acached descriptor.
The
date=object-key segment was computed by float divisionpipeline.py:296dividedcaptured_at_nsby1e9, which rounds. The C++port truncates. Now
// 1_000_000_000.The divergence window is 596 ns per day boundary, measured rather than
estimated:
…999999403agrees,…999999404diverges. A capture landing inthat window is written under tomorrow's
date=partition by Python andtoday's by the port — the two implementations disagree about where the object
lives, which is exactly the class of drift the conformance suite exists to
catch and could not see here.
A manual flush could orphan its barrier at the failure seam
If the worker failed between a flush installing its barrier and the worker
reaching it, the barrier was never woken and the flush waited forever. Now the
failure handler wakes it.
The six pins
These change no production behaviour. Each was verified red-then-green by
reverting the source under test and confirming exactly one test fails.
check and its own message —
grepfound that message only in the stub andthe
match=asserting it, so the test passed with the product's guarddeleted. The fix was to make the stub delegate.
Falseinto anunbounded hang.
Trueon afailure.
left unpinned, because its new test parks flush before the barrier exists.
Checks
214bdb29aaa871make -C native cpu-goals(forced-B)pytest -m cpupytest -m "clickhouse and manual and not garage"compileall src/dmiThe one live failure is pre-existing and environmental:
test_a_role_that_cannot_see_one_object_is_told_to_grant_it_not_to_rebuildneeds
CREATE USER, and the local standalone ClickHouse has no accessmanagement. It passes in CI.
The build check forces
-Bon purpose — without it, cached objects suppressre-emission and a warning count is vacuous.
Threading stability: 30 runs, zero flakes. Twenty full-file runs plus ten
under 48-way CPU saturation on 32 cores. Latency rose ~1.7x with no
behavioural change, which is the actual proof: the interleavings are forced by
synchronisation primitives rather than sleeps, so they hold under load.
Per-test JUnit parsing confirmed all five interleaving tests ran in all twenty
runs and none silently deselected.
make test-packageis NOT MEASURABLE on this host, and is not counted asgreen: it dies inside the stdlib (
venv.EnvBuilder(with_pip=True),ensurepiprc=127), reproduced with zero repo files involved. The parts that could
regress from source changes do pass — the wheel builds, and
_validate_archivereturns clean. Only the venv smoke-install is unmeasured.
Overlap
No production file here is touched by #132, #134 or #135. Three test files
collide and will need a trivial merge whichever lands second:
tests/test_engine_runtime_api.py— all four branchestests/test_native_pack_sink.py— Native capture decoder parity, and the tests that hid the gaps #134, Dynamic-dim shape inference, upload conflict reporting, and libcurl lifetime #135tests/test_capture_record_adapter.py,tests/test_record_cpu_direct.py— Capture-path correctness fixes, and tests for the paths that hid them #132Known, not fixed
record_adapter.py:240-241is dead code, shadowed unconditionally by theis_runningcheck at 231. Found during verification. It is a code changerather than a missing test, so it is out of scope for a behaviour-preserving
pass and wants its own PR.
-Wall -Wextraauthority has a hole. Two of the nine cpu goalscompile without it —
bench_builderand_dmi_native_sink— so "0 warnings"says nothing about
bench_builder.cpp,native_pack_sink.cpporbindings_sink.cpp. Pre-existing and untouched by this Python-only diff.optional, each with its evidence.They are a maintainer's call, not the pass's.
How to read the ledger
56 findings reported, 26 verified, 9 fixed: 5 refuted outright, 13
confirmed-then-downgraded, 1 deduped. The gap is the point — every refutation
and downgrade would otherwise have been a fix applied on a reviewer's say-so.
Reachability was the sole reason for the downgrades: the mechanism was real and
demonstrated, and no production caller could reach it.