fabro/apps/fabro-web/app/lib/sse.ts
Bryan Helmkamp bc31b8d23a
Serve the run stream as the only run event API
Step 4 of the legacy executor deletion, third commit: the legacy event
API and every reader of it go, so that the next commits can delete the
event log, its reducer and the types beneath them.

The API:
- `GET /runs/{id}/events` pages the run stream only
  (`PaginatedRunStreamList` by `after`); the legacy `since_seq`,
  `before_seq` and `order` cursors, the `oneOf` envelope, the legacy
  `EventEnvelope`, `PaginatedEventList`, `RunEvent`, `EventSeq`,
  `AppendEventResponse` and `RunEventDetailResponse` schemas,
  `POST /runs/{id}/events`, `GET /runs/{id}/events/{seq}` and
  `GET /runs/{id}/stages/{stageId}/events` are deleted. `GET
  /runs/{id}/attach` and `GET /attach` frame `RunStreamItem`s only.
- The Rust and TypeScript clients regenerate; the removed models leave
  the TypeScript package.

The readers:
- `fabro-client` drops the legacy run event listing, tail and attach
  methods and `RunEventStream`; `list_run_stream_until` bounds a stream
  read.
- `fabro-tool`'s `fabro_run_events` lists, searches and details the run
  stream: `after` is the exclusive `stream_seq` cursor, `event_id` the
  item's id, filters match the item's name and `recorded_at`.
- `fabro-dump` writes the stream to `events.jsonl`; `fabro dump` reads
  it.
- The CLI's progress renderer keeps only what the run stream drives:
  the legacy event conversion, the sandbox and setup displays and their
  styles go. `fabro system events` prints stream items.
- The server's demo mode folds its agent fixture straight into the
  session projection and answers the attach stub with a stream item;
  the demo stage events endpoint is gone.
- The web app: every run is a Petri run. The legacy event hooks,
  renderer props, stage popover summary, run phases derivation and
  live-event payload handling are deleted or ported to `RunStreamItem`;
  toasts and board refreshes read the stream's platform records.
- Tests: the legacy API round trips and pagination tests are deleted;
  the CLI's MCP, attach and system event mocks serve stream pages; the
  CLI test helpers read stream items.

Still failing until the later commits: the CLI tests seeded through
`POST /runs/{id}/events`, the server tests over the legacy store, and
the legacy type tests.

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

203 lines
5.5 KiB
TypeScript

