Repository navigation
Release idle shared-storage topics - #4016
Conversation
Topics were never removed, so every page load that streamed kept its last 32 frames for the life of the process (or forever in Redis). All three backends now take topic_ttl (default 300s, None disables): a topic nobody publishes to or reads for that long is dropped. The local engine sweeps idle topics not held by any call; Redis PEXPIREs the counter and stream on publish and poll; diskcache expires messages and touches the counter. A released topic restarts at seq 1, which would reset a page that streams again after idling with its old cursor. The renderer now picks a fresh downlinkId per run of streams and restarts its cursor at 0, so each run reads its own topic (<end_id>:<downlinkId>). The downlink lifecycle record stays keyed on the page.
Dash performance benchmarks✅ all within thresholds
growth = late-third / early-third per-op time; ~1 is flat, a large value means the per-op cost scales with accumulated state. machine scale vs baseline: 0.98x - divided out of the baseline ratios so they compare like for like (the absolute warn/fail ceilings are left un-scaled); calibrated on |
| # 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 |
There was a problem hiding this comment.
This default TTL seems to apply to every user of the shared storage, not only streaming callbacks. Before this PR, a user topic lived forever, but now a slow reader would lose their messages after expiry.
Since shared storage is meant to be user facing, I think we should either:
- Carefully document how expiry works
- Have a per-topic TTL? Streaming callback topics would use a short one, and user topics a custom length. Not sure if you think this is a good idea or not.
There was a problem hiding this comment.
Agreed, a store-wide default was the wrong place for it. Went with your option 2: the ttl is now per publish, publish(topic, message, ttl=None), same shape as set(key, value, ttl=None). The constructor topic_ttl is gone. Without a ttl a topic lives as long as the store, like before this PR. Streaming publishes with its own STREAM_TOPIC_TTL (300s), so only stream topics expire. The publish docstring says what idle means and what a returning reader gets (SharedStorageGap).
| 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) |
There was a problem hiding this comment.
The DiskCache backend now behaves slightly differently than Local or Redis and I don't think that's intended. When a topic stays active but its messages are older than topic_ttl:
- Local and Redis keep the whole buffer and drop it once the whole topic has been idle for
topic_ttl. - DiskCache expires individual messages even as readers keep polling. A slow reader gets a
SharedStorageGapexception.
This clock app shows it in action.
I think the backends should expire content consistently. And we should add a test around this too: we can run the same cases against all 3 backends to catch this kind of inconsistency going forward.
There was a problem hiding this comment.
Good catch, thanks for the repro app. Diskcache now expires a topic as a whole like the other two. The ttl lives in a key next to the counter, and on publish or poll, once less than one ttl is left, every key of the topic (counter, ttl key, buffered messages) is renewed to 1.25 ttl. So a topic that is being read keeps all its buffered messages, and an idle one goes between 1 and 1.25 ttl later. Redis now keeps the same 1.25 ttl, so a reader polling a short-ttl topic can't lose it between two polls.
Added tests/shared_storage/test_topic_ttl.py, which runs the same cases on Local, Diskcache and Redis. test_a_read_topic_keeps_every_buffered_message fails on the old diskcache logic with the same SharedStorageGap you saw.
| 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}" | ||
| ) |
There was a problem hiding this comment.
nit: should we also guard against very small topic_ttl values? There would be some value at which you start to get quirky behaviour. Example: a reconnection alone can easily take 1s so a ttl of 1s means messages would expire before they ever have a chance to be read.
There was a problem hiding this comment.
Yes, a ttl under 1s (MIN_TOPIC_TTL) is now rejected with SharedStorageError on every backend, along with NaN and infinity. Tests lower the floor to use short ttls.
Topics only expire when published with a ttl, so user topics live as long as the store again; streaming passes its own STREAM_TOPIC_TTL. Diskcache renews every key of a topic together instead of expiring messages one by one, Redis keeps 1.25 ttl of slack between polls, and a ttl under 1s is rejected. A shared test runs the same expiry cases on every backend.
KoolADE85
left a comment
There was a problem hiding this comment.
Found a few nits, but overall I'm good with this one.
| idle = [ | ||
| name | ||
| for name, t in self._topics.items() | ||
| if t.ttl is not None and t.users == 0 and t.touched < now - t.ttl | ||
| ] |
There was a problem hiding this comment.
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?
| # 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 |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
This has now drifted slightly with the latest changes. See also Line 853 for the same issue.
A Redis poll blocked in XREAD for up to 5s, past the 1.25s lifetime of a minimum-ttl topic, so a waiting reader lost it. Polls now block for at most a third of the topic's lifetime (half a ttl on diskcache). The local engine also drops empty, unheld topics, such as the one a returning reader's poll recreates after a release.
|



Fixes #4010.
Shared-storage topics were never removed. Streaming uses one topic per page load, so every page that streamed kept its last 32 frames for the life of the process (and forever in Redis).
Topic lifetime.
publish(topic, message, ttl=None)andapublishtake an optional ttl (seconds, at least 1s, NaN and infinity rejected). A topic with a ttl is released, buffer and sequence both, once nobody has published to or read from it for that long. Without a ttl it lives as long as the store, so user topics behave as before. Streaming publishes withSTREAM_TOPIC_TTL(300s). Every backend expires a topic as a whole, never message by message while it is read:PEXPIREs all three keys to 1.25 ttl, a poll script renews them, and a poll blocks inXREADfor at most a third of that lifetime so a waiting reader renews in time.tests/shared_storage/test_topic_ttl.pyruns the same expiry cases on all three backends.Per-run downlink id. A released topic restarts at seq 1, but the browser kept its cursor for the whole page load. A page that idled past the ttl and streamed again would resume from its old cursor, get
{reset: true}and fail the new stream. The renderer now picks a freshdownlinkId(and restarts its cursor at 0) each time it pins a connection from idle, and sends it on every stream request.get_stream_connection_idreturns<end_id>:<downlinkId>, so each run has its own topic. The id only partitions the page's own signed space, so it is not signed, just validated (403 if malformed). The downlink lifecycle record stays keyed on the page. This also fixes a SharedWorker switching to another tab'sendIdwith a stale cursor.Proof. The issue's repro (20 page loads, ~50 KB frames), with a short ttl: 0 topics and 0 frames after idle, versus 20 topics / 640 frames before.
test_stst012_idle_stream_topics_are_releaseddrives this in a browser and fails with either half reverted (without the per-run id, the page's next run after idle never renders).Contributor Checklist
tests/shared_storage,tests/streamingincl. Redis and the streaming browser tests, renderer karma suite)optionals
CHANGELOG.md(streaming and shared storage are unreleased, covered by the existing Shared storage: backend-agnostic state manager + pub/sub #3930/Streaming callbacks with multiplexed transport #3931 entries)