Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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
6 changes: 5 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

* `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.
* `reactive.lock()` no longer pauses reactive processing: Shiny itself no longer takes it (see Deprecations). To change reactive state from another `asyncio` task, set the value directly. A flush is scheduled automatically, so there's no need to `await reactive.flush()`.

* `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.

Expand All @@ -37,6 +37,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

* Navigation links that target a tab panel (e.g. in `ui.navset_tab()`) now carry `aria-controls` pointing at the panel's `id`, alongside the existing `href` (rstudio/bslib#1355). (#2526)

### Deprecations

* `reactive.lock()` is deprecated and now emits a `ShinyDeprecationWarning`. It does nothing else: it returns an `asyncio.Lock` that never blocks, so code that holds it no longer keeps out other code that holds it. Code that relies on it for mutual exclusion should create its own `asyncio.Lock`. Remove `async with reactive.lock():` and the `await reactive.flush()` that usually follows it, and set the reactive value directly. To apply a change only once a session's running effects have finished, use `session.run_once_when_idle()`. (#2520)

### 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)
Expand Down
2 changes: 1 addition & 1 deletion docs/_quartodoc-core.yml
Original file line number Diff line number Diff line change
Expand Up @@ -237,7 +237,6 @@ quartodoc:
- reactive.flush
- reactive.poll
- reactive.file_reader
- reactive.lock
- req
- title: Create and run applications
desc: ""
Expand Down Expand Up @@ -415,6 +414,7 @@ quartodoc:
- render.download
- render.transformer.output_transformer
- render.transformer.resolve_value_fn
- reactive.lock
# `ui.update_navs()` is deprecated (superseded by `ui.update_navset()`) but
# still exists. It is deliberately left undocumented; uncomment to publish it.
# - ui.update_navs
Expand Down
2 changes: 1 addition & 1 deletion docs/_quartodoc-express.yml
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,6 @@ quartodoc:
- reactive.flush
- reactive.poll
- reactive.file_reader
- reactive.lock
- req
- title: Reusable Express code
desc: Create reusable Express code.
Expand Down Expand Up @@ -255,3 +254,4 @@ quartodoc:
desc: ""
contents:
- express.render.download
- reactive.lock
2 changes: 1 addition & 1 deletion shiny/.agents/skills/shiny-for-python/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ reference file before writing code** for that area.

| Topic | Use when | Reference |
|---|---|---|
| Reactivity | A value should recompute or an output update as inputs change; choosing between calc / effect / value; `req`, `isolate`, timers, polling | `references/reactivity.md` |
| Reactivity | A value should recompute or an output update as inputs change; choosing between calc / effect / value; `req`, `isolate`, timers, polling, setting values from background tasks | `references/reactivity.md` |
| Express mode | Writing or converting an Express app (`from shiny.express import ...`); context-manager layout; `page_opts`, `@expressify` | `references/express.md` |
| Modules (Core) | A reusable, repeatable UI+server component in a Core app; avoiding input/output id collisions across copies | `references/modules-core.md` |
| Modules (Express) | The same reusable-component need in an Express app, via the single `@module` decorator | `references/modules-express.md` |
Expand Down
46 changes: 41 additions & 5 deletions shiny/.agents/skills/shiny-for-python/references/reactivity.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,10 +144,42 @@ def data2():
Declare `poll`/`file_reader` at module top level to share one cache across
sessions.

To change reactive state from a background task or a custom timer, the way
input changes and `invalidate_later` do, hand the change to the session. It runs
once, after the session's running effects finish, so they keep seeing stable
values:
To change reactive state from a background `asyncio` task, call `.set()`
directly. Setting a value schedules a flush, so there's no lock to take and no
`reactive.flush()` to await. (From another thread, use
`loop.call_soon_threadsafe(v.set, x)`.) For a producer shared by all sessions
(e.g. one polling an API), define it at module top level and start it from the
first session, once the event loop is running. Keep a reference to its task,
since the event loop holds tasks only weakly:

```python
import asyncio
import contextvars
import logging

from shiny import reactive

latest = reactive.value(None)
_tasks: set[asyncio.Task[None]] = set()

async def _produce():
while True:
try:
latest.set(await fetch()) # your async data source
except Exception:
logging.exception("fetch failed; retrying") # keep the producer alive
await asyncio.sleep(10)

def server(input, output, session):
if not _tasks:
# A fresh context keeps this session's context out of the shared task.
_tasks.add(contextvars.Context().run(asyncio.create_task, _produce()))
```

