"""Fault-tolerance coverage for Deep Research progress events.""" import pytest from deerflow.agents.deep_research.adapters.event_sink import PersistingEventSink class _UnavailableEventRepository: async def append(self, **kwargs): # noqa: ANN003 raise RuntimeError("database temporarily unavailable") @pytest.mark.asyncio async def test_event_persistence_failure_keeps_job_live_and_emits_unique_fallback_event(): published: list[dict] = [] async def publish(event: dict) -> None: published.append(event) sink = PersistingEventSink( _UnavailableEventRepository(), session_id="drs_1", job_id="drj_1", on_event=publish, ) first = await sink.emit("phase_changed", phase="planning", payload={"to": "planning"}) second = await sink.emit("queries_planned", phase="collecting", payload={"queries": ["AI"]}) # A degraded emit consumes two negative sequences: the original frame plus # the copyable diagnostic that follows it, so the second emit lands on -3. assert (first.seq, second.seq) == (-1, -3) assert first.payload["eventPersistenceDegraded"] is True assert published[0]["liveId"] == "event-persistence-fallback:drj_1:1" assert published[1]["liveId"] == "event-persistence-fallback:drj_1:2" @pytest.mark.asyncio async def test_live_publish_failure_does_not_fail_event_delivery(): class _EventRepository: async def append(self, **kwargs): # noqa: ANN003 return { "seq": 1, "event_type": kwargs["event_type"], "phase": kwargs["phase"], "payload": kwargs["payload"], } async def broken_publisher(event: dict) -> None: # noqa: ARG001 raise RuntimeError("browser disconnected") sink = PersistingEventSink( _EventRepository(), session_id="drs_1", job_id="drj_1", on_event=broken_publisher, ) event = await sink.emit("phase_changed", phase="planning") assert event.seq == 1 class _RecordingEventRepository: def __init__(self) -> None: self.rows: list[dict] = [] async def append(self, **kwargs): # noqa: ANN003 self.rows.append(kwargs) return { "seq": len(self.rows), "event_type": kwargs["event_type"], "phase": kwargs["phase"], "payload": kwargs["payload"], } async def _noop_publish(event: dict) -> None: # noqa: ARG001 return None def _mirrored(repo: _RecordingEventRepository, event_type: str) -> list[dict]: return [row for row in repo.rows if row["event_type"] == event_type] @pytest.mark.asyncio async def test_reasoning_frames_are_mirrored_durably_when_no_viewer_is_on_this_worker(): """Multi-worker guard: thinking must not vanish into a subscriber-less hub. A job runs on the worker that won its lease while the SSE connection is served by whichever worker accepted it. Live-only frames published into an empty in-process hub are lost with no durable row to replay them from, which is what silently removed the step bar's thinking text. """ repo = _RecordingEventRepository() sink = PersistingEventSink( repo, session_id="drs_1", job_id="drj_1", on_event=_noop_publish, live_subscriber_probe=lambda _job_id: False, ) for _ in range(3): await sink.publish_live( "report_thinking", phase="writing", payload={"delta": "思" * 240, "target": "report", "reportVersion": 1}, ) mirrored = _mirrored(repo, "report_thinking") assert mirrored, "reasoning must be persisted when the hub has no subscriber" assert "".join(row["payload"]["delta"] for row in mirrored) == "思" * 720 assert mirrored[0]["payload"]["liveMirror"] is True assert mirrored[0]["payload"]["target"] == "report" # Coalesced, never one row per provider token. assert len(mirrored) == 1 @pytest.mark.asyncio async def test_a_locally_attached_viewer_keeps_reasoning_off_the_database(): repo = _RecordingEventRepository() sink = PersistingEventSink( repo, session_id="drs_1", job_id="drj_1", on_event=_noop_publish, live_subscriber_probe=lambda _job_id: True, ) for _ in range(5): await sink.publish_live("report_thinking", phase="writing", payload={"delta": "思" * 240}) assert _mirrored(repo, "report_thinking") == [] @pytest.mark.asyncio async def test_report_text_is_never_mirrored_and_the_tail_flushes_before_the_next_event(): repo = _RecordingEventRepository() sink = PersistingEventSink( repo, session_id="drs_1", job_id="drj_1", on_event=_noop_publish, live_subscriber_probe=lambda _job_id: False, ) # Report prose already has a durable counterpart (``report_chunk``); a # mirror would duplicate it in the assembled report. await sink.publish_live("report_delta", phase="writing", payload={"delta": "正文", "offset": 0}) assert _mirrored(repo, "report_delta") == [] # A short reasoning tail stays buffered, then lands ahead of the next # durable milestone instead of waiting for the job to end. await sink.publish_live("report_thinking", phase="writing", payload={"delta": "短思考"}) assert _mirrored(repo, "report_thinking") == [] await sink.emit("report_completed", phase="writing", payload={}) assert [row["event_type"] for row in repo.rows] == ["report_thinking", "report_completed"] @pytest.mark.asyncio async def test_mirroring_is_off_without_a_probe(): """Single-worker deployments and tests keep zero extra database writes.""" repo = _RecordingEventRepository() sink = PersistingEventSink(repo, session_id="drs_1", job_id="drj_1", on_event=_noop_publish) await sink.publish_live("report_thinking", phase="writing", payload={"delta": "思" * 900}) assert repo.rows == [] @pytest.mark.asyncio async def test_document_rewrite_milestone_is_a_supported_durable_event(): """Guard the artifact-rewrite job path against EventType regressions.""" class _EventRepository: async def append(self, **kwargs): # noqa: ANN003 return { "seq": 1, "event_type": kwargs["event_type"], "phase": kwargs["phase"], "payload": kwargs["payload"], } sink = PersistingEventSink(_EventRepository(), session_id="drs_1", job_id="drj_1") event = await sink.emit( "document_rewrite", phase="planning", payload={"type": "requirements_started"}, ) assert event.type == "document_rewrite"