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
19 changes: 19 additions & 0 deletions .ai/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -925,6 +925,16 @@ 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.** Every backend takes `topic_ttl` (default 300s,
`DEFAULT_TOPIC_TTL`; `None` keeps topics forever). A topic nobody publishes
to or reads for that long is released, buffer and sequence both; a later
publish starts it over at 1, so a consumer returning with an old cursor gets
`SharedStorageGap`. Local: the engine sweeps idle topics (no publish, head or
poll, and no call holding it) at most every quarter ttl, on the next pub/sub
call. Redis: the publish script `PEXPIRE`s the counter and stream, and each
poll renews both. Diskcache: each message expires `topic_ttl` after publish,
the counter after the last publish or poll.

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

Expand Down Expand Up @@ -1383,6 +1393,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 `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
106 changes: 78 additions & 28 deletions dash/_shared_storage/_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,18 @@
``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 nobody publishes to, polls or subscribes to for ``topic_ttl`` seconds
is dropped, buffer and sequence both, 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 +33,11 @@
# Deployments that need a wider reconnect window set buffer_size explicitly.
DEFAULT_BUFFER = 32

# How long an untouched topic is kept. Long enough to outlast a streaming
# downlink's reconnect and poll grace windows many times over, short enough that
# a busy app does not hold every finished session's frames for hours.
DEFAULT_TOPIC_TTL = 300.0


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


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

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


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


class StoreEngine:
def __init__(self, buffer_size: int = DEFAULT_BUFFER, persistence: Any = None):
def __init__(
self,
buffer_size: int = DEFAULT_BUFFER,
persistence: Any = None,
topic_ttl: Optional[float] = DEFAULT_TOPIC_TTL,
):
self._buffer_size = buffer_size
self._topic_ttl = topic_ttl
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 +165,56 @@
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
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 the ttl. Runs at most
every quarter ttl, so a topic lives between 1 and 1.25 ttl idle."""
if self._topic_ttl is None or now < self._next_sweep:
return
self._next_sweep = now + self._topic_ttl / 4
cutoff = now - self._topic_ttl
idle = [
name
for name, t in self._topics.items()
if t.users == 0 and t.touched < cutoff
]
for name in idle:
del self._topics[name]

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
with self._use(topic) as t:
with t.cond:
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:
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:
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
7 changes: 7 additions & 0 deletions dash/_shared_storage/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,13 @@ class SharedStorageGap(SharedStorageError):
"""


def check_topic_ttl(topic_ttl: Optional[float]) -> None:
if topic_ttl is not None and topic_ttl <= 0:
raise SharedStorageError(
f"topic_ttl must be positive or None, got {topic_ttl!r}"
)


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

Expand Down
22 changes: 17 additions & 5 deletions dash/_shared_storage/diskcache.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,15 +13,17 @@
sequence numbers, each message is stored under its sequence and old sequences
are trimmed to a bounded window, and subscribers poll for sequences past their
cursor. A consumer that falls farther behind than the window gets a
``SharedStorageGap``.
``SharedStorageGap``. With ``topic_ttl``, each message expires that long after
it was published and the counter that long after the topic's last publish or
poll, so abandoned topics leave the cache.
"""
import time
from typing import Any, Optional

from ._codec import decode, encode
from ._engine import DEFAULT_BUFFER, PollResult
from ._engine import DEFAULT_BUFFER, DEFAULT_TOPIC_TTL, PollResult
from ._polling import PollingSubscription
from .base import BaseSharedStorage, Subscription
from .base import BaseSharedStorage, Subscription, check_topic_ttl

# Poll cycle: short so a subscription's close() stays responsive; diskcache has
# no server-side blocking wait, so this is a sleep-poll loop.
Expand All @@ -46,15 +48,19 @@ class DiskcacheSharedStorage(BaseSharedStorage):

Pass an existing ``cache`` (e.g. the one a ``DiskcacheManager`` already
holds) to share one store, or a ``directory`` to open/create one. Values and
published messages must be JSON-compatible.
published messages must be JSON-compatible. ``topic_ttl`` (seconds) releases
a topic nobody has published to or read from for that long; ``None`` keeps
topics until the cache evicts them.
"""

def __init__(
self,
cache: Any = None,
directory: Optional[str] = None,
buffer_size: int = DEFAULT_BUFFER,
topic_ttl: Optional[float] = DEFAULT_TOPIC_TTL,
):
check_topic_ttl(topic_ttl)
diskcache = _require_diskcache()
if cache is not None:
if not isinstance(cache, (diskcache.Cache, diskcache.FanoutCache)):
Expand All @@ -67,6 +73,7 @@ def __init__(
self._cache = diskcache.Cache(directory)
self._owns_cache = True
self._buffer_size = buffer_size
self._topic_ttl = topic_ttl

def close(self) -> None:
if self._owns_cache:
Expand Down Expand Up @@ -101,7 +108,9 @@ def publish(self, topic: str, message: Any) -> None:
payload = encode(message)
with self._cache.transact():
seq = int(self._cache.incr(self._seq(topic)))
self._cache.set(self._msg(topic, seq), payload)
if self._topic_ttl is not None:
self._cache.touch(self._seq(topic), expire=self._topic_ttl)
self._cache.set(self._msg(topic, seq), payload, expire=self._topic_ttl)
evicted = seq - self._buffer_size
if evicted >= 1:
self._cache.delete(self._msg(topic, evicted))
Expand All @@ -110,6 +119,9 @@ def _head(self, topic: str) -> int:
return int(self._cache.get(self._seq(topic), 0))

def _poll(self, topic: str, after_seq: int, timeout: float) -> PollResult:
if self._topic_ttl is not None:
# A reader keeps its topic alive; a missing key is left alone.
self._cache.touch(self._seq(topic), expire=self._topic_ttl)
deadline = time.monotonic() + timeout
while True:
head = self._head(topic)
Expand Down
Loading
Loading