You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
A consumer parked on AsyncQueue::pop resumes on the executor it was running on when it
parked, not on the executor the queue was constructed over -- on a push, a close() and a stop
request alike (see Added, the current-executor context). The queue's executor is now the
fallback, used only for a consumer that parked while no executor was current. The old behaviour
sent a coroutine that parked on a Strand back off it, onto the queue's executor, where it raced
the state the strand serialises (found by morph, PR #806).
Migration: a consumer that runs on the queue's own executor -- the usual shape -- sees no
change. One that parks while running on a different executor now comes back there; if it
relied on arriving on the queue's executor, hop explicitly after the pop with co_await ResumeOn { queueExecutor }. An IExecutor of your own that resumes coroutines
should state itself with core::async::ExecutorScope around the resumption (as testing::ManualExecutor does), or its coroutines keep the old behaviour.
Consumers: fastcached uses AsyncQueue in five files (RaftPeerTransport.hpp and .cpp,
and fastcache-cli's LiveEventSource.cpp, LiveSourceRig.hpp and ScriptedStopSignal.hpp)
and is unaffected: every queue there is built over the reactor, and every pop() parks either
on that reactor or outside any executor, so the executor it comes back to is the one it came back
to before. morph's handlers, which parked on a strand, come back to it now, which is the fix.
core::net::resumeSoonOn is core::net::detail::resumeSoonOn, and takes the chain's
unowned root instead of a work-item factory. It is ResultAwaitable's out-of-line completion
hook and has no other caller.
Migration: none expected -- fastcached, endo, contour, tuidu, dbtool and morph call it
nowhere. A transport completes an operation through ResultAwaitable::complete.
Added
core::async::Strand (<core/async/Strand.hpp>): an IExecutor over any IExecutor that
runs what it is given one task at a time, FIFO -- a task being one resumption, from the submit to
the next suspension. co_await ResumeOn { strand } hops onto it; runningHere() asks whether
the calling thread is inside one of its tasks. One coroutine pump per strand, queued on the base
once however many tasks arrive while it is busy, runs at most StrandOptions::batch (32) tasks
per turn before it hands the base back. A task that throws out of resume() propagates to
whoever resumed the pump and does not wedge the strand or lose the tasks behind it -- except
under MSVC's cl, where it ends the process with a message: an exception crossing the strand's
coroutine frames was measured corrupting the thread's executor scopes on cl-release
(core-cpp.strand-throw-canary). An allocation that fails leaves the strand as it was: submit
throws with nothing queued, and a key's strand whose replacement pump cannot be made is retired
if nothing is queued on it. A refused hand-off abandons the strand's queued work: a base
whose submit throws when the strand hands it its pump makes that submit throw, and the rest
of the queue -- work other threads queued meanwhile included -- is dropped as destruction drops
it; a KeyedStrands key's strand is retired with it. A key whose first submit cannot allocate
no longer leaves its new strand registered with nothing to retire it. What a frame freed by that drop submits to
the same strand (or KeyedStrands) from its destructor is dropped as well, not handed to the
refusing base. ~Strand waits for a hand-off still inside the base's submit. Destroying a strand drops what is
queued -- a chain rooted in a DetachedTask is freed, a Task-owned coroutine is left to its
owner -- and waits for a task running on another thread, but not for the task it is called
from: a task may release the last reference to the strand's owner. The strand's state outlives
it, so a coroutine that parked on it and is handed back later (an AsyncQueue push after the
strand died) is dropped the same way rather than reaching freed storage; inside a task, currentExecutor() is that state, not the Strand's address -- ask runningHere(). The base
must outlive the strand and run what it queued; an EventLoop destroyed with the strand's pump
still in its inbound queue drops it, leaking the strand's state. Written after morph's StrandExecutor, which consumers should replace with it.
core::async::KeyedStrands<Key, Hash, KeyEqual> (<core/async/KeyedStrands.hpp>): one
strand per key over a shared base, made when a key gets work and reclaimed when it runs out. submit(key, ...), co_await strands.resumeOn(key), runningHere(key), runningAnyHere(), size(), and waitIdle(), which blocks until no key has work, asserts when called from one of
its own tasks, and is declared only where threads exist. Destroying it from one of its own tasks
does not wait for that task. A coroutine that parked while its key's
strand was reclaimed comes back to the key, never to a second strand beside it.
Callables, a closed strand's answer, an around-task hook and idle() on both strands -- what
morph's switch from its own StrandExecutor found missing:
post(fn) / post(key, fn) run a callable as one task, held by value in one allocation (the
census in StrandAllocation_test.cpp: one per post and none per submit, to a busy key and to
an idle one alike, in the steady state -- KeyedStrands keeps up to 32 retired strands, with
their pump frames, queue room and map nodes, and gives them to the next key that needs one;
without that a post to an idle key cost five). What it throws takes
a task's way out, and ends the process under MSVC's cl as a task's throw does.
tryPost(fn), trySubmit(work) and their keyed forms return false once the strand is closed
and leave the work with the caller, which can then run it itself. close() is public on both,
idempotent, and what the destructors call.
StrandOptions::aroundTask (an AroundTask) and KeyedStrands' KeyedAroundTask<Key>, which
is given the key: a hook called around every task, a coroutine that came back through the
strand from another executor included, to install per-task ambient context such as a
session. A reference set at construction; unset, it costs one branch per task. Task carries
no context of its own, which would cost every co_await for every consumer.
idle(): nothing queued or running. On the single-threaded WebAssembly build, where waitIdle() does not exist, a host pumps its base until idle(); destroying or closing a
strand there drops what is queued without waiting, since nothing else can be running.
The current-executor context (<core/async/ExecutorContext.hpp>): ExecutorScope marks the
calling thread as running a task of an executor, nests, and restores on every exit; currentExecutor() answers the innermost; ResumeTarget is what an awaitable holds across a
suspension to resume there, with ResumeTarget::currentOr(fallback). core::net::EventLoop
states itself once per turn (and around its teardown drains), ThreadPoolExecutor once per
worker thread, a strand once per batch. The scope is two thread-local stores and no allocation;
on the loop it is paid per turn, not per resumption, because G2 puts every resumption inside one
turn's step 2. Measured on the drain path against the same tree without it (gcc-release, one
pinned core, 4M resumptions, median of five interleaved runs): 35.9 to 36.3 ns per resumption
with one resumption per turn -- the two stores, paid once per turn -- and 19.6 to 19.7 with a
turn of 64. The program is tests/bench/ExecutorContextBench.cpp (core-cpp-bench-executor-context).
core::net's socket operations and timers do not read it and keep resuming on their EventLoop (G2); a strand-bound coroutine hops back with co_await ResumeOn { strand }.
core::async::testing::ManualExecutor (<core/async/testing/ManualExecutor.hpp>): an
executor a test drains by hand (runOne, drain, pending) that states itself as the current
executor while it does. drain(bound) stops after DrainBound (2^20) resumptions by default and
throws std::length_error if work is still queued, so work that requeues itself for ever fails
a case instead of hanging it. The first public test double of core::async.
A design note, Strands and the resume context (docs/design/strands.md).
Changed
A socket completion costs less on the loop's thread.ResultAwaitable::complete hands its
waiter to the drain step (G2) through EventLoop's ready queue, and for a chain rooted in a DetachedTask -- every connection a server spawns -- the queue entry was a refcounted work item
whose claim cost five atomic operations per completion; it is now a detail::CountedClaim
costing two, with the same count, arm and teardown behaviour. The waiter a readiness callback
queues reaches the callback position by a swap of two vectors instead of a range insert and an
erase, and the drain resumes an entry in place rather than moving it out first. Measured with the
new [bench] cases (core-cpp-net_backend-test "[bench]", gcc-release, median of five on an
idle host): a completion through a drain-step callback went from 120.4 to 83.4 ns for a
detached chain, and from 122.6 to 79.3 ns for a chain a caller owns. The ordering (a waiter
resumes in its callback's position), teardown (a detached chain is freed, a borrowed one resumed)
and take-back paths are the same, and each is a case in CompletionClaim_test.cpp. Two things a
caller can observe did change:
a frameless park holding both of its handle's watch slots gets one callback per report that
flags both, not two;
at teardown, a readiness callback still queued is dropped with its park still marked queued,
so a requeue for the same reason during teardown is suppressed rather than queued again.
A readiness callback that completes one waiter, the common case, has that waiter held in a
slot and resumed by the drain straight after it returns, with no queue entry made for it; a
readiness report reaches its park through the handle's watch instead of probing the park
table; and EventLoop's queue entry is back to 56 bytes from 64.
A loop in its steady state allocates nothing per completion or per timer firing. Measured by
fastcached as allocations per request, and now asserted over 256 turns in CallbackAllocation_test.cpp. Four allocations are gone: the ready queue is a ring that keeps its
capacity (detail::RingQueue) rather than a std::deque, which allocated a node every few
entries of a FIFO that never grows; a parked operation no longer spells its never-completed
answer into a heap string when it is armed, only when that answer is read; a fired or cancelled
timer's park is recycled like every other; and a turn collects its due deadlines into a vector
it keeps.
Against fastcached's own reactor, measured by fastcached on this release's core::net
(release/next-031 at b0c82a4) at 64 connections: fewer syscalls per request than the reactor,
6.0 allocations per request against its 12.6, and total CPU per request at parity within the
noise. User CPU is not yet at parity: its medians are still about 1 µs per request higher, with
the runs' ranges overlapping, and a profile puts core-cpp at 7.6% of samples where the reactor's
own code was 3.6%. What remains is tracked as core-cpp#52 (reusing a socket's park across its
operations).
Fixed
A loop with more ready connections than dispatchBatch spent about 1,460 dispatches per round
trip; it spends four. A frameless readiness park -- every parked socket operation -- stays filed
across its wakes, and a level-triggered backend reports its handle on every wait until it is read,
which it was not while the owner's callback waited behind EventLoopOptions::dispatchBatch. Each
report queued the callback again, and the copies took the bound's slots doing nothing: 64 socket
pairs ping-ponging on one loop at the default bound of 64 made 47,424 round trips in 10 s, at
about 1,460 dispatches each, where 32 pairs made 64,000 in 0.28 s. A park is now queued once per
reason (a cancel or an abandonment behind a queued readiness still queues), and 64 pairs make
128,000 round trips in 0.65 s.
Known issues
DetachedTask under clang-cl at -O0 can read its return object back from a frame that
has already been freed, when its first suspension hands it to something that runs it to its end
before await_suspend returns -- a pool thread, or an inline executor (core-cpp#51). Debug
builds with clang-cl only: cl, clang and GCC elsewhere, and clang-cl with optimisation, are
unaffected.