import type { Engine } from '../../core/engine.ts'; import type { FleetEventFeed } from '../fleet-event-feed.ts'; import type { ServerContext } from './context.ts'; /** * Result of wiring up engine-to-WebSocket event broadcasting. * * - `dispose`: removes all listeners (abort signal). Called on server shutdown. * - `cleanupWorkflow`: drops the per-workflow sequence state for the given * workflow id. Should be invoked when a workflow reaches a terminal state * so the bookkeeping maps do not grow unbounded over the server's lifetime. * * @example * ```ts * import { Engine, MemoryStorage } from '@lostgradient/weft'; * import { wireEventBroadcasting, type EventBroadcastingHandle } from '@lostgradient/weft/server'; * * await using engine = new Engine({ storage: new MemoryStorage() }); * const bunServer = Bun.serve({ fetch: () => new Response('ok') }); * const handle: EventBroadcastingHandle = wireEventBroadcasting(engine, bunServer); * // Later, on shutdown: * handle.dispose(); * ``` */ export interface EventBroadcastingHandle { dispose: () => void; cleanupWorkflow: (workflowId: string) => void; } /** * Extract a `workflowId` from a DOM `Event` when the concrete event carries * one. All workflow, activity, token, signal, attribute, and update events * in `core/events.ts` expose a `workflowId: string` field, but the `Event` * base type does not know about it — so a runtime structural check narrows * the value before we use it to key bookkeeping maps. Returns `undefined` * for events without a string `workflowId` property. */ export declare function getWorkflowIdFromEvent(event: Event): string | undefined; export declare function registerWorkflowEventLifecycle(engine: Engine, context: ServerContext, broadcastingHandle: EventBroadcastingHandle): () => void; /** * Attach event listeners to the engine that broadcast events via WebSocket * and persist each event to storage so `GET /v1/workflows/:id/events` returns data. * Returns a handle exposing a cleanup function and a per-workflow eviction hook. * * @param engine - The engine whose events will be listened to. * @param server - The Bun server used to `server.publish()` WebSocket messages. * @param options.publishTokenMessage - Optional override for token-event delivery. * When provided, this callback is called instead of `server.publish()` for * token messages, enabling per-workflow stream sockets to be used in * place of the default pub/sub channel. Leave unset unless you manage stream * sockets separately (as `serve()` does internally). * @param options.publishWatchMessage - Optional override for watch-event delivery. * When provided, this callback is called instead of `server.publish()` for * watch messages, enabling per-socket replay buffering during catch-up. * * @example * ```ts * import { Engine, MemoryStorage } from '@lostgradient/weft'; * import { wireEventBroadcasting } from '@lostgradient/weft/server'; * * await using engine = new Engine({ storage: new MemoryStorage() }); * const bunServer = Bun.serve({ fetch: () => new Response('ok') }); * const handle = wireEventBroadcasting(engine, bunServer); * * // Wire a terminal-event listener to clean up per-workflow bookkeeping. * engine.addEventListener('workflow:completed', (e) => { * const workflowId = (e as { workflowId?: string }).workflowId; * if (workflowId) handle.cleanupWorkflow(workflowId); * }); * * // On server shutdown, remove all event listeners. * handle.dispose(); * ``` */ export declare function wireEventBroadcasting(engine: Engine, server: ReturnType, options?: { publishTokenMessage?: (workflowId: string, sequence: number, message: string) => void; publishWatchMessage?: (workflowId: string, sequence: number, message: string) => void; fleetEventFeed?: FleetEventFeed; }): EventBroadcastingHandle;