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 @@ -912,7 +912,7 @@ void SchedulerContext::deinit() {
drain_state_.sync_start_pending.store(0, std::memory_order_release);
drain_state_.drain_worker_elected.store(0, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_release);
drain_state_.pending_task = nullptr;
drain_state_.pending_task.store(nullptr, std::memory_order_release);

// Reset task counters and orchestrator state
completed_tasks_.store(0, std::memory_order_release);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,7 @@ bool SchedulerContext::enter_drain_mode(PTO2TaskSlotState *slot_state, int32_t b
return false; // Another thread already holds the drain slot.
}
// We own the drain slot. Store the task and reset election flag before making it visible.
drain_state_.pending_task = slot_state;
drain_state_.pending_task.store(slot_state, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_relaxed);
drain_state_.drain_worker_elected.store(0, std::memory_order_relaxed);
// Release store: all stores above are now visible to any thread that
Expand All @@ -344,7 +344,7 @@ int32_t SchedulerContext::count_global_available(PTO2ResourceShape shape) {
// Called only when global resources >= block_num, so one pass always suffices.
// All other threads are spinning -- the drain worker has exclusive tracker access.
void SchedulerContext::drain_worker_dispatch(int32_t block_num) {
PTO2TaskSlotState *slot_state = drain_state_.pending_task;
PTO2TaskSlotState *slot_state = drain_state_.pending_task.load(std::memory_order_acquire);
if (!slot_state) {
drain_state_.sync_start_pending.store(0, std::memory_order_release);
return;
Expand All @@ -363,7 +363,7 @@ void SchedulerContext::drain_worker_dispatch(int32_t block_num) {
// Release fence ensures tracker mutations are visible to threads that
// acquire-load sync_start_pending == 0 and resume normal operation.
std::atomic_thread_fence(std::memory_order_release);
drain_state_.pending_task = nullptr;
drain_state_.pending_task.store(nullptr, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_relaxed);
drain_state_.drain_worker_elected.store(0, std::memory_order_relaxed);
drain_state_.sync_start_pending.store(0, std::memory_order_release);
Expand Down Expand Up @@ -419,7 +419,17 @@ void SchedulerContext::handle_drain_mode(int32_t thread_idx) {
}

// Elected: check if global resources are sufficient.
PTO2TaskSlotState *slot_state = drain_state_.pending_task;
PTO2TaskSlotState *slot_state = drain_state_.pending_task.load(std::memory_order_acquire);
if (slot_state == nullptr) {
// pending_task is observed null only when a concurrent drain completion
// already cleared it (drain_worker_dispatch nulls it before reopening the
// gate). That drain is done and this is a stale-elected thread, so just
// release the election lock and return. Do NOT clear drain_ack_mask or
// sync_start_pending: a *new* drain run may already be active and
// accumulating acks, and zeroing them would corrupt it into a hang.
drain_state_.drain_worker_elected.store(0, std::memory_order_release);
return;
}
Comment thread
ChaoWao marked this conversation as resolved.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
PTO2ResourceShape shape = slot_state->active_mask.to_shape();
int32_t available = count_global_available(shape);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -383,7 +383,7 @@ struct alignas(64) SyncStartDrainState {
std::atomic<int32_t> sync_start_pending{0}; // 0=normal; -1=initializing; >0=active (value=block_num)
std::atomic<int32_t> drain_worker_elected{0}; // 0=none; >0: elected thread's (thread_idx+1)
std::atomic<uint32_t> drain_ack_mask{0}; // bit per thread; all-set = all threads reached ack barrier
PTO2TaskSlotState *pending_task{nullptr}; // held task (not re-queued)
std::atomic<PTO2TaskSlotState *> pending_task{nullptr}; // held task (not re-queued)
int32_t _pad[10];
};
static_assert(sizeof(SyncStartDrainState) == 64);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -920,7 +920,7 @@ void SchedulerContext::deinit() {
drain_state_.sync_start_pending.store(0, std::memory_order_release);
drain_state_.drain_worker_elected.store(0, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_release);
drain_state_.pending_task = nullptr;
drain_state_.pending_task.store(nullptr, std::memory_order_release);

// Reset task counters and orchestrator state
completed_tasks_.store(0, std::memory_order_release);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -341,7 +341,7 @@ bool SchedulerContext::enter_drain_mode(PTO2TaskSlotState *slot_state, int32_t b
return false; // Another thread already holds the drain slot.
}
// We own the drain slot. Store the task and reset election flag before making it visible.
drain_state_.pending_task = slot_state;
drain_state_.pending_task.store(slot_state, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_relaxed);
drain_state_.drain_worker_elected.store(0, std::memory_order_relaxed);
// Release store: all stores above are now visible to any thread that
Expand All @@ -363,7 +363,7 @@ int32_t SchedulerContext::count_global_available(PTO2ResourceShape shape) {
// Called only when global resources >= block_num, so one pass always suffices.
// All other threads are spinning -- the drain worker has exclusive tracker access.
void SchedulerContext::drain_worker_dispatch(Runtime *runtime, int32_t block_num) {
PTO2TaskSlotState *slot_state = drain_state_.pending_task;
PTO2TaskSlotState *slot_state = drain_state_.pending_task.load(std::memory_order_acquire);
if (!slot_state) {
drain_state_.sync_start_pending.store(0, std::memory_order_release);
return;
Expand All @@ -382,7 +382,7 @@ void SchedulerContext::drain_worker_dispatch(Runtime *runtime, int32_t block_num
// Release fence ensures tracker mutations are visible to threads that
// acquire-load sync_start_pending == 0 and resume normal operation.
std::atomic_thread_fence(std::memory_order_release);
drain_state_.pending_task = nullptr;
drain_state_.pending_task.store(nullptr, std::memory_order_release);
drain_state_.drain_ack_mask.store(0, std::memory_order_relaxed);
drain_state_.drain_worker_elected.store(0, std::memory_order_relaxed);
drain_state_.sync_start_pending.store(0, std::memory_order_release);
Expand Down Expand Up @@ -438,7 +438,17 @@ void SchedulerContext::handle_drain_mode(Runtime *runtime, int32_t thread_idx) {
}

// Elected: check if global resources are sufficient.
PTO2TaskSlotState *slot_state = drain_state_.pending_task;
PTO2TaskSlotState *slot_state = drain_state_.pending_task.load(std::memory_order_acquire);
if (slot_state == nullptr) {
// pending_task is observed null only when a concurrent drain completion
// already cleared it (drain_worker_dispatch nulls it before reopening the
// gate). That drain is done and this is a stale-elected thread, so just
// release the election lock and return. Do NOT clear drain_ack_mask or
// sync_start_pending: a *new* drain run may already be active and
// accumulating acks, and zeroing them would corrupt it into a hang.
drain_state_.drain_worker_elected.store(0, std::memory_order_release);
return;
}
Comment thread
ChaoWao marked this conversation as resolved.
PTO2ResourceShape shape = slot_state->active_mask.to_shape();
int32_t available = count_global_available(shape);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -385,7 +385,7 @@ struct alignas(64) SyncStartDrainState {
std::atomic<int32_t> sync_start_pending{0}; // 0=normal; -1=initializing; >0=active (value=block_num)
std::atomic<int32_t> drain_worker_elected{0}; // 0=none; >0: elected thread's (thread_idx+1)
std::atomic<uint32_t> drain_ack_mask{0}; // bit per thread; all-set = all threads reached ack barrier
PTO2TaskSlotState *pending_task{nullptr}; // held task (not re-queued)
std::atomic<PTO2TaskSlotState *> pending_task{nullptr}; // held task (not re-queued)
int32_t _pad[10];
};
static_assert(sizeof(SyncStartDrainState) == 64);
Expand Down
Loading