Skip to content
Merged
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
38 changes: 36 additions & 2 deletions .ai/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -850,7 +850,7 @@ outputs. Messages are encoded with `msgspec` (msgpack), a hard dependency.
| `get(key, default=None)` | Read a value |
| `set(key, value, ttl=None)` | Write a value; `ttl` = optional lifetime in seconds |
| `delete(key)` | Remove a key |
| `publish(topic, message)` | Append a message to a topic |
| `publish(topic, message, ttl=None)` | Append a message to a topic; `ttl` = release the topic after that long idle |
| `subscribe(topic, replay_from=None)` | Return a `Subscription` |

**Key expiry (TTL).** `set(key, value, ttl=<seconds>)` gives a key a bounded
Expand Down Expand Up @@ -925,6 +925,27 @@ deployment is its own owner and pays no socket overhead. Knobs: `namespace`
purpose, since each topic retains that many arbitrary payloads; raise it for a
wider reconnect window).

**Topic lifetime.** `publish(topic, message, ttl=None)` takes an optional ttl
(seconds, at least `MIN_TOPIC_TTL` = 1s, so a reconnecting reader is not
outrun). A topic with a ttl is released, buffer and sequence both, once nobody
has published to or read from it for that long; a later publish starts it over
at 1, so a consumer returning with an old cursor gets `SharedStorageGap`.
Without a ttl a topic lives as long as the store. The latest publish's ttl
wins. Streaming publishes with `STREAM_TOPIC_TTL` (300s); user topics default
to no ttl. Every backend expires a topic as a whole, never message by message
while it is read (`tests/shared_storage/test_topic_ttl.py` runs the same cases
on all three). Local: the engine records the ttl on the topic and sweeps idle
ones (no publish, head or poll, and no call holding it) at most once a second,
on the next pub/sub call. The sweep also drops empty, unheld topics, such as
the one a returning reader's poll recreates after a release. Redis: the ttl sits in a third key; the publish
script `PEXPIRE`s all three to 1.25 ttl, and a poll script renews them. A poll
on such a topic blocks in `XREAD` for at most a third of that lifetime, not the
usual 5s, so a waiting reader renews it in time. Diskcache: the ttl
sits in a key; on publish or poll, once less than one ttl is left, every key of
the topic (counter, ttl, buffered messages) is renewed to 1.25 ttl, and at
once when a publish changes the ttl. A poll on such a topic waits at most half
a ttl before it renews again.

*Durability* is controlled by `mode` (the key/value store only — pub/sub is
always transient):

Expand Down Expand Up @@ -982,7 +1003,11 @@ election only reaches processes in the same network + filesystem namespace.
in-tree or as a **separate package** — implements:

- `get(key, default)` / `set(key, value)` / `delete(key)` — JSON-compatible values
- `publish(topic, message)` and `subscribe(topic, replay_from=None) -> Subscription`
- `publish(topic, message, ttl=None)` and `subscribe(topic, replay_from=None) -> Subscription`
- topic expiry: with a `ttl`, release the topic as a whole once nobody has
published to or read from it for that long, never message by message while
it is read; reject a ttl under `MIN_TOPIC_TTL` (run `test_topic_ttl.py`
against a new backend)
- optional `start()` / `close()` (idempotent, called once per worker)

and returns a `Subscription` (`__iter__` / `__aiter__` / `close`) that raises
Expand Down Expand Up @@ -1383,6 +1408,15 @@ forged -- otherwise a client could read or inject into another page's topic.
Across worker processes every worker must resolve the same signing secret
(`secret_key`).

Each run of streams (from the first stream after idle until none is in flight)
also carries `&downlinkId=`, picked fresh by the client, and the connection id
is `<end_id>:<downlinkId>`: every run gets its own topic and the client's
cursor restarts at 0. Without it a page that sat idle past `STREAM_TOPIC_TTL` would
resume a cursor into a topic the store released, get `{reset: true}`, and fail
its new stream. The id only partitions the page's own space, so it is not
signed (just checked against `[A-Za-z0-9_-]{1,64}`). The downlink lifecycle
record (`connection_key`) stays keyed on the `end_id` alone.

