"""Deep Research event protocol. Every observable thing that happens during a research job — a search query fired, a source added, a report chunk written, a model invocation billed — becomes an event. Events are **persisted first** (to ``deep_research_events``) then pushed to the SSE stream; this ordering lets a client reconnect with ``?after=`` and replay anything missed (§15.3). This module defines: - :data:`EventType` — the closed set of event ``type`` strings (§15.2). - :data:`ResearchPhase` — the phase enum (§14.3). - :class:`DeepResearchEvent` — the envelope persisted + serialised over SSE. - :class:`EventSink` — the adapter protocol runners call. The concrete :class:`PersistingEventSink` (adapter) lives in ``adapters/event_sink.py`` and bridges to the event-store repository. """ from __future__ import annotations from datetime import UTC, datetime from typing import Any, Literal, Protocol, runtime_checkable from pydantic import BaseModel, Field # ── phase enum (§14.3) ─────────────────────────────────────────────────────── ResearchPhase = Literal[ "initializing", "planning", "collecting", "curating", "compressing", "deepening", "writing", "summarizing", "reviewing", "illustrating", "exporting", "done", ] ALL_PHASES: tuple[ResearchPhase, ...] = ( "initializing", "planning", "collecting", "curating", "compressing", "deepening", "writing", "summarizing", "reviewing", "illustrating", "exporting", "done", ) # ── event type enum (§15.2) ────────────────────────────────────────────────── # # Kept as a Literal union (not an enum.IntEnum) so the persisted string and the # SSE JSON ``type`` field are identical and stable across versions. EventType = Literal[ "job_status", # Terminal events are kept separate from job_status for the historical # event timeline. The executor persists these after the DB state has # been committed, so a reconnecting browser receives an explicit end. "job_completed", "job_failed", "phase_changed", "plan_created", "queries_planned", "query_started", "query_completed", "source_added", "sources_curated", "context_compressed", "deep_branch_started", "deep_branch_completed", "report_reset", "report_chunk", # Model reasoning normally travels as a live-only frame. It becomes a # durable event when no SSE connection for the job is attached to the # worker running it, because the in-process live hub cannot reach a # viewer served by a different worker (see PersistingEventSink). "report_thinking", "plan_reasoning", "summary_thinking", "report_completed", "summary_completed", "collection_completed", "artifact_created", "image_generation_started", "image_generated", "awaiting_input", "model_usage", # Durable whole-document rewrite milestones reuse the deep-research job # event log, but have their own frontend SSE payload contract. "document_rewrite", "warning", "error", "heartbeat", ] class DeepResearchEvent(BaseModel): """A single observable event in a research job. ``seq`` is assigned by the event store (monotonic per job); callers leave it ``0`` when emitting and read it back from the persisted row. """ seq: int = 0 session_id: str = "" job_id: str = "" type: EventType phase: ResearchPhase = "initializing" timestamp: datetime = Field(default_factory=lambda: datetime.now(UTC)) payload: dict[str, Any] = Field(default_factory=dict) model_config = {"extra": "allow"} def to_sse_dict(self) -> dict[str, Any]: """Serialise for the ``data:`` line of the SSE stream (§15.1).""" return { "seq": self.seq, "sessionId": self.session_id, "jobId": self.job_id, "type": self.type, "phase": self.phase, "timestamp": self.timestamp.isoformat() if self.timestamp else None, "payload": self.payload, } @runtime_checkable class EventSink(Protocol): """Adapter protocol for emitting research events. Implementations persist the event (assigning ``seq``) before returning, so a crash after ``emit()`` never loses an already-observed event. """ async def emit( self, type: EventType, # noqa: A002 — intentional arg name for readability *, phase: ResearchPhase = "initializing", payload: dict[str, Any] | None = None, ) -> DeepResearchEvent: """Persist + return an event. ``session_id``/``job_id`` are bound by the sink.""" ... async def publish_live( self, type: str, # noqa: A002 - intentionally outside the durable EventType union *, phase: ResearchPhase = "initializing", payload: dict[str, Any] | None = None, ) -> None: """Publish a best-effort, non-persisted frame. Provider text deltas use this path. Bounded ``report_chunk`` events remain the durable recovery checkpoints. """ ... __all__ = [ "ALL_PHASES", "DeepResearchEvent", "EventSink", "EventType", "ResearchPhase", ]