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 feac37904..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, @@ -829,11 +830,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..62744726d 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 // --------------------------------------------------------------------------- @@ -565,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"); @@ -620,13 +634,11 @@ 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); + if (reserve) drain.reserve(step_total_bytes, num_hooks); return STEP_RING_OK; // fast path -- no CUDA or thread interaction } @@ -635,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; } @@ -891,22 +903,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 @@ -917,17 +936,41 @@ 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 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 // 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"); - impl_->engine.drain_thread().reserve( - ring::align_up(nbytes, ring::PAYLOAD_ALIGN), 1); + 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: 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"); + } + 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 a465f7c33..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 @@ -268,13 +274,23 @@ 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. Saturates at 0 like available_capacity(). + 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. 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 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/src/dmi/configuration/estimate.py b/src/dmi/configuration/estimate.py index 1e55b8259..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,10 +816,11 @@ 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 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..20c7efd13 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 @@ -226,6 +228,194 @@ 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); +} + +// 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 +// 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 { @@ -1171,6 +1361,9 @@ 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_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(); test_prefix_force_flush(); 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 new file mode 100644 index 000000000..0535280fd --- /dev/null +++ b/tests/test_eager_safety_net_gpu.py @@ -0,0 +1,169 @@ +"""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`` 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 +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() + + +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_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..ef3d8bab9 --- /dev/null +++ b/tests/test_hook_point_eager_task_slots.py @@ -0,0 +1,231 @@ +"""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, 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 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.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 + + +class _FakeCudaTensor(torch.Tensor): + @property + def is_cuda(self) -> bool: + return True + + +class _LegacyRingEngine: + """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 + + 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 max(0, self.task_cap - (self.task_head - self.task_tail)) + + 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.published_tasks + self.payload_tail += self.published_bytes + self.published_tasks = self.published_bytes = 0 + + +def _activate(monkeypatch, engine, dispatched): + from dmi.transport import ring as ring_transport + from dmi.transport.ring import RingTransport + + transport = RingTransport(engine) + monkeypatch.setattr(ring_transport, "_active_transport", transport) + + 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 + 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"] + + +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_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 diff --git a/tests/test_review_findings_runtime.py b/tests/test_review_findings_runtime.py index a0e1a1733..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 @@ -858,6 +859,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):