From 4cc49f87009a0c94770d2a2aa057a6b30f520044 Mon Sep 17 00:00:00 2001 From: philippe Date: Wed, 30 Sep 2026 14:30:19 -0400 Subject: [PATCH 1/2] Add DASH_SECRET_KEY and DASH_SHARED_STORAGE platform hooks A hosting platform can now give every worker and pod the same signing secret (DASH_SECRET_KEY) and pick the shared-storage backend (DASH_SHARED_STORAGE) without editing the app. Without a shared secret, multi-worker stream requests 403 silently; the first failure in each process now logs a warning. --- .ai/ARCHITECTURE.md | 24 +++- CHANGELOG.md | 3 + dash/_callback.py | 43 ++++-- dash/_configs.py | 2 + dash/_shared_storage/_env.py | 61 +++++++++ dash/dash.py | 62 ++++++--- tests/shared_storage/test_dash_integration.py | 21 ++- tests/shared_storage/test_env_config.py | 127 ++++++++++++++++++ tests/streaming/test_stream_transport.py | 34 +++++ tests/streaming/test_stream_wsgi.py | 104 ++++++++++++-- .../unit/test_background_callback_signing.py | 35 +++++ 11 files changed, 474 insertions(+), 42 deletions(-) create mode 100644 dash/_shared_storage/_env.py create mode 100644 tests/shared_storage/test_env_config.py diff --git a/.ai/ARCHITECTURE.md b/.ai/ARCHITECTURE.md index d4bd2ace56..cd21890888 100644 --- a/.ai/ARCHITECTURE.md +++ b/.ai/ARCHITECTURE.md @@ -752,6 +752,8 @@ Special handling for Colab: | `DASH_PRUNE_ERRORS` | Simplify tracebacks | | `HOST` | Server host | | `PORT` | Server port | +| `DASH_SECRET_KEY` | Signing secret for page, background and stream tokens when `server.secret_key` is unset | +| `DASH_SHARED_STORAGE` | Shared-storage backend when `shared_storage=` is not passed (see Shared Storage) | ## Stores and Client-Side State @@ -976,6 +978,16 @@ election only reaches processes in the same network + filesystem namespace. This works at 1 pod and fragments silently once it scales — use `RedisSharedStorage` (one Redis shared by all pods) instead. +A hosting platform can switch the backend without editing the app through +`DASH_SHARED_STORAGE` (`_shared_storage/_env.py`), read only when the app did +not pass `shared_storage=` (the default is a sentinel, so an explicit argument, +`None` included, always wins): `local`, `none`, `diskcache:///abs/path`, or a +`redis://` / `rediss://` URL. `cluster://` is reserved and raises; anything +else raises `InvalidConfig` at construction. The value becomes a zero-argument +factory, so nothing is built or connected until `app.shared_storage` is first +read, but a missing extra (`dash[redis]`, `dash[diskcache]`) fails at +construction with the backend's own `ImportError`. + ### Custom and Out-of-Tree Backends `BaseSharedStorage` is the stable, public extension point. A backend — shipped @@ -1010,6 +1022,7 @@ same way as a built-in: `Dash(shared_storage=PostgresSharedStorage(...))`. | `_shared_storage/_polling.py` | `PollingSubscription`: shared poll-loop subscription for the diskcache/Redis backends | | `_shared_storage/_transport.py` | Length-prefixed, token-gated socket transport (local backend) | | `_shared_storage/_codec.py` | msgspec msgpack codec (data-only) | +| `_shared_storage/_env.py` | `DASH_SHARED_STORAGE` parsing | | `dash.py` | `shared_storage` constructor arg + lazy `app.shared_storage` property | | `_callback_context.py` | `dash.ctx.shared_storage` accessor | @@ -1380,8 +1393,15 @@ The connection id is never chosen by the client: every stream request rides on `?endId=`, the server-signed per-page-load token, and the backend derives the id from it (`get_stream_connection_id`), answering 403 when it is missing or 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`). +Across worker processes every worker must resolve the same signing secret. +`_get_signing_secret` resolves, in order: `server.secret_key`, then +`DASH_SECRET_KEY` (for Dash's signing only, never copied onto +`server.secret_key`, so Flask sessions are untouched), then a secret persisted +in the background-callback store, then a per-process random one. With the last, +tokens only verify on the worker that issued them: stream requests 403 on the +other workers, and the first failure in each process logs a warning pointing at +`DASH_SECRET_KEY` (`_warn_unverified_stream_token`). A request with no token at +all is not logged. The downlink is hosted in a SharedWorker (`dash-stream-worker.js`, served like the WebSocket worker; `config.stream.worker_url`) so **one connection per diff --git a/CHANGELOG.md b/CHANGELOG.md index fd04ee0987..ea8897c7e7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,8 @@ This project adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] ### Added +- Add `DASH_SECRET_KEY` so a hosting platform can give every worker and pod the same signing secret for page, background and stream tokens without editing the app (`server.secret_key` still wins and is not changed). +- Add `DASH_SHARED_STORAGE` to pick the shared-storage backend when the app does not pass `shared_storage=`: `local`, `none`, `diskcache:///abs/path`, or a `redis://` / `rediss://` URL. - [#3976](https://github.com/plotly/dash/pull/3976) Add a new `scrollToTop` prop to `dcc.Link` to control whether the page scrolls to the top after client-side navigation. It defaults to `True` to preserve the existing behavior. Fixes [#3974](https://github.com/plotly/dash/issues/3974). - [#3947](https://github.com/plotly/dash/pull/3947) Make `plotly-cloud` a default install dependency of Dash instead of an optional extra, so the `plotly` CLI and Dash's cloud integration work out of the box. The `dash[cloud]` extra is kept for backward compatibility. - [#3930](https://github.com/plotly/dash/pull/3930) Add shared storage: a backend-agnostic cross-process state manager (key/value with optional TTL, plus ordered replayable pub/sub) on every app via `dash.ctx.shared_storage`, started lazily and disabled with `shared_storage=None`. Ships `LocalSharedStorage` (default, in-memory with optional disk persistence), `DiskcacheSharedStorage`, and `RedisSharedStorage` for horizontally-scaled deployments; see `.ai/ARCHITECTURE.md`. @@ -21,6 +23,7 @@ This project adheres to [Semantic Versioning](https://semver.org/). - [#3646](https://github.com/plotly/dash/pull/3646) Remove React 16 support (`16.14.0` is no longer an accepted value for `REACT_VERSION` / `_set_react_version`). ### Changed +- Log a warning, once per process, when a streaming request is refused because its stream token failed verification, which on multi-worker deployments usually means the workers do not share a signing secret. - [#3987](https://github.com/plotly/dash/pull/3987) Forward FastAPI reload scope options (`reload_dirs`, `reload_excludes`, and `reload_includes`) to Uvicorn when reloading. - [#3986](https://github.com/plotly/dash/pull/3986) Adjust `_run_before_hooks` in the `fastapi` backend to honor a response returned by a `before_request` function, matching the `flask` backend's behavior. diff --git a/dash/_callback.py b/dash/_callback.py index 1b92907df1..6854c9720e 100644 --- a/dash/_callback.py +++ b/dash/_callback.py @@ -2,6 +2,7 @@ import hashlib import inspect import logging +import threading import warnings from functools import wraps from typing import Callable, Optional, Any, List, Tuple, Union, Dict, TypeVar, cast @@ -464,11 +465,14 @@ def get_request_end_id(secret: bytes): request; this verifies the signature and returns the underlying end_id so background handles can be checked against it. """ + return _callback_signing.unsign( + secret, _callback_signing.END_SCOPE, _request_end_token() + ) + + +def _request_end_token(): adapter = get_app().backend.request_adapter() - if not adapter: - return None - token = adapter.args.get("endId") - return _callback_signing.unsign(secret, _callback_signing.END_SCOPE, token) + return adapter.args.get("endId") if adapter else None def get_stream_connection_id() -> "str | None": @@ -482,11 +486,34 @@ def get_stream_connection_id() -> "str | None": forged token yields ``None``, and the backend refuses the request (403). ``end_id`` is signed with the server secret, so across worker processes every - 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. + worker must resolve the same secret: set a ``secret_key`` on the server or + the ``DASH_SECRET_KEY`` environment variable, or cross-worker stream requests + will not verify. Single-process apps are fine with no configuration. """ - return get_request_end_id(_get_signing_secret()) + connection_id = get_request_end_id(_get_signing_secret()) + if connection_id is None and _request_end_token(): + _warn_unverified_stream_token() + return connection_id + + +_stream_token_warning_lock = threading.Lock() +_stream_token_warned = False + + +def _warn_unverified_stream_token(): + # Once per process: a secret mismatch fails every cross-worker request. + global _stream_token_warned # pylint: disable=global-statement + with _stream_token_warning_lock: + if _stream_token_warned: + return + _stream_token_warned = True + get_app().logger.warning( + "A streaming request was refused (403): its stream token failed " + "verification. On multi-worker or multi-pod deployments this usually " + "means the workers do not share a signing secret (or the server " + "restarted since the page loaded); set server.secret_key or the " + "DASH_SECRET_KEY environment variable." + ) def _get_signing_secret() -> bytes: diff --git a/dash/_configs.py b/dash/_configs.py index 0e1ab75505..aa3399e0df 100644 --- a/dash/_configs.py +++ b/dash/_configs.py @@ -35,6 +35,8 @@ def load_dash_env_vars(): "DASH_COMPRESS", "DASH_MCP_ENABLED", "DASH_MCP_PATH", + "DASH_SECRET_KEY", + "DASH_SHARED_STORAGE", "HOST", "PORT", ) diff --git a/dash/_shared_storage/_env.py b/dash/_shared_storage/_env.py new file mode 100644 index 0000000000..61001f8dd6 --- /dev/null +++ b/dash/_shared_storage/_env.py @@ -0,0 +1,61 @@ +"""Pick a shared-storage backend from ``DASH_SHARED_STORAGE``. + +Lets a hosting platform switch backends without editing the app's ``Dash(...)`` +call. Only used when the app did not pass ``shared_storage=``. +""" + +import functools +import re +from typing import Any, Optional +from urllib.parse import urlparse + +from ..exceptions import InvalidConfig +from .diskcache import DiskcacheSharedStorage, _require_diskcache +from .local import LocalSharedStorage +from .redis import RedisSharedStorage, _require_redis + +ENV_VAR = "DASH_SHARED_STORAGE" + + +def _redact(value: str) -> str: + # Keep credentials in a URL out of the error message. + return re.sub(r"(://)[^/@]*@", r"\1***@", value) + + +def _invalid(value: str, reason: str) -> InvalidConfig: + return InvalidConfig( + f"{ENV_VAR}={_redact(value)!r} is not valid: {reason}. Use 'local', 'none', " + "'diskcache:///absolute/path', or a redis:// or rediss:// URL." + ) + + +def storage_from_env(value: Optional[str]) -> Any: + """Turn a ``DASH_SHARED_STORAGE`` value into a ``shared_storage`` argument. + + Returns ``None`` (disabled), or a zero-argument callable that builds the + backend. Nothing is built or connected here, so startup stays lazy; missing + optional dependencies still fail now, with the backend's own error. + """ + raw = (value or "").strip() + lowered = raw.lower() + if lowered in ("", "local"): + return LocalSharedStorage + if lowered == "none": + return None + + scheme = urlparse(raw).scheme.lower() + if scheme in ("redis", "rediss"): + _require_redis() + return functools.partial(RedisSharedStorage, url=raw) + if scheme == "diskcache": + parsed = urlparse(raw) + if parsed.netloc or not parsed.path.startswith("/"): + raise _invalid(raw, "diskcache needs an absolute path (three slashes)") + _require_diskcache() + return functools.partial(DiskcacheSharedStorage, directory=parsed.path) + if scheme == "cluster": + raise InvalidConfig( + f"{ENV_VAR}={_redact(raw)!r}: the cluster:// backend is not supported in this " + "version of Dash." + ) + raise _invalid(raw, "unknown backend") diff --git a/dash/dash.py b/dash/dash.py index 19892085a0..101f743597 100644 --- a/dash/dash.py +++ b/dash/dash.py @@ -48,7 +48,12 @@ ) from .backends import get_backend from .version import __version__ -from ._configs import get_combined_config, pathname_configs, pages_folder_config +from ._configs import ( + get_combined_config, + load_dash_env_vars, + pathname_configs, + pages_folder_config, +) from ._utils import ( AttributeDict, format_tag, @@ -76,11 +81,8 @@ from . import backends from ._get_app import with_app_context, with_app_context_factory -from ._shared_storage import ( - BaseSharedStorage, - LocalSharedStorage, - SharedStorageError, -) +from ._shared_storage import BaseSharedStorage, SharedStorageError +from ._shared_storage._env import storage_from_env from ._grouping import map_grouping, grouping_len, update_args_group from ._obsolete import ObsoleteChecker from ._callback_context import callback_context @@ -154,6 +156,9 @@ _ID_DUMMY = "_pages_dummy" _UNINITIALIZED = object() # Sentinel for tracking init_app state +# Default for ``shared_storage`` so "not passed" (DASH_SHARED_STORAGE may pick +# the backend) can be told apart from an explicit argument. +_SHARED_STORAGE_DEFAULT: Any = object() DASH_VERSION_URL = "https://dash-version.plotly.com:8080/current_version" @@ -467,6 +472,20 @@ class Dash(ObsoleteChecker): takes a thread for milliseconds. ASGI backends (Quart, FastAPI) keep one open connection per browser instead and ignore this. :type stream_poll_interval: int + + :param shared_storage: Backend for ``dash.ctx.shared_storage`` and the + streaming-callback transport: a ``BaseSharedStorage`` subclass or + instance, or ``None`` to disable. Default ``LocalSharedStorage`` (one + machine or pod). When not passed, the ``DASH_SHARED_STORAGE`` + environment variable can choose it: ``local``, ``none``, + ``diskcache:///absolute/path``, or a ``redis://`` / ``rediss://`` URL. + An explicit argument, including ``None``, always wins. + + Stream requests are signed per page load, so with several workers or + pods they all need the same signing secret: set ``server.secret_key`` + or the ``DASH_SECRET_KEY`` environment variable. ``DASH_SECRET_KEY`` is + used for Dash's own signing only and does not set ``server.secret_key``. + :type shared_storage: BaseSharedStorage subclass or instance, or None """ _plotlyjs_url: str @@ -531,7 +550,7 @@ def __init__( # pylint: disable=too-many-statements, too-many-branches stream_poll_interval: int = 100, shared_storage: Optional[ Union[Type[BaseSharedStorage], BaseSharedStorage] - ] = LocalSharedStorage, + ] = _SHARED_STORAGE_DEFAULT, enable_mcp: Optional[bool] = None, mcp_path: Optional[str] = None, **obsolete, @@ -714,6 +733,10 @@ def __init__( # pylint: disable=too-many-statements, too-many-branches # access so it costs nothing until used and never binds in a gunicorn # preload master or the Flask reloader parent -- only in the worker that # actually touches it. + if shared_storage is _SHARED_STORAGE_DEFAULT: + shared_storage = storage_from_env( + load_dash_env_vars().get("DASH_SHARED_STORAGE") + ) self._shared_storage_arg = shared_storage self._shared_storage_instance: Optional[BaseSharedStorage] = None self._shared_storage_lock = threading.Lock() @@ -743,9 +766,8 @@ def __init__( # pylint: disable=too-many-statements, too-many-branches self._pages_lock = threading.Lock() self._pages_async_lock: Optional[asyncio.Lock] = None - # Secret used to sign background-callback handles (see _callback_signing). - # Prefer the Flask/Quart secret_key (shared across workers when the - # operator sets one); otherwise fall back to a per-process random secret. + # Fallback signing secret, used when neither server.secret_key nor + # DASH_SECRET_KEY is set. See _get_signing_secret. self._generated_signing_secret: Optional[bytes] = None if server: @@ -1012,7 +1034,7 @@ def shared_storage(self) -> BaseSharedStorage: with self._shared_storage_lock: if self._shared_storage_instance is None: storage = self._shared_storage_arg - if isinstance(storage, type): + if isinstance(storage, (type, functools.partial)): storage = storage() storage.start() self._shared_storage_instance = storage @@ -1061,21 +1083,29 @@ def serve_layout(self): ) def _get_signing_secret(self) -> bytes: - """Return the secret used to sign background-callback handles. + """Return the secret used to sign page tokens (``end_id``), background + callback handles and stream connections. Resolution order: 1. The server's ``secret_key`` if set (shared across workers when the operator configures one, e.g. for Flask-Login). - 2. Otherwise a random secret persisted in the background-callback result + 2. ``DASH_SECRET_KEY`` from the environment, so a hosting platform can + give every worker and pod the same key without editing the app. Used + for Dash's own signing only, never assigned to ``server.secret_key`` + (that would change Flask session behavior). + 3. Otherwise a random secret persisted in the background-callback result store, so every worker reads back the same value. This is exactly as shared as the callback results themselves, so it works cross-worker whenever the deployment is set up for multi-worker background callbacks (an explicitly shared cache / broker). - 3. Finally, if no background manager is available, a per-process random - secret (there are no background handles to verify in that case). + 4. Finally, if no background manager is available, a per-process random + secret. Stream tokens then only verify on the worker that issued + them, so multi-worker streaming apps need 1 or 2. """ - key = getattr(self.server, "secret_key", None) + key = getattr(self.server, "secret_key", None) or load_dash_env_vars().get( + "DASH_SECRET_KEY" + ) if key: return key.encode("utf-8") if isinstance(key, str) else key if self._generated_signing_secret is None: diff --git a/tests/shared_storage/test_dash_integration.py b/tests/shared_storage/test_dash_integration.py index e74feb571b..d915085ccd 100644 --- a/tests/shared_storage/test_dash_integration.py +++ b/tests/shared_storage/test_dash_integration.py @@ -34,10 +34,14 @@ def _redis_available(): return False -def _counter_app(storage): +_FROM_ENV = object() + + +def _counter_app(storage=_FROM_ENV): """A two-callback app: 'bump' increments a shared counter, 'read' (a separate callback) shows the current shared value.""" - app = Dash(__name__, shared_storage=storage) + kwargs = {} if storage is _FROM_ENV else {"shared_storage": storage} + app = Dash(__name__, **kwargs) app.layout = html.Div( [ html.Button("bump", id="bump"), @@ -124,6 +128,19 @@ def test_counter_shared_across_callbacks_redis(dash_duo): _drive_counter(dash_duo) +def test_redis_selected_by_env(dash_duo, monkeypatch): + if not _redis_available(): + pytest.skip("no Redis reachable at REDIS_URL") + monkeypatch.setenv("DASH_SHARED_STORAGE", REDIS_URL) + app = _counter_app() + storage = app.shared_storage + assert isinstance(storage, RedisSharedStorage) + # The env var gives the default key prefix, so clear what earlier runs left. + storage.delete("count") + dash_duo.start_server(app) + _drive_counter(dash_duo) + + def test_disabled_shared_storage_errors_the_callback(dash_duo): """With shared_storage=None, a callback touching ctx.shared_storage fails; the app surfaces a callback error rather than updating the output.""" diff --git a/tests/shared_storage/test_env_config.py b/tests/shared_storage/test_env_config.py new file mode 100644 index 0000000000..7fb37f0a85 --- /dev/null +++ b/tests/shared_storage/test_env_config.py @@ -0,0 +1,127 @@ +"""DASH_SHARED_STORAGE picks the backend when the app does not pass one.""" +# pylint: disable=protected-access +import sys + +import pytest + +from dash import Dash +from dash._shared_storage import ( + DiskcacheSharedStorage, + LocalSharedStorage, + RedisSharedStorage, +) +from dash.exceptions import InvalidConfig + + +def _app(monkeypatch, value, **kwargs): + monkeypatch.setenv("DASH_SHARED_STORAGE", value) + return Dash(__name__, **kwargs) + + +def _built(app): + # Build the backend without start(), so nothing connects. + storage = app._shared_storage_arg + return storage if storage is None else storage() + + +def test_unset_defaults_to_local(monkeypatch): + monkeypatch.delenv("DASH_SHARED_STORAGE", raising=False) + app = Dash(__name__) + assert app._shared_storage_arg is LocalSharedStorage + + +@pytest.mark.parametrize("value", ["local", "LOCAL", " local ", ""]) +def test_local(monkeypatch, value): + app = _app(monkeypatch, value) + assert app._shared_storage_arg is LocalSharedStorage + + +@pytest.mark.parametrize("value", ["none", "None"]) +def test_none_disables(monkeypatch, value): + app = _app(monkeypatch, value) + assert not app.shared_storage_enabled + + +def test_diskcache_is_lazy(monkeypatch, tmp_path): + pytest.importorskip("diskcache") + directory = tmp_path / "ss" + app = _app(monkeypatch, f"diskcache://{directory}") + assert app.shared_storage_enabled + assert not directory.exists() + storage = app.shared_storage + assert isinstance(storage, DiskcacheSharedStorage) + assert directory.is_dir() + storage.close() + + +@pytest.mark.parametrize("value", ["diskcache://relative/path", "diskcache:"]) +def test_diskcache_needs_absolute_path(monkeypatch, value): + pytest.importorskip("diskcache") + with pytest.raises(InvalidConfig, match="DASH_SHARED_STORAGE"): + _app(monkeypatch, value) + + +@pytest.mark.parametrize( + "url", ["redis://localhost:6399/3", "rediss://user:pw@example.com:6380/0"] +) +def test_redis(monkeypatch, url): + pytest.importorskip("redis") + # A port nothing listens on: construction must not connect. + storage = _built(_app(monkeypatch, url)) + assert isinstance(storage, RedisSharedStorage) + kwargs = storage._redis.connection_pool.connection_kwargs + assert kwargs["port"] == int(url.rsplit(":", 1)[1].split("/")[0]) + storage.close() + + +def test_cluster_is_reserved(monkeypatch): + with pytest.raises(InvalidConfig, match="not supported in this version"): + _app(monkeypatch, "cluster://nodes") + + +@pytest.mark.parametrize("value", ["memcached://x", "redis", "/tmp/cache", "true"]) +def test_garbage_names_the_variable_and_value(monkeypatch, value): + with pytest.raises(InvalidConfig) as err: + _app(monkeypatch, value) + assert "DASH_SHARED_STORAGE" in str(err.value) + assert repr(value) in str(err.value) + + +@pytest.mark.parametrize( + "value", ["valkey://:hunter2@host:6379", "cluster://admin:hunter2@nodes"] +) +def test_error_hides_url_credentials(monkeypatch, value): + with pytest.raises(InvalidConfig) as err: + _app(monkeypatch, value) + assert "hunter2" not in str(err.value) + assert "***@" in str(err.value) + + +def test_explicit_argument_beats_env(monkeypatch): + storage = LocalSharedStorage() + app = _app(monkeypatch, "none", shared_storage=storage) + assert app._shared_storage_arg is storage + + +def test_explicit_none_beats_env(monkeypatch): + assert not _app(monkeypatch, "local", shared_storage=None).shared_storage_enabled + + +def test_explicit_argument_skips_env_validation(monkeypatch): + _app(monkeypatch, "garbage", shared_storage=LocalSharedStorage) + + +@pytest.mark.parametrize( + "value, module, backend", + [ + ("redis://localhost:6379", "redis", RedisSharedStorage), + ("diskcache:///tmp/dash-ss", "diskcache", DiskcacheSharedStorage), + ], +) +def test_missing_extra_matches_explicit_error(monkeypatch, value, module, backend): + monkeypatch.setitem(sys.modules, module, None) + with pytest.raises(ImportError) as explicit: + backend() + with pytest.raises(ImportError) as from_env: + _app(monkeypatch, value) + assert str(from_env.value) == str(explicit.value) diff --git a/tests/streaming/test_stream_transport.py b/tests/streaming/test_stream_transport.py index 458cb93fda..ebc413e3f2 100644 --- a/tests/streaming/test_stream_transport.py +++ b/tests/streaming/test_stream_transport.py @@ -127,6 +127,40 @@ def test_flask_downlink_rejects_missing_end_id(): storage.close() +def test_flask_forged_stream_tokens_warn_once(monkeypatch, caplog): + # A token that fails verification is almost always a secret mismatch + # between workers, so say so, but once per process, not per poll. + from dash import _callback + + monkeypatch.setattr(_callback, "_stream_token_warned", False) + app, storage = _streaming_app() + client = app.server.test_client() + forged = "/_dash-update-component?endId=forged~deadbeef" + with caplog.at_level("WARNING", logger="dash.dash"): + # A missing token is not a verification failure: no warning. + assert ( + client.post( + "/_dash-update-component", json={"streamDownlink": {"from": 0}} + ).status_code + == 403 + ) + assert not caplog.records + for _ in range(3): + assert ( + client.post(forged, json={"streamDownlink": {"from": 0}}).status_code + == 403 + ) + assert client.post(forged, json=_uplink_body("r1")).status_code == 403 + assert ( + client.post(forged, json={"streamCancel": {"requestId": "r1"}}).status_code + == 403 + ) + warnings = [r for r in caplog.records if "DASH_SECRET_KEY" in r.getMessage()] + assert len(warnings) == 1 + assert "failed verification" in warnings[0].getMessage() + storage.close() + + def test_flask_downlink_resets_a_stale_cursor(): # A downlink resuming from a cursor the fresh topic never reached (the page's # server restarted, so the topic is back at seq 0) gets a reset line, not a diff --git a/tests/streaming/test_stream_wsgi.py b/tests/streaming/test_stream_wsgi.py index bf5b616a6c..1d51f0c630 100644 --- a/tests/streaming/test_stream_wsgi.py +++ b/tests/streaming/test_stream_wsgi.py @@ -1,16 +1,24 @@ -"""Streaming over the multiplexed transport on a single-threaded WSGI worker. +"""Streaming over the multiplexed transport on gunicorn sync workers. gunicorn's default sync worker serves one request at a time. The client must open its long-lived downlink only once the uplink is acknowledged, otherwise the downlink holds the only worker and the stream never starts. + +With several workers, each must verify the stream token another one signed, +so they need a shared signing secret (``server.secret_key`` or +``DASH_SECRET_KEY``). """ +import contextlib +import json import os +import re import signal import socket import subprocess import sys import textwrap import time +from concurrent.futures import ThreadPoolExecutor import pytest import requests @@ -23,10 +31,17 @@ APP = textwrap.dedent( """ import asyncio + import os from dash import Dash, Input, Output, html app = Dash(__name__) server = app.server + + @server.after_request + def tag_worker(response): + response.headers["X-Worker-Pid"] = str(os.getpid()) + return response + app.layout = html.Div([html.Button("go", id="btn", n_clicks=0), html.Div(id="out")]) @app.callback(Output("out", "children"), Input("btn", "n_clicks")) @@ -46,9 +61,12 @@ def _free_port(): return sock.getsockname()[1] -def test_stwg001_single_sync_worker_streams_promptly(dash_br, tmp_path): +@contextlib.contextmanager +def _gunicorn(tmp_path, workers, env=None): (tmp_path / "wsgi_app.py").write_text(APP) port = _free_port() + proc_env = {k: v for k, v in os.environ.items() if k.lower() != "dash_secret_key"} + proc_env.update(env or {}) proc = subprocess.Popen( # pylint: disable=consider-using-with [ sys.executable, @@ -58,35 +76,28 @@ def test_stwg001_single_sync_worker_streams_promptly(dash_br, tmp_path): "--bind", f"127.0.0.1:{port}", "--workers", - "1", + str(workers), "--timeout", "30", ], cwd=tmp_path, + env=proc_env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True, ) + url = f"http://127.0.0.1:{port}" try: deadline = time.monotonic() + 30 while time.monotonic() < deadline: try: - if requests.get(f"http://127.0.0.1:{port}", timeout=1).ok: + if requests.get(url, timeout=1).ok: break except requests.RequestException: time.sleep(0.2) else: raise AssertionError("gunicorn never came up") - - dash_br.server_url = f"http://127.0.0.1:{port}" - dash_br.find_element("#btn").click() - started = time.monotonic() - # Well under gunicorn's worker timeout: the stream must not need the - # worker to be killed and respawned before it starts. - dash_br.wait_for_contains_text("#out", "token 1", timeout=10) - assert time.monotonic() - started < 10 - dash_br.wait_for_contains_text("#out", "token 19", timeout=15) - assert dash_br.get_logs() == [] + yield url finally: os.killpg(proc.pid, signal.SIGTERM) try: @@ -94,3 +105,68 @@ def test_stwg001_single_sync_worker_streams_promptly(dash_br, tmp_path): except subprocess.TimeoutExpired: os.killpg(proc.pid, signal.SIGKILL) proc.wait() + + +def _downlink_polls(url, polls=40): + """Poll the downlink with one page load's token from many connections at + once, so the polls spread over the workers (one idle sync worker would take + them all if sent one by one). Returns the pid of the worker that served the + page and a ``(pid, status)`` per poll, from at least two workers.""" + page = requests.get(url, timeout=5) + config = json.loads( + re.search( + r'', page.text, re.S + ).group(1) + ) + + def poll(_): + resp = requests.post( + f"{url}/_dash-update-component", + params={"endId": config["end_id"]}, + json={"streamDownlink": {"from": 0}}, + headers={"Connection": "close"}, + timeout=5, + ) + return resp.headers["X-Worker-Pid"], resp.status_code + + results = [] + with ThreadPoolExecutor(max_workers=8) as pool: + for _ in range(5): + results += pool.map(poll, range(polls)) + if len({pid for pid, _ in results}) > 1: + break + assert len({pid for pid, _ in results}) > 1, "polls never left one worker" + return page.headers["X-Worker-Pid"], results + + +def _stream_completes(dash_br, url): + dash_br.server_url = url + dash_br.find_element("#btn").click() + started = time.monotonic() + # Well under gunicorn's worker timeout: the stream must not need the + # worker to be killed and respawned before it starts. + dash_br.wait_for_contains_text("#out", "token 1", timeout=10) + assert time.monotonic() - started < 10 + dash_br.wait_for_contains_text("#out", "token 19", timeout=15) + assert dash_br.get_logs() == [] + + +def test_stwg001_single_sync_worker_streams_promptly(dash_br, tmp_path): + with _gunicorn(tmp_path, workers=1) as url: + _stream_completes(dash_br, url) + + +def test_stwg002_workers_share_dash_secret_key(dash_br, tmp_path): + with _gunicorn(tmp_path, workers=4, env={"DASH_SECRET_KEY": "shared"}) as url: + _, polls = _downlink_polls(url) + assert {status for _, status in polls} == {200} + _stream_completes(dash_br, url) + + +def test_stwg003_workers_without_shared_secret_refuse_streams(tmp_path): + # Keeps the multi-worker failure visible: each worker makes up its own + # secret, so a token only verifies on the worker that signed it. + with _gunicorn(tmp_path, workers=4) as url: + page_pid, polls = _downlink_polls(url) + for pid, status in polls: + assert status == (200 if pid == page_pid else 403) diff --git a/tests/unit/test_background_callback_signing.py b/tests/unit/test_background_callback_signing.py index 0ed17689b0..efff782480 100644 --- a/tests/unit/test_background_callback_signing.py +++ b/tests/unit/test_background_callback_signing.py @@ -229,3 +229,38 @@ def test_mcp_parse_task_id_rejects_forged_handles(): # A raw pid + cache key with no valid signature (what an attacker sends). with pytest.raises(MCPError): parse_task_id("mytool:9999:operator-secret-key:0") + + +def test_server_secret_key_beats_env(monkeypatch): + monkeypatch.setenv("DASH_SECRET_KEY", "env-secret") + app, _ = _make_app() + app.server.secret_key = "configured-secret" + assert app._get_signing_secret() == b"configured-secret" + + +def test_env_secret_beats_background_store(monkeypatch): + shared_dir = tempfile.mkdtemp() + stored = _make_app(cache_dir=shared_dir)[0]._get_signing_secret() + monkeypatch.setenv("DASH_SECRET_KEY", "env-secret") + app, _ = _make_app(cache_dir=shared_dir) + assert app._get_signing_secret() == b"env-secret" != stored + assert not app.server.secret_key + + +def test_env_secret_without_background_manager(monkeypatch): + from dash import Dash + + monkeypatch.setenv("DASH_SECRET_KEY", "env-secret") + assert Dash(__name__)._get_signing_secret() == b"env-secret" + + +def test_empty_env_secret_is_ignored(monkeypatch): + from dash import Dash + + monkeypatch.setenv("DASH_SECRET_KEY", "") + shared_dir = tempfile.mkdtemp() + app_a, _ = _make_app(cache_dir=shared_dir) + app_b, _ = _make_app(cache_dir=shared_dir) + assert app_a._get_signing_secret() == app_b._get_signing_secret() != b"" + # No store and no key: a random per-process secret, different per app. + assert Dash(__name__)._get_signing_secret() != Dash(__name__)._get_signing_secret() From 9a2f54d6202e9553c748c98892e1209f73a0a428 Mon Sep 17 00:00:00 2001 From: philippe Date: Wed, 30 Sep 2026 14:30:47 -0400 Subject: [PATCH 2/2] Add PR number to CHANGELOG --- CHANGELOG.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ea8897c7e7..7e92f699c1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,8 +5,8 @@ This project adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] ### Added -- Add `DASH_SECRET_KEY` so a hosting platform can give every worker and pod the same signing secret for page, background and stream tokens without editing the app (`server.secret_key` still wins and is not changed). -- Add `DASH_SHARED_STORAGE` to pick the shared-storage backend when the app does not pass `shared_storage=`: `local`, `none`, `diskcache:///abs/path`, or a `redis://` / `rediss://` URL. +- [#4026](https://github.com/plotly/dash/pull/4026) Add `DASH_SECRET_KEY` so a hosting platform can give every worker and pod the same signing secret for page, background and stream tokens without editing the app (`server.secret_key` still wins and is not changed). +- [#4026](https://github.com/plotly/dash/pull/4026) Add `DASH_SHARED_STORAGE` to pick the shared-storage backend when the app does not pass `shared_storage=`: `local`, `none`, `diskcache:///abs/path`, or a `redis://` / `rediss://` URL. - [#3976](https://github.com/plotly/dash/pull/3976) Add a new `scrollToTop` prop to `dcc.Link` to control whether the page scrolls to the top after client-side navigation. It defaults to `True` to preserve the existing behavior. Fixes [#3974](https://github.com/plotly/dash/issues/3974). - [#3947](https://github.com/plotly/dash/pull/3947) Make `plotly-cloud` a default install dependency of Dash instead of an optional extra, so the `plotly` CLI and Dash's cloud integration work out of the box. The `dash[cloud]` extra is kept for backward compatibility. - [#3930](https://github.com/plotly/dash/pull/3930) Add shared storage: a backend-agnostic cross-process state manager (key/value with optional TTL, plus ordered replayable pub/sub) on every app via `dash.ctx.shared_storage`, started lazily and disabled with `shared_storage=None`. Ships `LocalSharedStorage` (default, in-memory with optional disk persistence), `DiskcacheSharedStorage`, and `RedisSharedStorage` for horizontally-scaled deployments; see `.ai/ARCHITECTURE.md`. @@ -23,7 +23,7 @@ This project adheres to [Semantic Versioning](https://semver.org/). - [#3646](https://github.com/plotly/dash/pull/3646) Remove React 16 support (`16.14.0` is no longer an accepted value for `REACT_VERSION` / `_set_react_version`). ### Changed -- Log a warning, once per process, when a streaming request is refused because its stream token failed verification, which on multi-worker deployments usually means the workers do not share a signing secret. +- [#4026](https://github.com/plotly/dash/pull/4026) Log a warning, once per process, when a streaming request is refused because its stream token failed verification, which on multi-worker deployments usually means the workers do not share a signing secret. - [#3987](https://github.com/plotly/dash/pull/3987) Forward FastAPI reload scope options (`reload_dirs`, `reload_excludes`, and `reload_includes`) to Uvicorn when reloading. - [#3986](https://github.com/plotly/dash/pull/3986) Adjust `_run_before_hooks` in the `fastapi` backend to honor a response returned by a `before_request` function, matching the `flask` backend's behavior.