A session's effects that are paused at an `await` see a direct `.set()` right
away. To apply the change the way input changes and `invalidate_later` do, hand
it to the session. It runs once, after the session's running effects finish, so
they keep seeing stable values:

```python
session.run_once_when_idle(lambda: latest.set(new_value))
Expand All @@ -167,7 +199,7 @@ session.run_once_when_idle(lambda: latest.set(new_value))
| Wait for / validate a value | `req(x)`, `req(x, cancel_output=True)` |
| Read without depending | `with reactive.isolate():` |
| Timer | `reactive.invalidate_later(secs)` |
| Change state between a session's cycles | `session.run_once_when_idle(fn)` |
| Set a value from a background task | `v.set(x)`; between cycles: `session.run_once_when_idle(fn)` |
| Poll a data source / file | `@reactive.poll(...)` / `@reactive.file_reader(...)` |
| Force sync execution in tests | `reactive.flush()` |

Expand All @@ -186,3 +218,7 @@ session.run_once_when_idle(lambda: latest.set(new_value))
`isolate()` -> self-invalidating loop; wrap the read in `reactive.isolate()`.
- `while`/`sleep` loop to watch a DB or file -> use `reactive.poll` or
`reactive.file_reader`.
- `async with reactive.lock():` around a `.set()`, followed by
`await reactive.flush()` -> `reactive.lock()` is deprecated and does nothing,
and the flush adds nothing but a wait on every session's effects. Call
`.set()` alone.
100 changes: 70 additions & 30 deletions shiny/reactive/_core.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,15 @@
import typing
import warnings
from contextvars import ContextVar
from typing import TYPE_CHECKING, Awaitable, Callable, Generator, Optional, TypeVar
from typing import (
TYPE_CHECKING,
Awaitable,
Callable,
Generator,
Literal,
Optional,
TypeVar,
)

from .. import _utils
from .._datastructures import PriorityQueueFIFO
Expand Down Expand Up @@ -170,7 +178,6 @@ def __init__(self) -> None:
)
self._next_id: int = 0
self._effect_queue: PriorityQueueFIFO[Context] = PriorityQueueFIFO()
self._lock: Optional[asyncio.Lock] = None
self._round_finished_callbacks = _utils.AsyncCallbacks()
self._round_requested: bool = False
# The event loop Shiny runs on, for requests made from other threads.
Expand All @@ -186,21 +193,6 @@ def __init__(self) -> None:
# Tasks stranded on an event loop that stopped (see `_adopt_loop()`).
self._abandoned_tasks: set[asyncio.Task[None]] = set()

@property
def lock(self) -> asyncio.Lock:
"""
Lock that protects this ReactiveEnvironment. It must be lazily created, because
at the time the module is loaded, there generally isn't a running asyncio loop
yet. This causes the asyncio.Lock to be created with a different loop than it
will be invoked from later; when that happens, acquire() will succeed if there's
no contention, but throw a "hey you're on the wrong loop" error if there is.
"""
if self._lock is None:
# Ensure we have a loop; get_running_loop() throws an error if we don't
asyncio.get_running_loop()
self._lock = asyncio.Lock()
return self._lock

def next_id(self) -> int:
"""Return the next available id"""
id = self._next_id
Expand Down Expand Up @@ -472,10 +464,13 @@ async def flush() -> None:
-------
You shouldn't ever need to call this function inside of a Shiny app. It's only
useful for testing and running reactive code interactively in the console.
Setting a :class:`~shiny.reactive.value` already schedules a flush, so code that
sets one (from a background task, for example) doesn't need to call this.

Runs rounds until the reactive environment is idle: the reactive effect queue is
empty, every started effect (including its async part) has finished, and each
session's resulting outputs have been sent.
session's resulting outputs have been sent. That covers every session, not only
the caller's, so in an app it can wait on other sessions' slow effects.

Note
----
Expand Down Expand Up @@ -521,22 +516,67 @@ def on_flushed(
return _reactive_environment.on_round_finished(func, once)


@no_example()
def lock() -> asyncio.Lock:
class _NoOpLock(asyncio.Lock):
"""
What :func:`~shiny.reactive.lock` returns: an :class:`asyncio.Lock` that never
blocks, so any number of holders can hold it at once.
"""
A process-wide lock, kept for backward compatibility.

Shiny no longer takes this lock itself, so holding it doesn't pause reactive
processing or serialize effects. Code that holds it only excludes other code
that also holds it.
async def acquire(self) -> Literal[True]:
return True

To change reactive state from a different :class:`~asyncio.Task` than the one
running the Shiny :class:`~shiny.Session`, set the :class:`~reactive.value`
directly; a round is scheduled automatically. Await
:func:`~shiny.reactive.flush` to wait until the resulting reactive work has
finished.
def release(self) -> None:
pass

def locked(self) -> bool:
return False


@no_example()
def lock() -> asyncio.Lock:
"""
Deprecated. Set a reactive value directly instead.

Apart from emitting a deprecation warning, this function does nothing. It
returns an :class:`asyncio.Lock` that never blocks, so holding it neither pauses
reactive processing nor keeps out other code that holds it. Code that needs
mutual exclusion of its own should create its own :class:`asyncio.Lock`.

To change reactive state from outside a reactive context (for example, from a
background :class:`asyncio.Task`), set the :class:`~shiny.reactive.value`
directly. Setting it schedules a round, so there's no need to call
:func:`~shiny.reactive.flush` afterwards:

```python
# Deprecated
async with reactive.lock():
current_query.set(query)
await reactive.flush()

