Eager safety net: flush when the task ring is full, not only the payload - #160
Merged
Merged
Conversation
A legacy step with more firing hooks than task_ring_entries gets STEP_OVERSIZED and runs every hook through the eager safety net in HookPoint.forward. The net checked only payload bytes before reserve_one, and reserve_one advanced the task head with no task-capacity check, so hook task_cap reserved a sequence whose producer publishes into slot 0 while hook 0's READY word is still unread. The drain then pairs the wrong size with a TensorMeta, or clears the new word and waits at that sequence forever while flush_and_wait reports success. Found by TLA+ model checking of the payload ring. - RingEnginePy::available_task_slots() (bound to Python) reports free task entries, computed like available_capacity(). - The safety net reserves without a flush only when bytes AND a task entry are free; otherwise it takes the flush-then-reserve branch. - reserve_one throws std::logic_error, reserving nothing, when no task entry is free, so no caller can overwrite an unread slot silently. - estimate.py no longer claims the over-task-cap case goes CPU-direct.
The eager safety net now flushes when the task ring is full, but the only native coverage was the reserve_one refusal. This drives nine hooks of a STEP_OVERSIZED step through a legacy RingEnginePy with 1, 2 and 4 task entries, using the same check-and-reserve as HookPoint.forward, and checks the consumer receives every payload byte for byte, in order, paired with its own TensorMeta, with both rings fully released afterwards. Row counts differ per hook, so a payload paired with its neighbour's meta fails the consumer's shape check and goes missing.
check_ring_fit kept the busiest rank's hook count in over_task_cap whether or not it exceeded task_entries, and then branched on its truthiness. Any bytes-only overflow with task_entries passed therefore got the task-entry message, which told the user their 10 hooks were "above the 1024 task entries". The configurator always passes task_entries, so it showed this for every bytes overflow. The flag now records whether the hook count actually exceeds the cap, and the count keeps its own name.
available_task_slots() returned task_cap - (head - tail) in uint64_t, and reserve_one refused only when that was exactly 0. The CPU accounting can hold more than the ring: a flush releases only what producers published, so a reservation nothing publishes stays (HF's driver swallows a step that fails after prepare_step, a leak adapter.py already documents), and prepare_step's flushed path reserves on top of it without checking again. Past task_cap the difference wraps, HookPoint's "slots > 0" check and reserve_one's guard both pass, and the overwrite the task check exists to stop comes back. On a 2-entry ring with 3 leaked entries, 12 entries ended up outstanding and 1 of 9 captures arrived, mispaired, while flush_and_wait returned normally. The payload side had the same wrap in available_capacity() and in prepare_step's own copies of both sums. One helper now computes the room for all four and reads 0 once the reservations reach the ring, so a leaked ring flushes every step and a full one refuses loudly instead of publishing over unread slots. Releasing a failed step's reservation is the root fix and is left for a follow-up. New native case test_ring_room_saturates_past_capacity leaks 3 entries on a 2-entry ring; it fails 6 checks before this change and passes after.
Once a legacy drain fails it parks, and DrainThread::force_flush_and_wait returns at once without raising, as the legacy path always has (the error goes to stderr). With task entries checked, the eager safety net then fills the ring, calls flush_and_wait, which frees nothing, and reserve_one raises "no free task-ring entry; flush_and_wait first" out of the model forward. It tells a caller that has just flushed to flush, and hides the cause. reserve_one now rethrows the drain's own failure first, as the record paths already do before they reserve. With no drain failure, the message says what a full ring right after a flush means: reservations no producer publishes, such as a step that failed after prepare_step. The eager forward still fails loudly after a drain failure; the non-eager prepare_step path still loses capture quietly, as before. No test: the drain fails only on a CUDA error, there is no fault-injection seam, and forcing a sticky device fault on a shared GPU is not acceptable here.
The real-GPU regression kept for the reviewer copies HookPoint.forward's check-and-reserve into C++ and calls RingEnginePy directly, and the Python side of the check runs only against fake engines. Nothing committed drove HookPoint.forward -> engine.available_task_slots() (pybind) -> a real legacy ring, so binding available_task_slots to available_capacity passed the whole GPU/CPU set while the model forward raised from reserve_one on the hook after the task entries ran out. tests/test_eager_safety_net_gpu.py builds a MonitoringEngine legacy ring with 1, 2 and 4 task entries, pushes metadata for nine hooks, gets STEP_OVERSIZED from prepare_step and fires nine real HookPoints on CUDA tensors under force_eager. Each hook must return without raising and leave the ring's free counts within its capacity; after flush_and_wait the ring must be empty again.
_spec_needs_eager() is ORed into force_eager even when the step fits, so commit_step had already reserved the whole step through prepare_step (RESERVED or FLUSHED) before every hook's safety net reserved its own entry again with reserve_one. Producers publish one entry per hook, so each such step left hook_count task entries and total_bytes in the CPU accounting that no flush ever frees. The eager path never looked at task entries before this branch, so the leak went unnoticed; with the task check, the forward raises once the leak reaches task_cap (on a real 4-entry ring with 2 hooks per step, on the second step), or the free count wraps and captures are lost silently. prepare_step takes reserve (default true). With reserve=false it makes the same decision, flushing when the ring is full, and returns the same code, but reserves nothing. commit_step passes reserve=not plan.needs_eager, so the eager safety net's per-hook reservations are the only ones, and each hook reserves its real size. RESERVED and FLUSHED keep their meaning for such a step: whether the ring had room for it. Test-first: - tests/test_hook_point_eager_task_slots.py: the fake ring's flush now releases only published entries, as the native drain does. Three needs_eager steps through commit_step and real HookPoints must leave nothing outstanding; before this change step 0 leaked (2, 256). - tests/test_eager_safety_net_gpu.py: the same on a real legacy ring; before this change 2 of 4 entries stayed reserved after step 0's flush. - The prepare_step fakes accept reserve, and test_adapter_protocol pins reserve=False for a needs_eager plan and True otherwise.
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.
Replaces #155, which the Claude GitHub app opened. The original commit is re-landed on current main (8b7991d) and authored by Alan Liu, so the squash commit is his. Follow-up commits from an independent review are added on top. #155 was approved by XbzOnGit and by zaoxing at the original commit; that commit's changes (19f6492) are identical to the approved 3763122.
Original description (#155)
Requested by Alan Liu · Slack thread
Before: In the legacy ring, a step that fires more hooks than the task ring has entries could quietly lose or mislabel captured tensors. The eager fallback only checked that there was enough payload space. It never checked for a free task slot, so the hook after the ring filled wrote over a slot the drain had not read yet.
flush_and_waitstill reported success.After: The eager fallback flushes the ring when either payload space or task slots run out, so every hook gets a free slot.
reserve_oneitself now refuses (raises) instead of overwriting when no task slot is free. Record-mode rings were never affected and behave the same.How. Found by TLA+ model checking of the payload ring (
CapacityRespectedand thenNoSlotOverwritefail in 12–15 steps). The failure path, confirmed againstmain @ a987dfe:prepare_stepseesnum_hooks > task_cap, flushes, and returnsSTEP_OVERSIZED. The adapter setsforce_eager.HookPoint.forwardchecks onlyavailable_capacity()(payload bytes) and callsreserve_one. That callsDrainThread::reserve(n, 1), which advances the task head with no task-capacity check.task_caprelease-stores itsREADY|sizeword into slot 0 while hook 0's word is still unread (producers never read the tails). The drain then either reads the wrong size, shifting the TensorMeta pairing by one, or clears the new word inflush_state_updateand waits at that sequence forever.The fix adds
RingEnginePy::available_task_slots(), computed and bound the same way asavailable_capacity(). The safety net takes the no-flush branch only when the bytes fit and at least one task entry is free. Otherwise it flushes and then reserves, and a flush frees both. As a second guard,reserve_onethrowsstd::logic_errorand reserves nothing when no task entry is free. The ring's only otherreserve_onecaller is this safety net, and it now checks first, so no correct path hits the throw. Theestimate.pymessage that said the over-task-cap case "falls back to eager CPU-direct dispatch" is corrected, and the check-and-reserve comments inring_engine_py.cunow cover task slots.Tests
tests/test_hook_point_eager_task_slots.py: a fake legacy ring with 4 task entries and plenty of payload fires 5 hooks through the eager safety net. It asserts that no more than 4 tasks are ever outstanding and that the 5th hook flushes before it reserves. Fails onmain(assert 5 <= 4), passes with the fix.test_reserve_one_refuses_without_a_free_task_slotintests/native/ring/test_ring_engine.cu.test_hook_point_eager_cap_cache.pyandtest_producer_chunked_schema.pygainavailable_task_slots().python -m pytest -m cpu -q: 1944 passed, 324 skipped (the skips are native binaries not built here).ruff check --select F821 src/dmiis clean.bindings.cpppassedg++ -fsyntax-onlyagainst the torch/pybind headers. A host-sideg++parse ofring_engine_py.cuandtest_ring_engine.cushows the same errors asmain(only__CUDACC__-gated launcher declarations) and none in the touched code. CI's native compile is the real build check.Re-landed on main, plus review fixes
Composition with #152. #152 rewrote the record-ring paths in the same files,
ring_engine_py.{cu,h}and the bindings. The cherry-pick applied cleanly. Reading the merged code confirmed the two paths stay separate:prepare_step,reserve_oneandflush_and_waitare unchanged by Sink refusals and stalls never kill the forward: failure policy, per-step stall budget, bounded admission #152.reserve_onestill refuses on a record ring before the new task-slot check.DrainThreadandP2PThreadare untouched.Added commits:
layer_no, shapes paired with the right hooks, and both rings fully released. It fails with the pre-fix check, mispairing shapes and losing deliveries.needs_eagerstep whose reservation already succeeded is reserved once, per hook, instead of also as a whole.prepare_step(reserve=False)makes the decision without reserving. Before, every such step leakedhook_counttask entries.reserve_onereports a failed drain first, instead of a bare "no free slot".HookPointand pybind on a realMonitoringEnginelegacy ring.check_ring_fitblames task entries only when the busiest rank's hook count exceeds them. This was already on main since DMI-configurator: structured configuration, YAML, and web UI #122.Evidence:
test_ring_engine: 306 passed, 0 failed.