Skip to content

Commit c53c2a9

Browse files
Address latest Copilot review on PR #193.
Correct the scheduler_data comment (atomic stores are not a data race), cache Pool::current() on the submit hot path, and reflow scheduler.md with semantic line breaks. Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 4591e31 commit c53c2a9

3 files changed

Lines changed: 57 additions & 29 deletions

File tree

‎docs/explanation/scheduler.md‎

Lines changed: 50 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -2,18 +2,21 @@
22

33
This page explains how NUClear's task scheduler works internally — the lock-free queues, thread pools, group tokens, and the path from `emit()` to a running reaction callback.
44

5-
For the user-facing view of pools, priorities, groups, and idle tasks, see [Threading Model](threading.md). For DSL usage, see the [Scheduling](../reference/dsl/index.md) reference words.
5+
For the user-facing view of pools, priorities, groups, and idle tasks, see [Threading Model](threading.md).
6+
For DSL usage, see the [Scheduling](../reference/dsl/index.md) reference words.
67

78
## Role in the system
89

9-
Every reaction execution is a **task** (`ReactionTask`) submitted to the scheduler. The `PowerPlant` owns a single `Scheduler` instance and forwards all work to it:
10+
Every reaction execution is a **task** (`ReactionTask`) submitted to the scheduler.
11+
The `PowerPlant` owns a single `Scheduler` instance and forwards all work to it:
1012

1113
1. A trigger (message emit, timer, IO event, etc.) creates a `ReactionTask`.
1214
1. `PowerPlant::submit()` calls `Scheduler::submit()`.
1315
1. The scheduler resolves the target **pool**, acquires any required **group** tokens, and enqueues the task.
1416
1. A pool worker dequeues the task, runs the callback, and releases group locks when the callback returns.
1517

16-
`PowerPlant::start()` calls `Scheduler::start()`, which starts worker pools and then blocks the calling thread in the **MainThread** pool until shutdown. `PowerPlant::shutdown()` emits the shutdown event and calls `Scheduler::stop()`.
18+
`PowerPlant::start()` calls `Scheduler::start()`, which starts worker pools and then blocks the calling thread in the **MainThread** pool until shutdown.
19+
`PowerPlant::shutdown()` emits the shutdown event and calls `Scheduler::stop()`.
1720

