From be1741c2f84b7102a7560e65a91c0c7164e58621 Mon Sep 17 00:00:00 2001 From: chick Date: Fri, 17 Jul 2026 18:30:58 +0800 Subject: [PATCH] feat(web): add durable control-plane event updates --- apps/web/package.json | 1 + apps/web/src/app-shell.tsx | 19 +- .../event-invalidation-controller.test.ts | 55 ++++ .../events/event-invalidation-controller.ts | 78 +++++ .../src/events/event-stream-client.test.ts | 231 +++++++++++++++ apps/web/src/events/event-stream-client.ts | 266 ++++++++++++++++++ .../src/events/use-control-plane-events.ts | 70 +++++ apps/web/src/fleet/fleet-page.tsx | 7 +- apps/web/src/index.ts | 19 ++ pnpm-lock.yaml | 3 + 10 files changed, 745 insertions(+), 4 deletions(-) create mode 100644 apps/web/src/events/event-invalidation-controller.test.ts create mode 100644 apps/web/src/events/event-invalidation-controller.ts create mode 100644 apps/web/src/events/event-stream-client.test.ts create mode 100644 apps/web/src/events/event-stream-client.ts create mode 100644 apps/web/src/events/use-control-plane-events.ts diff --git a/apps/web/package.json b/apps/web/package.json index f63ecab..3c456e9 100644 --- a/apps/web/package.json +++ b/apps/web/package.json @@ -12,6 +12,7 @@ "typecheck": "tsc -p tsconfig.json" }, "dependencies": { + "@multi-simadmin/contracts": "workspace:*", "react": "19.2.4", "react-dom": "19.2.4" }, diff --git a/apps/web/src/app-shell.tsx b/apps/web/src/app-shell.tsx index 0895999..4649c5f 100644 --- a/apps/web/src/app-shell.tsx +++ b/apps/web/src/app-shell.tsx @@ -1,4 +1,8 @@ import type { ReactNode } from 'react'; +import { useMemo } from 'react'; + +import { useControlPlaneEvents } from './events/use-control-plane-events.js'; +import { createEventStreamClient, type EventStreamClient } from './events/event-stream-client.js'; import { FleetPage, type FleetDataSource, type FleetSnapshot } from './fleet/fleet-page.js'; import { createFleetApiDataSource } from './fleet/fleet-api-data-source.js'; @@ -57,6 +61,7 @@ export interface AppShellProps { instanceDataSource?: InstanceDataSource; capabilities?: InstanceCapabilityMap; capabilityDataSource?: CapabilityDataSource; + eventStreamClient?: EventStreamClient; } const MODULE_LABELS = INSTANCE_MODULE_LABELS; @@ -123,6 +128,7 @@ function Page({ instanceDataSource, capabilities, capabilityDataSource, + fleetRefreshSignal, }: { route: ResolvedRoute; fleetDataSource: FleetDataSource | undefined; @@ -131,12 +137,14 @@ function Page({ instanceDataSource: InstanceDataSource | undefined; capabilities: InstanceCapabilityMap | undefined; capabilityDataSource: CapabilityDataSource | undefined; + fleetRefreshSignal: number; }): ReactNode { if (route.kind === 'fleet') return ( ); if (route.kind === 'instance-new') @@ -201,9 +209,17 @@ export function AppShell({ instanceDataSource, capabilities, capabilityDataSource, + eventStreamClient, }: AppShellProps) { + const defaultEventStreamClient = useMemo(() => createEventStreamClient(), []); + const defaultFleetDataSource = useMemo(() => createFleetApiDataSource(), []); const resolved = resolveRoute(pathname); const route = resolved.kind === 'redirect' ? resolveRoute(resolved.to ?? '/fleet') : resolved; + const routeInstanceId = route.params?.instanceId; + const refresh = useControlPlaneEvents( + routeInstanceId, + eventStreamClient ?? defaultEventStreamClient, + ); const currentSection = section(route); return (
@@ -239,12 +255,13 @@ export function AppShell({
diff --git a/apps/web/src/events/event-invalidation-controller.test.ts b/apps/web/src/events/event-invalidation-controller.test.ts new file mode 100644 index 0000000..f74eac2 --- /dev/null +++ b/apps/web/src/events/event-invalidation-controller.test.ts @@ -0,0 +1,55 @@ +import { describe, expect, it, vi } from 'vitest'; + +import { createEventInvalidationController } from './event-invalidation-controller.js'; + +const base = { id: 'event', occurredAt: '2026-07-17T10:00:00.000Z', requestId: 'request' } as const; + +describe('event invalidation controller', () => { + it('coalesces bursts into the correct global invalidation domains', () => { + const callbacks: Array<() => void> = []; + const emit = vi.fn(); + const controller = createEventInvalidationController(emit, { + schedule: (callback) => callbacks.push(callback), + }); + controller.push({ ...base, kind: 'job', jobId: 'job-1' }); + controller.push({ ...base, kind: 'item', jobId: 'job-1', itemId: 'item-1' }); + controller.push({ ...base, kind: 'audit', auditEventId: 'audit-1' }); + expect(callbacks).toHaveLength(1); + callbacks[0]?.(); + expect(emit).toHaveBeenCalledWith( + expect.objectContaining({ fleet: false, jobs: true, audit: true }), + ); + }); + + it('invalidates Fleet for every instance event but detail only for the latest route owner', () => { + const callbacks: Array<() => void> = []; + const emit = vi.fn(); + const controller = createEventInvalidationController(emit, { + schedule: (callback) => callbacks.push(callback), + }); + controller.setActiveInstance('old'); + controller.push({ ...base, kind: 'instance', instanceId: 'old' }); + controller.setActiveInstance('new'); + controller.push({ ...base, kind: 'instance', instanceId: 'new' }); + callbacks[0]?.(); + const batch = emit.mock.calls[0]?.[0]; + expect(batch.fleet).toBe(true); + expect([...batch.instanceIds]).toEqual(['new']); + }); + + it('cancels pending work and ignores events after disposal', () => { + const callback = vi.fn(); + const cancel = vi.fn(); + const emit = vi.fn(); + const controller = createEventInvalidationController(emit, { + schedule: () => callback, + cancel, + }); + controller.push({ ...base, kind: 'audit', auditEventId: 'audit-1' }); + controller.dispose(); + controller.push({ ...base, kind: 'job', jobId: 'job-1' }); + callback(); + expect(cancel).toHaveBeenCalledWith(callback); + expect(emit).not.toHaveBeenCalled(); + }); +}); diff --git a/apps/web/src/events/event-invalidation-controller.ts b/apps/web/src/events/event-invalidation-controller.ts new file mode 100644 index 0000000..a0172ba --- /dev/null +++ b/apps/web/src/events/event-invalidation-controller.ts @@ -0,0 +1,78 @@ +import type { EventEnvelope } from '@multi-simadmin/contracts'; + +export interface InvalidationBatch { + readonly fleet: boolean; + readonly instanceIds: ReadonlySet; + readonly jobs: boolean; + readonly audit: boolean; +} + +export interface EventInvalidationController { + push(event: EventEnvelope): void; + setActiveInstance(instanceId: string | undefined): void; + dispose(): void; +} + +export interface EventInvalidationOptions { + readonly delayMs?: number; + readonly schedule?: (callback: () => void, delayMs: number) => unknown; + readonly cancel?: (handle: unknown) => void; +} + +export function createEventInvalidationController( + emit: (batch: InvalidationBatch) => void, + options: EventInvalidationOptions = {}, +): EventInvalidationController { + const schedule = options.schedule ?? ((callback, delay) => setTimeout(callback, delay)); + const cancel = + options.cancel ?? ((handle) => clearTimeout(handle as ReturnType)); + let activeInstanceId: string | undefined; + let handle: unknown; + let disposed = false; + let fleet = false; + let jobs = false; + let audit = false; + const instanceIds = new Set(); + + const flush = (): void => { + handle = undefined; + if (disposed) return; + const batch: InvalidationBatch = { + fleet, + instanceIds: new Set(instanceIds), + jobs, + audit, + }; + fleet = false; + jobs = false; + audit = false; + instanceIds.clear(); + emit(batch); + }; + const arm = (): void => { + if (handle === undefined) handle = schedule(flush, Math.max(0, options.delayMs ?? 50)); + }; + + return { + push(event) { + if (disposed) return; + if (event.kind === 'instance') { + fleet = true; + if (event.instanceId === activeInstanceId) instanceIds.add(event.instanceId); + } else if (event.kind === 'job' || event.kind === 'item' || event.kind === 'attempt') + jobs = true; + else if (event.kind === 'audit') audit = true; + arm(); + }, + setActiveInstance(instanceId) { + activeInstanceId = instanceId; + for (const id of instanceIds) if (id !== instanceId) instanceIds.delete(id); + }, + dispose() { + disposed = true; + if (handle !== undefined) cancel(handle); + handle = undefined; + instanceIds.clear(); + }, + }; +} diff --git a/apps/web/src/events/event-stream-client.test.ts b/apps/web/src/events/event-stream-client.test.ts new file mode 100644 index 0000000..b8a16e2 --- /dev/null +++ b/apps/web/src/events/event-stream-client.test.ts @@ -0,0 +1,231 @@ +// @vitest-environment jsdom +import { describe, expect, it, vi } from 'vitest'; + +import { + createEventStreamClient, + createSseParser, + parseEventEnvelope, + type EventStreamState, +} from './event-stream-client.js'; + +const instanceEvent = { + kind: 'instance', + id: 'event-1', + occurredAt: '2026-07-17T10:00:00.000Z', + requestId: 'request-1', + instanceId: 'instance-1', +} as const; + +function streamResponse(chunks: readonly string[], status = 200): Response { + const encoder = new TextEncoder(); + return new Response( + new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(encoder.encode(chunk)); + controller.close(); + }, + }), + { status, headers: { 'content-type': 'text/event-stream; charset=utf-8' } }, + ); +} + +function deferred(): { promise: Promise; resolve: () => void } { + let resolve!: () => void; + return { promise: new Promise((done) => (resolve = done)), resolve }; +} + +describe('strict EventEnvelope parser', () => { + it('accepts every frozen discriminated-union shape', () => { + expect(parseEventEnvelope(instanceEvent)).toEqual(instanceEvent); + expect( + parseEventEnvelope({ ...instanceEvent, occurredAt: '2026-07-17T10:00:00Z' }), + ).toMatchObject({ + occurredAt: '2026-07-17T10:00:00Z', + }); + expect( + parseEventEnvelope({ ...instanceEvent, occurredAt: '2026-07-17T18:00:00+08:00' }), + ).toMatchObject({ + occurredAt: '2026-07-17T18:00:00+08:00', + }); + const base = { + kind: instanceEvent.kind, + id: instanceEvent.id, + occurredAt: instanceEvent.occurredAt, + requestId: instanceEvent.requestId, + }; + expect( + parseEventEnvelope({ ...base, kind: 'item', jobId: 'job-1', itemId: 'item-1' }), + ).toMatchObject({ kind: 'item', jobId: 'job-1', itemId: 'item-1' }); + expect( + parseEventEnvelope({ ...base, kind: 'attempt', jobId: 'job-1', attemptId: 'try-1' }), + ).toMatchObject({ kind: 'attempt', attemptId: 'try-1' }); + expect(parseEventEnvelope({ ...base, kind: 'job', jobId: 'job-1' })).toMatchObject({ + kind: 'job', + }); + expect(parseEventEnvelope({ ...base, kind: 'audit', auditEventId: 'audit-1' })).toMatchObject({ + kind: 'audit', + }); + }); + + it.each([ + null, + [], + { ...instanceEvent, kind: 'unknown' }, + { ...instanceEvent, id: '' }, + { ...instanceEvent, occurredAt: 'yesterday' }, + { ...instanceEvent, instanceId: 1 }, + { ...instanceEvent, secret: 'must-not-pass' }, + { ...instanceEvent, kind: 'job' }, + { ...instanceEvent, kind: 'instance', jobId: 'wrong-union-member' }, + ])('rejects invalid or non-exact envelopes: %j', (value) => { + expect(() => parseEventEnvelope(value)).toThrow('Invalid event envelope.'); + }); + + it.each(['2026-02-30T10:00:00Z', '2026-04-31T10:00:00Z', '2026-01-01T24:00:00Z'])( + 'rejects normalized invalid RFC3339 timestamp %s', + (occurredAt) => { + expect(() => parseEventEnvelope({ ...instanceEvent, occurredAt })).toThrow( + 'Invalid event envelope.', + ); + }, + ); +}); + +describe('SSE framing', () => { + it('handles arbitrary chunk boundaries, CRLF, comments, and multiline data', () => { + const records: Array<{ id?: string; data: string }> = []; + const parser = createSseParser((record) => records.push(record)); + parser.push(': keepalive\r'); + parser.push('\nid: cursor-1\r\nda'); + parser.push('ta: {"one":\r\ndata: 2}\r\n\r'); + parser.push('\n'); + parser.finish(); + expect(records).toEqual([{ id: 'cursor-1', data: '{"one":\n2}' }]); + }); + + it('dispatches a final unterminated event and ignores id-only events and NUL ids', () => { + const records: Array<{ id?: string; data: string }> = []; + const parser = createSseParser((record) => records.push(record)); + parser.push('id: ignored\n\nid: bad\0id\ndata: final'); + parser.finish(); + expect(records).toEqual([{ data: 'final' }]); + }); +}); + +describe('fetch event stream client', () => { + it('uses protected same-origin fetch, retains cursor, and reconnects with Last-Event-ID', async () => { + const gate = deferred(); + const fetcher = vi + .fn() + .mockResolvedValueOnce( + streamResponse([ + `id: cursor-7\ndata: ${JSON.stringify(instanceEvent).slice(0, 20)}`, + `${JSON.stringify(instanceEvent).slice(20)}\n\n`, + ]), + ) + .mockImplementationOnce(async (_input, init) => { + await gate.promise; + init?.signal?.throwIfAborted(); + return streamResponse([]); + }); + const sleep = vi.fn(async () => undefined); + const events: unknown[] = []; + const states: EventStreamState[] = []; + const unsubscribe = createEventStreamClient({ fetch: fetcher, sleep }).subscribe( + (event) => events.push(event), + (state) => states.push(state), + ); + + await vi.waitFor(() => expect(fetcher).toHaveBeenCalledTimes(2)); + expect(events).toEqual([instanceEvent]); + expect(fetcher.mock.calls[0]?.[0]).toBe('/api/v1/events'); + expect(fetcher.mock.calls[0]?.[1]).toMatchObject({ + credentials: 'same-origin', + headers: { Accept: 'text/event-stream' }, + }); + expect(fetcher.mock.calls[1]?.[1]?.headers).toEqual({ + Accept: 'text/event-stream', + 'Last-Event-ID': 'cursor-7', + }); + expect(states).toContain('open'); + expect(states).toContain('reconnecting'); + unsubscribe(); + gate.resolve(); + }); + + it('applies bounded exponential reconnect delay and never exposes response bodies', async () => { + const gate = deferred(); + const fetcher = vi + .fn() + .mockResolvedValueOnce(new Response('password=top-secret', { status: 401 })) + .mockResolvedValueOnce(new Response('another-secret', { status: 503 })) + .mockImplementationOnce(async (_input, init) => { + await gate.promise; + init?.signal?.throwIfAborted(); + return streamResponse([]); + }); + const delays: number[] = []; + const errors: string[] = []; + const client = createEventStreamClient({ + fetch: fetcher, + initialRetryMs: 10, + maxRetryMs: 20, + sleep: async (ms) => void delays.push(ms), + }); + const unsubscribe = client.subscribe( + () => undefined, + (_state, error) => error && errors.push(error.message), + ); + + await vi.waitFor(() => expect(fetcher).toHaveBeenCalledTimes(3)); + expect(delays.slice(0, 2)).toEqual([10, 20]); + expect(errors.join(' ')).not.toMatch(/secret|password/i); + expect(errors[0]).toBe('The event stream connection failed.'); + unsubscribe(); + gate.resolve(); + }); + + it('clears an expired cursor and reconnects without Last-Event-ID', async () => { + const gate = deferred(); + const fetcher = vi + .fn() + .mockResolvedValueOnce( + streamResponse([`id: cursor-old\ndata: ${JSON.stringify(instanceEvent)}\n\n`]), + ) + .mockResolvedValueOnce(new Response(null, { status: 409 })) + .mockImplementationOnce(async (_input, init) => { + await gate.promise; + init?.signal?.throwIfAborted(); + return streamResponse([]); + }); + const states: EventStreamState[] = []; + const unsubscribe = createEventStreamClient({ + fetch: fetcher, + sleep: async () => undefined, + }).subscribe( + () => undefined, + (state) => states.push(state), + ); + await vi.waitFor(() => expect(fetcher).toHaveBeenCalledTimes(3)); + expect(fetcher.mock.calls[1]?.[1]?.headers).toMatchObject({ 'Last-Event-ID': 'cursor-old' }); + expect(fetcher.mock.calls[2]?.[1]?.headers).toEqual({ Accept: 'text/event-stream' }); + expect(states).toContain('resetting'); + unsubscribe(); + gate.resolve(); + }); + + it('rejects cross-origin endpoints before fetching and aborts on unsubscribe', async () => { + const fetcher = vi.fn(); + const errors: Error[] = []; + const unsubscribe = createEventStreamClient({ + url: 'https://attacker.example/events', + fetch: fetcher, + }).subscribe( + () => undefined, + (_state, error) => error && errors.push(error), + ); + await vi.waitFor(() => expect(errors).toHaveLength(1)); + expect(fetcher).not.toHaveBeenCalled(); + unsubscribe(); + }); +}); diff --git a/apps/web/src/events/event-stream-client.ts b/apps/web/src/events/event-stream-client.ts new file mode 100644 index 0000000..3978ea5 --- /dev/null +++ b/apps/web/src/events/event-stream-client.ts @@ -0,0 +1,266 @@ +import type { EventEnvelope, EventKind } from '@multi-simadmin/contracts'; + +export type EventStreamState = 'connecting' | 'open' | 'reconnecting' | 'resetting' | 'closed'; +export type EventStreamStateListener = (state: EventStreamState, error?: Error) => void; +export type EventStreamListener = (event: EventEnvelope) => void; + +export interface SseRecord { + readonly id?: string; + readonly data: string; +} + +export interface SseParser { + push(chunk: string): void; + finish(): void; +} + +export interface EventStreamClient { + subscribe(onEvent: EventStreamListener, onState?: EventStreamStateListener): () => void; +} + +export interface EventStreamClientOptions { + readonly url?: string; + readonly fetch?: typeof globalThis.fetch; + readonly initialRetryMs?: number; + readonly maxRetryMs?: number; + readonly sleep?: (milliseconds: number, signal: AbortSignal) => Promise; + readonly onPositionReset?: () => void; +} + +const KINDS: readonly EventKind[] = ['instance', 'job', 'item', 'attempt', 'audit']; +const COMMON_KEYS = ['kind', 'id', 'occurredAt', 'requestId'] as const; +const KIND_KEYS: Readonly> = { + instance: ['instanceId'], + job: ['jobId'], + item: ['jobId', 'itemId'], + attempt: ['jobId', 'attemptId'], + audit: ['auditEventId'], +}; +const CONNECTION_ERROR = 'The event stream connection failed.'; +const INVALID_EVENT_ERROR = 'Invalid event envelope.'; + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === 'object' && !Array.isArray(value); +} + +function isNonEmptyString(value: unknown): value is string { + return typeof value === 'string' && value.length > 0; +} + +function isIsoTimestamp(value: unknown): value is string { + if (!isNonEmptyString(value)) return false; + const match = + /^(\d{4})-(\d{2})-(\d{2})T(\d{2}):(\d{2}):(\d{2})(?:\.\d+)?(Z|[+-]\d{2}:\d{2})$/u.exec(value); + if (!match) return false; + const [, year, month, day, hour, minute, second] = match; + const parts = [year, month, day, hour, minute, second].map(Number); + if (parts.some((part) => !Number.isInteger(part))) return false; + const [y, m, d, h, min, s] = parts as [number, number, number, number, number, number]; + if (m < 1 || m > 12 || d < 1 || h > 23 || min > 59 || s > 59) return false; + const daysInMonth = new Date(Date.UTC(y, m, 0)).getUTCDate(); + if (d > daysInMonth) return false; + return Number.isFinite(Date.parse(value)); +} + +/** Validates the exact, frozen EventEnvelope union without retaining unknown fields. */ +export function parseEventEnvelope(value: unknown): EventEnvelope { + if ( + !isRecord(value) || + typeof value.kind !== 'string' || + !KINDS.includes(value.kind as EventKind) + ) + throw new Error(INVALID_EVENT_ERROR); + + const kind = value.kind as EventKind; + const kindKeys = KIND_KEYS[kind]; + if (kindKeys === undefined) throw new Error(INVALID_EVENT_ERROR); + const expected = [...COMMON_KEYS, ...kindKeys]; + const actual = Object.keys(value); + if (actual.length !== expected.length || actual.some((key) => !expected.includes(key))) + throw new Error(INVALID_EVENT_ERROR); + if ( + !isNonEmptyString(value.id) || + !isIsoTimestamp(value.occurredAt) || + !isNonEmptyString(value.requestId) || + kindKeys.some((key) => !isNonEmptyString(value[key])) + ) + throw new Error(INVALID_EVENT_ERROR); + + return value as unknown as EventEnvelope; +} + +/** Incremental SSE field parser implementing LF, CRLF, CR, comments, and multiline data. */ +export function createSseParser(onRecord: (record: SseRecord) => void): SseParser { + let buffer = ''; + let data: string[] = []; + let id: string | undefined; + + const dispatch = (): void => { + if (data.length > 0) + onRecord(id === undefined ? { data: data.join('\n') } : { id, data: data.join('\n') }); + data = []; + id = undefined; + }; + const line = (value: string): void => { + if (value === '') { + dispatch(); + return; + } + if (value.startsWith(':')) return; + const colon = value.indexOf(':'); + const field = colon < 0 ? value : value.slice(0, colon); + let fieldValue = colon < 0 ? '' : value.slice(colon + 1); + if (fieldValue.startsWith(' ')) fieldValue = fieldValue.slice(1); + if (field === 'data') data.push(fieldValue); + else if (field === 'id' && !fieldValue.includes('\0')) id = fieldValue; + }; + const consume = (final: boolean): void => { + let start = 0; + for (let index = 0; index < buffer.length; index += 1) { + const character = buffer[index]; + if (character !== '\r' && character !== '\n') continue; + if (character === '\r' && index + 1 === buffer.length && !final) break; + line(buffer.slice(start, index)); + if (character === '\r' && buffer[index + 1] === '\n') index += 1; + start = index + 1; + } + buffer = buffer.slice(start); + if (final && buffer.length > 0) { + line(buffer); + buffer = ''; + } + }; + return { + push(chunk) { + buffer += chunk; + consume(false); + }, + finish() { + consume(true); + dispatch(); + }, + }; +} + +function sameOriginUrl(input: string): string { + const base = globalThis.location?.href ?? 'http://localhost/'; + const url = new URL(input, base); + const origin = new URL(base).origin; + if (url.origin !== origin || url.username !== '' || url.password !== '') + throw new Error('The event stream URL must be same-origin.'); + return input.startsWith('/') ? `${url.pathname}${url.search}` : url.href; +} + +function defaultSleep(milliseconds: number, signal: AbortSignal): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout(resolve, milliseconds); + signal.addEventListener( + 'abort', + () => { + clearTimeout(timer); + reject(signal.reason); + }, + { once: true }, + ); + }); +} + +function safeConnectionError(): Error { + return new Error(CONNECTION_ERROR); +} + +export function createEventStreamClient(options: EventStreamClientOptions = {}): EventStreamClient { + const fetcher = options.fetch ?? globalThis.fetch; + const initialRetryMs = Math.max(1, options.initialRetryMs ?? 1_000); + const maxRetryMs = Math.max(initialRetryMs, options.maxRetryMs ?? 30_000); + const sleep = options.sleep ?? defaultSleep; + + return { + subscribe(onEvent, onState = () => undefined) { + const controller = new AbortController(); + let cursor: string | undefined; + let retryMs = initialRetryMs; + + const run = async (): Promise => { + let url: string; + try { + url = sameOriginUrl(options.url ?? '/api/v1/events'); + } catch (error) { + onState('closed', error instanceof Error ? error : safeConnectionError()); + return; + } + + onState('connecting'); + while (!controller.signal.aborted) { + try { + const headers: Record = { Accept: 'text/event-stream' }; + if (cursor !== undefined) headers['Last-Event-ID'] = cursor; + const response = await fetcher(url, { + credentials: 'same-origin', + headers, + signal: controller.signal, + }); + if (response.status === 409 && cursor !== undefined) { + cursor = undefined; + options.onPositionReset?.(); + onState( + 'resetting', + new Error('The event stream position expired; refreshing current state.'), + ); + continue; + } + if (!response.ok || response.body === null) throw safeConnectionError(); + const contentType = response.headers.get('content-type')?.toLowerCase() ?? ''; + if (!contentType.startsWith('text/event-stream')) throw safeConnectionError(); + + onState('open'); + retryMs = initialRetryMs; + const decoder = new TextDecoder(); + let recordFailure: Error | undefined; + const parser = createSseParser((record) => { + if (recordFailure !== undefined) return; + try { + const event = parseEventEnvelope(JSON.parse(record.data) as unknown); + onEvent(event); + if (record.id !== undefined) cursor = record.id; + } catch { + recordFailure = new Error(INVALID_EVENT_ERROR); + } + }); + const reader = response.body.getReader(); + try { + while (!controller.signal.aborted) { + const result = await reader.read(); + if (result.done) break; + parser.push(decoder.decode(result.value, { stream: true })); + if (recordFailure !== undefined) throw recordFailure; + } + parser.push(decoder.decode()); + parser.finish(); + if (recordFailure !== undefined) throw recordFailure; + } finally { + if (controller.signal.aborted) await reader.cancel().catch(() => undefined); + reader.releaseLock(); + } + } catch { + if (controller.signal.aborted) break; + } + if (controller.signal.aborted) break; + // Deliberately replace fetch/parser errors: response bodies and transport details may be sensitive. + onState('reconnecting', safeConnectionError()); + const delay = retryMs; + retryMs = Math.min(maxRetryMs, retryMs * 2); + try { + await sleep(delay, controller.signal); + } catch { + break; + } + } + onState('closed'); + }; + + void run(); + return () => controller.abort(); + }, + }; +} diff --git a/apps/web/src/events/use-control-plane-events.ts b/apps/web/src/events/use-control-plane-events.ts new file mode 100644 index 0000000..f8f6293 --- /dev/null +++ b/apps/web/src/events/use-control-plane-events.ts @@ -0,0 +1,70 @@ +import { useEffect, useRef, useState } from 'react'; +import type { EventEnvelope } from '@multi-simadmin/contracts'; + +import { createEventInvalidationController } from './event-invalidation-controller.js'; +import { + createEventStreamClient, + type EventStreamClient, + type EventStreamState, +} from './event-stream-client.js'; + +export interface EventRefreshState { + readonly fleet: number; + readonly jobs: number; + readonly audit: number; + readonly detail: number; + readonly stream: EventStreamState; + readonly positionReset: number; +} + +export function useControlPlaneEvents( + activeInstanceId?: string, + client: EventStreamClient = createEventStreamClient(), +): EventRefreshState { + const [state, setState] = useState({ + fleet: 0, + jobs: 0, + audit: 0, + detail: 0, + stream: 'connecting', + positionReset: 0, + }); + + const activeInstanceRef = useRef(activeInstanceId); + activeInstanceRef.current = activeInstanceId; + + useEffect(() => { + const controller = createEventInvalidationController((batch) => { + setState((current) => ({ + ...current, + fleet: current.fleet + (batch.fleet ? 1 : 0), + jobs: current.jobs + (batch.jobs ? 1 : 0), + audit: current.audit + (batch.audit ? 1 : 0), + detail: current.detail + (batch.instanceIds.size > 0 ? 1 : 0), + })); + }); + controller.setActiveInstance(activeInstanceRef.current); + const unsubscribe = client.subscribe( + (event: EventEnvelope) => { + controller.setActiveInstance(activeInstanceRef.current); + controller.push(event); + }, + (stream) => + setState((current) => ({ + ...current, + stream, + fleet: current.fleet + (stream === 'resetting' ? 1 : 0), + jobs: current.jobs + (stream === 'resetting' ? 1 : 0), + audit: current.audit + (stream === 'resetting' ? 1 : 0), + detail: current.detail + (stream === 'resetting' && activeInstanceRef.current ? 1 : 0), + positionReset: current.positionReset + (stream === 'resetting' ? 1 : 0), + })), + ); + return () => { + unsubscribe(); + controller.dispose(); + }; + }, [client]); + + return state; +} diff --git a/apps/web/src/fleet/fleet-page.tsx b/apps/web/src/fleet/fleet-page.tsx index 551e015..c9aab42 100644 --- a/apps/web/src/fleet/fleet-page.tsx +++ b/apps/web/src/fleet/fleet-page.tsx @@ -20,6 +20,7 @@ export interface FleetDataSource { export interface FleetPageProps { readonly dataSource?: FleetDataSource; readonly initialData?: FleetSnapshot; + readonly refreshSignal?: number; } const EMPTY_SNAPSHOT: FleetSnapshot = { instances: [], statuses: new Map() }; @@ -54,7 +55,7 @@ export function canonicalHttpOrigin(value: string): string | null { } } -export function FleetPage({ dataSource, initialData }: FleetPageProps) { +export function FleetPage({ dataSource, initialData, refreshSignal = 0 }: FleetPageProps) { const [snapshot, setSnapshot] = useState(initialData ?? null); const [error, setError] = useState(null); const [attempt, setAttempt] = useState(0); @@ -76,7 +77,7 @@ export function FleetPage({ dataSource, initialData }: FleetPageProps) { const selectAllRef = useRef(null); useEffect(() => { - if (initialData) { + if (initialData && refreshSignal === 0) { setSnapshot(initialData); setError(null); return; @@ -100,7 +101,7 @@ export function FleetPage({ dataSource, initialData }: FleetPageProps) { return () => { active = false; }; - }, [attempt, dataSource, initialData]); + }, [attempt, dataSource, initialData, refreshSignal]); const model = useMemo( () => diff --git a/apps/web/src/index.ts b/apps/web/src/index.ts index 008bb9c..9ea447b 100644 --- a/apps/web/src/index.ts +++ b/apps/web/src/index.ts @@ -34,4 +34,23 @@ export { type PasswordUpdate, } from './instances/instance-api-data-source.js'; +export { + createEventStreamClient, + createSseParser, + parseEventEnvelope, + type EventStreamClient, + type EventStreamClientOptions, + type EventStreamListener, + type EventStreamState, + type EventStreamStateListener, + type SseParser, + type SseRecord, +} from './events/event-stream-client.js'; +export { + createEventInvalidationController, + type EventInvalidationController, + type EventInvalidationOptions, + type InvalidationBatch, +} from './events/event-invalidation-controller.js'; + export const webWorkspaceReady = true; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index a6ed6c2..1f550ec 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -61,6 +61,9 @@ importers: apps/web: dependencies: + '@multi-simadmin/contracts': + specifier: workspace:* + version: link:../../packages/contracts react: specifier: 19.2.4 version: 19.2.4