Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion docs/integration-api-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
8 changes: 6 additions & 2 deletions native/csrc/bindings.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<py::gil_scoped_release>())
.def("reserve_record",
&ring_py::RingEnginePy::reserve_record,
Expand Down Expand Up @@ -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"))
Expand Down
83 changes: 63 additions & 20 deletions native/csrc/ring/ring_engine_py.cu
Original file line number Diff line number Diff line change
Expand Up @@ -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

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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
}

Expand All @@ -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;
}

Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
24 changes: 20 additions & 4 deletions native/csrc/ring/ring_engine_py.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down
8 changes: 7 additions & 1 deletion src/dmi/adapters/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = (
Expand Down
16 changes: 9 additions & 7 deletions src/dmi/configuration/estimate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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."
Expand Down
9 changes: 7 additions & 2 deletions src/dmi/hooks/point.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
3 changes: 2 additions & 1 deletion src/dmi/transport/ring.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading