From 19f6492a06f5a3b38a14dba34b4058b26c585506 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 19:52:46 +0000 Subject: [PATCH 1/7] Eager safety net: flush when the task ring is full, not only the payload 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. --- native/csrc/bindings.cpp | 7 +- native/csrc/ring/ring_engine_py.cu | 38 ++++++-- native/csrc/ring/ring_engine_py.h | 8 +- src/dmi/configuration/estimate.py | 3 +- src/dmi/hooks/point.py | 9 +- src/dmi/transport/ring.py | 3 +- tests/native/ring/test_ring_engine.cu | 29 ++++++ tests/test_hook_point_eager_cap_cache.py | 3 + tests/test_hook_point_eager_task_slots.py | 108 ++++++++++++++++++++++ tests/test_producer_chunked_schema.py | 3 + 10 files changed, 195 insertions(+), 16 deletions(-) create mode 100644 tests/test_hook_point_eager_task_slots.py diff --git a/native/csrc/bindings.cpp b/native/csrc/bindings.cpp index feac37904..a84622d57 100644 --- a/native/csrc/bindings.cpp +++ b/native/csrc/bindings.cpp @@ -829,11 +829,14 @@ PYBIND11_MODULE(TORCH_EXTENSION_NAME, m) { .def("staging_cap", &ring_py::RingEnginePy::staging_cap) .def("task_cap", &ring_py::RingEnginePy::task_cap) .def("payload_tensor", &ring_py::RingEnginePy::payload_tensor) - // Safety-net surface (eager only). available_capacity() and - // reserve_one() are CPU-only and fast -- no GIL release needed. + // Safety-net surface (eager only). available_capacity(), + // available_task_slots() and reserve_one() are CPU-only and fast -- + // no GIL release needed. // flush_and_wait() blocks on cudaStreamSynchronize + drain flush -- // GIL released so other Python threads aren't blocked. .def("available_capacity", &ring_py::RingEnginePy::available_capacity) + .def("available_task_slots", + &ring_py::RingEnginePy::available_task_slots) .def("reserve_one", &ring_py::RingEnginePy::reserve_one, py::arg("nbytes")) diff --git a/native/csrc/ring/ring_engine_py.cu b/native/csrc/ring/ring_engine_py.cu index 541486193..609aabf7c 100644 --- a/native/csrc/ring/ring_engine_py.cu +++ b/native/csrc/ring/ring_engine_py.cu @@ -891,22 +891,29 @@ at::Tensor RingEnginePy::payload_tensor() const { // // Thread safety of the check-and-reserve pattern used by the safety net: // -// if nbytes <= available_capacity(): +// if nbytes <= available_capacity() and available_task_slots() > 0: // reserve_one(nbytes) // +// Both halves are needed. A step with more hooks than task entries gets +// STEP_OVERSIZED and runs through the safety net with plenty of payload +// room, so a bytes-only check lets the (task_cap+1)-th producer publish +// over slot 0 while its READY word is still unread. +// // The main thread (this thread) is the only writer of cpu_payload_head_ -// (it advances only through reserve / reserve_one calls). The drain -// thread only ever advances cpu_payload_tail_committed_ forward as it -// frees ring space. Between the check and the reserve: +// and cpu_task_head_ (they advance only through reserve / reserve_one +// calls). The drain thread only ever advances the committed tails +// forward as it frees ring space. Between the check and the reserve: // - tail may move forward (drain freed more): actual available at // reserve time is >= what we observed. // - head is unchanged (single-threaded writer). // So the check's "fits" decision remains valid at reserve time. No extra -// locking around the pair is required. +// locking around the pair is required. The same holds for the task +// head and tail. // -// Within available_capacity(), the two accessor calls happen under -// separate mutex acquires (drain.cpu_payload_head() and -// drain.cpu_payload_tail_committed() each take mgmt_mu_ internally). +// Within available_capacity() (and likewise available_task_slots()), the +// two accessor calls happen under separate mutex acquires +// (drain.cpu_payload_head() and drain.cpu_payload_tail_committed() each +// take mgmt_mu_ internally). // The observed snapshot is non-atomic: if drain advances tail between // the two reads, available_observed = pcap - head + tail_later, which // is >= the true available at the time of the head read. That is, the @@ -921,11 +928,24 @@ uint64_t RingEnginePy::available_capacity() const { return pcap - (drain.cpu_payload_head() - drain.cpu_payload_tail_committed()); } +uint64_t RingEnginePy::available_task_slots() const { + auto& drain = impl_->engine.drain_thread(); + const uint64_t tcap = impl_->engine.task_cap(); + return tcap - (drain.cpu_task_head() - drain.cpu_task_tail_committed()); +} + // Per-hook reservation: claim nbytes of payload + 1 task entry for an // upcoming producer kernel launch. Caller must have checked -// available_capacity() first. drain.reserve takes mgmt_mu_ internally. +// available_capacity() and available_task_slots() first. The task check +// is repeated here because producers never read the tails: an entry past +// task_cap silently overwrites an unread slot rather than failing. +// drain.reserve takes mgmt_mu_ internally. void RingEnginePy::reserve_one(uint64_t nbytes) { refuse_on_record_ring(impl_->record_mode, "legacy per-hook reservation"); + if (available_task_slots() == 0) { + throw std::logic_error( + "reserve_one: no free task-ring entry; flush_and_wait first"); + } impl_->engine.drain_thread().reserve( ring::align_up(nbytes, ring::PAYLOAD_ALIGN), 1); } diff --git a/native/csrc/ring/ring_engine_py.h b/native/csrc/ring/ring_engine_py.h index a465f7c33..c388d5236 100644 --- a/native/csrc/ring/ring_engine_py.h +++ b/native/csrc/ring/ring_engine_py.h @@ -271,10 +271,16 @@ class RingEnginePy { // pending drain. CPU-only read. uint64_t available_capacity() const; + // Free task-ring entries not currently reserved and not pending + // drain. CPU-only read. + uint64_t available_task_slots() const; + // Per-hook reservation: claim `nbytes` of payload ring + 1 task entry // for an upcoming producer kernel launch. Used by the safety net // when force_eager is on and the spec is dynamic-shape. Advances - // cpu_payload_head/cpu_task_head atomically. + // cpu_payload_head/cpu_task_head atomically. Throws std::logic_error, + // reserving nothing, when no task entry is free: the producer would + // overwrite an unconsumed READY word. Callers flush_and_wait() first. void reserve_one(uint64_t nbytes); // Synchronise the current CUDA stream + force drain to process all diff --git a/src/dmi/configuration/estimate.py b/src/dmi/configuration/estimate.py index 1e55b8259..78beb52ee 100644 --- a/src/dmi/configuration/estimate.py +++ b/src/dmi/configuration/estimate.py @@ -818,7 +818,8 @@ def check_ring_fit( f"The busiest rank fires {over_task_cap} hooks in one step, " f"above the {task_entries} task entries the ring was configured " "with. prepare_step returns STEP_OVERSIZED regardless of bytes " - "and the adapter falls back to eager CPU-direct dispatch -- " + "and the adapter falls back to eager per-hook dispatch, which " + "syncs and flushes the ring each time its task entries fill -- " "capture keeps working, but the serving path pays for it. " "Raise ring task entries, narrow the layer range, or deselect " "observations." diff --git a/src/dmi/hooks/point.py b/src/dmi/hooks/point.py index 89ac20d20..52b7a50e5 100644 --- a/src/dmi/hooks/point.py +++ b/src/dmi/hooks/point.py @@ -317,9 +317,14 @@ def forward(self, x: Tensor) -> Tensor: # min is cached on the transport and shared by every hook -- # querying the pair per hook would repeat it for the whole # active set on the first eager forward. + # A free task entry is needed too: a step with more hooks + # than task entries lands here with the payload ring nearly + # empty, and reserving past task_cap makes a producer + # overwrite an unread slot. A flush frees both. effective_cap = transport.effective_cap - if transport_bytes <= min(engine.available_capacity(), - effective_cap): + if (transport_bytes <= min(engine.available_capacity(), + effective_cap) + and engine.available_task_slots() > 0): engine.reserve_one(nbytes) dispatch_producer(ring_payload, x_cont, strip_t, strip_rb, self._ring_hook_type, self._ring_hook_id) diff --git a/src/dmi/transport/ring.py b/src/dmi/transport/ring.py index 7a57ccc78..b92b07caf 100644 --- a/src/dmi/transport/ring.py +++ b/src/dmi/transport/ring.py @@ -226,7 +226,8 @@ def __init__(self, ring_engine: Any) -> None: # When True, HookPoint.forward takes the runtime safety-net branch # instead of the fast path: - # 1. fits in current slack -> reserve_one + ring + # 1. fits in current slack (bytes AND a free task entry) + # -> reserve_one + ring # 2. fits after flushing the ring -> flush_and_wait + reserve_one + ring # 3. single tensor > ring -> flush_and_wait + submit_cpu_direct # Owned by adaptor_base.before_forward (per-batch reassignment based diff --git a/tests/native/ring/test_ring_engine.cu b/tests/native/ring/test_ring_engine.cu index 23284a108..d50162321 100644 --- a/tests/native/ring/test_ring_engine.cu +++ b/tests/native/ring/test_ring_engine.cu @@ -226,6 +226,34 @@ static void test_native_reservation_uses_transport_alignment() { EXPECT(before - engine.available_capacity() == 32); } +static void test_reserve_one_refuses_without_a_free_task_slot() { + banner("reserve_one refuses once every task entry is reserved"); + // Producers never read the tails, so a reservation past task_cap would + // publish over an unread READY word. Payload room is plentiful here: + // the task ring is the only limit, as for a STEP_OVERSIZED step whose + // hook count exceeds task_ring_entries. + ring_py::RingEnginePy engine(make_py_config(), ring_py::SubmitFn{}); + engine.init(); + const uint64_t task_cap = engine.task_cap(); + EXPECT(engine.available_task_slots() == task_cap); + for (uint64_t i = 0; i < task_cap; ++i) { + engine.reserve_one(16); + } + EXPECT(engine.available_task_slots() == 0); + + const uint64_t bytes = engine.available_capacity(); + EXPECT(bytes >= 16); + bool refused = false; + try { + engine.reserve_one(16); + } catch (const std::logic_error& error) { + refused = std::strstr(error.what(), "task-ring entry") != nullptr; + } + EXPECT(refused); + EXPECT(engine.available_task_slots() == 0); + EXPECT(engine.available_capacity() == bytes); +} + template static bool refuses_as_legacy_on_record_ring(Fn&& call) { try { @@ -1171,6 +1199,7 @@ int main() { std::printf("test_ring_engine (current drain pipeline)\n"); test_ring_geometry_requires_payload_alignment(); test_native_reservation_uses_transport_alignment(); + test_reserve_one_refuses_without_a_free_task_slot(); test_record_ring_refuses_every_legacy_producer_entry(); test_static_force_flush(); test_prefix_force_flush(); diff --git a/tests/test_hook_point_eager_cap_cache.py b/tests/test_hook_point_eager_cap_cache.py index 034ad7549..db85ddb5d 100644 --- a/tests/test_hook_point_eager_cap_cache.py +++ b/tests/test_hook_point_eager_cap_cache.py @@ -37,6 +37,9 @@ def payload_tensor(self) -> torch.Tensor: def available_capacity(self) -> int: return self.available + def available_task_slots(self) -> int: + return 1 # task capacity is not under test here + def payload_cap(self) -> int: self.payload_cap_calls += 1 return self.capacity diff --git a/tests/test_hook_point_eager_task_slots.py b/tests/test_hook_point_eager_task_slots.py new file mode 100644 index 000000000..5832db9f7 --- /dev/null +++ b/tests/test_hook_point_eager_task_slots.py @@ -0,0 +1,108 @@ +"""The eager safety net must not reserve more tasks than the task ring holds. + +A legacy step with more firing hooks than ``task_ring_entries`` gets +``STEP_OVERSIZED`` from ``prepare_step`` (after a flush, so the ring is +empty), and the adapter sets ``force_eager``. Each hook then took the safety +net in ``HookPoint.forward``, which checked only payload BYTES before +``reserve_one`` -- and ``reserve_one`` advanced the task head with no +task-capacity check. With plenty of payload room, hook ``task_cap`` reserved +sequence ``task_cap``, whose producer release-stores into slot 0 while hook +0's READY word there is still unread (producers never read the tails). The +drain then either pairs the wrong size with hook 0's 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. + +The fake below keeps a legacy ring's head/tail counters. ``HookPoint.forward`` +only takes the eager branch for a CUDA tensor, so a CPU tensor subclass poses +as one (as in tests/test_record_runtime.py) and ``dispatch_producer`` is +monkeypatched; the transport is a real ``RingTransport``. +""" +from __future__ import annotations + +import pytest +import torch + +from dmi.hooks.specs import align_up_py + +pytestmark = pytest.mark.cpu + + +class _FakeCudaTensor(torch.Tensor): + @property + def is_cuda(self) -> bool: + return True + + +class _LegacyRingEngine: + """Payload and task heads/tails of a legacy ring; flush drains both.""" + + def __init__(self, task_cap: int, payload_cap: int): + self.task_cap = task_cap + self.capacity = payload_cap + self.task_head = self.task_tail = 0 + self.payload_head = self.payload_tail = 0 + self.events: list[str] = [] + self.max_outstanding_tasks = 0 + + def payload_tensor(self) -> torch.Tensor: + return torch.empty(64, dtype=torch.uint8) + + def payload_cap(self) -> int: + return self.capacity + + def staging_cap(self) -> int: + return self.capacity + + def available_capacity(self) -> int: + return self.capacity - (self.payload_head - self.payload_tail) + + def available_task_slots(self) -> int: + return self.task_cap - (self.task_head - self.task_tail) + + def reserve_one(self, nbytes: int) -> None: + self.events.append("reserve") + self.payload_head += align_up_py(nbytes, 16) + self.task_head += 1 + self.max_outstanding_tasks = max( + self.max_outstanding_tasks, self.task_head - self.task_tail) + + def flush_and_wait(self) -> None: + self.events.append("flush") + self.task_tail = self.task_head + self.payload_tail = self.payload_head + + +def _eager_hook(monkeypatch, engine, dispatched): + from dmi.hooks.point import HookPoint + from dmi.transport import ring as ring_transport + from dmi.transport.ring import RingTransport + + transport = RingTransport(engine) + transport.force_eager = True + monkeypatch.setattr(ring_transport, "_active_transport", transport) + monkeypatch.setattr( + "dmi.hooks.point.dispatch_producer", + lambda *args: dispatched.append(args), + ) + hook = HookPoint() + hook._ring_hook_type = 1 + hook._ring_hook_id = 2 + hook._ring_payload = transport._ring_payload + return hook + + +def test_a_step_with_more_hooks_than_task_entries_flushes_before_overflowing( + monkeypatch): + task_cap = 4 + engine = _LegacyRingEngine(task_cap=task_cap, payload_cap=1 << 20) + dispatched = [] + hook = _eager_hook(monkeypatch, engine, dispatched) + + value = torch.arange(32, dtype=torch.uint8).as_subclass(_FakeCudaTensor) + for _ in range(task_cap + 1): + hook(value) + + assert len(dispatched) == task_cap + 1, "every hook still takes the ring" + assert engine.max_outstanding_tasks <= task_cap, ( + "a reservation past task_cap overwrites an unread task slot") + assert engine.events == ["reserve"] * task_cap + ["flush", "reserve"] diff --git a/tests/test_producer_chunked_schema.py b/tests/test_producer_chunked_schema.py index dc46c00f0..c2b9596fe 100644 --- a/tests/test_producer_chunked_schema.py +++ b/tests/test_producer_chunked_schema.py @@ -172,6 +172,9 @@ def payload_tensor(self) -> torch.Tensor: def available_capacity(self) -> int: return self.available + def available_task_slots(self) -> int: + return 1 # task capacity is not under test here + def payload_cap(self) -> int: return self.capacity From 53b5a8c8d8e2510b3a7f65cd0e1be701c111e3ab Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 16:26:44 -0400 Subject: [PATCH 2/7] Keep a real-GPU delivery regression for the eager task-ring flush 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. --- tests/native/ring/test_ring_engine.cu | 123 ++++++++++++++++++++++++++ 1 file changed, 123 insertions(+) diff --git a/tests/native/ring/test_ring_engine.cu b/tests/native/ring/test_ring_engine.cu index d50162321..f81245545 100644 --- a/tests/native/ring/test_ring_engine.cu +++ b/tests/native/ring/test_ring_engine.cu @@ -11,6 +11,7 @@ #include #include +#include #include #include #include @@ -19,6 +20,7 @@ #include #include #include +#include #include #include #include @@ -254,6 +256,126 @@ static void test_reserve_one_refuses_without_a_free_task_slot() { EXPECT(engine.available_capacity() == bytes); } +// A step that fires more hooks than the task ring has entries gets +// STEP_OVERSIZED, and every hook goes through HookPoint.forward's eager +// safety net: reserve without a flush only while the bytes AND a task entry +// are free, otherwise flush_and_wait first. The task ring wraps several +// times inside the one step, so each payload must still reach the consumer +// byte for byte and paired with its own TensorMeta. Row counts differ per +// hook, so a payload paired with its neighbour's meta fails the shape check +// in the consumer and goes undelivered. +static void check_eager_delivery_past_the_task_ring(uint32_t task_slots) { + constexpr int hooks = 9; + constexpr int64_t cols = 24; // not a multiple of PAYLOAD_ALIGN + ring_py::RingConfig cfg = make_py_config(); + cfg.task_ring_entries = task_slots; + + struct Delivery { + std::string act_name; + int32_t layer_no; + std::vector shape; + std::vector bytes; + }; + std::mutex delivered_mu; + std::vector delivered; + ring_py::RingEnginePy engine( + cfg, + [&](const std::string&, int32_t, const std::string&, + const std::string& act_name, int32_t layer_no, int32_t, int32_t, + at::Tensor slice) { + const uint8_t* data = slice.data_ptr(); + Delivery delivery{act_name, layer_no, slice.sizes().vec(), + std::vector(data, data + slice.nbytes())}; + std::lock_guard lock(delivered_mu); + delivered.push_back(std::move(delivery)); + }); + engine.init(); + engine.start(); + + std::vector> sources; + std::vector devices; + std::vector metas; + uint64_t step_bytes = 0; + for (int hook = 0; hook < hooks; ++hook) { + const int64_t rows = 1 + hook % 3; + sources.push_back(pattern(static_cast(rows * cols), + static_cast(17 * hook + task_slots))); + uint8_t* device = nullptr; + CUDA_CHECK(cudaMalloc(&device, sources.back().size())); + CUDA_CHECK(cudaMemcpy(device, sources.back().data(), + sources.back().size(), cudaMemcpyHostToDevice)); + devices.push_back(device); + ring_py::TensorMeta meta; + meta.hook_type = ring_py::HOOK_TYPE_RESID_PRE; + meta.layer_no = hook; + meta.shape = {1, rows, cols}; + meta.dtype = static_cast(at::kByte); + meta.last_in_step = hook == hooks - 1; + metas.push_back(std::move(meta)); + step_bytes += ring::align_up(sources.back().size(), ring::PAYLOAD_ALIGN); + } + auto* context = new ring_py::StepContext(); + context->model_id = "model"; + context->requests.push_back({"request", 0, 4, 0, 0}); + engine.push_step(context, metas); + + EXPECT(engine.prepare_step(step_bytes, hooks) == + ring_py::RingEnginePy::STEP_OVERSIZED); + const uint64_t effective_cap = + std::min(engine.payload_cap(), engine.staging_cap()); + uint64_t most_outstanding = 0; + bool refused = false; + for (int hook = 0; hook < hooks; ++hook) { + const uint64_t nbytes = sources[hook].size(); + const uint64_t transport_bytes = + ring::align_up(nbytes, ring::PAYLOAD_ALIGN); + try { + if (transport_bytes > std::min(engine.available_capacity(), + effective_cap) || + engine.available_task_slots() == 0) { + engine.flush_and_wait(); + } + engine.reserve_one(nbytes); + } catch (const std::logic_error&) { + refused = true; + break; + } + most_outstanding = std::max( + most_outstanding, + engine.task_cap() - engine.available_task_slots()); + engine.hook_no_notify(reinterpret_cast(devices[hook]), + nbytes, ring_py::HOOK_TYPE_RESID_PRE, 0); + } + engine.flush_and_wait(); + EXPECT(!refused); + EXPECT(most_outstanding <= task_slots); + EXPECT(engine.available_task_slots() == task_slots); + EXPECT(engine.available_capacity() == engine.payload_cap()); + // stop() lets the consumer finish every task already drained. + engine.stop(); + + EXPECT(delivered.size() == static_cast(hooks)); + for (size_t hook = 0; hook < delivered.size() && hook < sources.size(); + ++hook) { + const Delivery& delivery = delivered[hook]; + const int64_t rows = 1 + static_cast(hook) % 3; + EXPECT(delivery.layer_no == static_cast(hook)); + EXPECT(delivery.act_name == "blocks.hook_resid_pre"); + EXPECT((delivery.shape == std::vector{rows, cols})); + EXPECT(delivery.bytes == sources[hook]); + } + for (uint8_t* device : devices) { + CUDA_CHECK(cudaFree(device)); + } +} + +static void test_eager_safety_net_delivers_past_the_task_ring() { + banner("eager safety net delivers nine hooks through 1, 2 and 4 task slots"); + for (uint32_t task_slots : {1u, 2u, 4u}) { + check_eager_delivery_past_the_task_ring(task_slots); + } +} + template static bool refuses_as_legacy_on_record_ring(Fn&& call) { try { @@ -1200,6 +1322,7 @@ int main() { test_ring_geometry_requires_payload_alignment(); test_native_reservation_uses_transport_alignment(); test_reserve_one_refuses_without_a_free_task_slot(); + test_eager_safety_net_delivers_past_the_task_ring(); test_record_ring_refuses_every_legacy_producer_entry(); test_static_force_flush(); test_prefix_force_flush(); From 68d63cd3bf774c0c6d86a22d7d1fb6ae7c18f9e2 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 17:19:10 -0400 Subject: [PATCH 3/7] Blame bytes, not task entries, when only the bytes overflow 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. --- src/dmi/configuration/estimate.py | 13 +++++++------ tests/test_review_findings_runtime.py | 25 +++++++++++++++++++++++++ 2 files changed, 32 insertions(+), 6 deletions(-) diff --git a/src/dmi/configuration/estimate.py b/src/dmi/configuration/estimate.py index 78beb52ee..14ffe9b25 100644 --- a/src/dmi/configuration/estimate.py +++ b/src/dmi/configuration/estimate.py @@ -799,13 +799,14 @@ def check_ring_fit( # step than task_ring_entries returns OVERSIZED regardless of bytes. # The hook counts are already on every rank; the cap is the caller's # ring configuration, so it arrives as an argument. - over_task_cap = 0 + busiest_hooks = 0 + over_task_cap = False if task_entries is not None and task_entries >= 1: for load in estimate.ranks: - worst_hooks = max(load.prefill_hooks, load.decode_hooks) - if worst_hooks > over_task_cap: - over_task_cap = worst_hooks - if over_task_cap > task_entries: + busiest_hooks = max(busiest_hooks, load.prefill_hooks, + load.decode_hooks) + over_task_cap = busiest_hooks > task_entries + if over_task_cap: fits = False if fits: @@ -815,7 +816,7 @@ def check_ring_fit( ) elif over_task_cap: detail = ( - f"The busiest rank fires {over_task_cap} hooks in one step, " + f"The busiest rank fires {busiest_hooks} hooks in one step, " f"above the {task_entries} task entries the ring was configured " "with. prepare_step returns STEP_OVERSIZED regardless of bytes " "and the adapter falls back to eager per-hook dispatch, which " diff --git a/tests/test_review_findings_runtime.py b/tests/test_review_findings_runtime.py index a0e1a1733..f65b9c962 100644 --- a/tests/test_review_findings_runtime.py +++ b/tests/test_review_findings_runtime.py @@ -858,6 +858,31 @@ def test_hooks_within_the_cap_fit_on_bytes_alone(self): fit = check_ring_fit(est, payload_bytes=1024 * 1024, task_entries=1024) assert fit.fits is True + def test_a_bytes_overflow_within_the_task_cap_blames_the_bytes(self): + """The configurator always passes task_entries, so a step that + overflows on bytes alone must still get the bytes explanation, not + one claiming its 10 hooks exceed 1024 task entries.""" + from dmi.configuration.estimate import Estimate, RankLoad, check_ring_fit + + est = Estimate( + peak_step_bytes=4096, + peak_step_rank="pp0/tp0", + decode_step_bytes=4096, + bytes_per_request=4096, + aggregate_peak_step_bytes=4096, + sustained_bytes_per_second=None, + bytes_per_day=None, + ranks=( + RankLoad(label="pp0/tp0", pp_stage=0, tp_rank=0, + prefill_step_bytes=4096, decode_step_bytes=4096, + prefill_hooks=10, decode_hooks=10), + ), + ) + fit = check_ring_fit(est, payload_bytes=1024, task_entries=1024) + assert fit.fits is False + assert "effective ring" in fit.detail + assert "task entries" not in fit.detail + class TestPPSplitMatchesVLLM: def test_remainder_layers_land_on_the_first_stages(self): From 4729abd3a385a7f08c70ec483b803351822152ab Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 17:32:42 -0400 Subject: [PATCH 4/7] Read a legacy ring past its capacity as full, not as ~2^64 free 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. --- native/csrc/ring/ring_engine_py.cu | 29 +++++++++++++------ native/csrc/ring/ring_engine_py.h | 6 ++-- tests/native/ring/test_ring_engine.cu | 41 +++++++++++++++++++++++++++ 3 files changed, 66 insertions(+), 10 deletions(-) diff --git a/native/csrc/ring/ring_engine_py.cu b/native/csrc/ring/ring_engine_py.cu index 609aabf7c..d6969fd2f 100644 --- a/native/csrc/ring/ring_engine_py.cu +++ b/native/csrc/ring/ring_engine_py.cu @@ -130,6 +130,19 @@ void refuse_on_record_ring(bool record_mode, const char* what) { } } +// Free room in a legacy ring of `cap` entries (or bytes) whose CPU +// accounting has reserved up to `head` and released up to `tail`, saturated +// at 0. The accounting can hold more than the ring: a flush releases only +// what producers published, so a reservation nothing publishes (a step that +// fails after prepare_step) stays, and prepare_step's flushed path reserves +// on top of it without checking again. Unsaturated, cap - (head - tail) then +// wraps to ~2^64, every check admits every reservation, and the producers +// publish over unread slots. +uint64_t legacy_ring_room(uint64_t cap, uint64_t head, uint64_t tail) { + const uint64_t outstanding = head - tail; + return outstanding >= cap ? 0 : cap - outstanding; +} + } // namespace // --------------------------------------------------------------------------- @@ -620,10 +633,8 @@ int RingEnginePy::prepare_step(uint64_t step_total_bytes, } // Case A: step fits. Check available space for BOTH payload AND tasks. - const uint64_t payload_avail = pcap - - (drain.cpu_payload_head() - drain.cpu_payload_tail_committed()); - const uint64_t task_avail = tcap - - (drain.cpu_task_head() - drain.cpu_task_tail_committed()); + const uint64_t payload_avail = available_capacity(); + const uint64_t task_avail = available_task_slots(); if (step_total_bytes <= payload_avail && num_hooks <= task_avail) { drain.reserve(step_total_bytes, num_hooks); @@ -924,14 +935,16 @@ at::Tensor RingEnginePy::payload_tensor() const { uint64_t RingEnginePy::available_capacity() const { auto& drain = impl_->engine.drain_thread(); - const uint64_t pcap = impl_->engine.payload_cap(); - return pcap - (drain.cpu_payload_head() - drain.cpu_payload_tail_committed()); + const uint64_t head = drain.cpu_payload_head(); + return legacy_ring_room(impl_->engine.payload_cap(), head, + drain.cpu_payload_tail_committed()); } uint64_t RingEnginePy::available_task_slots() const { auto& drain = impl_->engine.drain_thread(); - const uint64_t tcap = impl_->engine.task_cap(); - return tcap - (drain.cpu_task_head() - drain.cpu_task_tail_committed()); + const uint64_t head = drain.cpu_task_head(); + return legacy_ring_room(impl_->engine.task_cap(), head, + drain.cpu_task_tail_committed()); } // Per-hook reservation: claim nbytes of payload + 1 task entry for an diff --git a/native/csrc/ring/ring_engine_py.h b/native/csrc/ring/ring_engine_py.h index c388d5236..09fc71b9b 100644 --- a/native/csrc/ring/ring_engine_py.h +++ b/native/csrc/ring/ring_engine_py.h @@ -268,11 +268,13 @@ class RingEnginePy { // CUDA-graph capture or replay. // Free bytes in the payload ring not currently reserved and not - // pending drain. CPU-only read. + // pending drain. CPU-only read. Reads 0, never a wrapped value, once + // reservations no producer publishes have pushed the CPU accounting + // past the ring. uint64_t available_capacity() const; // Free task-ring entries not currently reserved and not pending - // drain. CPU-only read. + // drain. CPU-only read. Saturates at 0 like available_capacity(). uint64_t available_task_slots() const; // Per-hook reservation: claim `nbytes` of payload ring + 1 task entry diff --git a/tests/native/ring/test_ring_engine.cu b/tests/native/ring/test_ring_engine.cu index f81245545..20c7efd13 100644 --- a/tests/native/ring/test_ring_engine.cu +++ b/tests/native/ring/test_ring_engine.cu @@ -256,6 +256,46 @@ static void test_reserve_one_refuses_without_a_free_task_slot() { EXPECT(engine.available_capacity() == bytes); } +// A step that fails after prepare_step reserved it (HF's driver swallows the +// error on the legacy ring) leaves a reservation no producer publishes, and +// a flush frees only published entries. prepare_step's flushed path then +// reserves on top of it without checking again, so the CPU accounting can +// hold more than the ring. The free counts must read 0 there, not wrap to +// ~2^64: a wrapped count admits every reservation, the eager safety net +// stops flushing, and producers publish over unread slots. +static void test_ring_room_saturates_past_capacity() { + banner("free task entries and bytes read 0 once reservations pass the ring"); + ring_py::RingConfig cfg = make_py_config(); + cfg.task_ring_entries = 2; + ring_py::RingEnginePy engine(cfg, ring_py::SubmitFn{}); + engine.init(); + engine.start(); + const uint64_t payload_cap = engine.payload_cap(); + + // Two failed steps: the first takes the whole ring, the second flushes + // (freeing nothing) and reserves on top. + EXPECT(engine.prepare_step(payload_cap, 2) == + ring_py::RingEnginePy::STEP_RING_OK); + EXPECT(engine.prepare_step(16, 1) == + ring_py::RingEnginePy::STEP_RING_FLUSHED); + EXPECT(engine.available_task_slots() == 0); + EXPECT(engine.available_capacity() == 0); + + bool refused = false; + try { + engine.reserve_one(16); + } catch (const std::logic_error& error) { + refused = std::strstr(error.what(), "task-ring entry") != nullptr; + } + EXPECT(refused); + // Not the fast path: by its own accounting the ring has no room. + EXPECT(engine.prepare_step(16, 1) == + ring_py::RingEnginePy::STEP_RING_FLUSHED); + EXPECT(engine.available_task_slots() == 0); + EXPECT(engine.available_capacity() == 0); + engine.stop(); +} + // A step that fires more hooks than the task ring has entries gets // STEP_OVERSIZED, and every hook goes through HookPoint.forward's eager // safety net: reserve without a flush only while the bytes AND a task entry @@ -1322,6 +1362,7 @@ int main() { test_ring_geometry_requires_payload_alignment(); test_native_reservation_uses_transport_alignment(); test_reserve_one_refuses_without_a_free_task_slot(); + test_ring_room_saturates_past_capacity(); test_eager_safety_net_delivers_past_the_task_ring(); test_record_ring_refuses_every_legacy_producer_entry(); test_static_force_flush(); From 2c68a4aac8a58ecd043089c2690d0326af4ab654 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 17:33:11 -0400 Subject: [PATCH 5/7] Name a failed drain when reserve_one finds no free task entry 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. --- native/csrc/ring/ring_engine_py.cu | 15 ++++++++++++--- native/csrc/ring/ring_engine_py.h | 8 +++++--- 2 files changed, 17 insertions(+), 6 deletions(-) diff --git a/native/csrc/ring/ring_engine_py.cu b/native/csrc/ring/ring_engine_py.cu index d6969fd2f..08bc90d45 100644 --- a/native/csrc/ring/ring_engine_py.cu +++ b/native/csrc/ring/ring_engine_py.cu @@ -955,12 +955,21 @@ uint64_t RingEnginePy::available_task_slots() const { // drain.reserve takes mgmt_mu_ internally. void RingEnginePy::reserve_one(uint64_t nbytes) { refuse_on_record_ring(impl_->record_mode, "legacy per-hook reservation"); + auto& drain = impl_->engine.drain_thread(); if (available_task_slots() == 0) { + // A failed drain never frees an entry again, and on a legacy ring + // its flush_and_wait returns at once without raising (the failure + // went to stderr), so a caller that just flushed lands here. Name + // that failure rather than asking for the flush it already made. + drain.rethrow_drain_failure(); throw std::logic_error( - "reserve_one: no free task-ring entry; flush_and_wait first"); + "reserve_one: every task-ring entry is reserved; call " + "flush_and_wait first. A flush frees only entries a producer " + "published, so right after one the entries belong to " + "reservations nothing publishes, such as a step that failed " + "after prepare_step"); } - impl_->engine.drain_thread().reserve( - ring::align_up(nbytes, ring::PAYLOAD_ALIGN), 1); + drain.reserve(ring::align_up(nbytes, ring::PAYLOAD_ALIGN), 1); } // Synchronise the current CUDA stream so all queued producer kernels diff --git a/native/csrc/ring/ring_engine_py.h b/native/csrc/ring/ring_engine_py.h index 09fc71b9b..838ec7d76 100644 --- a/native/csrc/ring/ring_engine_py.h +++ b/native/csrc/ring/ring_engine_py.h @@ -280,9 +280,11 @@ class RingEnginePy { // Per-hook reservation: claim `nbytes` of payload ring + 1 task entry // for an upcoming producer kernel launch. Used by the safety net // when force_eager is on and the spec is dynamic-shape. Advances - // cpu_payload_head/cpu_task_head atomically. Throws std::logic_error, - // reserving nothing, when no task entry is free: the producer would - // overwrite an unconsumed READY word. Callers flush_and_wait() first. + // cpu_payload_head/cpu_task_head atomically. Reserves nothing and + // throws when no task entry is free, since the producer would + // overwrite an unconsumed READY word: the drain's failure when the + // drain has failed, std::logic_error otherwise. Callers + // flush_and_wait() first. void reserve_one(uint64_t nbytes); // Synchronise the current CUDA stream + force drain to process all From 88d1507f7fdc276bf0ddaf8f235a6c402447653c Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 17:33:37 -0400 Subject: [PATCH 6/7] Drive the eager safety net through HookPoint and pybind on a real ring 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. --- tests/test_eager_safety_net_gpu.py | 101 +++++++++++++++++++++++++++++ 1 file changed, 101 insertions(+) create mode 100644 tests/test_eager_safety_net_gpu.py diff --git a/tests/test_eager_safety_net_gpu.py b/tests/test_eager_safety_net_gpu.py new file mode 100644 index 000000000..cf95c64fb --- /dev/null +++ b/tests/test_eager_safety_net_gpu.py @@ -0,0 +1,101 @@ +"""The eager safety net on a real legacy ring, through HookPoint and pybind. + +tests/test_hook_point_eager_task_slots.py runs ``HookPoint.forward``'s +check-and-reserve against a fake engine, and +test_eager_safety_net_delivers_past_the_task_ring in +tests/native/ring/test_ring_engine.cu runs a C++ copy of it against the +native ring. Neither goes through point.py and the pybind surface together +(``available_task_slots``, ``reserve_one``, ``flush_and_wait``), so binding +``available_task_slots`` to the wrong method passed every committed test +while the model forward raised. These drive the real composition: a +``MonitoringEngine``'s legacy ring, its ``RingTransport``, and real +``HookPoint`` modules on CUDA tensors. + +Needs CUDA and the full native backend. The ring has no host engine, so its +consumer drops what it drains; delivering each payload byte for byte with +its own metadata is the native test's job. +""" +from __future__ import annotations + +import pytest +import torch + +from tests._requirements import require_cuda, require_native_backend + +pytestmark = [ + pytest.mark.gpu, + pytest.mark.native_backend, + require_cuda(), + require_native_backend(), +] + +RING_BYTES = 64 * 1024 + + +def _legacy_engine(task_entries): + from dmi.api.v1 import RingConfig + from dmi.engine import MonitoringEngine + + config = RingConfig() + config.task_ring_entries = task_entries + config.payload_ring_bytes = RING_BYTES + config.pinned_staging_bytes = RING_BYTES + engine = MonitoringEngine(model_id="eager-safety-net", ring_config=config) + assert not engine._record_mode + return engine + + +def _arm(hook, transport, hook_type, layer_no): + """What install_ring_hooks does to a HookPoint.""" + hook._ring_hook_type = hook_type + hook._ring_hook_id = layer_no + hook._ring_payload = transport._ring_payload + return hook + + +@pytest.mark.parametrize("task_entries", [1, 2, 4]) +def test_an_oversized_step_fires_every_hook_through_a_small_task_ring( + task_entries): + """Nine hooks through 1, 2 and 4 task entries: the step is OVERSIZED on + hook count alone, and the safety net must flush each time the entries + run out rather than reserve past them.""" + from dmi.adapters.base import StepReservation + from dmi.hooks.point import HookPoint + from dmi.hooks.specs import HOOK_TYPE_RESID_PRE, align_up_py + + hooks = 9 + engine = _legacy_engine(task_entries) + transport = engine._ring_transport + ring = engine._ring_engine + try: + assert ring.task_cap() == task_entries + # Row counts differ per hook and 24 is not a multiple of 16, as in + # the native delivery test. + tensors = [ + torch.full((1 + layer % 3, 24), layer, dtype=torch.uint8, + device="cuda") + for layer in range(hooks) + ] + ring.push_all_metas( + [HOOK_TYPE_RESID_PRE] * hooks, list(range(hooks)), + [list(t.shape) for t in tensors], [torch.uint8] * hooks, + [0] * hooks, "eager-safety-net", 0, 0, 0, 0, False, + ["0:0"], [(0, 1)], [0], [0]) + step_bytes = sum(align_up_py(t.nbytes, 16) for t in tensors) + assert (ring.prepare_step(step_bytes, hooks) + == StepReservation.OVERSIZED) + + transport.force_eager = True + for layer, tensor in enumerate(tensors): + hook = _arm(HookPoint(), transport, HOOK_TYPE_RESID_PRE, layer) + assert hook(tensor) is tensor + assert 0 <= ring.available_task_slots() <= task_entries + assert ring.available_capacity() <= ring.payload_cap() + + ring.flush_and_wait() + assert ring.available_task_slots() == task_entries + assert ring.available_capacity() == ring.payload_cap() + finally: + transport.force_eager = False + engine.close() + From 74aecc3c95395cb5efb4e37facb99c6ee3261fd9 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Mon, 28 Sep 2026 17:34:22 -0400 Subject: [PATCH 7/7] Reserve a needs_eager step once: per hook, not also as a whole _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. --- docs/integration-api-v1.md | 5 +- native/csrc/bindings.cpp | 1 + native/csrc/ring/ring_engine_py.cu | 7 +- native/csrc/ring/ring_engine_py.h | 10 +- src/dmi/adapters/base.py | 8 +- tests/test_adapter_base_pipeline.py | 3 +- tests/test_adapter_protocol.py | 8 +- tests/test_eager_safety_net_gpu.py | 78 ++++++++++- tests/test_hf_capture_refusal.py | 2 +- tests/test_hook_point_eager_task_slots.py | 159 +++++++++++++++++++--- tests/test_review_findings_runtime.py | 3 +- 11 files changed, 250 insertions(+), 34 deletions(-) diff --git a/docs/integration-api-v1.md b/docs/integration-api-v1.md index bcd7d1706..b8d63bb75 100644 --- a/docs/integration-api-v1.md +++ b/docs/integration-api-v1.md @@ -673,7 +673,10 @@ hook metadata is emitted. Exceptions propagate; the driver does not roll back partial state. For each firing spec, `_spec_needs_eager()` is ORed into the step's eager -decision even when reservation succeeds. For an oversized step the callback +decision even when the step fits. Such a step is not reserved as a whole: +`commit_step()` still returns `RESERVED` or `FLUSHED` to say whether the ring +had room for it, but the eager safety net reserves each hook as it fires. For +an oversized step the callback order is `adapt_for_cpu_direct()`, `on_capacity_exceeded()`, then `_warn_once_capacity()`, and metadata is built from the adapted context. diff --git a/native/csrc/bindings.cpp b/native/csrc/bindings.cpp index a84622d57..f5ca7be95 100644 --- a/native/csrc/bindings.cpp +++ b/native/csrc/bindings.cpp @@ -751,6 +751,7 @@ PYBIND11_MODULE(TORCH_EXTENSION_NAME, m) { &ring_py::RingEnginePy::prepare_step, py::arg("step_total_bytes"), py::arg("num_hooks"), + py::arg("reserve") = true, py::call_guard()) .def("reserve_record", &ring_py::RingEnginePy::reserve_record, diff --git a/native/csrc/ring/ring_engine_py.cu b/native/csrc/ring/ring_engine_py.cu index 08bc90d45..62744726d 100644 --- a/native/csrc/ring/ring_engine_py.cu +++ b/native/csrc/ring/ring_engine_py.cu @@ -578,7 +578,8 @@ void RingEnginePy::notify_drain() { // drain thread to flush all pending entries. // --------------------------------------------------------------------------- int RingEnginePy::prepare_step(uint64_t step_total_bytes, - uint32_t num_hooks) + uint32_t num_hooks, + bool reserve) { // Record rings reserve through reserve_record. refuse_on_record_ring(impl_->record_mode, "legacy step reservation"); @@ -637,7 +638,7 @@ int RingEnginePy::prepare_step(uint64_t step_total_bytes, const uint64_t task_avail = available_task_slots(); if (step_total_bytes <= payload_avail && num_hooks <= task_avail) { - drain.reserve(step_total_bytes, num_hooks); + if (reserve) drain.reserve(step_total_bytes, num_hooks); return STEP_RING_OK; // fast path -- no CUDA or thread interaction } @@ -646,7 +647,7 @@ int RingEnginePy::prepare_step(uint64_t step_total_bytes, cudaStream_t ms = at::cuda::getCurrentCUDAStream().stream(); cudaStreamSynchronize(ms); drain.force_flush_and_wait(); - drain.reserve(step_total_bytes, num_hooks); + if (reserve) drain.reserve(step_total_bytes, num_hooks); return STEP_RING_FLUSHED; } diff --git a/native/csrc/ring/ring_engine_py.h b/native/csrc/ring/ring_engine_py.h index 838ec7d76..94bd93239 100644 --- a/native/csrc/ring/ring_engine_py.h +++ b/native/csrc/ring/ring_engine_py.h @@ -209,7 +209,12 @@ class RingEnginePy { // Synced + flushed. All hooks must use .cpu() path. // // For cases 0 and 1, advances cpu_payload_head_ and cpu_task_head_ - // under mgmt_mu_ to pre-allocate ring space for this step's producers. + // under mgmt_mu_ to pre-allocate ring space for this step's producers, + // unless `reserve` is false: then it only makes room. A step whose + // hooks all take the eager safety net although it fits (the adapter's + // needs_eager) passes false: each hook reserves its own entry with + // reserve_one, and a step reservation on top would never be released, + // since the producers publish one entry per hook. // Also resets the internal hook index counter for hook_no_notify. static constexpr int STEP_RING_OK = 0; static constexpr int STEP_RING_FLUSHED = 1; @@ -218,7 +223,8 @@ class RingEnginePy { // safety net in HookPoint.forward via transport.force_eager = True. static constexpr int STEP_OVERSIZED = 2; - int prepare_step(uint64_t step_total_bytes, uint32_t num_hooks); + int prepare_step(uint64_t step_total_bytes, uint32_t num_hooks, + bool reserve = true); // Each item is (aligned upper bound, needs actual-byte reconciliation). // When the ring has no room, waits for the drain within the step's diff --git a/src/dmi/adapters/base.py b/src/dmi/adapters/base.py index 0b1a79011..795155cd6 100644 --- a/src/dmi/adapters/base.py +++ b/src/dmi/adapters/base.py @@ -387,9 +387,15 @@ def commit_step( # their internal loops produce nothing to push when the # active-spec list is empty or every shape was empty. if plan.hook_count > 0: + # A needs_eager step runs every hook through the eager safety + # net, which reserves each hook's entry itself (reserve_one), so + # prepare_step only makes room for it. Reserving the step as + # well would hold hook_count entries and total_bytes that no + # producer publishes and no flush frees, every such step. reservation = StepReservation( self.ring_engine.prepare_step( - plan.total_bytes, plan.hook_count + plan.total_bytes, plan.hook_count, + reserve=not plan.needs_eager, ) ) self.transport.force_eager = ( diff --git a/tests/test_adapter_base_pipeline.py b/tests/test_adapter_base_pipeline.py index 83a7e504c..b3b4344ee 100644 --- a/tests/test_adapter_base_pipeline.py +++ b/tests/test_adapter_base_pipeline.py @@ -93,7 +93,8 @@ def __init__(self, prepare_step_result: int = 2) -> None: self._result = prepare_step_result self.prepare_step_calls: list = [] - def prepare_step(self, total_bytes: int, n_hooks: int) -> int: + def prepare_step(self, total_bytes: int, n_hooks: int, + reserve: bool = True) -> int: self.prepare_step_calls.append((total_bytes, n_hooks)) return self._result diff --git a/tests/test_adapter_protocol.py b/tests/test_adapter_protocol.py index 83ce63d4e..104c6c95c 100644 --- a/tests/test_adapter_protocol.py +++ b/tests/test_adapter_protocol.py @@ -55,9 +55,12 @@ class FakeRingEngine: def __init__(self, prepare_step_result: int = 0) -> None: self._result = prepare_step_result self.prepare_step_calls: list = [] + self.reserve_flags: list = [] - def prepare_step(self, total_bytes: int, n_hooks: int) -> int: + def prepare_step(self, total_bytes: int, n_hooks: int, + reserve: bool = True) -> int: self.prepare_step_calls.append((total_bytes, n_hooks)) + self.reserve_flags.append(reserve) return self._result @@ -168,6 +171,8 @@ def test_commit_step_uses_supplied_plan_without_replanning(): assert reservation is StepReservation.FLUSHED assert a.call_order == [] assert a.engine._ring_engine.prepare_step_calls == [(2048, 4)] + # needs_eager: the safety net reserves each hook, so the step is not. + assert a.engine._ring_engine.reserve_flags == [False] assert a.transport.force_eager is True assert len(a.transport.set_step_context_calls) == 1 assert len(a.transport.pre_push_all_metas_calls) == 1 @@ -181,6 +186,7 @@ def test_commit_step_plans_once_when_plan_is_omitted(): assert reservation is StepReservation.RESERVED assert a.call_order == ["plan_step"] assert a.engine._ring_engine.prepare_step_calls == [(1024, 3)] + assert a.engine._ring_engine.reserve_flags == [True] def test_commit_step_returns_skipped_without_hooks_but_publishes_context(): diff --git a/tests/test_eager_safety_net_gpu.py b/tests/test_eager_safety_net_gpu.py index cf95c64fb..0535280fd 100644 --- a/tests/test_eager_safety_net_gpu.py +++ b/tests/test_eager_safety_net_gpu.py @@ -5,11 +5,11 @@ test_eager_safety_net_delivers_past_the_task_ring in tests/native/ring/test_ring_engine.cu runs a C++ copy of it against the native ring. Neither goes through point.py and the pybind surface together -(``available_task_slots``, ``reserve_one``, ``flush_and_wait``), so binding -``available_task_slots`` to the wrong method passed every committed test -while the model forward raised. These drive the real composition: a -``MonitoringEngine``'s legacy ring, its ``RingTransport``, and real -``HookPoint`` modules on CUDA tensors. +(``available_task_slots``, ``reserve_one``, ``flush_and_wait`` and +``prepare_step``'s ``reserve``), so binding ``available_task_slots`` to the +wrong method passed every committed test while the model forward raised. +These drive the real composition: a ``MonitoringEngine``'s legacy ring, its +``RingTransport``, and real ``HookPoint`` modules on CUDA tensors. Needs CUDA and the full native backend. The ring has no host engine, so its consumer drops what it drains; delivering each payload byte for byte with @@ -99,3 +99,71 @@ def test_an_oversized_step_fires_every_hook_through_a_small_task_ring( transport.force_eager = False engine.close() + +def test_needs_eager_steps_that_fit_leave_no_reservation_behind(): + """_spec_needs_eager turns the safety net on for steps that fit, and + each hook then reserves its own entry. commit_step used to reserve the + whole step as well: each step kept its hook count of task entries after + the flush, and the second step's forward raised from reserve_one.""" + from types import SimpleNamespace + + from dmi.adapters.base import BackendAdapter, StepReservation + from dmi.adapters.types import StepContext + from dmi.hooks.point import HookPoint + from dmi.hooks.specs import ( + HOOK_TYPE_RESID_PRE, HookSpec, ModelShapeConfig) + + class DynamicShapeAdapter(BackendAdapter): + def detect_model_shape(self, model): + return self._cfg + + def detect_parallel_ranks(self): + return (0, 0, 0, 0) + + def is_pp_first(self): + return True + + def is_pp_last(self): + return True + + def build_step_context(self, *raw): + return None + + def on_capacity_exceeded(self, ctx): + pass + + def _spec_needs_eager(self, spec): + return True + + task_entries = 4 + engine = _legacy_engine(task_entries) + transport = engine._ring_transport + ring = engine._ring_engine + try: + adapter = DynamicShapeAdapter(engine, "eager-safety-net") + adapter._cfg = ModelShapeConfig( + hidden_dim=16, num_heads=4, num_kv_heads=4, head_dim=4, + dtype=torch.float16, vocab_size=32, intermediate_dim=0) + hooks = [HookPoint(), HookPoint()] + model = SimpleNamespace(get_hook_specs=lambda: [ + HookSpec(HOOK_TYPE_RESID_PRE, hook, layer) + for layer, hook in enumerate(hooks)]) + adapter.attach_model(model) + assert [spec.module for spec in adapter.active_specs] == hooks + ctx = StepContext( + model_id="eager-safety-net", flattened=False, req_ids=["0:0"], + token_ranges=[(0, 4)], dim0_offsets=[0], kv_offsets=[0], + batch=1, q_len=4, kv_dim=4) + value = torch.ones(1, 4, 16, dtype=torch.float16, device="cuda") + + for step in range(3): + assert adapter.commit_step(ctx) is StepReservation.RESERVED + assert transport.force_eager + for hook in hooks: + hook(value) + ring.flush_and_wait() + assert ring.available_task_slots() == task_entries, step + assert ring.available_capacity() == ring.payload_cap(), step + finally: + transport.force_eager = False + engine.close() diff --git a/tests/test_hf_capture_refusal.py b/tests/test_hf_capture_refusal.py index 6f08a05e3..45270d327 100644 --- a/tests/test_hf_capture_refusal.py +++ b/tests/test_hf_capture_refusal.py @@ -84,7 +84,7 @@ def payload_cap(self): def staging_cap(self): return 1 << 20 - def prepare_step(self, total_bytes, num_hooks): + def prepare_step(self, total_bytes, num_hooks, reserve=True): self.prepare_step_calls.append((total_bytes, num_hooks)) if self.prepare_step_error is not None: raise self.prepare_step_error diff --git a/tests/test_hook_point_eager_task_slots.py b/tests/test_hook_point_eager_task_slots.py index 5832db9f7..ef3d8bab9 100644 --- a/tests/test_hook_point_eager_task_slots.py +++ b/tests/test_hook_point_eager_task_slots.py @@ -12,17 +12,28 @@ new word and waits at that sequence forever while ``flush_and_wait`` reports success. Found by TLA+ model checking of the payload ring. -The fake below keeps a legacy ring's head/tail counters. ``HookPoint.forward`` +The fake below keeps a legacy ring's head/tail counters, and like the real +drain a flush releases only what producers published. ``HookPoint.forward`` only takes the eager branch for a CUDA tensor, so a CPU tensor subclass poses as one (as in tests/test_record_runtime.py) and ``dispatch_producer`` is -monkeypatched; the transport is a real ``RingTransport``. +monkeypatched to publish; the transport is a real ``RingTransport``. The real +native ring runs the same checks in tests/test_eager_safety_net_gpu.py. """ from __future__ import annotations +from types import SimpleNamespace + import pytest import torch -from dmi.hooks.specs import align_up_py +from dmi.adapters.base import BackendAdapter, StepReservation +from dmi.adapters.types import StepContext +from dmi.hooks.specs import ( + HOOK_TYPE_RESID_PRE, + HookSpec, + ModelShapeConfig, + align_up_py, +) pytestmark = pytest.mark.cpu @@ -34,13 +45,19 @@ def is_cuda(self) -> bool: class _LegacyRingEngine: - """Payload and task heads/tails of a legacy ring; flush drains both.""" + """Payload and task heads/tails of a legacy ring. + + Reservations advance the heads; a flush advances the tails by what + producers published since the last one, as the native drain does, so a + reservation nothing publishes stays outstanding. + """ def __init__(self, task_cap: int, payload_cap: int): self.task_cap = task_cap self.capacity = payload_cap self.task_head = self.task_tail = 0 self.payload_head = self.payload_tail = 0 + self.published_tasks = self.published_bytes = 0 self.events: list[str] = [] self.max_outstanding_tasks = 0 @@ -57,33 +74,71 @@ def available_capacity(self) -> int: return self.capacity - (self.payload_head - self.payload_tail) def available_task_slots(self) -> int: - return self.task_cap - (self.task_head - self.task_tail) + return max(0, self.task_cap - (self.task_head - self.task_tail)) - def reserve_one(self, nbytes: int) -> None: - self.events.append("reserve") - self.payload_head += align_up_py(nbytes, 16) - self.task_head += 1 + def outstanding(self) -> tuple[int, int]: + return (self.task_head - self.task_tail, + self.payload_head - self.payload_tail) + + def _reserve(self, nbytes: int, tasks: int) -> None: + self.payload_head += nbytes + self.task_head += tasks self.max_outstanding_tasks = max( self.max_outstanding_tasks, self.task_head - self.task_tail) + def reserve_one(self, nbytes: int) -> None: + self.events.append("reserve") + self._reserve(align_up_py(nbytes, 16), 1) + + def prepare_step(self, step_total_bytes: int, num_hooks: int, + reserve: bool = True) -> int: + """RingEnginePy::prepare_step's decision, on the same counters.""" + if step_total_bytes > self.capacity or num_hooks > self.task_cap: + self.flush_and_wait() + return int(StepReservation.OVERSIZED) + result = StepReservation.RESERVED + if (step_total_bytes > self.available_capacity() + or num_hooks > self.available_task_slots()): + self.flush_and_wait() + result = StepReservation.FLUSHED + if reserve: + self._reserve(step_total_bytes, num_hooks) + return int(result) + + def push_all_metas(self, *args) -> None: + pass + + def publish(self, nbytes: int) -> None: + self.published_tasks += 1 + self.published_bytes += align_up_py(nbytes, 16) + def flush_and_wait(self) -> None: self.events.append("flush") - self.task_tail = self.task_head - self.payload_tail = self.payload_head + self.task_tail += self.published_tasks + self.payload_tail += self.published_bytes + self.published_tasks = self.published_bytes = 0 -def _eager_hook(monkeypatch, engine, dispatched): - from dmi.hooks.point import HookPoint +def _activate(monkeypatch, engine, dispatched): from dmi.transport import ring as ring_transport from dmi.transport.ring import RingTransport transport = RingTransport(engine) - transport.force_eager = True monkeypatch.setattr(ring_transport, "_active_transport", transport) - monkeypatch.setattr( - "dmi.hooks.point.dispatch_producer", - lambda *args: dispatched.append(args), - ) + + def dispatch(payload, x, *args): + dispatched.append((payload, x, *args)) + engine.publish(x.nbytes) + + monkeypatch.setattr("dmi.hooks.point.dispatch_producer", dispatch) + return transport + + +def _eager_hook(monkeypatch, engine, dispatched): + from dmi.hooks.point import HookPoint + + transport = _activate(monkeypatch, engine, dispatched) + transport.force_eager = True hook = HookPoint() hook._ring_hook_type = 1 hook._ring_hook_id = 2 @@ -106,3 +161,71 @@ def test_a_step_with_more_hooks_than_task_entries_flushes_before_overflowing( assert engine.max_outstanding_tasks <= task_cap, ( "a reservation past task_cap overwrites an unread task slot") assert engine.events == ["reserve"] * task_cap + ["flush", "reserve"] + + +class _DynamicShapeAdapter(BackendAdapter): + """Every hook needs eager dispatch, as for a dynamic-shape backend.""" + + def detect_model_shape(self, model): + raise NotImplementedError + + def detect_parallel_ranks(self): + return (0, 0, 0, 0) + + def is_pp_first(self): + return True + + def is_pp_last(self): + return True + + def build_step_context(self, *raw): + return None + + def on_capacity_exceeded(self, ctx): + pass + + def _spec_needs_eager(self, spec): + return True + + +def test_a_needs_eager_step_that_fits_is_reserved_once(monkeypatch): + """_spec_needs_eager turns force_eager on even when the step fits, and + each hook then reserves its own entry in the safety net. Had + commit_step reserved the whole step too, every such step would hold + hook_count entries (and its bytes) that no producer publishes and no + flush frees, until the task check made a forward raise.""" + from dmi.hooks.point import HookPoint + + engine = _LegacyRingEngine(task_cap=4, payload_cap=1 << 20) + dispatched = [] + transport = _activate(monkeypatch, engine, dispatched) + adapter = _DynamicShapeAdapter( + SimpleNamespace(_ring_transport=transport, _ring_engine=engine), + "m") + cfg = ModelShapeConfig(hidden_dim=16, num_heads=4, num_kv_heads=4, + head_dim=4, dtype=torch.float16, vocab_size=32, + intermediate_dim=0) + adapter.model_cfg = cfg + transport.set_model_cfg(cfg) + hooks = [HookPoint(), HookPoint()] + specs = [HookSpec(HOOK_TYPE_RESID_PRE, hook, layer) + for layer, hook in enumerate(hooks)] + for spec in specs: + spec.module._ring_hook_type = spec.hook_type + spec.module._ring_hook_id = spec.layer_no + spec.module._ring_payload = transport._ring_payload + adapter.active_specs = transport._active_specs = specs + ctx = StepContext(model_id="m", flattened=False, req_ids=["0:0"], + token_ranges=[(0, 4)], dim0_offsets=[0], + kv_offsets=[0], batch=1, q_len=4, kv_dim=4) + value = torch.zeros(1, 4, 16, dtype=torch.float16).as_subclass( + _FakeCudaTensor) + + for step in range(3): + assert adapter.commit_step(ctx) is StepReservation.RESERVED + assert transport.force_eager + for hook in hooks: + hook(value) + engine.flush_and_wait() + assert engine.outstanding() == (0, 0), f"step {step} leaked" + assert len(dispatched) == 3 * len(hooks) diff --git a/tests/test_review_findings_runtime.py b/tests/test_review_findings_runtime.py index f65b9c962..8b0f126d4 100644 --- a/tests/test_review_findings_runtime.py +++ b/tests/test_review_findings_runtime.py @@ -61,7 +61,8 @@ class FakeRingEngine: def __init__(self) -> None: self.prepare_step_calls: list = [] - def prepare_step(self, total_bytes: int, n_hooks: int) -> int: + def prepare_step(self, total_bytes: int, n_hooks: int, + reserve: bool = True) -> int: self.prepare_step_calls.append((total_bytes, n_hooks)) return 0