The downlink is hosted in a SharedWorker (`dash-stream-worker.js`, served like
the WebSocket worker; `config.stream.worker_url`) so **one connection per
browser** serves every tab: browsers cap HTTP/1.1 connections per host at
Expand Down
19 changes: 18 additions & 1 deletion dash/_callback.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
import hashlib
import inspect
import logging
import re
import warnings
from functools import wraps
from typing import Callable, Optional, Any, List, Tuple, Union, Dict, TypeVar, cast
Expand Down Expand Up @@ -485,8 +486,24 @@ def get_stream_connection_id() -> "str | None":
worker must resolve the same secret: set a ``secret_key`` on the server, or
cross-worker stream requests will not verify. Single-process apps are fine
with no configuration.

The renderer also sends a ``downlinkId`` it picks fresh for each run of
streams, giving every run its own topic (``<end_id>:<downlinkId>``). A run
then never resumes a cursor into a topic the store released while the page
sat idle. It only partitions the page's own space, so it needs no signing.
"""
return get_request_end_id(_get_signing_secret())
end_id = get_request_end_id(_get_signing_secret())
if end_id is None:
return None
downlink_id = get_app().backend.request_adapter().args.get("downlinkId")
if not downlink_id:
return end_id
if not _DOWNLINK_ID_RE.fullmatch(downlink_id):
return None
return f"{end_id}:{downlink_id}"


_DOWNLINK_ID_RE = re.compile(r"[A-Za-z0-9_-]{1,64}")


def _get_signing_secret() -> bytes:
Expand Down
110 changes: 80 additions & 30 deletions dash/_shared_storage/_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,19 @@
``poll`` blocks a thread; ``apoll`` parks an asyncio task on a future that
``publish`` resolves from whichever thread it runs on, so an ASGI server can
hold thousands of subscriptions without an executor thread each.

A topic published with a ``ttl`` is dropped, buffer and sequence both, once
nobody has published to, polled or subscribed to it for that long, so
per-session topics do not pile up for the life of the process. A later publish
starts it over at sequence 1.
"""

import asyncio
import contextlib
import threading
import time
from collections import deque
from typing import Any, Deque, Dict, List, NamedTuple, Optional, Tuple
from typing import Any, Deque, Dict, Iterator, List, NamedTuple, Optional, Tuple

# Per-topic replay buffer size. Kept small by default because messages are
# arbitrary user payloads and every topic retains up to this many -- unbounded
Expand All @@ -28,6 +34,10 @@
# Deployments that need a wider reconnect window set buffer_size explicitly.
DEFAULT_BUFFER = 32

# How often idle topics are looked for, so a topic outlives its ttl by at most
# this much.
_SWEEP_INTERVAL = 1.0


class PollResult(NamedTuple):
messages: List[Any]
Expand All @@ -39,14 +49,19 @@


class _Topic: # pylint: disable=too-few-public-methods
__slots__ = ("seq", "buf", "cond", "waiters")
__slots__ = ("seq", "buf", "cond", "waiters", "users", "touched", "ttl")

def __init__(self, maxlen: int):
def __init__(self, maxlen: int, now: float):
self.seq = 0
self.buf: Deque[Tuple[int, Any]] = deque(maxlen=maxlen)
self.cond = threading.Condition()
# asyncio tasks parked in apoll(), woken by the next publish/close.
self.waiters: List[_Waiter] = []
# Calls currently holding this topic (a blocked poll among them), and
# when the last one let go. Both guarded by the engine's _topics_lock.
self.users = 0
self.touched = now
self.ttl: Optional[float] = None


def _wake(fut: "asyncio.Future[None]") -> None:
Expand All @@ -55,8 +70,13 @@


class StoreEngine:
def __init__(self, buffer_size: int = DEFAULT_BUFFER, persistence: Any = None):
def __init__(
self,
buffer_size: int = DEFAULT_BUFFER,
persistence: Any = None,
):
self._buffer_size = buffer_size
self._next_sweep = 0.0
# key -> (value, expiry). expiry is a monotonic deadline, or None for
# no TTL. Expired entries are dropped lazily on the next read.
self._data: Dict[str, Tuple[Any, Optional[float]]] = {}
Expand Down Expand Up @@ -144,30 +164,57 @@
return out