1821
```mermaid
1922
flowchart LR
@@ -46,7 +49,8 @@ flowchart LR
4649

4750
### Scheduler
4851

49-
The scheduler is the central coordinator. It:
52+
The scheduler is the central coordinator.
53+
It:
5054

5155
- **Owns pools** — lazily created from `ThreadPoolDescriptor` values (default pool, `MainThread`, custom `Pool<T>`, etc.).
5256
- **Owns groups** — lazily created from `GroupDescriptor` values (`Sync<T>`, `Group<T>`, etc.).
@@ -67,11 +71,13 @@ Each pool is a set of worker threads (or a single thread for `MainThread`) plus:
6771

6872
Workers loop in `Pool::run()`: dequeue a task, call `ReactionTask::run()`, repeat until shutdown.
6973

70-
The default pool's thread count comes from `Configuration::default_pool_concurrency` (typically hardware concurrency). Other pools use the `concurrency` value from their descriptor.
74+
The default pool's thread count comes from `Configuration::default_pool_concurrency` (typically hardware concurrency).
75+
Other pools use the `concurrency` value from their descriptor.
7176

7277
### Group
7378

74-
A group limits how many tasks sharing the same descriptor may run concurrently. `Sync<T>` is a group with concurrency 1.
79+
A group limits how many tasks sharing the same descriptor may run concurrently.
80+
`Sync<T>` is a group with concurrency 1.
7581

7682
Groups maintain:
7783

@@ -117,31 +123,38 @@ sequenceDiagram
117123

118124
### Pool resolution cache
119125

120-
The first submit for a reaction calls `get_pool()` under `pools_mutex`. The resulting `Pool*` is stored in `Reaction::scheduler_data` — a plain `std::atomic<Pool*>` rather than `atomic<shared_ptr>` to avoid libstdc++'s hashed mutex pool for atomic shared pointers, which would contend on hot paths.
126+
The first submit for a reaction calls `get_pool()` under `pools_mutex`.
127+
The resulting `Pool*` is stored in `Reaction::scheduler_data` — a plain `std::atomic<Pool*>` rather than `atomic<shared_ptr>` to avoid libstdc++'s hashed mutex pool for atomic shared pointers, which would contend on hot paths.
121128

122-
Subsequent submits load the cached pointer with acquire semantics. Concurrent first submits may both resolve the pool; they store the same pointer, so the race is benign.
129+
Subsequent submits load the cached pointer with acquire semantics.
130+
Concurrent first submits may both resolve the pool; they store the same pointer, so the race is benign.
123131

124132
### Inline execution
125133

126-
If a reaction is bound with `Inline` and belongs to a single group, the scheduler tries to acquire a group token and run the callback on the submitting thread without enqueueing. This avoids queue overhead for synchronous emit paths.
134+
If a reaction is bound with `Inline` and belongs to a single group, the scheduler tries to acquire a group token and run the callback on the submitting thread without enqueueing.
135+
This avoids queue overhead for synchronous emit paths.
127136

128137
## Thread pools and queue selection
129138

130-
Each pool holds an array of five `Queue<Task>` instances — one per priority bucket. At construction time the pool chooses the concrete queue type:
139+
Each pool holds an array of five `Queue<Task>` instances — one per priority bucket.
140+
At construction time the pool chooses the concrete queue type:
131141

132142
| Pool kind | Queue type | Why |
133143
| ---------------------------------------------------------- | ------------------ | -------------------------------------------------------------------------------------------------- |
134144
| Default pool (`Pool<>`) | `TaskQueue` (MPMC) | Concurrency may differ from the descriptor's nominal value; multiple workers dequeue concurrently. |
135145
| `MainThread`, Trace pool, any pool with `concurrency == 1` | `MPSCQueue` (MPSC) | Exactly one consumer; simpler and cheaper than MPMC. |
136146
| Custom pools with `concurrency > 1` | `TaskQueue` (MPMC) | Multiple workers compete for tasks. |
137147

138-
The virtual `Queue` interface lets `Pool` store both implementations in one `std::array` without templating the entire pool. The virtual call cost is negligible compared to the atomic operations inside enqueue and dequeue.
148+
The virtual `Queue` interface lets `Pool` store both implementations in one `std::array` without templating the entire pool.
149+
The virtual call cost is negligible compared to the atomic operations inside enqueue and dequeue.
139150

140-
Workers identify themselves via a thread-local `Pool::current_pool` pointer, set when `run()` starts. `Pool::current()` returns a `shared_ptr` to the active pool, or `nullptr` off-scheduler threads.
151+
Workers identify themselves via a thread-local `Pool::current_pool` pointer, set when `run()` starts.
152+
`Pool::current()` returns a `shared_ptr` to the active pool, or `nullptr` off-scheduler threads.
141153

142154
## Priority buckets
143155

144-
Tasks are not kept in one monolithic priority queue. Instead, each pool has **five fixed buckets** scanned from highest to lowest priority:
156+
Tasks are not kept in one monolithic priority queue.
157+
Instead, each pool has **five fixed buckets** scanned from highest to lowest priority:
145158

146159
| Bucket | Priority range | DSL level |
147160
| -------- | -------------- | ---------------------------- |
@@ -151,13 +164,18 @@ Tasks are not kept in one monolithic priority queue. Instead, each pool has **fi
151164
| LOW | ≥ 250 | `Priority::LOW` |
152165
| IDLE | < 250 | `Priority::IDLE` |
153166

154-
`Pool::try_dequeue_task()` walks buckets 0→4 and returns the first available task. Within a bucket, ordering is **FIFO** (per-producer FIFO in the MPMC queue; strict FIFO in MPSC). Priority therefore dominates bucket order; tie-breaking within a bucket follows enqueue order, not reaction ID.
167+
`Pool::try_dequeue_task()` walks buckets 0→4 and returns the first available task.
168+
Within a bucket, ordering is **FIFO** (per-producer FIFO in the MPMC queue; strict FIFO in MPSC).
169+
Priority therefore dominates bucket order; tie-breaking within a bucket follows enqueue order, not reaction ID.
155170

156-
Priority affects **queuing order only**. Running tasks are never preempted.
171+
Priority affects **queuing order only**.
172+
Running tasks are never preempted.
157173

158174
## Lock-free queues
159175

160-
Both queue implementations use a **block-based** design: fixed-size blocks of 64 slots linked in a list. Producers claim slots with `write.fetch_add(1)`, construct the payload in place, then set a `committed` flag. Consumers read committed slots and advance head/tail as blocks drain.
176+
Both queue implementations use a **block-based** design: fixed-size blocks of 64 slots linked in a list.
177+
Producers claim slots with `write.fetch_add(1)`, construct the payload in place, then set a `committed` flag.
178+
Consumers read committed slots and advance head/tail as blocks drain.
161179

162180
### TaskQueue (MPMC)
163181

@@ -173,17 +191,20 @@ Cross-producer ordering is not guaranteed; per-producer FIFO is preserved.
173191

174192
Used for single-consumer pools (`MainThread`, concurrency-1 custom pools).
175193

176-
The producer side matches `TaskQueue`. The consumer side is simpler: a plain (non-atomic) read index, no CAS on dequeue, and immediate block retirement to the graveyard when advancing.
194+
The producer side matches `TaskQueue`.
195+
The consumer side is simpler: a plain (non-atomic) read index, no CAS on dequeue, and immediate block retirement to the graveyard when advancing.
177196

178-
`try_dequeue` must only be called from the designated consumer thread. Force shutdown from another thread delegates queue draining to that consumer via `discard_queues_requested`.
197+
`try_dequeue` must only be called from the designated consumer thread.
198+
Force shutdown from another thread delegates queue draining to that consumer via `discard_queues_requested`.
179199

180200
### Shared block helpers
181201

182202
`queue/detail/block_ops.hpp` provides `link_next_block` and `retire_block` shared by both queues.
183203

184204
### Lock-free vs wait-free
185205

186-
The queues are **lock-free** at the algorithm level: no mutexes, and the system makes progress under contention. They are **not wait-free end-to-end**:
206+
The queues are **lock-free** at the algorithm level: no mutexes, and the system makes progress under contention.
207+
They are **not wait-free end-to-end**:
187208

188209
- Block allocation uses `operator new`.
189210
- Overflow paths use CAS loops on list pointers.
@@ -195,25 +216,29 @@ The hot-path slot claim via `fetch_add` is wait-free within a non-full block.
195216

196217
### Single-group fast path
197218

198-
Most reactions belong to at most one group (including `Sync<T>`). For these, `Group::try_submit()`:
219+
Most reactions belong to at most one group (including `Sync<T>`).
220+
For these, `Group::try_submit()`:
199221

200222
1. Tries to decrement `tokens` with a compare-exchange.
201223
1. On success, submits to the pool immediately with a `RunningLock` that calls `release_token()` on destruction.
202224
1. On failure, **parks** the task in priority-ordered waiter buckets via `park_publish()` / `park_reconcile()`.
203225

204-
The token counter can go **negative** when waiters reserve slots they have not yet consumed. This signed counter, combined with per-waiter **arbiter slots** (`atomic<bool>`), ensures no lost wakeups and exact accounting when multiple waiters race with draining threads.
226+
The token counter can go **negative** when waiters reserve slots they have not yet consumed.
227+
This signed counter, combined with per-waiter **arbiter slots** (`atomic<bool>`), ensures no lost wakeups and exact accounting when multiple waiters race with draining threads.
205228

206229
When a running task finishes, `release_token()` increments `tokens` and drains at most one parked waiter into the pool — keeping running count bounded by the group's concurrency.
207230

208231
### Multi-group slow path
209232

210-
Tasks bound to multiple groups (`Sync<A>` and `Sync<B>`, etc.) use `CombinedLock`: each group gets a `GroupLock` backed by a mutex-protected sorted queue. `slow_pending` on each group prevents fast-path submitters from jumping ahead of older multi-group waiters.
233+
Tasks bound to multiple groups (`Sync<A>` and `Sync<B>`, etc.) use `CombinedLock`: each group gets a `GroupLock` backed by a mutex-protected sorted queue.
234+
`slow_pending` on each group prevents fast-path submitters from jumping ahead of older multi-group waiters.
211235

212236
When a `GroupLock` is released, the group may drain a fast-path waiter even if slow-path waiters exist, if the pre-release token count indicates a committed fast waiter is owed a slot — avoiding deadlocks between fast and slow paths.
213237

214238
### External waiters
215239

216-
When a task is parked in a group's wait buckets (not yet in the pool queue), the destination pool must not go idle as if it had no work. `Pool::register_external_waiter()` increments `external_waiters`, keeping workers alive until the parked task is drained or the registration is destroyed.
240+
When a task is parked in a group's wait buckets (not yet in the pool queue), the destination pool must not go idle as if it had no work.
241+
`Pool::register_external_waiter()` increments `external_waiters`, keeping workers alive until the parked task is drained or the registration is destroyed.
217242

218243
If idle reactions are registered for that pool (or globally), a `pending_idle` latch ensures one idle epoch fires before the next dequeue — preserving the invariant that parking a non-runnable task triggers idle detection, even if the worker is preempted and a runnable task arrives in the queue before the worker resumes.
219244

@@ -245,7 +270,8 @@ When a pool worker finds no runnable task:
245270
| `FINAL` | Used after the main thread exits `start()`; even persistent pools stop once their queues empty. |
246271
| `FORCE` | Clears queues and wakes all threads; used for forced test timeouts. MPSC pools require the consumer thread to perform the drain. |
247272

248-
`Scheduler::start()` starts worker pools first, then blocks in `MainThread::start()`. When the main thread pool exits (after shutdown), pools are stopped in order — non-persistent pools before persistent ones — then joined.
273+
`Scheduler::start()` starts worker pools first, then blocks in `MainThread::start()`.
274+
When the main thread pool exits (after shutdown), pools are stopped in order — non-persistent pools before persistent ones — then joined.
249275

250276
Persistent pools (`ThreadPoolDescriptor::persistent`) continue accepting tasks during a normal shutdown so networking or logging reactors can finish in-flight work.
251277

‎src/threading/Reaction.hpp‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -148,9 +148,9 @@ namespace threading {
148148
/// outlive scheduler-side resources because PowerPlant tears reactors down before the
149149
/// scheduler. The first submit resolves the pool and stores it here (release); later submits
150150
/// just load it (acquire). The write is a plain store rather than a CAS: every writer
151-
/// resolves the same pool for a given reaction, so concurrent first-submit stores are still
152-
/// a data race (though they publish identical values) and a reader either sees nullptr
153-
/// (and re-resolves) or the one pointer.
151+
/// resolves the same pool for a given reaction. Concurrent stores/loads are well-defined
152+
/// on the atomic (no data race); a reader either sees nullptr (and re-resolves) or the
153+
/// cached pointer.
154154
std::atomic<scheduler::Pool*> scheduler_data{nullptr};
155155
friend class scheduler::Scheduler; /// Let the scheduler mess with reaction objects
156156
};

‎src/threading/scheduler/Scheduler.cpp‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -203,7 +203,8 @@ namespace threading {
203203
auto lock = std::make_unique<CombinedLock>();
204204
for (const auto& desc : descs) {
205205
lock->add(get_group(desc)->lock(task_id, priority, [pool] {
206-
const bool current_pool_idle = Pool::current() != nullptr && Pool::current()->is_idle();
206+
const auto current_pool = Pool::current();
207+
const bool current_pool_idle = current_pool != nullptr && current_pool->is_idle();
207208
pool->notify(!current_pool_idle);
208209
}));
209210
}
@@ -245,7 +246,8 @@ namespace threading {
245246
pool = get_pool(task->pool_descriptor).get();
246247
}
247248

248-
const bool current_pool_idle = Pool::current() != nullptr && Pool::current()->is_idle();
249+
const auto current_pool = Pool::current();
250+
const bool current_pool_idle = current_pool != nullptr && current_pool->is_idle();
249251

250252
// Fast path for a single group: lock-free token acquisition and waiter buckets
251253
if (task->group_descriptors.size() == 1) {

0 commit comments

Comments
 (0)