fabro/apps/fabro-web/app/lib/session-stream.ts
Bryan Helmkamp 34996d630f
Give Ask Fabro sessions their own event log
Step 4 of the legacy executor deletion, second commit. Ask Fabro's
sessions were the last writer of `run_events`: a session's creation, its
turns and their messages, tool calls and endings went into the run's
legacy event log, keyed by the run's sequence. They now have a log of
their own.

- `run_session_events` (migration `2026091802`): one row per session
  event, numbered per session from 1, with the owning run, the turn, the
  event name and its properties. `RunSessionEventStore` appends under the
  write lock, lists a session from a sequence, names a session's owner
  from its creation event, deletes a run's sessions with the run, and
  publishes each committed event to its subscribers.
- `fabro_types::SessionEvent`: `seq`, `session_id`, `run_id`, `ts` and a
  flattened body (`event` naming the kind, `properties` its fields), with
  the same event names and property shapes the legacy events carried,
  so the web app and the CLI read the same JSON. The property structs
  move to `session_event`; `run_event::session` re-exports them under
  their old names until the legacy event log goes.
- The API: `GET /sessions/{id}/events` pages `PaginatedSessionEventList`
  by the session's own sequence, `GET /sessions/{id}/attach` replays and
  streams `SessionEvent` frames (subscribed before the replay, so no
  event falls between the two), the turn stream carries the same frames,
  and an interrupt answers with the recorded event. The session
  projection folds `SessionEvent`s; the legacy `find_session_owner` over
  `run_events` is gone.
- The CLI's `run ask` and the web app's session stream read
  `SessionEvent`; the web runtime no longer accepts the nested legacy
  envelope shape.

The two session resume tests in the server keep failing for a reason
this commit does not touch: Ask Fabro reconnects to the run's sandbox
from the projection's sandbox instance, which the Petri projection does
not carry yet (`VIEWS.md`, the `scope.acquired` gap).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-18 12:47:43 -04:00

142 lines
3.8 KiB
TypeScript

import {
SessionsApiAxiosParamCreator,
type SessionEvent,
type SubmitTurnRequest,
} from "@qltysh/fabro-api-client";
import {
apiErrorFromFetchResponse,
generatedApiConfiguration,
} from "./api-client";
export type SessionStreamEvent = SessionEvent;
type FetchLike = (
input: string,
init?: RequestInit,
) => Promise<Response>;
interface SessionStreamOptions {
sessionId: string;
signal?: AbortSignal;
fetchImpl?: FetchLike;
onEvent: (event: SessionStreamEvent) => void;
}
export interface StreamSessionTurnOptions extends SessionStreamOptions {
input: string;
turnId?: string;
}
export interface StreamSessionTurnResult {
turnId: string | null;
}
export interface AttachSessionEventsOptions extends SessionStreamOptions {
sinceSeq?: number;
}
export async function streamSessionTurn({
sessionId,
input,
turnId,
signal,
fetchImpl = fetch,
onEvent,
}: StreamSessionTurnOptions): Promise<StreamSessionTurnResult> {
const body: SubmitTurnRequest = { input };
if (turnId) body.turn_id = turnId;
const request = await SessionsApiAxiosParamCreator(
generatedApiConfiguration,
).submitSessionTurn(sessionId, body, { signal });
const response = await fetchImpl(request.url, fetchInitFromAxiosRequest(request.options));
await throwIfApiError(response);
await readEventStream(response, onEvent);
return { turnId: response.headers.get("x-fabro-turn-id") };
}
export async function attachSessionEvents({
sessionId,
sinceSeq,
signal,
fetchImpl = fetch,
onEvent,
}: AttachSessionEventsOptions): Promise<void> {
const request = await SessionsApiAxiosParamCreator(
generatedApiConfiguration,
).attachSessionEvents(sessionId, sinceSeq, { signal });
const response = await fetchImpl(request.url, fetchInitFromAxiosRequest(request.options));
await throwIfApiError(response);
await readEventStream(response, onEvent);
}
async function throwIfApiError(response: Response): Promise<void> {
const error = await apiErrorFromFetchResponse(response);
if (error) throw error;
}
function fetchInitFromAxiosRequest(options: {
method?: string;
headers?: unknown;
data?: unknown;
signal?: unknown;
}): RequestInit {
const init: RequestInit = {
method: options.method,
credentials: "same-origin",
headers: options.headers as HeadersInit,
signal: options.signal as AbortSignal | undefined,
};
if (options.data !== undefined) {
init.body = typeof options.data === "string"
? options.data
: JSON.stringify(options.data);
}
return init;
}
async function readEventStream(
response: Response,
onEvent: (event: SessionStreamEvent) => void,
): Promise<void> {
if (!response.body) return;
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
// react-doctor-disable-next-line react-doctor/async-await-in-loop -- Streaming readers must consume chunks sequentially to preserve SSE order and decoder state.
const { value, done } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
buffer = drainSseBuffer(buffer, onEvent);
}
buffer += decoder.decode();
drainSseBuffer(`${buffer}\n\n`, onEvent);
}
function drainSseBuffer(
buffer: string,
onEvent: (event: SessionStreamEvent) => void,
): string {
let cursor = 0;
while (true) {
const match = /\r?\n\r?\n/g.exec(buffer.slice(cursor));
if (!match) return buffer.slice(cursor);
const next = cursor + match.index;
const frame = buffer.slice(cursor, next);
cursor = next + match[0].length;
const data = frame
.split(/\r?\n/)
.filter((line) => line.startsWith("data:"))
.map((line) => line.slice("data:".length).trimStart())
.join("\n");
if (!data) continue;
onEvent(JSON.parse(data) as SessionStreamEvent);
}
}