import type { Key, MutatorCallback } from "swr";
export type KeyMatcher = (key: Key) => boolean;
export type KeyOrMatcher = Key | KeyMatcher;
export type MutateFn = (key: KeyOrMatcher) => ReturnType<MutatorCallback>;
/** A parsed SSE frame: a run stream item, read through `isStreamItemPayload`. */
export interface EventPayload {
[key: string]: unknown;
}
export interface EventSourceLike {
onmessage: ((event: { data: string }) => void) | null;
close(): void;
}
export interface EventInvalidation {
keys: KeyOrMatcher[];
close?: boolean;
immediate?: boolean;
}
type EventResolver = (payload: EventPayload) => EventInvalidation;
export interface SharedEventSubscription {
source: EventSourceLike;
refcount: number;
mutators: Map<MutateFn, number>;
resolvers: Map<symbol, EventResolver>;
pendingKeys: Map<string, KeyOrMatcher>;
debounceTimer: ReturnType<typeof setTimeout> | null;
}
export function sseKeyDedupeId(key: KeyOrMatcher): string {
return typeof key === "function" ? `fn:${key.toString()}` : stringifyKeyValue(key);
}
export function createBrowserEventSource(url: string): EventSourceLike {
return new EventSource(url);
}
export function subscribeToSharedEventSource<TPayload extends EventPayload>({
subscriptions,
subscriptionKey,
url,
mutate,
resolveInvalidation,
eventSourceFactory = createBrowserEventSource,
debounceMs = 300,
}: {
subscriptions: Map<string, SharedEventSubscription>;
subscriptionKey: string;
url: string;
mutate: MutateFn;
resolveInvalidation: (payload: TPayload) => EventInvalidation;
eventSourceFactory?: (url: string) => EventSourceLike;
debounceMs?: number;
}): () => void {
let subscription = subscriptions.get(subscriptionKey);
if (!subscription) {
const source = eventSourceFactory(url);
subscription = {
source,
refcount: 0,
mutators: new Map(),
resolvers: new Map(),
pendingKeys: new Map(),
debounceTimer: null,
};
subscriptions.set(subscriptionKey, subscription);
source.onmessage = (message) => {
const current = subscriptions.get(subscriptionKey);
if (!current) return;
let payload: TPayload;
try {
payload = JSON.parse(message.data) as TPayload;
} catch {
return;
}
const keys = new Map<string, KeyOrMatcher>();
let close = false;
let immediate = false;
for (const resolver of current.resolvers.values()) {
const invalidation = resolver(payload);
for (const key of invalidation.keys) {
keys.set(sseKeyDedupeId(key), key);
}
close ||= Boolean(invalidation.close);
immediate ||= Boolean(invalidation.immediate);
}
queueInvalidations(current, [...keys.values()], { debounceMs, immediate });
if (close) {
closeSharedEventSource(subscriptions, subscriptionKey, { flushPending: true });
}
};
}
const resolverId = Symbol(subscriptionKey);
subscription.resolvers.set(
resolverId,
resolveInvalidation as EventResolver,
);
subscription.refcount += 1;
subscription.mutators.set(mutate, (subscription.mutators.get(mutate) ?? 0) + 1);
return () => {
const current = subscriptions.get(subscriptionKey);
if (!current) return;
current.resolvers.delete(resolverId);
const mutateCount = current.mutators.get(mutate) ?? 0;
if (mutateCount <= 1) {
current.mutators.delete(mutate);
} else {
current.mutators.set(mutate, mutateCount - 1);
}
current.refcount -= 1;
if (current.refcount <= 0) {
closeSharedEventSource(subscriptions, subscriptionKey);
}
};
}
function queueInvalidations(
subscription: SharedEventSubscription,
keys: KeyOrMatcher[],
{
debounceMs,
immediate,
}: {
debounceMs: number;
immediate?: boolean;
},
) {
if (keys.length === 0) return;
for (const key of keys) {
subscription.pendingKeys.set(sseKeyDedupeId(key), key);
}
if (immediate || debounceMs <= 0) {
flushInvalidations(subscription);
return;
}
if (subscription.debounceTimer) {
clearTimeout(subscription.debounceTimer);
}
subscription.debounceTimer = setTimeout(() => {
subscription.debounceTimer = null;
flushInvalidations(subscription);
}, debounceMs);
}
function flushInvalidations(subscription: SharedEventSubscription) {
if (subscription.pendingKeys.size === 0) return;
const keys = [...subscription.pendingKeys.values()];
subscription.pendingKeys.clear();
for (const mutator of subscription.mutators.keys()) {
for (const key of keys) {
void mutator(key);
}
}
}
function closeSharedEventSource(
subscriptions: Map<string, SharedEventSubscription>,
subscriptionKey: string,
{ flushPending = false }: { flushPending?: boolean } = {},
) {
const subscription = subscriptions.get(subscriptionKey);
if (!subscription) return;
if (flushPending) {
flushInvalidations(subscription);
}
if (subscription.debounceTimer) {
clearTimeout(subscription.debounceTimer);
}
subscription.source.close();
subscriptions.delete(subscriptionKey);
}
function stringifyKeyValue(value: unknown): string {
if (Array.isArray(value)) {
return `[${value.map((item) => stringifyKeyValue(item)).join(",")}]`;
}
if (value && typeof value === "object") {
const record = value as Record<string, unknown>;
return `{${Object.keys(record)
.sort()
.map((key) => `${JSON.stringify(key)}:${stringifyKeyValue(record[key])}`)
.join(",")}}`;
}
return JSON.stringify(value) ?? String(value);
}