Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
36 commits
Select commit Hold shift + click to select a range
a10db83
test: Add concurrency tests for cross-session and busy-session behavior
schloerke Sep 25, 2026
3c889f5
feat: Run reactive effects concurrently, R Shiny-style
schloerke Sep 25, 2026
332edff
test(chat): Don't assert message order in chat stream test
schloerke Sep 25, 2026
98b8a8c
test: Add two-tab e2e test for cross-session concurrency
schloerke Sep 25, 2026
feeeee2
docs: Move flush reentrancy note to a code comment; drop R references
schloerke Sep 25, 2026
3b9bd39
test(chat): Restore message order check in chat stream test
schloerke Sep 25, 2026
00f74e0
fix(reactive): Share in-flight async calc runs and serialize effect runs
schloerke Sep 30, 2026
156c5bf
fix(reactive): Make reactive.flush() wait reliably without hanging
schloerke Sep 30, 2026
7869400
fix(session): Flush only sessions that request it, each on its own
schloerke Sep 30, 2026
dd72820
docs: Update extended-task skill and otel example for concurrent effects
schloerke Sep 30, 2026
d7b444a
test: Add regression tests from the PR review's repros
schloerke Sep 30, 2026
2f33e40
Merge remote-tracking branch 'origin/main' into schloerke/async-pr-21…
schloerke Sep 30, 2026
7b16aa0
docs: Add changelog entry for concurrent reactive effects
schloerke Sep 30, 2026
dc42c0c
fix(reactive): Don't wait in reactive.flush() under a running flush o…
schloerke Sep 30, 2026
5bd513c
fix(reactive): Cache only an async calc's latest, still-valid run
schloerke Sep 30, 2026
0bf0ebe
test: Cover flushing through child tasks and stale async calc runs
schloerke Sep 30, 2026
f283caf
refactor: Drop two checks no behavior depends on
schloerke Sep 30, 2026
fd6b68e
test: Give every concurrency fix a test that fails without it
schloerke Sep 30, 2026
33d1711
docs: Say when reactive.flush() returns without waiting
schloerke Sep 30, 2026
a00ce1a
feat(session): Send each session's outputs in a task of its own
schloerke Sep 30, 2026
f910a1f
feat(reactive): Accept flush requests from other threads
schloerke Sep 30, 2026
fc959b3
docs: Describe per-session output and thread-safe flush requests
schloerke Sep 30, 2026
8e4587e
test: Cover per-session output tasks and flush requests from threads
schloerke Sep 30, 2026
3430e1c
test: Catch Exception, not BaseException, in closed-loop test (flake8…
schloerke Sep 30, 2026
b02a18d
test: Cover reactive values set from a download handler (#1785)
schloerke Oct 1, 2026
f50c245
test: Cover values set and input handled while a download streams
jat255 Oct 2, 2026
0b89ab7
test: Fail any unit test that runs longer than 2 minutes
schloerke Oct 2, 2026
d0865c0
fix(reactive): Discard flush state left by an event loop that has sto…
schloerke Oct 2, 2026
4c7eabc
test: Make the flushed-callback busy test deterministic
schloerke Oct 2, 2026
d481337
fix(reactive): Keep tasks from a stopped event loop alive
schloerke Oct 2, 2026
a9fe446
Merge remote-tracking branch 'origin/main' into schloerke/async-pr-21…
schloerke Oct 2, 2026
7d27a41
test: Wait for the flush task in the dead-event-loop test
jat255 Oct 3, 2026
4be9d36
refactor(reactive): Name rounds, the reactive effect queue, and outpu…
jat255 Oct 8, 2026
4a70f30
Merge remote-tracking branch 'origin/main' into schloerke/async-pr-21…
schloerke Oct 8, 2026
8721519
Merge branch 'schloerke/async-pr-2182-reimpl' into jat255/2508-flush-…
schloerke Oct 8, 2026
38d74df
refactor(reactive): Name the round methods by what they start
jat255 Oct 8, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions .claude/references/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,10 @@ The reactive system is based on a **push-pull** model with three core abstractio
sources (Value) it reads from.
- **Dependents**: Each reactive source maintains a list of downstream consumers
that depend on it. When a source invalidates, it notifies all dependents.
- **ReactiveEnvironment**: Global singleton managing the reactive graph,
execution queue, and flush cycles.
- **ReactiveEnvironment**: Global singleton managing the reactive graph, the
reactive effect queue, and rounds. Its class docstring in
`shiny/reactive/_core.py` defines the terms used for the reactive system:
reactive effect queue, round, idle, cycle, and output flush.

Key implementation details:

Expand All @@ -25,8 +27,9 @@ Key implementation details:
- `Effect_()` is a side-effect that re-executes when dependencies change
- `event()` decorator suppresses reactive dependencies for specific reads
- The reactive graph is built automatically through the Context's dependency tracking
- Execution uses a priority queue to ensure correct invalidation ordering
- In tests, `reactive.flush()` forces a synchronous flush of the reactive graph
- Each round takes effects from the reactive effect queue in priority order
- In tests, `await reactive.flush()` runs rounds until the reactive environment is
idle

## Session Hierarchy

Expand Down
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,24 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Breaking changes

* Reactive effects now run concurrently, following Shiny for R's model. A reactive flush starts each invalidated effect and no longer waits for an `async` effect's awaited part, so a slow `async` effect, calc, or render function no longer delays other sessions. Within a session, input changes and `reactive.invalidate_later()` still wait until all of that session's effects have finished, and outputs are sent together once they have. What app authors may notice:

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do not mention R in public notes or comments. Only mention R in internal comments where we deviate from R shinys's implementation


* `async` effects interleave at each `await` instead of running one after another. `priority` orders when effects start, not when they finish. Runs of the same effect still don't overlap: a re-run waits for the previous run.

* `reactive.lock()` no longer pauses reactive processing: Shiny itself no longer takes it. To change reactive state from another `asyncio` task, set the value directly (a flush is scheduled automatically) and `await reactive.flush()` to wait until the resulting reactive work has finished.

* `await reactive.flush()` returns right away, without waiting for dependents, when called from within an effect or from a task started by an effect that is still running, since waiting there could deadlock. A task that outlives the effect that started it, such as an extended task's body, still waits.

* Message handlers (`session.set_message_handler()`) run in their own task, so a slow handler no longer holds up other messages from the client.

* Each session sends its outputs in a task of its own, once all of its effects have finished. A slow client, or a slow `session.on_flush()` callback, delays only its own session's output, not other sessions' or the next reactive flush.

* Setting a `reactive.value` schedules a flush, including when `set()` is called from another thread. The reactive graph itself still isn't thread-safe, so from another thread use `loop.call_soon_threadsafe(value.set, new_value)`.

(#2508)

### New features

* `ui.sidebar()` gains a `role` parameter (`"form"`, `"search"`, `"complementary"`, or `"region"`) for opt-in ARIA landmark markup (rstudio/bslib#1359). `"complementary"` renders an `<aside>`, other roles render a `<div>` with the corresponding `role` attribute, and landmark roles require an accessible name (from `title`, `aria_label`, or `aria_labelledby`). Note that the default (`role=None`) now renders a neutral `<div>` instead of an `<aside>`. In addition, `ui.page_sidebar()` now places the whole sidebar layout inside the page's `<main>` landmark. (#2526)
Expand All @@ -17,6 +35,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Bug fixes

* Setting a reactive value from a download handler (`@render.download_button`) now updates the outputs and effects that read it right away, including while a streamed download is still sending. The session also keeps handling input during a streamed download. Before, they updated only after the next message from the client. (#1785)

* `Progress.set(0)` now places the bar at `0` relative to `min` and `max`, instead of sending `0` unnormalized. With a negative `min`, such as `Progress(min=-10, max=10)`, a value of `0` showed an empty bar rather than a half-full one. (#2518)

* `ui.input_date()`'s `datesdisabled` now works when `format` is not the default `yyyy-mm-dd`. The dates are now converted on the client the same way `min`/`max` are, instead of being parsed by bootstrap-datepicker with the display `format`. The `data-date-dates-disabled` attribute is replaced by `data-dates-disabled` (and omitted when `datesdisabled` is `None`), and `controller.InputDate.expect_datesdisabled()` checks the new attribute. Requires the updated vendored `shiny.js` (rstudio/shiny#4434). (#2523)
Expand Down
11 changes: 11 additions & 0 deletions pytest.ini
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,14 @@ asyncio_default_fixture_loop_scope=function
testpaths=tests/pytest/
; Note: Browsers are set within `./Makefile`
addopts = --strict-markers --durations=6 --durations-min=5.0 --numprocesses auto
; No test may take more than 2 minutes, so a hang fails instead of running until
; CI cancels the job. The `thread` method kills the hung xdist worker (which
; xdist reports as a failure of that test); its own stack dump is lost in the
; worker's IPC pipe, so faulthandler dumps every thread's stack to stderr first.
; (Same approach as the Playwright targets in the Makefile.)
timeout = 120
timeout_method = thread
faulthandler_timeout = 110
; Tests must not emit stray warnings; assert expected ones with `pytest.warns()`.
; Third-party deprecations we can't act on go in the ignore list below.
filterwarnings =
Expand All @@ -24,4 +32,7 @@ filterwarnings =
; test at fault -- i.e. they make the suite flaky instead of strict.
ignore::ResourceWarning
ignore::pytest.PytestUnraisableExceptionWarning
; When a timeout kills an xdist worker, syrupy warns that the killed worker's
; tests are missing; as an error, that hides the summary naming the hung test.
ignore:.*selected test\(s\) missing from collected items:UserWarning
; verbosity_test_cases=2
18 changes: 12 additions & 6 deletions shiny/.agents/skills/shiny-for-python/references/extended-tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,18 @@
## Overview

A slow computation inside a `@reactive.calc`, `@reactive.effect`, or `@render.*`
blocks the reactive flush - even an `async def` one - so the whole session (and,
for shared resources, other sessions) freezes until it returns. Do NOT put a
slow API call, model inference, or long loop directly in a calc/effect.
keeps its session busy until it returns: that session applies no input changes
and sends no outputs in the meantime. A sync one also blocks every session in the
process, since it holds the event loop. (Other sessions keep running while an
`async def` one awaits.) Do NOT put a slow API call, model inference, or long loop
directly in a calc/effect.

`@reactive.extended_task` runs an **async** function in a background asyncio
task, off the reactive flush. Reactivity keeps processing while it runs; the
result flows back into the reactive graph when it finishes.
task, outside the session's reactive cycle. The session keeps processing inputs
and outputs while it runs; the result flows back into the reactive graph when it
finishes.

Because it runs outside the reactive flush, the task **cannot read reactive
Because it runs outside the reactive graph, the task **cannot read reactive
sources** (`input.x()`, a `reactive.value`, a calc). Read those values *before*
invoking and pass them in as arguments; reading one inside raises an error.

Expand Down Expand Up @@ -121,6 +124,9 @@ current one finishes.
- Long synchronous work (a blocking library call) inside the async task still
blocks the event loop - offload it with `asyncio.to_thread(...)` or an async
client.
- Calling `value.set(...)` from that worker thread -> the reactive graph isn't
thread-safe and it can raise. Return the result to the task instead, or hand the
call to the event loop: `loop.call_soon_threadsafe(value.set, new_value)`.
- Calling `task.result()` outside a reactive context -> no dependency is tracked
and it errors; read it inside a `@render.*`, calc, or effect.
- Expecting a fresh `.invoke()` to interrupt the running task -> it queues
Expand Down
33 changes: 26 additions & 7 deletions shiny/_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
)
from .html_dependencies import _page_deps
from .http_staticfiles import FileResponse, StaticFiles
from .reactive._core import _reactive_environment
from .session._session import AppSession, Inputs, Outputs, Session, session_context
from .types import MISSING, MISSING_TYPE
from .ui._page import DEPS_PLACEHOLDER, PageHtmlDocument, page_html
Expand Down Expand Up @@ -240,7 +241,10 @@ def __init__(

self._sessions: dict[str, AppSession] = {}

# self._sessions_needing_flush: dict[int, AppSession] = {}
# Sessions with outputs, errors, or input messages to send after the next
# round (R's `appsNeedingFlush`).
self._sessions_needing_output_flush: dict[str, AppSession] = {}
self._unregister_output_flush_hook: Optional[Callable[[], None]] = None

self._registered_dependencies: dict[str, HTMLDependency] = {}
self._dependency_handler = starlette.routing.Router()
Expand Down Expand Up @@ -469,13 +473,28 @@ async def _on_session_request_cb(self, request: Request) -> ASGIApp:
return JSONResponse({"detail": "Not Found"}, status_code=404)

# ==========================================================================
Comment thread
schloerke marked this conversation as resolved.
# Flush
# Output flush
# ==========================================================================
def _request_flush(self, session: AppSession) -> None:
# TODO: Until we have reactive domains, because we can't yet keep track
# of which sessions need a flush.
pass
# self._sessions_needing_flush[session.id] = session
def _request_output_flush(self, session: AppSession) -> None:
self._sessions_needing_output_flush[session.id] = session
if self._unregister_output_flush_hook is None:
self._unregister_output_flush_hook = (
_reactive_environment.on_round_finished(self._start_output_flushes)
)
_reactive_environment.request_round()

async def _start_output_flushes(self) -> None:
"""
After each round, start an output flush for each requesting session.

Each session flushes in its own task (see `AppSession._start_output_flush()`),
so a slow client, or a slow `on_flush` callback, delays only its own session:
other sessions' output and the next round don't wait for it.
"""
sessions = list(self._sessions_needing_output_flush.values())
self._sessions_needing_output_flush.clear()
for session in sessions:
session._start_output_flush()

# ==========================================================================
# HTML Dependency stuff
Expand Down
4 changes: 2 additions & 2 deletions shiny/otel/_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -283,8 +283,8 @@ def detached_otel_context() -> Generator[None, None, None]:
from shiny.otel._core import detached_otel_context

with detached_otel_context():
ctx.invalidate()
await flush()
with tracer.start_as_current_span("timer_tick"): # a root span
...
```
"""
token = otel_context.attach(otel_context.Context())
Expand Down
Loading
Loading