# --- pub/sub -----------------------------------------------------------
def _topic(self, name: str) -> _Topic:
@contextlib.contextmanager
def _use(self, name: str) -> Iterator[_Topic]:
"""Hold a topic for one call. A held topic is never swept, so a
publish cannot land in a topic that was just dropped from the map."""
now = time.monotonic()
with self._topics_lock:
self._sweep(now)
topic = self._topics.get(name)
if topic is None:
topic = self._topics[name] = _Topic(self._buffer_size)
return topic

def publish(self, topic: str, message: Any) -> int:
t = self._topic(topic)
with t.cond:
t.seq += 1
t.buf.append((t.seq, message))
t.cond.notify_all()
waiters, t.waiters = t.waiters, []
seq = t.seq
topic = self._topics[name] = _Topic(self._buffer_size, now)
topic.users += 1
try:
yield topic
finally:
with self._topics_lock:
topic.users -= 1
topic.touched = time.monotonic()

def _sweep(self, now: float) -> None:
"""Under ``_topics_lock``: drop topics idle past their ttl, and empty
ones a read left behind, which are no different from a missing one."""
if now < self._next_sweep:
return
self._next_sweep = now + _SWEEP_INTERVAL
idle = [
name
for name, t in self._topics.items()
if t.users == 0
and (t.seq == 0 or (t.ttl is not None and t.touched < now - t.ttl))
]
Comment on lines +191 to +196

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

A placeholder topic created by a read operation does not have any ttl, so a cleanup sweep will never release it. This can happen in ordinary use:
A user loses their connection mid-stream (or closes the laptop) for longer than the grace window plus the ttl. The pump gives up and publishes its last frame, and the topic is released after the ttl. When the browser returns it polls the old topic, poll creates it again empty with no ttl, and the client gets {reset: true}. That topic is never released.
Since it should always be safe to drop an empty topic, maybe the sweep can also drop topics with no messages and no readers after a short idle time?

for name in idle:
del self._topics[name]

def publish(self, topic: str, message: Any, ttl: Optional[float] = None) -> int:
with self._use(topic) as t:

Check warning on line 201 in dash/_shared_storage/_engine.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Combine these nested "with" statements into a single "with" with multiple contexts.

See more on https://sonarcloud.io/project/issues?id=plotly_dash&issues=AaEXFzrJdqjdxOhcRso_&open=AaEXFzrJdqjdxOhcRso_&pullRequest=4016
with t.cond:
t.ttl = ttl
t.seq += 1
t.buf.append((t.seq, message))
t.cond.notify_all()
waiters, t.waiters = t.waiters, []
seq = t.seq
for loop, fut in waiters:
loop.call_soon_threadsafe(_wake, fut)
return seq

def head_seq(self, topic: str) -> int:
"""Current highest sequence -- where a fresh subscription starts."""
t = self._topic(topic)
with t.cond:
return t.seq
with self._use(topic) as t:

Check warning on line 215 in dash/_shared_storage/_engine.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Combine these nested "with" statements into a single "with" with multiple contexts.

See more on https://sonarcloud.io/project/issues?id=plotly_dash&issues=AaEXFzrJdqjdxOhcRspA&open=AaEXFzrJdqjdxOhcRspA&pullRequest=4016
with t.cond:
return t.seq

