Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 24 additions & 24 deletions docs/rfcs/session-sse-contract-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,8 @@ against current source before any route is added.
- This RFC does **not** implement `GET /api/sessions/{session_id}/events`. No
route, handler, or related code is added in this PR.
- This RFC does **not** modify `GET /api/sessions/events` (the existing global
session-list invalidation stream routed at `api/routes.py:12345-12346` and
implemented by `_handle_session_events_stream()` at `api/routes.py:16177`).
session-list invalidation stream routed in `api/routes.py` and
implemented by `_handle_session_events_stream()` in `api/routes.py`).
- This RFC does **not** replace or modify existing streams: `/api/chat/stream`,
`/api/approval/stream`, or `/api/clarify/stream`.
- This RFC does **not** introduce Android, iOS, or PWA client code.
Expand All @@ -53,42 +53,42 @@ against current source before any route is added.
### Existing global session-list stream

`GET /api/sessions/events` is a **different endpoint** from the one this RFC
proposes. It is routed at `api/routes.py:12345-12346` and implemented by
`_handle_session_events_stream()` at `api/routes.py:16177`. It emits bare
proposes. It is routed in `api/routes.py` and implemented by
`_handle_session_events_stream()` in `api/routes.py`. It emits bare
`sessions_changed` events and keepalives for any change to the session list. It
is a global invalidation signal, not a per-session lifecycle stream. The proposed
`GET /api/sessions/{session_id}/events` is per-session and path-distinct.

### Heartbeat

`_SSE_HEARTBEAT_INTERVAL_SECONDS = 5` (`api/routes.py:1018-1030`) is the current
`_SSE_HEARTBEAT_INTERVAL_SECONDS = 5` (defined in `api/routes.py`) is the current
heartbeat interval for SSE streams. Phase 1 reuses this constant rather than
adding a separate configurable knob.

### Run-journal cursor and replay

Current replay identity is run/stream-scoped:

Line ranges in this inventory were verified against WebUI `master` when this
RFC was written. Function, constant, and endpoint names are the stable anchors
if source layout moves later.
Symbols in this inventory were verified against WebUI `master` when this RFC
was written. Function, constant, and endpoint **names** are the stable anchors:
this RFC deliberately cites them by name (not by line number) so a source-layout
shift in `api/routes.py` cannot invalidate the doc or its contract test.

