diff --git a/backend/packages/harness/deerflow/runtime/stream_bridge/memory.py b/backend/packages/harness/deerflow/runtime/stream_bridge/memory.py index 0a03be8af..2690c8f4c 100644 --- a/backend/packages/harness/deerflow/runtime/stream_bridge/memory.py +++ b/backend/packages/harness/deerflow/runtime/stream_bridge/memory.py @@ -86,6 +86,10 @@ class MemoryStreamBridge(StreamBridge): ) return stream.start_offset + async def stream_exists(self, run_id: str) -> bool: + """Return whether the in-process event log still has data for *run_id*.""" + return run_id in self._streams + # -- StreamBridge API ------------------------------------------------------ async def publish(self, run_id: str, event: str, data: Any) -> None: diff --git a/backend/tests/test_stream_bridge.py b/backend/tests/test_stream_bridge.py index bf710b436..b5c413e53 100644 --- a/backend/tests/test_stream_bridge.py +++ b/backend/tests/test_stream_bridge.py @@ -176,6 +176,23 @@ async def test_cleanup(bridge: MemoryStreamBridge): assert run_id not in bridge._counters +@pytest.mark.anyio +async def test_stream_exists_reports_cleanup(bridge: MemoryStreamBridge): + """Callers can detect when the in-process event log has been cleaned up. + + Before cleanup a completed run's retained history still exists; after + cleanup ``stream_exists`` reports False so a reconnecting subscriber does + not hang waiting on a stream whose data is already gone. + """ + run_id = "run-post-cleanup" + await bridge.publish(run_id, "event-1", {"n": 1}) + await bridge.publish_end(run_id) + + assert await bridge.stream_exists(run_id) is True + await bridge.cleanup(run_id) + assert await bridge.stream_exists(run_id) is False + + @pytest.mark.anyio async def test_history_is_bounded(): """Retained history should be bounded by queue_maxsize."""