Repository navigation
Conversation
Five in-process tests for the R Shiny-style concurrency model: - a slow async effect in one session doesn't block another session - a shared value invalidating another session's slow effect doesn't make the invalidating session's flush wait for it - while an effect is busy, the session still answers dispatch messages, keeps input values stable, and applies queued updates once idle - an output set while a flush is awaiting its send is kept - an invalidate_later timer firing while an update is queued doesn't hang the session (the #2182 busy_count=1 case) All five fail on main. Co-authored-by: Josh Taillon <joshua.taillon@posit.co>
Reimplements the concurrency model from #2182 so that a flush starts effects and never waits for their async parts: - Remove the global reactive lock from the session loop, invalidate_later and ExtendedTask completion. The receive loop never awaits a flush. - Invalidations call request_flush(), which schedules one coalesced flush on the next loop pass with a fresh contextvars context. - The internal flush starts each pending context as a task (priority orders starts only), keeps strong references to those tasks, and, like R's flushReact(), doesn't start a second flush while one is active. - Input updates and invalidate_later invalidations are sync cycle-start actions, run only from the session's own output flush: idle -> send outputs -> start the next action. Dispatch messages run as tasks. - AppSession._flush sends once per cycle (held while the session is busy), always including the empty message the client needs to settle output progress, and copies+resets the queues before awaiting the send. - Exceptions escaping an effect task are reported to its session. - Public reactive.flush() waits until settled, but returns right away when called from within a flush or when nothing is pending or running. Co-authored-by: Josh Taillon <joshua.taillon@posit.co>
Effects now run concurrently, so a streamed message can land after one appended later (posit-dev/shinychat#417). Check each message is present; TODO restore the order check once that issue is fixed.
A 3s async effect in one tab must not stop the other tab's timer. Fails on main (the global lock freezes the other tab).
With the flush reentrancy guard, the stream's task starts before later effects append, so the ordered assertion holds again.
Contributor
|
For reference, I worked up some architectural diagrams to help me understand this a bit more: Before this PR
After this PR (as of Sep 30)Outdated, so removed Comparison to R Shiny approach
|
jat255
previously requested changes
Sep 30, 2026
Now that a flush doesn't wait for effects, two effects can read an async calc mid-run, and an effect can be re-invalidated mid-run. - Calc: a reader from another task waits for the run in flight instead of starting a second one (R caches the promise). A second run used to leave `_running` stuck True, so the calc never cached again. A run invalidated partway through no longer leaves its stale value in front of the fresh one. - Effect: a re-run waits for the previous run to finish, so a slower, older run can't finish last and win. (R lets them overlap.) - `priority` docs: priority orders starts; an async effect's awaited part doesn't hold back lower-priority effects.
- Track "inside a flush or effect run" by task identity (a ContextVar holding the owning task), not an inherited bool. Tasks an effect starts (ExtendedTask bodies, chat tool calls, background loops) now wait for dependents instead of returning at once, and don't wait on the effect that started them. - Add flush_pass(): run, or wait for, one complete flush that starts after the call. reactive.flush() always does at least one pass and no longer spins re-running flushes while one is active. - ExtendedTask waits for a full flush pass before starting a queued invocation, so dependents see every result. - Rewrite the reactive.lock() docstring: nothing takes it any more.
Fill in the App._sessions_needing_flush placeholder (R's appsNeedingFlush). A session asks for a flush when it has something to send: the idle transition, a queued cycle action, send_input_message(), on_flush() registration, and a finished message handler. After each reactive flush the app flushes just those sessions, and an error in one closes that session instead of skipping the rest. - Message handlers run in session-owned tasks that reactive.flush() doesn't wait on, so handlers that await reactive.flush() can't deadlock. - Requests made by a session's own flush callbacks don't schedule another flush (their messages go out in the one being sent), avoiding a loop. - Only downloads take the busy count; an upload or dynamic route no longer pauses the session's cycle. - `_flush` skips straight to a queued action only when there is one.
Covers every case from the review's repro scripts, plus edge cases: async calc sharing (incl. invalidated mid-run and shared errors), effect runs not overlapping, reactive.flush() from effects, tasks, flush callbacks and message handlers, ExtendedTask queued results, per-session flush requests and error isolation, one message per cycle (even empty), timers waiting while busy, busy only for downloads, and session teardown. The timer test uses a longer interval so it no longer flakes under load.
…r effect reactive.flush() awaited through a child task of a running flush or effect (`asyncio.gather()`, or `asyncio.wait_for()` before Python 3.12) waited for work that couldn't finish until that ancestor did, e.g. a flush pass that can't start while the current flush is in a flushed callback, or an effect's queued re-run. It now returns right away whenever the owning flush or effect run is still going. A task whose starting effect has finished (an extended task's body) still waits as usual.
A calc invalidated mid-run can be re-run by the next reader while the stale run still awaits. If the stale run finished last it overwrote the fresh value, and the calc then served it from the cache. Now only the most recent run, if it wasn't invalidated, stores its result and ends the running state; a superseded run returns its result to its own caller only.
A finished message handler no longer requests a flush (its ui.update_*() calls already do, via send_input_message()), and the latest calc run no longer skips caching when invalidated (_invalidated already forces the next read to recompute). Found by mutation testing: reverting either changed no observable behavior.
Mutation-tested all 34 fixes (revert one, run the suite): these tests cover the ones that were only covered jointly or not at all -- the two busy gates in _flush, the _start_cycle busy guard, the latest-run calc bookkeeping, the idle transition ending the cycle, on_flush registration requesting a flush, busy ending when a run is cancelled, no concurrent flushes, the fresh flush context, and slow message handlers. The recording connection also asserts that no output message is sent while busy.
From within an effect, or a task started by one that is still running, it returns right away. Previously this was only in the internal flush_settled() docstring.
After a reactive flush, the app used to await each requesting session's _flush(), including its websocket send and on_flush callbacks, one after another inside the global flush. A slow client or callback in one session delayed every other session's output, the next reactive flush, and anyone awaiting reactive.flush(). Each session now flushes in its own task, one at a time per session (a request made mid-send runs once more afterwards, so messages stay in order). An error closes that session only. reactive.flush() still waits for these tasks, so its "settled" meaning is unchanged. A flush callback that awaits reactive.flush() returns right away rather than waiting on its own session's task. Also skip the send when no cycle has ended and nothing is queued (R's hasPendingUpdates()): each input update now yields one output message instead of a value message plus an empty one.
request_flush() used get_running_loop(), so a value set from another thread silently scheduled no flush, and calling loop.call_soon() from that thread would not have been safe anyway. The environment now remembers the loop Shiny runs on and hands requests from other threads (including a thread running its own loop) to it with call_soon_threadsafe(); requests are ignored if there is no loop yet or it has closed. Only the request is thread-safe: Value.set()'s docs now say to call it on the loop's thread, or via loop.call_soon_threadsafe(value.set, x).
Slow clients and slow on_flush callbacks delaying only their own session; the next reactive flush not waiting on a stuck send; per-session sends never overlapping and staying in order across cycles; reactive.flush() waiting for pending sends; concurrent flush callbacks in two sessions; a session closing mid-send with a re-run pending; a failed send closing only its session; one output message per update. For threads: set() via call_soon_threadsafe and direct, a thread with its own loop, merged requests, no loop yet, a closed loop, and a new loop after the old one closed. Every new code path was mutation-tested.
schloerke
marked this pull request as ready for review
October 1, 2026 19:17
schloerke
added this pull request to stack #2516
October 1, 2026 19:36
Contributor
Issue sweep for this PR stack (#2508, #2515)I searched the open issues for problems that this stack fixes or changes. For #1785 and #1889, I ran the app from the issue in a browser with Playwright, on Fixed by this PR
Partly fixed
Not changed by this PRThese issues ask that async work not block the same session. This PR keeps the R behavior: the inputs and timers of a session wait while the effects of that session run.
Regression risk
Checked and not related#14, #365, #771, #1271, #1763, #1776, #2124, #2194, #2443, #2483. No open issue is about the area of #2515 (a module |
The download's busy count covers only the handler call, so a streamed body runs with the session idle: values it sets reach the client mid-stream, and the session keeps handling input. Lock that in with unit tests (sync and async handlers) and an e2e test, and say so in the code and changelog.
…pped The reactive environment is process-wide, but a flush runs on one event loop. If that loop stops mid-flush (a test's loop closing, or a previous `test_server()` run), the flush task is destroyed without running its `finally`, leaving `_in_flush` set forever. Every later flush then returned as if one were active, and `flush_pass()` waited forever: ExtendedTask completion hung. A pending flush request could be lost the same way, handed to the dead loop. When a flush, flush pass, or flush request runs on a different loop than the recorded one, and that loop is closed or no longer running, drop the stale state and adopt the new loop. This hung the Windows/Python <= 3.12 CI jobs for six hours: an otel test stranded a flush, and an ExtendedTask test that later landed on the same xdist worker waited forever.
It assumed that after `sleep(0.005)` the `on_flushed` callback was still in its `sleep(0.02)`. With Windows' ~15.6 ms timer resolution (and on loaded machines) the callback could finish first, so the queued update ran before the session turned busy. The callback now signals when it is running and waits for the test.
3 of 5 tasks
This was referenced Oct 2, 2026
The previous commit dropped the environment's references to a dead loop's tasks. Garbage collection then finalized those coroutines during whichever test was running: an effect's `finally` reset a ContextVar in the wrong context, and the resulting "Error in Effect" warning failed an unrelated test (seen on macOS / Python 3.12 CI). Keep them in a separate set that no flush waits on, as they were kept before.
…82-reimpl # Conflicts: # shiny/reactive/_reactives.py
This was referenced Oct 3, 2026
A 10 ms sleep ends on the next loop turn on Windows, whose 15.6 ms clock tick makes the timer due at once. The flush task's done callback had not yet removed it from the environment's tasks.
This was referenced Oct 5, 2026
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.


Fixes #1785.
Fixes #21.
Reimplements #2182 (supersedes it). Related: #837, #39, #907.
Summary
Adopts R Shiny's concurrency model: a flush starts effects and never waits for their async parts. A slow
asynceffect in one session no longer delays other sessions, and a busy session keeps handling messages while its inputs stay stable.Key changes:
invalidate_later, andExtendedTaskcompletion no longer takereactive.lock(). The receive loop never awaits a flush.request_flush(). Every invalidation requests a flush. The request schedules one merged flush on the next loop pass (loop.call_soon), using a freshcontextvarscontext so the requester's session and OTel span don't leak into it.ReactiveEnvironment.flush()starts each pending context as a task, in priority order (priority orders starts only, as in R). It keeps strong references to those tasks. Like R'sflushReact(), it won't start a second flush while one is active; the active one runs again afterwards.invalidate_laterinvalidations are sync "cycle start actions". They run only from the session's own output flush, in the order idle → send outputs → start the next action. They never run inline in whichever task queued them, which is the root cause of feat: Implement R Shiny-style concurrency model #2182'sinvalidate_laterhang. Dispatch messages run as their own tasks.AppSession._flushholds everything while the session is busy. At the end of each cycle it sends one message, even an empty one, which the client needs to settle output progress (seeoutputProgress.ts). It copies and resets the queues before awaiting the send, so a value set during the send isn't lost.reactive.flush()still waits until everything settles, so tests and thelock()+flush()idiom keep working. It returns immediately when called from within a flush (from an effect or a flushed callback), or when nothing is pending or running.Behavior changes app authors may notice:
awaitinstead of running one after another.append_message_stream()jump ahead of the stream (append_message_stream() doesn't reserve its place before its task starts, so later messages can jump ahead shinychat#417, fixed in fix(py): append_message_stream() reserves its place before its task starts shinychat#418). The chat stream test keeps its order check and passes.TODOs
Open (other repos, merge after the release)
reactive.lock()is deprecated and does nothing, and setting a value no longer needsawait reactive.flush().docs/genai-tools.qmd, whose tool example wrapscurrent_query.set(query)inasync with reactive.lock(). Merge it after the release that includes this PR and feat(reactive): Deprecatereactive.lock()#2520.on_response()docstring in shinychat'spkg-py/src/shinychat/_history.py. It says that the effect is safe because Shiny serializes flushes behind a process-widereactive.lock().Done
reactive.lock()+await reactive.flush(), including for code that runs without a session (Design pattern for global async reactive #698).reactive.lock()#2520:reactive.lock()is deprecated and does nothing. Set the value directly, or usesession.run_once_when_idle().ExtendedTask. Design pattern for global async reactive #698 was closed with an example.invalidate_later's privatesession._cycle_start_action.session.run_once_when_idle()#2517:session.run_once_when_idle().destroy()of effects that are mid-run.reactive_updatespans (one per flush, or one per session). They lostsession.idafter init, because a requested flush runs in a fresh context and can serve several sessions.reactive_updatespan #2522: one span per session cycle, as in R.on_restorecallbacks (priority 1000000).windows-latest(run). Onmain, none of four outputs got its progress message before its value. On this PR, the first output did and the other three did not. So Some busy indicators not firing #1381 still occurs onmain, and this PR does not make it worse.request_flush.call_soon_threadsafe. TheValue.set()docs say to call it on the loop's thread (loop.call_soon_threadsafe(value.set, x)).on_flushcallback) delays only its own session.test_value_set_in_download_handler_updates_outputs_and_effects(async and sync handlers) fails onmainand passes here.priorityandreactive.lock()docstrings, and CHANGELOG.Verification
tests/playwright/shiny/async/concurrent_sessions/: it opens the app in two tabs and starts a 3s async effect in the first. The second tab's 100ms timer must keep ticking while that effect runs. The test fails onmain, where the global lock freezes the other tab.tests/pytest/test_concurrency.pyadds five in-process tests, all of which fail onmain:invalidate_later+ queued-update hang doesn't occur.tests/playwright/shinyin chromium.uv run make format check-lint check-typesandmake check-pyrightare clean.