- `_parse_run_journal_event_id()` (`api/routes.py:15673-15686`) and
`_parse_run_journal_after_seq()` (`api/routes.py:15688-15701`) parse the replay
cursor from the `after_event_id` / `after_seq` **query params** (not the
- `_parse_run_journal_event_id()` and `_parse_run_journal_after_seq()` (both in
`api/routes.py`) parse the replay cursor from the `after_event_id` /
`after_seq` **query params** (not the
`Last-Event-ID` header — that header is the *proposed* new-endpoint contract
below, §Reconnect).
- `_runner_event_id()` at `api/routes.py:15765-15772` constructs the event `id`
- `_runner_event_id()` (in `api/routes.py`) constructs the event `id`
field as `stream_id:seq`.
- SSE frames carry their `id:` via the `_sse_with_id()` helper, emitted on the
live `/api/chat/stream` path at `api/routes.py:15918`, on the runner-observe
path at `api/routes.py:15811`, and during journal replay at
`api/routes.py:15721` / `15734`.
- `_replay_run_journal()` reads events by `(session_id, stream_id)` at
`api/routes.py:15703-15735`.
- `api/streaming.py:6265-6285` writes current live agent streams to
live `/api/chat/stream` path, on the runner-observe path, and during journal
replay — all in `api/routes.py`.
- `_replay_run_journal()` (in `api/routes.py`) reads events by
`(session_id, stream_id)`.
- `api/streaming.py` writes current live agent streams to
`STREAMS[stream_id]`.
- `api/streaming.py:6620-6634` appends SSE events to the run journal and carries
- `api/streaming.py` appends SSE events to the run journal and carries
per-item `event_id` into the live queue.

The existing run journal represents `session_id`, `stream_id`, `seq`, and
Expand Down Expand Up @@ -164,8 +164,8 @@ maintainer review before implementation.
position.

**`event_id` is opaque to clients.** Its current source-compatible form is
`stream_id:seq`, as constructed by `_runner_event_id()` at
`api/routes.py:15765-15772`. Clients must treat it as an opaque string and must
`stream_id:seq`, as constructed by `_runner_event_id()` in `api/routes.py`.
Clients must treat it as an opaque string and must
not parse or construct cursor values.

**`seq` is monotonic within a stream/run.** It is not a session-global counter
Expand All @@ -179,11 +179,11 @@ events, clients use `event_id` to detect and skip duplicates.
## Replay source

Phase 1 uses the **durable run journal** as the replay source for replayable
events. The live `STREAMS[stream_id]` queue (`api/streaming.py:6265-6285`) is
events. The live `STREAMS[stream_id]` queue (in `api/streaming.py`) is
not a reliable replay source because it holds only recent in-memory state.

A future implementation must replay from the run journal via the existing
`_replay_run_journal()` path (`api/routes.py:15703-15735`) and fall back to the
`_replay_run_journal()` path (in `api/routes.py`) and fall back to the
snapshot mechanism when journal entries are unavailable for a given cursor.

## Snapshot fallback
Expand All @@ -201,7 +201,7 @@ resync from the snapshot payload.

## Heartbeat

Phase 1 reuses `_SSE_HEARTBEAT_INTERVAL_SECONDS` (`api/routes.py:1018-1030`) for
Phase 1 reuses `_SSE_HEARTBEAT_INTERVAL_SECONDS` (defined in `api/routes.py`) for
heartbeat cadence. A new per-session configurable heartbeat knob is **not** added
in Phase 1. The implementation PR must follow whatever value the constant holds
at implementation time; it must not hard-code a separate interval.
Expand Down
111 changes: 53 additions & 58 deletions tests/test_issue4812_session_sse_contract_rfc.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,73 +92,68 @@ def test_rfc_states_global_endpoint_is_different(self):

def test_rfc_cites_current_global_endpoint_source(self):
"""The RFC's source anchors for the existing global stream must be
ACCURATE against current api/routes.py — verify the cited lines actually
contain what the RFC claims, rather than string-matching a fixed number
(which silently rots and codifies wrong anchors, #5513 gate finding)."""
import re
ACCURATE against current api/routes.py, verified by SYMBOL not by line
number. The RFC cites the route string and handler function by name;
this test confirms (a) each symbol still exists in api/routes.py and
(b) the RFC names that symbol. It deliberately does NOT check line
numbers: a routes.py line-shift must never break this test or the RFC
(#5513 gate finding, chronic brittle failure #5542)."""
text = _rfc()
routes = (REPO / "api" / "routes.py").read_text(encoding="utf-8").splitlines()

# Pull every `api/routes.py:<start>-<end>` or `api/routes.py:<line>`
# anchor the RFC cites and confirm the referenced span exists.
anchors = re.findall(r"api/routes\.py:(\d+)(?:-(\d+))?", text)
assert anchors, "RFC must cite at least one api/routes.py source anchor"
for start, end in anchors:
lo = int(start)
hi = int(end) if end else lo
assert 1 <= lo <= len(routes), f"RFC cites api/routes.py:{start} beyond EOF ({len(routes)} lines)"
assert 1 <= hi <= len(routes), f"RFC cites api/routes.py:{end} beyond EOF ({len(routes)} lines)"

# The two load-bearing anchors must land on the real definitions.
def _anchor_line(label):
m = re.search(r"%s.*?api/routes\.py:(\d+)" % re.escape(label), text, re.DOTALL)
assert m, f"RFC must cite a routes.py anchor near {label!r}"
return int(m.group(1))

route_line = _anchor_line("routed at")
assert "/api/sessions/events" in routes[route_line - 1], (
f"RFC's routed-at anchor api/routes.py:{route_line} must be the "
f"/api/sessions/events route; got: {routes[route_line - 1].strip()!r}"
)
handler_line = _anchor_line("_handle_session_events_stream()")
assert "def _handle_session_events_stream" in routes[handler_line - 1], (
f"RFC's handler anchor api/routes.py:{handler_line} must be the "
f"_handle_session_events_stream definition; got: {routes[handler_line - 1].strip()!r}"
)
routes_src = (REPO / "api" / "routes.py").read_text(encoding="utf-8")

# (RFC-cited symbol, existence probe in api/routes.py source)
checks = [
("/api/sessions/events", "/api/sessions/events"),
("_handle_session_events_stream", "def _handle_session_events_stream"),
]
for rfc_symbol, source_probe in checks:
assert source_probe in routes_src, (
f"api/routes.py must still define/route {source_probe!r}; the RFC "
f"cites {rfc_symbol!r} as a stable source anchor"
)
assert rfc_symbol in text, (
f"RFC must name the symbol {rfc_symbol!r} (symbol-based anchor, "
f"not a line number)"
)

def test_rfc_run_journal_anchors_land_on_real_source(self):
"""Every named-symbol / emission anchor the RFC cites in the run-journal
inventory must land on the actual source token (not just be in-bounds),
so a stale line number can't silently pass (#5513 gate finding 2)."""
import re
"""Every named-symbol anchor the RFC cites in the run-journal inventory
must be a REAL symbol in api/routes.py and be NAMED in the RFC prose.
This is verified by symbol, never by line number, so a routes.py
line-shift can't silently break it or the RFC (#5513 gate finding 2,
chronic brittle failure #5542)."""
text = _rfc()
routes = (REPO / "api" / "routes.py").read_text(encoding="utf-8").splitlines()

def _first_anchor_after(label):
m = re.search(r"%s.*?api/routes\.py:(\d+)" % re.escape(label), text, re.DOTALL)
assert m, f"RFC must cite a routes.py anchor near {label!r}"
return int(m.group(1))
routes_src = (REPO / "api" / "routes.py").read_text(encoding="utf-8")

# (label in RFC prose, token that must appear on the cited line)
# (RFC-cited symbol name, existence probe in api/routes.py source)
checks = [
("_parse_run_journal_event_id()", "def _parse_run_journal_event_id"),
("_parse_run_journal_after_seq()", "def _parse_run_journal_after_seq"),
("_runner_event_id()", "def _runner_event_id"),
("_replay_run_journal()", "def _replay_run_journal"),
("_parse_run_journal_event_id", "def _parse_run_journal_event_id"),
("_parse_run_journal_after_seq", "def _parse_run_journal_after_seq"),
("_runner_event_id", "def _runner_event_id"),
("_replay_run_journal", "def _replay_run_journal"),
("_sse_with_id", "_sse_with_id"),

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 _sse_with_id probe is inconsistent with the other four checks

The four sibling entries all use "def <symbol>" as source_probe, so they verify a definition exists in api/routes.py. The _sse_with_id entry uses just "_sse_with_id", which also matches call sites, string literals, or comments — meaning the check still passes if the helper is removed but its name survives in a comment or import string. If this is intentional (e.g., _sse_with_id is defined elsewhere and only called from routes.py), a brief inline comment explaining that would prevent future maintainers from "fixing" it to "def _sse_with_id" and introducing a spurious failure.

]
for label, token in checks:
line = _first_anchor_after(label)
assert 1 <= line <= len(routes), f"{label} anchor api/routes.py:{line} beyond EOF"
assert token in routes[line - 1], (
f"RFC's {label} anchor api/routes.py:{line} must contain {token!r}; "
f"got: {routes[line - 1].strip()!r}"
for rfc_symbol, source_probe in checks:
assert source_probe in routes_src, (
f"api/routes.py must still define {source_probe!r}; the RFC cites "
f"{rfc_symbol!r} as a stable source anchor"
)
assert rfc_symbol in text, (
f"RFC must name the symbol {rfc_symbol!r} (symbol-based anchor, "
f"not a line number)"
)

# The live-emission bullet must cite the real _sse_with_id call site.
emit_line = _first_anchor_after("live `/api/chat/stream` path at")
assert "_sse_with_id" in routes[emit_line - 1], (
f"RFC's live-emission anchor api/routes.py:{emit_line} must be an "
f"_sse_with_id() call; got: {routes[emit_line - 1].strip()!r}"
def test_rfc_uses_no_hardcoded_routes_line_numbers(self):
"""Guard against regression to line-number coupling: the RFC must not
cite `api/routes.py:<line>` (or streaming.py:<line>) anchors. Symbol
names are the durable anchor; absolute line numbers rot on any
source-layout shift and caused chronic brittle failures (#5513, #5542)."""
import re
text = _rfc()
stale = re.findall(r"\b\w+\.py:\d+(?:-\d+)?", text)
assert not stale, (
"RFC must not cite hardcoded source line numbers (they rot on any "
f"line-shift); found: {stale!r}. Reference symbols by name instead."
)

def test_contracts_distinguishes_both_endpoints(self):
Expand Down
Loading