def _ready(self, t: _Topic, after_seq: int) -> Optional[PollResult]:
"""Under ``t.cond``: the result available right now, or None to wait."""
Expand Down Expand Up @@ -196,22 +243,25 @@
elapsed (caller re-polls) or the engine closed. ``gap`` is True when the
next expected message was already evicted from the buffer.
"""
t = self._topic(topic)
deadline = time.monotonic() + timeout
with t.cond:
while True:
res = self._ready(t, after_seq)
if res is not None:
return res
remaining = deadline - time.monotonic()
if remaining <= 0:
return PollResult([], after_seq, False)
t.cond.wait(remaining)
with self._use(topic) as t:

Check warning on line 247 in dash/_shared_storage/_engine.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Combine these nested "with" statements into a single "with" with multiple contexts.

See more on https://sonarcloud.io/project/issues?id=plotly_dash&issues=AaEXFzrJdqjdxOhcRspB&open=AaEXFzrJdqjdxOhcRspB&pullRequest=4016
with t.cond:
while True:
res = self._ready(t, after_seq)
if res is not None:
return res
remaining = deadline - time.monotonic()
if remaining <= 0:
return PollResult([], after_seq, False)
t.cond.wait(remaining)

async def apoll(self, topic: str, after_seq: int, timeout: float) -> PollResult:
""":meth:`poll` for asyncio: parks the task on a future instead of
blocking a thread; ``publish`` (from any thread) or ``close`` wakes it."""
t = self._topic(topic)
with self._use(topic) as t:
return await self._apoll(t, after_seq, timeout)

async def _apoll(self, t: _Topic, after_seq: int, timeout: float) -> PollResult:

Check warning on line 264 in dash/_shared_storage/_engine.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Remove this "timeout" parameter and use a timeout context manager instead.

See more on https://sonarcloud.io/project/issues?id=plotly_dash&issues=AaDo3b6D0kONGPq9vh3p&open=AaDo3b6D0kONGPq9vh3p&pullRequest=4016
loop = asyncio.get_running_loop()
deadline = time.monotonic() + timeout
while True:
Expand Down
3 changes: 2 additions & 1 deletion dash/_shared_storage/_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,8 @@ def _dispatch(self, req): # pylint: disable=too-many-return-statements
self._engine.delete(req[1])
return ("ok", None)
if op == "publish":
return ("ok", self._engine.publish(req[1], req[2]))
ttl = req[3] if len(req) > 3 else None
return ("ok", self._engine.publish(req[1], req[2], ttl))
if op == "head":
return ("ok", self._engine.head_seq(req[1]))
if op == "poll":
Expand Down
31 changes: 27 additions & 4 deletions dash/_shared_storage/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import abc
import asyncio
import math
from typing import Any, AsyncIterator, Iterator, List, Optional, Tuple


Expand All @@ -35,6 +36,18 @@ class SharedStorageGap(SharedStorageError):
"""


# Shortest topic ttl a publish accepts. A reconnecting reader can easily be
# gone for a second, and a shorter ttl would drop its topic before it is back.
MIN_TOPIC_TTL = 1.0
Comment on lines +39 to +41

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

I suspect for Redis, 1s might be too low. Not sure if you think it will be a problem in practice, but I'll call it out here for your consideration.



def check_topic_ttl(ttl: Optional[float]) -> None:
if ttl is not None and not (math.isfinite(ttl) and ttl >= MIN_TOPIC_TTL):
raise SharedStorageError(
f"topic ttl must be None or a finite {MIN_TOPIC_TTL}s or more, got {ttl!r}"
)


class Subscription(abc.ABC):
"""A live, ordered view of a topic.

Expand Down Expand Up @@ -133,8 +146,16 @@ def delete(self, key: str) -> None:
...

@abc.abstractmethod
def publish(self, topic: str, message: Any) -> None:
"""Append ``message`` to ``topic``; delivered to every current subscriber."""
def publish(self, topic: str, message: Any, ttl: Optional[float] = None) -> None:
"""Append ``message`` to ``topic``; delivered to every current subscriber.

``ttl`` (seconds, at least ``MIN_TOPIC_TTL``) lets the store release the
topic, buffer and sequence both, once nobody has published to or read
from it for that long. A reader that keeps polling keeps it alive. A
later publish starts the topic over, so a consumer returning with an old
cursor gets ``SharedStorageGap``. ``None`` (the default) keeps the topic
for the life of the store. The latest publish's ``ttl`` applies.
"""

# --- asyncio variants --------------------------------------------------
# Code running on an event loop (ASGI request handlers, the streaming
Expand All @@ -155,9 +176,11 @@ async def aset(self, key: str, value: Any, ttl: Optional[float] = None) -> None:
async def adelete(self, key: str) -> None:
await asyncio.get_running_loop().run_in_executor(None, self.delete, key)

async def apublish(self, topic: str, message: Any) -> None:
async def apublish(
self, topic: str, message: Any, ttl: Optional[float] = None
) -> None:
await asyncio.get_running_loop().run_in_executor(
None, self.publish, topic, message
None, self.publish, topic, message, ttl
)

@abc.abstractmethod
Expand Down
Loading
Loading