# Use instead
current_query.set(query)
```

Effects of a session that are paused at an ``await`` see the new value as soon
as it is set. To apply the change only once the session's effects have
finished, as Shiny does for input changes from the client, use
:meth:`~shiny.Session.run_once_when_idle`:

```python
session.run_once_when_idle(lambda: current_query.set(query))
```
"""
return _reactive_environment.lock
# Imported here because `shiny._deprecated` imports `shiny.reactive`.
from .._deprecated import warn_deprecated

warn_deprecated(
"reactive.lock() is deprecated, does nothing, and will be removed in a future "
"version of shiny. To change reactive state from a background task, set the "
"reactive value directly: a flush is scheduled automatically, so "
"`await reactive.flush()` is not needed. To apply the change once the "
"session's effects have finished, use `session.run_once_when_idle()`."
)
return _NoOpLock()


_timer_tasks: set[asyncio.Task[None]] = set()
Expand Down
63 changes: 63 additions & 0 deletions tests/pytest/test_concurrency.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

from shiny import App, Inputs, Outputs, Session, module, reactive, render, ui
from shiny._connection import MockConnection
from shiny._deprecated import ShinyDeprecationWarning
from shiny.bookmark._bookmark import BookmarkApp
from shiny.bookmark._restore_state import RestoreContext
from shiny.reactive._core import (
Expand Down Expand Up @@ -2658,3 +2659,65 @@ def boom() -> None:
assert await wait_until(lambda: c.session._has_run_session_ended_tasks)
finally:
await c.close()


@pytest.mark.asyncio
async def test_lock_warns_on_every_call():
for _ in range(2):
with pytest.warns(ShinyDeprecationWarning, match="reactive.lock") as record:
reactive.lock()
# The warning points at the code that called `lock()`.
assert record[0].filename == __file__


@pytest.mark.asyncio
async def test_lock_does_not_exclude_other_holders():
with pytest.warns(ShinyDeprecationWarning):
lock = reactive.lock()
inside = 0
both_inside = asyncio.Event()

async def hold() -> None:
nonlocal inside
async with lock:
inside += 1
if inside == 2:
both_inside.set()
await asyncio.wait_for(both_inside.wait(), TIMEOUT)

await asyncio.gather(hold(), hold())
assert await lock.acquire() is True
assert await lock.acquire() is True
assert not lock.locked()
lock.release()
assert isinstance(lock, asyncio.Lock)

with pytest.raises(ValueError):
async with lock:
raise ValueError("raised inside the lock")


@pytest.mark.asyncio
async def test_value_set_from_background_task_reruns_session_effect():
# The replacement for `async with reactive.lock(): v.set(x); await flush()`:
# setting the value from a task outside any session is enough.
shared = reactive.value(0)
seen: list[int] = []

def server(input: Inputs, output: Outputs, session: Session) -> None:
@reactive.effect
def _():
seen.append(shared())

c = await started_client(server)
try:
assert await wait_until(lambda: seen == [0])

async def producer() -> None:
await asyncio.sleep(0)
shared.set(1)

await asyncio.create_task(producer())
assert await wait_until(lambda: seen == [0, 1])
finally:
await c.close()
Loading