diff --git a/docs/spec/concurrency_and_lifetimes.md b/docs/spec/concurrency_and_lifetimes.md index 21447c16a..526a04ca2 100644 --- a/docs/spec/concurrency_and_lifetimes.md +++ b/docs/spec/concurrency_and_lifetimes.md @@ -104,7 +104,11 @@ turns it into a set of per-key serial queues: may run concurrently on different pool threads. This is what removes the need for per-model mutexes. - Each strand is a `shared_ptr` in a map guarded by `_mapMtx`. When a - strand's queue drains, the map entry is erased. The invariant is: **at most one + strand's queue drains, the map entry is removed — `extract`ed into a + single-slot `_spare` the next miss re-keys, which recycles the node's memory + without changing when the entry leaves the map (see + [`core/executor.md`](core/executor.md), "Lifetime & ownership"). The + invariant is: **at most one live strand per `ModelId`, and any `running` strand is the one currently in the map** — that is what keeps a key's tasks from overlapping. Both sides that can break it hold `_mapMtx` across their *whole* decision: `post()` takes `_mapMtx`, diff --git a/docs/spec/core/executor.md b/docs/spec/core/executor.md index 0828f0537..444e34b4f 100644 --- a/docs/spec/core/executor.md +++ b/docs/spec/core/executor.md @@ -205,14 +205,18 @@ lifetime of the model. It supports three-way comparison and can be used as an Internally `StrandExecutor` maintains a map of `ModelId → shared_ptr` (shared state per key). A `Strand` holds a pointer to the base `IExecutor`, a mutex, a pending queue, and a `running` flag. The executor also tracks an -`_inFlight` counter (guarded by the map mutex) that the destructor waits on. +`_inFlight` counter (guarded by the map mutex) that the destructor waits on, +and a single-slot `_spare` node handle (also guarded by the map mutex) that +recycles one detached map entry — see +[Lifetime & ownership](#lifetime--ownership). The pending queue is `StrandExecutor::PendingQueue`, not `std::queue`. It is a FIFO with the head task stored **inside** the `Strand` and a lazily constructed `std::deque` behind it, and it exists purely to make the common case cheaper: -because the drain step below destroys the whole `Strand` as soon as the queue -empties, a serial workload (one action at a time, each waited out) rebuilds the -queue on every dispatch and puts exactly one task in it — and libstdc++'s +because the drain step below detaches the whole map entry as soon as the queue +empties, a serial workload (one action at a time, each waited out) starts from +an empty queue on every dispatch and puts exactly one task in it — and +libstdc++'s `std::deque` allocates its node map *and* a 512-byte first buffer in its default constructor, whether or not anything is ever pushed. See [Lifetime & ownership](#lifetime--ownership) below for the measurement. @@ -220,7 +224,7 @@ default constructor, whether or not anything is ever pushed. See under the owning `Strand::mtx`, exactly as the `std::queue` it replaced was, and it tracks occupancy with a flag rather than by testing the callable, so an empty `std::function` is queued and dispatched like any other. It changes no -lifetime or locking rule: the erase still fires when `empty()` becomes true, +lifetime or locking rule: the removal still fires when `empty()` becomes true, still under the `{_mapMtx, strand->mtx}` pair. **`_inFlight` is incremented with the *decision* to dispatch, not lazily.** @@ -336,16 +340,15 @@ data race the `_inFlight` wait exists to prevent. Callers must ensure all task sources are shut down before the `StrandExecutor` is destroyed. The strand map is self-cleaning: when a strand drains (its `pending` queue is -empty), `scheduleNext` clears `running` and erases the map entry under the +empty), `scheduleNext` clears `running` and removes the map entry under the combined `{_mapMtx, strand->mtx}` lock. Live memory therefore tracks the set of *currently active* models rather than every model ever seen — there is no per-model registration to leak. -**The cost is allocation churn, and it is bounded rather than removed.** A -model posted to serially — one action at a time, each waited out — never has a -task queued at the instant the previous one finishes, so it never keeps a -strand: every dispatch takes `post()`'s `if (!slot)` branch and rebuilds the -map node and the `Strand`. Measured on `7a343e6f` with +**The cost was allocation churn.** A model posted to serially — one action at +a time, each waited out — never has a task queued at the instant the previous +one finishes, so it never keeps a strand: every dispatch missed in the map and +rebuilt the map node and the `Strand`. Measured on `7a343e6f` with `tests/bench/bench_dispatch_allocations.cpp` (see [testing_strategy.md](../testing_strategy.md)), x86-64 Linux, GCC 16.2.1 / libstdc++, `-O2`, that came to **4 allocations and 760 of the 1990 bytes** a @@ -357,11 +360,35 @@ buffer in its default constructor. Replacing that container with allocations and 152 bytes** and the whole round trip to 18.9 allocations / 1396 bytes (morph#660). -The two allocations that remain — the map node and the `Strand` itself — are -inherent to the erase: removing them means keeping the slot alive across the -drain, which trades this churn for a per-model entry that nothing reclaims, -since `StrandExecutor` has no deregistration hook. That trade is deliberately -not made here; the erase is what bounds the map. +**The remaining two allocations — the map node and the `Strand` itself — are +recycled rather than removed (morph#670).** They looked inherent to the erase: +removing them appeared to mean keeping the slot alive across the drain, which +would trade the churn for a per-model entry nothing reclaims, since +`StrandExecutor` has no deregistration hook. That framing turned out to be +avoidable. The entry still leaves the map at exactly the same moment, under +exactly the same locks; the drain simply calls `extract` instead of `erase` and +parks the detached node in a single-slot `_spare` member, and the next +`post()` that misses re-keys that node and inserts it back. The map is still +bounded by the removal — `_spare` holds **at most one** node, is guarded by +`_mapMtx` like the map itself, and is freed with the executor. + +Reusing the parked node's `Strand` object is guarded additionally by +`use_count() == 1`: the recycled node is then the only owner, so no strand task +can still reach the object and reusing it is indistinguishable from +constructing a new one. When that guard fails — a finishing strand lambda still +holds its `shared_ptr` when the next `post()` looks — a fresh `Strand` is +constructed exactly as before and only the node is recycled. To make the guard +usually hold, the strand lambda drops its `shared_ptr` immediately after the +drain block rather than at its own destruction; nothing after that point +touches the strand. That timing affects *whether* the object is recycled, never +whether the recycling is safe. + +Re-measured on `7d4ca453` (this change's base) with the same instrument, +x86-64 Linux, **clang 22.1.8 / libstdc++ 16.2.1, Release**: the round trip went +from **18.90 allocations / 1394.8 bytes** to **16.95 / 1244.6** — the full 2 +allocations and ~150 bytes the strand had left. Six alternating runs of each +binary; spread within 0.1 allocations and 2 bytes per call. The magnitude is +libstdc++-specific, as it was for morph#660. ## Thread safety @@ -375,8 +402,9 @@ concurrently. - `MainThreadExecutor` guards its queue with `_m`. `post()` may be called from any thread, but `runFor()` must be called only from the single owning ("main") thread; concurrent `runFor()` calls are not supported. -- `StrandExecutor` uses two lock levels: `_mapMtx` protects the `_strands` map - and the `_inFlight` counter, and each `Strand::mtx` protects that strand's +- `StrandExecutor` uses two lock levels: `_mapMtx` protects the `_strands` map, + the `_spare` recycled node and the `_inFlight` counter, and each + `Strand::mtx` protects that strand's `pending` queue and `running` flag. Both operations that can break the per-key invariant hold `_mapMtx` across their whole decision: `post()` takes `_mapMtx`, does the slot lookup/create, and then — still holding `_mapMtx` — diff --git a/include/morph/core/strand.hpp b/include/morph/core/strand.hpp index 44fb25049..a658c7347 100644 --- a/include/morph/core/strand.hpp +++ b/include/morph/core/strand.hpp @@ -102,12 +102,11 @@ class StrandExecutor { // is uncontended; an existing strand's mtx can only be held // elsewhere under the same _mapMtx-first order, so no deadlock. std::scoped_lock const mapLock{_mapMtx}; - auto& slot = _strands[key]; - if (!slot) { - slot = std::make_shared(); - slot->base = _base; + auto slotIter = _strands.find(key); + if (slotIter == _strands.end()) { + slotIter = installStrand(key); } - strand = slot; + strand = slotIter->second; std::scoped_lock const strandLock{strand->mtx}; strand->pending.push(std::move(task)); if (!strand->running) { @@ -126,7 +125,7 @@ class StrandExecutor { } } if (schedule) { - scheduleNext(strand, key); + scheduleNext(std::move(strand), key); } } @@ -205,6 +204,64 @@ class StrandExecutor { bool running = false; }; + /// @brief The `ModelId` → strand map. Named so the recycled node type can be. + using StrandMap = std::unordered_map, ModelIdHash>; + + /// @brief Returns an iterator to the strand for @p key, creating the entry. + /// + /// **Precondition:** the caller holds `_mapMtx` and has already + /// established that `key` has no entry. + /// + /// This is a pure allocation optimisation and changes no lifetime or + /// locking rule. The drain step in `scheduleNext` removes the whole map + /// entry as soon as the queue empties, so a workload that dispatches one + /// action at a time against a model paid for a fresh map node *and* a + /// fresh `make_shared` on every call — the 2 allocations / 152 + /// bytes that were left after morph#660 took the container's share. + /// Rather than keep the slot alive across the drain (which would need a + /// deregistration hook and would trade this churn for a per-model entry + /// nothing reclaims), the drain `extract`s the node instead of erasing it + /// and parks it in `_spare`, and this re-keys and re-inserts that one + /// node. The entry still leaves the map at the same point under the same + /// locks, so the map is bounded exactly as before; `_spare` holds at most + /// one node and is freed with the executor. + /// + /// Reusing the parked node's `Strand` as well is guarded by sole + /// ownership. `use_count() == 1` means the recycled node holds the only + /// reference, so nothing else can reach the object and reusing it is + /// indistinguishable from constructing a new one. That is the whole + /// argument, and it is deliberately not "the previous owner makes no + /// further access": a strand lambda that is still finishing does hold a + /// reference, and when it does, this constructs a fresh `Strand` exactly + /// as before and recycles only the node. + /// + /// The parked strand needs no reset. It is only ever extracted from a + /// strand observed `!running` with an empty `pending` under + /// `{_mapMtx, strand->mtx}`, which is the state a fresh one is in. A + /// runtime re-check of that here would be an arm nothing can take, so it + /// is written down rather than branched on. + /// @param key Model identifier to install a strand for. + /// @return Iterator to the entry for @p key. + StrandMap::iterator installStrand(ModelId key) { + if (!_spare) { + auto const iter = _strands.emplace(key, std::make_shared()).first; + iter->second->base = _base; + return iter; + } + _spare.key() = key; + auto& reused = _spare.mapped(); + if (reused.use_count() != 1) { + reused = std::make_shared(); + } + reused->base = _base; + // `insert` consumes the node. Were the precondition ever violated it + // would instead hand the node back inside the returned object, which + // frees it — the same fate the old `erase` gave it — and `position` + // would name the existing entry, so the caller is right either way + // and there is nothing to branch on. + return _strands.insert(std::move(_spare)).position; + } + /// @brief Dispatches one strand lambda onto the base executor. /// /// **Precondition:** the caller must have already incremented `_inFlight` @@ -214,8 +271,15 @@ class StrandExecutor { /// with the *decision* (rather than here) closes the window where /// `~StrandExecutor` could observe `_inFlight == 0` between the decision and /// this dispatch and destroy `_strands` out from under us. - void scheduleNext(const std::shared_ptr& strand, ModelId key) { - strand->base->post([this, strand, key] { + void scheduleNext(std::shared_ptr strand, ModelId key) { + // Read `base` out before the capture list moves `strand` into the + // lambda: `strand->base` and the lambda's construction are + // unsequenced within one call expression, so reading through the + // moved-from pointer would be a real hazard rather than a stylistic + // one. The lambda's capture is non-const (hence `mutable` and the + // by-value parameter) so the drain below can release it early. + IExecutor* const base = strand->base; + base->post([this, strand = std::move(strand), key]() mutable { std::function task; { std::scoped_lock const lock{strand->mtx}; @@ -257,7 +321,19 @@ class StrandExecutor { strand->running = false; auto iter = _strands.find(key); if (iter != _strands.end() && iter->second == strand) { - _strands.erase(iter); + // `extract`, not `erase`: same removal, same moment, + // same locks — the entry leaves the map here exactly as + // before, and every reason the erase had to happen + // under {_mapMtx, strand->mtx} still applies unchanged. + // The difference is only that the detached node's + // memory is parked for `installStrand` to re-key + // instead of being returned to the allocator. Any node + // already parked is freed by this assignment, so at + // most one is ever held. Freeing it runs no user code + // under these locks: a parked strand was parked + // because its pending queue was empty, so there is no + // captured task left to destroy. + _spare = _strands.extract(iter); } } else { // Account for the re-armed dispatch *before* releasing @@ -271,6 +347,24 @@ class StrandExecutor { } if (more) { scheduleNext(strand, key); + } else { + // Drop this run's co-ownership here rather than leaving it to + // the lambda's destruction a few lines below. Nothing after + // this point touches the strand, and releasing it early is + // what lets `installStrand` see `use_count() == 1` on the + // node just parked in `_spare`: the next post() for this key + // is typically already blocked on _mapMtx when the block above + // releases it, so a reference held until the lambda dies would + // usually still be there when that post looks. This only + // affects *whether the object is recycled*, never whether the + // recycling is safe — a post that looks too early simply sees + // two owners and constructs a fresh Strand. + // + // Safe to be the last owner here: no lock is held (the block + // above released both), so this never destroys a mutex it is + // standing on. If the extract above did run, `_spare` owns the + // strand and this merely decrements. + strand.reset(); } // Decrement after all map access is done; wake destructor if it is waiting. { @@ -286,7 +380,14 @@ class StrandExecutor { std::mutex _mapMtx; std::condition_variable _cv; int _inFlight{0}; - std::unordered_map, ModelIdHash> _strands; + StrandMap _strands; + /// @brief The one detached map node kept for reuse. Guarded by `_mapMtx`. + /// + /// Declared after `_strands` so it is destroyed first: the node owns + /// storage obtained from the map's allocator, and returning it before the + /// container goes away keeps that ordering obvious even though the default + /// allocator is stateless. + StrandMap::node_type _spare; }; } // namespace morph::exec::detail diff --git a/scripts/branch_partial_allowlist.json b/scripts/branch_partial_allowlist.json index 473029fe1..22d79160f 100644 --- a/scripts/branch_partial_allowlist.json +++ b/scripts/branch_partial_allowlist.json @@ -74,9 +74,9 @@ }, { "file": "include/morph/core/strand.hpp", - "line": 259, + "line": 323, "source": "if (iter != _strands.end() && iter->second == strand) {", - "reason": "Unreachable by construction given this class's lock discipline (core audit finding ST1, resolved to (b) by a concurrency-focused review pass after an initial (a)/(b)-undecided pass). `_strands` has exactly two mutation sites: `post()`'s insert-if-absent (this file, `if (!slot) { slot = make_shared(); }`) and this exact block's own erase a few lines below, both under `_mapMtx`. At most one lambda per `Strand` runs at a time (`post()` only schedules when `!strand->running`, and re-arming happens only through this same lambda's own `more` branch), so dispatch for one `Strand` is strictly serial; and only a strand's own currently-running lambda can erase its map entry (the erase fires only in the `!more` branch for the entry this frame just found under `_mapMtx`, and a concurrent `post(key)` while this lambda runs can only push onto the existing `Strand`, never replace it). Together these force `_strands.find(key)` to yield this exact strand whenever this line runs, so `iter->second == strand` cannot be false. No stress test needed: one was considered, but given the strength of the lock-discipline argument it would spend CI time re-confirming an already-proven invariant rather than searching for an unknown one." + "reason": "Unreachable by construction given this class's lock discipline (core audit finding ST1, resolved to (b) by a concurrency-focused review pass after an initial (a)/(b)-undecided pass). `_strands` has exactly two mutation sites: the insert-if-absent `post()` reaches through `installStrand` (this file) and this exact block's own removal a few lines below -- an `extract` into `_spare` since morph#670, which detaches the entry at the same point the `erase` did -- both under `_mapMtx`. At most one lambda per `Strand` runs at a time (`post()` only schedules when `!strand->running`, and re-arming happens only through this same lambda's own `more` branch), so dispatch for one `Strand` is strictly serial; and only a strand's own currently-running lambda can erase its map entry (the erase fires only in the `!more` branch for the entry this frame just found under `_mapMtx`, and a concurrent `post(key)` while this lambda runs can only push onto the existing `Strand`, never replace it; a node parked in `_spare` is out of the map, so `find` cannot return it). Together these force `_strands.find(key)` to yield this exact strand whenever this line runs, so `iter->second == strand` cannot be false. No stress test needed: one was considered, but given the strength of the lock-discipline argument it would spend CI time re-confirming an already-proven invariant rather than searching for an unknown one." }, { "file": "include/morph/core/backend.hpp", diff --git a/tests/test_strand_race.cpp b/tests/test_strand_race.cpp index ace091859..870e9a448 100644 --- a/tests/test_strand_race.cpp +++ b/tests/test_strand_race.cpp @@ -1,10 +1,12 @@ // SPDX-License-Identifier: Apache-2.0 #include +#include #include #include #include #include +#include #include #include #include @@ -21,7 +23,7 @@ // both strands dispatched tasks for that key concurrently — breaking the // per-model serialisation guarantee. // -// This file holds two cases, and they cover different things. This first one is +// This file holds three cases, and they cover different things. This first one is // a *load* test: it hammers post() on a single key from many threads, and every // task bumps a per-key in-flight counter on entry and drops it on exit, so if // two tasks for the same key ever run concurrently the counter exceeds 1 and the @@ -39,7 +41,9 @@ // not better. // // The second case below produces the shape this one cannot, and is the one that -// fails against that mutant. Keep both: saturation and the drain boundary are +// fails against that mutant. The third covers the node the drain now recycles +// (morph#670), which neither of the first two can be wrong about. Keep all +// three: saturation, the drain boundary, and the recycled node's key are // different failure modes of the same invariant. TEST_CASE("StrandExecutor never runs two tasks for one key concurrently under contention", "[strand][race]") { constexpr int kThreads = 8; @@ -296,6 +300,138 @@ TEST_CASE("StrandExecutor keeps one strand per key when a post races the drain", } } +// The recycled map node (morph#670), which neither case above can be wrong +// about. +// +// When a strand drains, `scheduleNext` no longer `erase`s the map entry: it +// `extract`s it into a single-slot `_spare`, and the next `post()` that misses +// re-keys that node and inserts it back. Re-keying is the new step, and it is +// the one a functional test does not see. An entry left under the *previous* +// key still serialises every task that reaches it, still runs them in order, +// and still completes them all; what it corrupts is which key the map answers +// for, and that only becomes a serialisation failure two posts later: +// +// 1. Key A drains, parking a node still keyed A. +// 2. Key B misses and takes that node -- which, unkeyed, goes back into the +// map under A. B's first task starts running on a strand the map calls A, +// so `find(B)` still misses. +// 3. B's *next* post therefore misses too, and installs a second strand for +// B while the first is still running its task. Two strands for one key: +// the invariant the two cases above exist for, reached through a door +// neither of them opens. +// +// Step 3 is what the shape below is for, and it is why this case is not simply +// "post to several keys and let them drain". The mis-key is only observable +// while a task is *still running* on the mis-keyed strand, so each round posts +// a short burst back-to-back -- the second and third posts of a burst arrive +// while the first is running, which on the correct code is an ordinary re-arm +// and on the mutant is a second strand. Between rounds the key is allowed to +// go quiet, which is what produces the drain step 1 needs; several keys +// running this cycle out of phase is what carries a parked node from one key +// to another. Measured, not assumed: with `_spare.key() = key;` deleted, this +// case fails 10/10 under ThreadSanitizer, while a variant that drained between +// every single post (no burst) passed 10/10 against the same mutant. +// +// Same three detectors as the case above, kept per key: an in-flight counter +// for the symptom, plain (non-atomic) per-key state for the data race a +// sanitizer sees whether or not the two tasks overlap in wall clock, and a +// per-key FIFO check for an ordering break that needs no overlap at all. +TEST_CASE("StrandExecutor recycles a drained strand under the key that asked for it", "[strand][race]") { + constexpr std::size_t kKeys = 3; + constexpr int kBurst = 4; + constexpr int kRounds = 900; + constexpr int kIterations = 6; + constexpr std::size_t kCells = 24; + + for (int iter = 0; iter < kIterations; ++iter) { + // Same ordering rule as the cases above: the pool outlives the strand. + morph::exec::ThreadPoolExecutor pool{4}; + + // Declared outside the strand's scope so the drain in `~StrandExecutor` + // runs their final updates against live objects. + std::array, kKeys> inFlight{}; + std::array, kKeys> maxInFlight{}; + std::array, kKeys> completed{}; + std::array, kKeys> outOfOrder{}; + // Plain, per key, and touched only by that key's tasks: distinct + // objects, so a sanitizer report here means two tasks for the *same* + // key raced, never two keys sharing a cache line. Under the invariant + // this file exists for, the strand is the synchronisation and these + // accesses are data-race-free. + std::array, kKeys> cells{}; + std::array lastSeq{}; + + { + morph::exec::detail::StrandExecutor strand{pool}; + + std::vector producers; + producers.reserve(kKeys); + for (std::size_t slot = 0; slot < kKeys; ++slot) { + lastSeq.at(slot) = -1; + producers.emplace_back([&, slot] { + // Distinct, non-zero ids: 0 is `ModelId`'s reserved + // "unbound" sentinel. + morph::exec::detail::ModelId const key{100 + slot}; + for (int round = 0; round < kRounds; ++round) { + for (int post = 0; post < kBurst; ++post) { + int const seq = (round * kBurst) + post; + strand.post(key, [&, slot, seq] { + int const cur = inFlight.at(slot).fetch_add(1) + 1; + int prev = maxInFlight.at(slot).load(); + while (cur > prev && !maxInFlight.at(slot).compare_exchange_weak(prev, cur)) { + } + if (seq <= lastSeq.at(slot)) { + outOfOrder.at(slot).fetch_add(1); + } + lastSeq.at(slot) = seq; + // A little plain work rather than none: an + // empty task gives an overlap a window a few + // instructions wide, which is how a broken + // strand can still report `maxInFlight == 1`. + for (auto& cell : cells.at(slot)) { + cell += 1; + } + inFlight.at(slot).fetch_sub(1); + completed.at(slot).fetch_add(1, std::memory_order_release); + }); + } + // Let this key go quiet before the next burst: the + // drain that parks a node in `_spare` only fires from + // an empty pending queue, and without this gap every + // post after the first re-arms a strand that is + // already in the map and the install path is never + // reached at all. + while (completed.at(slot).load(std::memory_order_acquire) < (round + 1) * kBurst) { + std::this_thread::yield(); + } + } + }); + } + for (auto& producer : producers) { + producer.join(); + } + // Closing this scope runs `~StrandExecutor`, which blocks until + // `_inFlight == 0`; see the first case for why that is a complete + // drain and not a deadline. + } + + constexpr int kPerKey = kRounds * kBurst; + for (std::size_t slot = 0; slot < kKeys; ++slot) { + INFO("iteration " << iter << ", key slot " << slot << ": completed " << completed.at(slot).load() << " of " + << kPerKey); + // `CHECK` for the counts, `REQUIRE` for the invariant: they answer + // different questions and stopping at the first would hide the + // others (see the first case). + CHECK(completed.at(slot).load() == kPerKey); + auto const [lowest, highest] = std::ranges::minmax_element(cells.at(slot)); + CHECK(*lowest == kPerKey); + CHECK(*highest == kPerKey); + CHECK(outOfOrder.at(slot).load() == 0); + REQUIRE(maxInFlight.at(slot).load() == 1); + } + } +} + // Regression test for ThreadPoolExecutor(0): a zero-worker pool used to accept // tasks that could never run, hanging every post() forever. The constructor now // clamps the worker count to at least 1, so a pool built with 0 is still usable.