Skip to content

feat: Run reactive effects - #2508

Open
schloerke wants to merge 32 commits into
mainfrom
schloerke/async-pr-2182-reimpl
Open

schloerke wants to merge 32 commits into
mainfrom
schloerke/async-pr-2182-reimpl

Conversation

@schloerke

@schloerke schloerke commented Sep 25, 2026 •

Copy link
Copy Markdown
Collaborator

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 async effect in one session no longer delays other sessions, and a busy session keeps handling messages while its inputs stay stable.

Key changes:

  • No global lock. The session receive loop, invalidate_later, and ExtendedTask completion no longer take reactive.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 fresh contextvars context so the requester's session and OTel span don't leak into it.
  • Non-waiting flush. 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's flushReact(), it won't start a second flush while one is active; the active one runs again afterwards.
  • Cycles from the busy count. Input updates and invalidate_later invalidations 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's invalidate_later hang. Dispatch messages run as their own tasks.
  • Outputs once per cycle. AppSession._flush holds 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 (see outputProgress.ts). It copies and resets the queues before awaiting the send, so a value set during the send isn't lost.
  • Errors. An exception that escapes an effect task is reported to that effect's session, just as it was when flushes ran inline.
  • Public reactive.flush() still waits until everything settles, so tests and the lock() + 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:

TODOs

Open (other repos, merge after the release)

Done

Verification

  • New e2e test 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 on main, where the global lock freezes the other tab.
  • tests/pytest/test_concurrency.py adds five in-process tests, all of which fail on main:
    • a slow async effect in one session doesn't block another session;
    • a value shared across sessions doesn't make the setting session's flush wait on the other session's slow effect;
    • while busy, the session answers a dispatch message, keeps input values stable, and applies the queued update once idle;
    • an output set while a send is in flight is kept;
    • the feat: Implement R Shiny-style concurrency model #2182 invalidate_later + queued-update hang doesn't occur.
  • The full pytest suite passes, as does tests/playwright/shiny in chromium. uv run make format check-lint check-types and make check-pyright are clean.

schloerke and others added 6 commits September 25, 2026 11:42
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.
@jat255

jat255 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

For reference, I worked up some architectural diagrams to help me understand this a bit more:

Before this PR

before

After this PR (as of Sep 30)

Outdated, so removed

Comparison to R Shiny approach

r-shiny-arch

jat255
jat255 previously requested changes Sep 30, 2026
Comment thread shiny/_app.py
Comment thread shiny/reactive/_reactives.py
Comment thread shiny/reactive/_core.py Outdated
Comment thread shiny/reactive/_core.py Outdated
Comment thread shiny/session/_session.py
Comment thread shiny/reactive/_core.py
Comment thread shiny/reactive/_extended_task.py Outdated
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.
@jat255
jat255 dismissed their stale review October 1, 2026 19:17

PR has been updated

@schloerke
schloerke marked this pull request as ready for review October 1, 2026 19:17
@schloerke
schloerke added this pull request to stack #2516 October 1, 2026 19:36
@jat255

jat255 commented Oct 1, 2026 •

Copy link
Copy Markdown
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 main and on this branch.

Fixed by this PR

Partly fixed

Not changed by this PR

These 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. ExtendedTask is still the solution for them.

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 destroy() while one of its effects runs).

jat255 and others added 4 commits October 2, 2026 14:36
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.
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
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 branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Can't update reactive value from download handler Implement request_flush() at the app/session level

2 participants