--- name: river-ts-streaming description: Type-safe Server-Sent Events (SSE) and WebSocket communication using river.ts library. Use when working with this codebase to: (1) Define typed event schemas with RiverEvents builder, (2) Implement SSE streaming on server with RiverEmitter, (3) Consume SSE streams on client with RiverClient, (4) Handle WebSocket communication with RiverSocketAdapter, (5) Implement request/response RPC patterns over WebSocket, (6) Work with chunked/streamed data events. --- ## Quick Reference river.ts provides three main components: - `RiverEvents` - Type-safe event schema builder - `RiverEmitter` - Server-side SSE streaming - `RiverClient` - Client-side SSE consumption - `RiverSocketAdapter` - WebSocket message handling with request/response support ## Event Definition Define events using the builder pattern: ```typescript import { RiverEvents } from 'river.ts'; const events = new RiverEvents() .defineEvent('message', { message: 'Hello' }) .defineEvent('data', { data: {} as { id: number; name: string } }) .defineEvent('stream', { data: [] as string[], stream: true, chunkSize: 100 }) // Request/response pattern with explicit response type .defineEvent('rpc.call', { data: {} as { method: string; params: unknown }, response: {} as { result: unknown; error?: string } }) // Runtime validation: `data` is inferred from the schema's output .defineEvent('job.run', { schema: z.object({ id: z.string(), priority: z.number() }), responseSchema: z.object({ ok: z.boolean() }) // optional, for request() }) .build(); ``` Reserved event types: `close`, `error` - do not define these. `schema` and `responseSchema` accept any Standard Schema (zod, valibot, arktype; https://standardschema.dev). Events without one are type-checked only. ## Server-Side SSE (RiverEmitter) ```typescript import { RiverEmitter } from 'river.ts/server'; const emitter = RiverEmitter.init(events); // Create SSE stream for HTTP response const stream = emitter.stream({ callback: async (emit, clientId) => { await emit('message', { message: 'Connected' }); await emit('data', { data: { id: 1, name: 'test' } }); }, clientId: 'optional-custom-id', ondisconnect: (clientId) => console.log(`${clientId} disconnected`) }); return new Response(stream, { headers: emitter.headers() }); // Broadcast to all clients await emitter.broadcast('message', { message: 'Update' }); // Send to specific client await emitter.sendToClient('client-id', 'data', { data: { id: 2, name: 'specific' } }); ``` Resumable streams. All options are optional: ```typescript const stream = emitter.stream({ retry: 3000, // ms; written once as `retry:` when the stream opens keepAlive: 15_000, // ms; writes a `: keep-alive` comment so proxies keep the connection lastEventId: request.headers.get('Last-Event-ID'), // read it from the request yourself signal: request.signal, callback: async (emit, clientId, lastEventId) => { // lastEventId is undefined on a first connection; replay everything after it for (const entry of log.after(lastEventId)) { await emit('message', { message: entry.text }, entry.id); // 3rd arg = event id } } }); await emitter.broadcast('message', { message: 'Update' }, 42); // optional id await emitter.sendToClient('client-id', 'message', { message: 'x' }, 43); // optional id ``` Event ids are `string | number` and may be non-ASCII: `stream()` decodes the `Last-Event-ID` header as UTF-8, so the callback gets the id as it was emitted. For `stream: true` events the id is written after the last chunk. ## Client-Side SSE (RiverClient) ```typescript import { RiverClient } from 'river.ts/client'; const client = RiverClient.init(events, { reconnect: true }); client .prepare('http://localhost:3000/events', { method: 'GET' }) .on('message', (data) => console.log(data.message)) .on('data', (data) => console.log(data.id, data.name)) .stream(); // Close connection client.close(); // stream() can be called again after close() ``` Client options (all optional; reconnect is off by default): ```typescript const client = RiverClient.init(events, { reconnect: true, // or { initialDelay: 1000, maxDelay: 30_000 } lastEventId: savedId, // sent as Last-Event-ID on the first connection onInvalid: (type, issues, raw) => {}, // events that fail their schema; never dispatched fetchFn: fetch, headers: { Authorization: 'Bearer ...' } }); client.lastEventId; // read-only: id of the last event received; persist it to resume later client.addEventListener('open', () => {}); // each successful connection client.addEventListener('reconnect', (e) => {}); // (e as CustomEvent).detail = { attempt, delay, error } client.addEventListener('close', () => {}); // stopped for good ``` Reconnect rules: - Retries: network errors, 5xx, 429, and a stream that ends without a `close` event. - Stops for good: HTTP 204, other 4xx, `client.close()`, the server's `close` event. - Delay: `Retry-After` (429/503), else the server's `retry:` value, else exponential backoff with jitter, capped at `maxDelay`. - Every reconnect sends `Last-Event-ID` (fetch path). A GET with no headers at all (none in `init()`, none in `prepare()`) uses the browser `EventSource`, which reconnects and resumes by itself; any header, another method or an initial `lastEventId` selects `fetch`. The parser follows the WHATWG event-stream rules: CRLF/LF/CR line endings, `:` comments, `event`/`data`/`id`/`retry` fields, multiple `data:` lines joined with `\n`, default type `message`. Event data must be JSON. ## WebSocket Adapter (RiverSocketAdapter) ```typescript import { RiverSocketAdapter } from 'river.ts/websocket'; const adapter = new RiverSocketAdapter(events, { debug: false }); // Register event handlers adapter.on('message', (data) => console.log(data)); adapter.off('message', handler); // Unregister // Handle incoming messages (call from ws.onmessage) adapter.handleMessage(messageData); // Send messages adapter.send('data', { data: { id: 1, name: 'test' } }, (msg) => ws.send(msg)); ``` Runtime validation (events with a `schema`): ```typescript import { RiverSocketAdapter, InvalidMessageError } from 'river.ts/websocket'; const adapter = new RiverSocketAdapter(events, { onInvalid: (type, issues, raw) => console.warn(type, issues, raw) }); ``` - `handleMessage()` validates `data` against the event's `schema`; an invalid message goes to `onInvalid` and is not dispatched. Without `onInvalid` it is logged with `console.warn`. - `request()` validates the response against `responseSchema` (or `schema` when the event has neither `responseSchema` nor a `response` type) and rejects with `InvalidMessageError` (`.type`, `.issues`) when invalid. - Handlers receive the schema's output. Async schemas are supported and arrival order is kept. - `RiverClient` does the same for the `data` field of incoming SSE events. ## WebSocket Request/Response Pattern For RPC-style communication with automatic type inference: ```typescript import { RiverSocketAdapter, RequestTimeoutError, WebSocketClosedError } from 'river.ts/websocket'; // Events with explicit response types const events = new RiverEvents() .defineEvent('instance.spawn', { data: {} as { cwd: string }, response: {} as { instanceId: string; status: 'created' | 'error' } }) .build(); const adapter = new RiverSocketAdapter(events); // Route messages through adapter ws.onmessage = (e) => adapter.handleMessage(e.data); ws.onclose = () => adapter.clearPendingRequests(); // Make request - response type is inferred from event definition const response = await adapter.request( 'instance.spawn', { cwd: '/app' }, (msg) => ws.send(msg), 10000 // timeout in ms (default: 30000) ); // response is typed as { instanceId: string; status: 'created' | 'error' } ``` Wire format for request/response: ```json // Request (outgoing) { "type": "instance.spawn", "data": { "cwd": "/app" }, "id": "uuid" } // Response (incoming) - server echoes back the id { "type": "instance.spawn", "data": { "instanceId": "123", "status": "created" }, "id": "uuid" } ``` ## Key Types ```typescript import { EventData, ResponseData, EmitPayload } from 'river.ts'; // EventData - Extract data type for receiving/handling // ResponseData - Extract response type for request() return value // EmitPayload - Extract payload type for emitting (excludes type/stream/chunkSize/schema/responseSchema) // InvalidHandler - (type, issues, raw) => void, the `onInvalid` signature // InvalidMessageError - thrown by request() for a response that fails its schema ``` ## Project Structure ``` src/ ├── index.ts # Main exports (RiverEvents, types) ├── builder.ts # RiverEvents builder class ├── validate.ts # Standard Schema validation shared by client and websocket ├── client/ # RiverClient for SSE consumption ├── server/ # RiverEmitter for SSE streaming ├── websocket/ # RiverSocketAdapter for WebSocket └── types/ ├── core.ts # BaseEvent, EventMap, EventData, ResponseData └── http.ts # HTTPMethods type ``` ## Verification There are no unit tests. Check a change with `bunx tsc --noEmit`, `bun run build`, and a live run against a real server. ## Build and release Build with: `npm run build` (uses unbuild) Output goes to `dist/` with separate entry points for `/client`, `/server`, `/websocket`. Release by bumping the version in `package.json` and pushing to main: `.github/workflows/publish.yml` typechecks, builds and publishes to npm when the registry does not have that version yet. A prerelease version (one containing `-`, such as `1.3.0-test.1` from `npm run bump:test`) is published under the `test` dist-tag; any other version becomes `latest`. The workflow is the only publisher.