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
Original file line number Diff line number Diff line change
Expand Up @@ -219,8 +219,14 @@ struct PTO2TaskDescriptor {
enum PTO2EarlyDispatchState : uint8_t {
PTO2_EARLY_DISPATCH_NONE = 0, // not pre-staged
PTO2_EARLY_DISPATCH_STAGING = 1, // Hook 1 claimed it; staging in progress
PTO2_EARLY_DISPATCH_STAGED = 2, // staged on a core, gated; staged_* fields valid
PTO2_EARLY_DISPATCH_DISPATCHED = 3 // routed via the normal dispatch path (no pre-stage)
PTO2_EARLY_DISPATCH_STAGED = 2, // reserved
PTO2_EARLY_DISPATCH_DISPATCHED = 3 // producers released; staged blocks may still be gated
};

enum PTO2EarlyDispatchLaunchState : uint8_t {
PTO2_EARLY_DISPATCH_LAUNCH_NONE = 0,
PTO2_EARLY_DISPATCH_LAUNCH_RINGING = 1,
PTO2_EARLY_DISPATCH_LAUNCH_COMPLETE = 2,
};

// A pre-staged consumer occupies one core per gated subtask block. WHICH cores
Expand Down Expand Up @@ -250,8 +256,8 @@ struct PTO2TaskPayload {
// 576), so sizeof and tensors[] are unchanged.
//
// Bitmask of global core_ids this consumer is pre-staged (gated) on. Set with
// atomic fetch_or by concurrent stagers; read by release. (Re)initialized in
// PTO2TaskPayload::init before the slot can be staged again.
// atomic fetch_or by concurrent stagers, then destructively split between the
// release and late-stager paths. (Re)initialized in PTO2TaskPayload::init.
std::atomic<uint64_t> staged_core_mask[PTO2_EARLY_DISPATCH_CORE_MASK_WORDS]{};
// Early-dispatch CANDIDATE detection (event-driven, dual of fanin_refcount):
// seeded at wiring with producers already complete, then a flagged producer
Expand All @@ -270,11 +276,14 @@ struct PTO2TaskPayload {
// 3=DISPATCHED (2=STAGED is unused now). STAGING is the STABLE gated state —
// many threads stage blocks concurrently while it holds, each claiming a block
// via the atomic next_block_idx and OR-ing its cores into staged_core_mask.
// Release does STAGING->DISPATCHED then rings the mask; a thread that stages a
// block AFTER release flipped DISPATCHED rings that block's doorbell itself
// (self-ring), so no doorbell is ever missed.
// Release does STAGING->DISPATCHED and claims the current mask; a thread that
// stages a block after that flip claims and rings only its remaining bits.
std::atomic<uint8_t> early_dispatch_state{0};
std::atomic<uint8_t> dispatch_propagated{0}; // PRODUCER side: once-guard for fanout propagation
// The release owner publishes COMPLETE only after all doorbells it claimed
// are visible. Combined with published_block_count, this keeps fanout
// private until release-owned and late-owned blocks have both launched.
std::atomic<uint8_t> early_dispatch_launch_state{PTO2_EARLY_DISPATCH_LAUNCH_NONE};
// === Cache lines 9-72 (4096B) — tensors (alignas(64) forces alignment) ===
Tensor tensors[MAX_TENSOR_ARGS];
// === Cache lines 73-74 (128B) — scalars ===
Expand Down Expand Up @@ -350,14 +359,15 @@ struct PTO2TaskPayload {
// one of ITS producers is flagged (propagate_dispatch_fanin bumps
// dispatch_fanin and may CAS early_dispatch_state on any consumer, independent of the
// consumer's own hint). So they MUST be zeroed here unconditionally.
// published_block_count and dispatch_propagated are producer-side, but share
// this same per-submit lifetime and are reset here too.
// Publication, propagation, and launch fields share this same
// per-submit lifetime and are reset here too.
early_dispatch_state.store(PTO2_EARLY_DISPATCH_NONE, std::memory_order_relaxed);
for (int w = 0; w < PTO2_EARLY_DISPATCH_CORE_MASK_WORDS; w++)
staged_core_mask[w].store(0, std::memory_order_relaxed);
dispatch_fanin.store(0, std::memory_order_relaxed);
dispatch_propagated.store(0, std::memory_order_relaxed);
published_block_count.store(0, std::memory_order_relaxed);
early_dispatch_launch_state.store(PTO2_EARLY_DISPATCH_LAUNCH_NONE, std::memory_order_relaxed);
}
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
#include <atomic>

#include "common/core_type.h"
#include "common/memory_barrier.h"
#include "utils/device_arena.h"
#include "aicpu/platform_regs.h" // get_reg_ptr / RegId for the early-dispatch doorbell
#include "pto_async_wait.h"
Expand Down Expand Up @@ -657,8 +658,8 @@ struct PTO2SchedulerState {
// PARTIALLY pre-staged — the gated blocks are released by the doorbells rung
// here, and the remaining (next_block_idx .. logical_block_num) blocks
// dispatch normally off the ready queue. Lock-free claim shared with Hook 1
// (the stager): CAS NONE->DISPATCHED wins => not pre-staged; lose => STAGED
// (spin past the brief STAGING window so the mask is visible), then ring.
// (the stager): CAS NONE->DISPATCHED wins => not pre-staged; otherwise flip
// STAGING->DISPATCHED and destructively claim the published doorbell bits.

// Per-core early-dispatch doorbell table. Hook 1 records each gated core's
// (reg_addr, dispatch token) here at stage time; the completion-path release
Expand Down Expand Up @@ -698,6 +699,42 @@ struct PTO2SchedulerState {
*dmb = (tk << 32) | tk; // 64-bit STR: high=low=token releases the gated AICore
}

inline void ring_staged_doorbell_bits(int word, uint64_t bits) {
while (bits != 0) {
int core_id = word * 64 + __builtin_ctzll(bits);
bits &= bits - 1;
ring_one_doorbell(
early_dispatch_doorbell_table[core_id].addr, early_dispatch_doorbell_table[core_id].token
);
}
}

static inline uint64_t claim_all_staged_doorbell_bits(std::atomic<uint64_t> &mask) {
return mask.exchange(0, std::memory_order_seq_cst);
}

static inline uint64_t claim_late_staged_doorbell_bits(std::atomic<uint64_t> &mask, uint64_t candidates) {
return mask.fetch_and(~candidates, std::memory_order_seq_cst) & candidates;
}

static inline bool should_gate_early_dispatch(bool force_gate, uint8_t early_dispatch_state) {
return force_gate || early_dispatch_state == PTO2_EARLY_DISPATCH_STAGING;
}

static inline bool
ring_claimed_local_doorbell(uint64_t claimed_word, int core_id, uint64_t reg_addr, uint32_t token) {
if ((claimed_word & (1ULL << (core_id & 63))) == 0) return false;
ring_one_doorbell(reg_addr, token);
return true;
}

static inline bool try_claim_early_dispatch_launch(PTO2TaskPayload &payload) {
uint8_t expected = PTO2_EARLY_DISPATCH_LAUNCH_NONE;
return payload.early_dispatch_launch_state.compare_exchange_strong(
expected, PTO2_EARLY_DISPATCH_LAUNCH_RINGING, std::memory_order_seq_cst, std::memory_order_seq_cst
);
}

inline void record_published_blocks(PTO2TaskSlotState &slot_state, int32_t count) {
if (count <= 0 || !slot_state.allow_early_resolve) return;
slot_state.payload->published_block_count.fetch_add(static_cast<int16_t>(count), std::memory_order_seq_cst);
Expand All @@ -716,6 +753,10 @@ struct PTO2SchedulerState {
void propagate_dispatch_fanin(PTO2TaskSlotState &p) {
if (!p.allow_early_resolve) return; // only codegen-flagged (direct) producers propagate
if (p.payload->published_block_count.load(std::memory_order_seq_cst) < p.logical_block_num) return;
if (p.payload->early_dispatch_state.load(std::memory_order_seq_cst) == PTO2_EARLY_DISPATCH_STAGING) return;
if (p.payload->early_dispatch_launch_state.load(std::memory_order_seq_cst) ==
PTO2_EARLY_DISPATCH_LAUNCH_RINGING)
return;
if (p.payload->dispatch_propagated.exchange(1, std::memory_order_acq_rel) != 0)
return; // already propagated once
p.lock_fanout();
Expand Down Expand Up @@ -767,28 +808,25 @@ struct PTO2SchedulerState {
)) {
return false;
}
// Staged (STAGING). Flip STAGING->DISPATCHED, THEN read the mask. seq_cst
// gives a total order with the concurrent stagers, each of which OR-s its
// core into the mask and THEN loads early_dispatch_state: a stager whose bit lands
// before this CAS is read here and rung; a stager whose bit lands after
// sees DISPATCHED and rings that core itself (self-ring in
// stage_consumer_blocks). Either way every gated core's doorbell fires once
// (a double-ring is harmless — the AICore already matched). This replaces
// the old transient-STAGING spin: STAGING is now the stable gated state.
// route_ready_once admits one caller, but keep the helper defensive: a
// duplicate that observes or loses an in-progress launch must never
// route the same partial task a second time.
if (expect != PTO2_EARLY_DISPATCH_STAGING || !try_claim_early_dispatch_launch(*slot_state.payload)) return true;
expect = PTO2_EARLY_DISPATCH_STAGING;
slot_state.payload->early_dispatch_state.compare_exchange_strong(
expect, PTO2_EARLY_DISPATCH_DISPATCHED, std::memory_order_seq_cst, std::memory_order_seq_cst
);
// Destructively claim every published bit. A stager racing this pass
// can claim only bits that land after the exchange, so each gated core
// has exactly one doorbell writer.
for (int w = 0; w < PTO2_EARLY_DISPATCH_CORE_MASK_WORDS; w++) {
uint64_t bits = slot_state.payload->staged_core_mask[w].load(std::memory_order_seq_cst);
while (bits != 0) {
int core_id = w * 64 + __builtin_ctzll(bits);
bits &= bits - 1;
ring_one_doorbell(
early_dispatch_doorbell_table[core_id].addr, early_dispatch_doorbell_table[core_id].token
);
}
uint64_t owned = claim_all_staged_doorbell_bits(slot_state.payload->staged_core_mask[w]);
ring_staged_doorbell_bits(w, owned);
}
wmb();
slot_state.payload->early_dispatch_launch_state.store(
PTO2_EARLY_DISPATCH_LAUNCH_COMPLETE, std::memory_order_seq_cst
);
// This pre-staged consumer was just released by its doorbell — it starts
// running NOW, so propagate dispatch_fanin to ITS consumers (only if it is
// itself codegen-flagged; the gate inside no-ops otherwise). Defer it via the
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -230,7 +230,8 @@ class SchedulerContext {
);

void build_payload(
PTO2DispatchPayload &dispatch_payload, PTO2TaskSlotState &slot_state, PTO2SubtaskSlot subslot, int32_t block_idx
PTO2DispatchPayload &dispatch_payload, PTO2TaskSlotState &slot_state, PTO2SubtaskSlot subslot,
int32_t block_idx, bool force_gate
);

// Batched-dispatch primitives. prepare_* builds the payload and per-core
Expand All @@ -251,7 +252,7 @@ class SchedulerContext {

PublishHandle prepare_subtask_to_core(
int32_t thread_idx, int32_t core_offset, PTO2TaskSlotState &slot_state, PTO2SubtaskSlot subslot,
bool to_pending, int32_t block_idx
bool to_pending, int32_t block_idx, bool force_gate
);

inline void publish_subtask_to_core(const PublishHandle &h, uint64_t dispatch_ts) {
Expand Down Expand Up @@ -298,7 +299,7 @@ class SchedulerContext {
// caller-supplied handles buffer. Returns the number of handles written.
int prepare_block_for_dispatch(
int32_t thread_idx, int32_t core_offset, PTO2TaskSlotState &slot_state, PTO2ResourceShape shape,
bool to_pending, int32_t block_idx, PublishHandle *out_handles
bool to_pending, int32_t block_idx, PublishHandle *out_handles, bool force_gate = false
);

void dispatch_shape(
Expand All @@ -311,7 +312,7 @@ class SchedulerContext {
// is queued) and sets made_progress / try_pushed when it stages, so the caller
// is a single unconditional call like normal dispatch. After normal dispatch
// leaves idle cores spare, pre-stage the consumers of any RUNNING flagged
// producer onto those cores with not_ready=1 (gated). Touches no dependency
// producer onto those cores with a non-zero src_payload (gated). Touches no dependency
// state — the task is released by the doorbell at its normal ready-pop (Hook 2).
int32_t try_early_dispatch(
int32_t thread_idx, CoreTracker &tracker, bool pmu_active, bool &made_progress, bool &try_pushed
Expand Down
Loading
Loading