yatta/realtime
WebSockets, SSE, rooms, pub/sub and AI streaming.
35 exported symbols and 187 members, read from src/types/realtime.ts.
Construct
createRealtime
functionFactory function creating and registering a configured instance.
createRealtime<TData = Record<string, unknown>>(config?: RealtimeConfig<TData>): RealtimeServer<TData>- config
- Optional server configuration (authentication, rate limits, adapters, etc.).
Returns Configured instance.
import { createRealtime } from "./realtime"; export const realtime = createRealtime({ rateLimit: { messages: 100, windowMs: 10_000 },});createRealtimeClient
functionFactory helper for initializing a .
createRealtimeClient(options: ClientOptions | string): RealtimeClient- options
- Target endpoint string URL or client configuration object.
Returns Configured instance.
const client = createRealtimeClient("http://localhost:3000/realtime");client.on("open", () => console.log("Connected!"));getRealtime
functionRetrieves the global default instance, initializing one if not already created.
getRealtime(): RealtimeServer<any>Returns Global singleton.
sse
functionCreates an SSE `Response` for a given HTTP request, establishing a persistent SSE stream.
sse<TData = Record<string, unknown>>(req: Request, handler: (client: SSEClient<TData>) => void | Promise<void>, options?: SSEOptions, initialData?: TData, metrics?: MetricsCollector, logger?: Logger): Response- req
- Incoming HTTP `Request`.
- handler
- Callback invoked with the connected . May be async.
- options
- SSE transport configuration.
- initialData
- Initial per-connection data attached to `client.data`.
- metrics
- Optional metrics collector.
- logger
- Optional structured logger.
Returns An HTTP `Response` with `Content-Type: text/event-stream`.
export function GET(req: Request) { return sse(req, async (client) => { client.send("connected", { id: client.id }); for await (const event of eventSource) { if (client.signal.aborted) break; client.send("update", event); } });}InMemoryPubSubAdapter
classDefault single-process in-memory pub/sub adapter.
class InMemoryPubSubAdapter implements PubSubAdapter3 members
handlerspropertyhandlers:publishpublish(topic: string, envelope: RealtimeEnvelope): Promise<void>Publishes an envelope to all in-process subscribers of the given topic.
- topic
- Topic name.
- envelope
- Realtime message envelope.
subscribesubscribe(topic: string, handler: (env: RealtimeEnvelope) => void): Promise<() => void>Subscribes to a topic with a handler callback.
- topic
- Topic name.
- handler
- Callback invoked on each published envelope.
Returns Async cleanup function to unsubscribe.
JobTracker
classTracks and broadcasts the lifecycle of a long-running background job over SSE.
class JobTrackerconst tracker = realtime.job("reports");tracker.start("Generating report...");await generateReport((pct) => tracker.progress(pct, `${pct}% done`));tracker.done({ url: "/reports/latest.pdf" }, "Report ready");13 members
abortControllerpropertyabortController:cancelCallbackspropertycancelCallbacks: Array<(reason?: string) => void>statepropertystate: JobState<TResult>Current serializable job state snapshot.
signalgettersignal: AbortSignal`AbortSignal` that is triggered when the job is cancelled. Use to cooperatively stop work.
dispatchdispatch()startstart(message?: ): thisTransitions the job to `running` and broadcasts the update.
- message
- Status message displayed in client UIs.
Returns `this` for chaining.
progressprogress(percent: number, message?: string, extra?: Record<string, unknown>): thisUpdates the job's progress percentage and optional status message.
- percent
- Progress 0–100 (clamped automatically).
- message
- Optional status message.
- extra
- Optional extra metadata to merge into `state.extra`.
Returns `this` for chaining.
pausepause(message?: ): thisTransitions the job to `paused` and broadcasts the update.
- message
- Optional pause message.
Returns `this` for chaining.
resumeresume(message?: ): thisTransitions the job from `paused` back to `running`.
- message
- Optional resume message.
Returns `this` for chaining.
donedone(result?: TResult, message?: ): voidMarks the job as `completed`, setting progress to 100 and broadcasting the result.
- result
- Optional result payload delivered to subscribed clients.
- message
- Completion message.
failfail(error: string | Error): voidMarks the job as `failed` and broadcasts the error message.
- error
- Error instance or error message string.
cancelcancel(reason?: ): voidCancels the job, triggers the abort signal, and invokes all registered cancel callbacks.
- reason
- Human-readable cancellation reason.
onCancelonCancel(callback: (reason?: string) => void): thisRegisters a callback to be invoked when the job is cancelled.
- callback
- Function called with the cancellation reason.
Returns `this` for chaining.
MetricsCollector
classInternal metrics accumulator tracking connection counts, message throughput, and byte transfer totals.
class MetricsCollector16 members
wsCountpropertywsCount:sseCountpropertysseCount:msgSentpropertymsgSent:msgReceivedpropertymsgReceived:bytesOutpropertybytesOut:bytesInpropertybytesIn:activeJobsCountpropertyactiveJobsCount:incWSincWS()Increments the active WebSocket connection counter by 1.
decWSdecWS()Decrements the active WebSocket connection counter by 1 (floored at 0).
incSSEincSSE()Increments the active SSE connection counter by 1.
decSSEdecSSE()Decrements the active SSE connection counter by 1 (floored at 0).
recordSentrecordSent(bytes: number)Records an outbound message and increments cumulative byte totals.
- bytes
- Byte length of the dispatched frame.
recordReceivedrecordReceived(bytes: number)Records an inbound message and increments cumulative byte totals.
- bytes
- Byte length of the received frame.
incJobincJob()Increments the count of active background job trackers.
decJobdecJob()Decrements the count of active background job trackers (floored at 0).
snapshotsnapshot(topicCount: number): RealtimeStatsReturns a point-in-time snapshot copy of current operational metrics.
- topicCount
- Current number of active named SSE channels.
Returns Immutable snapshot.
RealtimeBroadcaster
classPublishes events to a designated topic across distributed pub/sub adapters,
class RealtimeBroadcaster3 members
sendsend<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K], id?: string): voidPublishes a typed event and payload to all topic subscribers across transports.
- event
- Registered event name.
- data
- Typed payload data.
- id
- Optional explicit UUID for message tracking.
sendsend(event: string, data: unknown, id?: string): voidPublishes an arbitrary named event and payload to all topic subscribers across transports.
- event
- Event name string.
- data
- Payload data.
- id
- Optional explicit UUID for message tracking.
sendsend(event: string, data: unknown, id?: string): void
RealtimeClient
classIsomorphic Realtime Client SDK supporting native WebSockets and Server-Sent Events (SSE).
class RealtimeClientconst client = new RealtimeClient({ url: "http://localhost:3000/realtime" }); client.on("chat.message", (msg) => { console.log("New message:", msg.text);}); client.subscribe("room:101");client.send("chat.message", { text: "Hello!" });17 members
wspropertyws: WebSocket | nullespropertyes: EventSource | nulllistenerspropertylisteners:subscribedTopicspropertysubscribedTopics:reconnectAttemptspropertyreconnectAttempts:isExplicitClosepropertyisExplicitClose:connectconnect(): voidinitWebSocketinitWebSocket(): voidinitSSEinitSSE(): voidscheduleReconnectscheduleReconnect(): voidonon<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], env?: RealtimeEnvelope) => void): () => voidRegisters an event listener for a typed registered event name.
- event
- Registered event name.
- handler
- Callback receiving the typed payload and optional raw envelope.
Returns Cleanup function to deregister this listener.
onon(event: string, handler: (data: any, env?: RealtimeEnvelope) => void): () => voidRegisters an event listener for an arbitrary event name or lifecycle event (`"open"`, `"close"`, `"error"`).
- event
- Event name string.
- handler
- Callback receiving the payload and optional raw envelope.
Returns Cleanup function to deregister this listener.
onon(event: string, handler: (data: any, env?: RealtimeEnvelope) => void): () => voidsubscribesubscribe(topic: string): () => voidSubscribes to a pub/sub topic on the server. If using WebSocket, automatically resubscribes on reconnect.
- topic
- Topic identifier to join.
Returns Cleanup function to unsubscribe from this topic.
sendsend(event: string, data: unknown): voidSends an outbound event to the server. Requires an active WebSocket transport.
- event
- Event name.
- data
- Payload to serialize and dispatch.
dispatchdispatch(event: string, data: any, envelope?: RealtimeEnvelope)closeclose(): voidCloses active WebSocket and EventSource connections and disables automatic reconnect attempts.
RealtimeError
classBase exception thrown by the Yatta Realtime engine for protocol, connection, and stream errors.
class RealtimeError extends Errortry { client.send("chat.message", undefined);} catch (error) { if (error instanceof RealtimeError) { console.error(error.message, error.code); }}RealtimeServer
classCore realtime engine orchestrating WebSocket and SSE transports, pub/sub routing,
Typically created via rather than instantiated directly.
class RealtimeServerimport { createRealtime } from "./realtime"; export const realtime = createRealtime({ handlers: { open: (client) => console.log("Connected:", client.id), message: (client, event, data) => client.send(event, data), },}); // In Bun.serve:Bun.serve({ fetch(req, server) { if (new URL(req.url).pathname === "/ws") { return realtime.connect(req, server); } return new Response("Not found", { status: 404 }); }, websocket: realtime.websocket,});23 members
sseChannelspropertysseChannels:jobspropertyjobs:rateLimitBucketspropertyrateLimitBuckets:messageValidatorspropertymessageValidators:metricspropertymetrics:Live connection and throughput metrics. Query via `.stats()`.
adapterpropertyadapter: PubSubAdapterActive pub/sub adapter (in-memory by default, swappable for Redis/NATS).
loggerpropertylogger: LoggerStructured logger instance used throughout the engine.
bindServerbindServer(server: Server<unknown>): voidBinds the live `Bun.Server` instance to this engine, enabling WebSocket pub/sub broadcasts.
- server
- The `Bun.Server` returned by `Bun.serve(...)`.
servergetterserver: Server<unknown>Returns the bound `Bun.Server` instance.
registerValidatorregisterValidator<K extends keyof RegisteredEvents>(event: K, validator: (data: RegisteredEvents[K]) => boolean): voidRegisters a type-safe per-event inbound message validator.
- event
- Registered event name to validate.
- validator
- Function receiving the typed payload; return `true` to accept, `false` to reject.
tsts realtime.registerValidator("chat.message", (data) => typeof data.text === "string" && data.text.length <= 2000);checkAuthorizeJoincheckAuthorizeJoin(client: SocketClient<TData>, topic: string): Promise<boolean>Runs the `authorize.join` guard from config against a client/topic pair.
- client
- The socket client attempting to subscribe.
- topic
- Topic name being joined.
checkAuthorizeSendcheckAuthorizeSend(client: SocketClient<TData>, event: string, data: unknown): Promise<boolean>Runs the `authorize.send` guard from config before delivering an inbound message.
- client
- The socket client sending the message.
- event
- Event name.
- data
- Event payload.
reportErrorreportError(error: unknown, context?: string): voidRoutes an internal engine error to `config.onError` if provided,
- error
- The caught error value.
- context
- Short label for the error site (e.g. `"adapter:publish"`, `"websocket:publish"`).
isRateLimitedisRateLimited(clientId: string): booleanChecks whether a client has exceeded the configured message rate limit within the current window.
upgradeupgrade(req: Request, server: Server<unknown>, customData?: Partial<TData>): Promise<boolean>Attempts to upgrade an incoming HTTP request to a WebSocket connection.
Prefer for a unified handler that auto-selects between WebSocket and SSE.
- req
- Incoming HTTP `Request`.
- server
- The `Bun.Server` instance.
- customData
- Optional extra data merged into the session after authentication.
Returns `true` if the upgrade was successful, `false` if authentication failed or upgrade was refused.
connectconnect(req: Request, server: Server<unknown>, customData?: Partial<TData>): Promise<Response | boolean>Unified connection entry point — auto-routes between WebSocket upgrade and SSE streaming.
Routing logic: - If the request has an `Upgrade: websocket` header → delegates to . - If the request has `Accept: text/event-stream` → returns an SSE `Response`. - Otherwise → returns a `400 Bad Request` listing supported transports. For SSE connections, the `topic` query parameter (default `"global"`) selects the channel, and `Last-Event-ID` / `lastEventId` query param enables reconnection replay.
- req
- Incoming HTTP `Request`.
- server
- The `Bun.Server` instance.
- customData
- Optional extra data merged into the session after authentication.
Returns A `Response` (for SSE or errors) or `true`/`false` (for WebSocket upgrades).
tsts Bun.serve({ fetch: (req, server) => realtime.connect(req, server), websocket: realtime.websocket,});websocketgetterwebsocket:WebSocket handler object to pass directly to `Bun.serve({ websocket })`.
tsts Bun.serve({ fetch: (req, server) => realtime.connect(req, server), websocket: realtime.websocket,});getChannelgetChannel(name: string): SSEChannel | undefinedReturns an existing named SSE channel without creating one.
- name
- Channel/topic name.
Returns The if it exists, or `undefined`.
channelchannel(name: string): SSEChannelReturns the named SSE channel, creating it if it doesn't exist yet.
- name
- Channel/topic name.
Returns The for the given name.
tsts const ch = realtime.channel("notifications");ch.broadcast("alert", { level: "info", message: "Deploy complete" });namespacenamespace(prefix: string)Fluent helper to scope rooms (WebSockets) and channels (SSE) under a common prefix.
- prefix
- Common namespace prefix (e.g. `"tenant:101"` or `"chat"`).
Returns Scoped builder exposing `.room(name)` and `.channel(name)`.
tsts const tenant = realtime.namespace("tenant:org_42");tenant.room("billing").send("invoice:paid", { id: "inv_123" });toto(topic: string): RealtimeBroadcasterReturns a targeting a specific topic across both
- topic
- Topic name to publish to.
Returns Topic-scoped broadcaster.
tsts realtime.to("orders").send("order:created", { id: "ord_99" });jobjob<T = unknown>(jobId: string): JobTracker<T>Retrieves or instantiates a tracked for a long-running background task.
- jobId
- Unique identifier for the job.
tsts const tracker = realtime.job("video-transcode-42");tracker.start("Transcoding 1080p...");tracker.progress(50, "Halfway done");tracker.done({ url: "/video.mp4" });statsstats(): RealtimeStatsCaptures a real-time operational metrics snapshot of active connections,
Returns Current telemetry snapshot.
RoomBroadcaster
classSends a realtime event to all WebSocket clients subscribed to a pub/sub topic.
class RoomBroadcaster1 member
sendsend(event: string, data: unknown): voidBroadcast a named event with payload to all topic subscribers.
- event
- Event name.
- data
- Event payload (JSON-serialized).
SSEChannel
classNamed SSE broadcast channel — a group of connections subscribed to the same topic.
class SSEChannelconst ch = realtime.channel("notifications");ch.broadcast("alert", { level: "info", message: "Deploy complete" });console.log(`${ch.size} clients listening`);9 members
subscriberspropertysubscribers:messageHistorypropertymessageHistory: Array<{ id: string; event: string; data: unknown; }>sizegettersize: numberReturns the number of active subscribers currently connected to this channel.
subscribesubscribe(client: SSEClient<any>): voidAdds a client to this channel and wires up an abort listener to auto-remove on disconnect.
- client
- SSE client to subscribe.
removeremove(client: SSEClient<any>): voidRemoves a client from this channel's subscriber set.
- client
- SSE client to remove.
replayHistoryreplayHistory(client: SSEClient<any>, lastEventId: string): voidReplays missed historical events to a reconnecting client, starting after `lastEventId`.
- client
- The reconnecting SSE client.
- lastEventId
- The `Last-Event-ID` header value from the client's reconnect request.
broadcastbroadcast<K extends keyof TEvents>(event: K, data: TEvents[K], id?: string): voidBroadcasts a typed event to all connected subscribers on this channel.
- event
- Event name from the registered event map.
- data
- Event payload.
- id
- Optional message ID; auto-generated UUID if omitted.
tsts channel.broadcast("order.updated", { orderId: "x", status: "shipped" });broadcastbroadcast(event: string, data: unknown, id?: string): voidbroadcastbroadcast(event: string, data: unknown, id?: string): void
SSEClient
classRepresents a single SSE connection to a browser or HTTP client.
class SSEClientsse(req, (client) => { client.send("welcome", { message: "connected!" }); const timer = setInterval(() => { client.send("tick", { time: Date.now() }); }, 1000); client.signal.addEventListener("abort", () => clearInterval(timer));});14 members
isClosedpropertyisClosed:heartbeatTimerpropertyheartbeatTimer?: ReturnType<typeof setInterval>bufferedEventCountpropertybufferedEventCount:signalpropertysignal: AbortSignal`AbortSignal` that fires when the client disconnects. Use to clean up resources.
idpropertyid: stringUnique UUID assigned to this SSE connection.
datapropertydata: TDataMutable per-connection session data. Useful for storing auth state or metadata.
writewrite(payload: string): booleansendsend<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K], id?: string): thisSend a typed named event with data payload and optional message ID.
- event
- Registered or arbitrary event name.
- data
- Event payload (must be JSON-serializable).
- id
- Optional event ID for `Last-Event-ID` reconnection support.
Returns `this` for chaining.
tsts client.send("user.joined", { userId: "abc", name: "Alice" });client.send("chat.message", { text: "Hello!" }, "msg-001");client.send({ raw: "anonymous data" }); // no event namesendsend(event: string, data: unknown, id?: string): thissendsend(data: unknown): thissendsend(arg1: unknown, arg2?: unknown, id?: string): thiscommentcomment(text: string): thisSend an SSE comment line (begins with `:`). Useful for keep-alive pings or debug annotations.
- text
- Comment text.
Returns `this` for chaining.
streamTextstreamText(source: AsyncIterable<string> | ReadableStream<string> | ((emit: (chunk: string) => void) => Promise<void>), options?: { eventName?: string; streamId?: string; usage?: () => { inputTokens?: number; outputTokens?: number; }; }): Promise<void>First-class LLM / AI token streaming protocol.
Handles client disconnect gracefully — will abort the stream and emit a `done/aborted` frame.
- source
- Token source: async iterable, readable stream, or a callback-based producer.
- options
- Streaming options including event name, stream ID override, and token usage reporter.
tsts await client.streamText(llm.stream("Tell me a joke"), { eventName: "ai:stream", streamId: "session-123", usage: () => ({ inputTokens: 10, outputTokens: 25 }),});closeclose(): voidCloses the SSE connection, stopping the heartbeat and terminating the readable stream.
TypedSocket
classConcrete implementation of wrapping a raw Bun `ServerWebSocket`.
class TypedSocket implements SocketClient<TData>13 members
idpropertyid: stringUnique UUID for this socket connection.
datagetterdata: TDatadatasetterdata:sendsend(event: string, data: unknown): voidSend a named event JSON envelope to this client.
sendJSONsendJSON(event: string, data: unknown): voidSerialize and send a named event as a JSON `RealtimeEnvelope` to this specific client.
- event
- Event name.
- data
- Event payload.
sendTextsendText(text: string): voidSend a raw UTF-8 text WebSocket frame to this client.
- text
- Raw text to send.
sendBinarysendBinary(data: ArrayBufferView | ArrayBuffer): voidSend a raw binary WebSocket frame to this client.
- data
- Binary buffer to send.
joinjoin(topic: string): Promise<boolean>Attempt to join a pub/sub topic, pending `authorize.join` validation.
- topic
- Topic name to join.
Returns `true` if authorized and subscribed; `false` if denied.
leaveleave(topic: string): voidUnsubscribe from a pub/sub topic.
- topic
- Topic name to leave.
isSubscribedisSubscribed(topic: string): booleanCheck whether this socket is subscribed to a given topic.
- topic
- Topic name to check.
toto(topic: string): RoomBroadcasterReturns a broadcaster that sends to all sockets on the topic (including this socket).
- topic
- Topic name.
broadcastbroadcast(topic: string): RoomBroadcasterReturns a broadcaster that sends to all sockets on the topic, **excluding** this socket.
- topic
- Topic name.
closeclose(code?: number, reason?: string): voidClose this WebSocket connection with an optional status code and reason.
- code
- WebSocket close code (default: 1000 normal closure).
- reason
- Human-readable close reason.
Types
ClientOptions
interfaceConfiguration options for the isomorphic SDK.
interface ClientOptions4 members
urlpropertyurl: stringTarget server endpoint URL (e.g. `"http://localhost:3000/realtime"` or `"ws://localhost:3000/realtime"`).
transportpropertytransport?: "auto" | "websocket" | "sse"Transport mode: `"websocket"`, `"sse"`, or `"auto"` (prefers WebSocket if available). Defaults to `"auto"`.
protocolspropertyprotocols?: string | string[]Optional sub-protocols for WebSocket handshakes.
reconnectpropertyreconnect?: { enabled?: boolean; minDelay?: number; maxDelay?: number; factor?: number; jitter?: boolean; }Exponential backoff and jitter configuration for automatic reconnection.
JobState
interfaceSerializable snapshot of a job's current lifecycle state, broadcast to clients via SSE.
interface JobState<TResult = unknown>8 members
jobIdpropertyjobId: stringUnique job identifier.
statuspropertystatus: JobStatusCurrent lifecycle status.
percentpropertypercent: numberProgress percentage (0–100).
messagepropertymessage: stringHuman-readable status message for display in UIs.
resultpropertyresult?: TResultResult value emitted on successful completion.
errorpropertyerror?: stringError message string emitted on failure.
updatedAtpropertyupdatedAt: numberUnix millisecond timestamp of the last state update.
extrapropertyextra?: Record<string, unknown>Optional arbitrary extra metadata for custom progress data.
Logger
interfaceStructured logger interface for realtime engine diagnostics.
interface Logger4 members
debugdebug(...args: unknown[]): voidEmit debug-level trace messages.
infoinfo(...args: unknown[]): voidEmit informational operational messages.
warnwarn(...args: unknown[]): voidEmit warning-level messages.
errorerror(...args: unknown[]): voidEmit error-level messages.
PubSubAdapter
interfacePluggable publish/subscribe adapter interface enabling horizontal scaling across multiple server nodes.
interface PubSubAdapterclass RedisPubSubAdapter implements PubSubAdapter { async publish(topic: string, envelope: RealtimeEnvelope) { await redis.publish(topic, JSON.stringify(envelope)); } async subscribe(topic: string, handler: (envelope: RealtimeEnvelope) => void) { const sub = redis.duplicate(); await sub.subscribe(topic, (msg) => handler(JSON.parse(msg))); return () => sub.unsubscribe(topic); }}3 members
publishpublish(topic: string, envelope: RealtimeEnvelope): Promise<void>Publish an envelope to all subscribers of a topic across the cluster.
- topic
- Topic name.
- envelope
- Realtime message envelope.
subscribesubscribe(topic: string, handler: (envelope: RealtimeEnvelope) => void): Promise<() => void>Subscribe to a topic, returning an unsubscribe cleanup function.
- topic
- Topic name.
- handler
- Callback invoked when an envelope is received for this topic.
Returns Async function that terminates the subscription when called.
closeclose(): Promise<void>?Optional cleanup hook for closing client connections.
RateLimitConfig
interfaceConfiguration for WebSocket inbound message rate limiting.
interface RateLimitConfigcreateRealtime({ rateLimit: { messages: 50, windowMs: 10_000, onRateLimit: (client) => client.send("error:rate_limit", { message: "Slow down!" }), },});3 members
messagespropertymessages: numberMaximum number of messages allowed per client within `windowMs`.
windowMspropertywindowMs: numberSliding window duration in milliseconds over which `messages` are counted.
onRateLimitpropertyonRateLimit?: (client: SocketClient<any>) => voidCustom handler invoked when a client exceeds the rate limit.
RealtimeConfig
interfaceTop-level configuration for / .
interface RealtimeConfig<TData = Record<string, unknown>>const realtime = createRealtime<{ userId: string }>({ authenticate: async (req) => { const user = await verifyToken(req.headers.get("authorization")); return user ? { userId: user.id } : null; }, authorize: { join: (client, topic) => topic.startsWith(`user:${client.data.userId}`), }, rateLimit: { messages: 100, windowMs: 10_000 }, handlers: { open: (client) => console.log("Connected:", client.data.userId), },});11 members
handlerspropertyhandlers?: RealtimeHandlers<TData>WebSocket lifecycle and message event handlers. See .
authenticatepropertyauthenticate?: (req: Request) => Promise<TData | null> | TData | nullAuthentication hook invoked on every incoming connection (WebSocket upgrade or SSE).
authorizepropertyauthorize?: { join?: (client: SocketClient<TData>, topic: string) => boolean | Promise<boolean>; send?: (client: SocketClient<TData>, event: string, data: unknown) => boolean | Promise<boolean>; }Fine-grained authorization guards.
validatepropertyvalidate?: (event: string, data: unknown) => boolean | { valid: boolean; error?: string; }Global inbound message payload validator.
rateLimitpropertyrateLimit?: RateLimitConfigWebSocket inbound message rate limiting. See .
adapterpropertyadapter?: PubSubAdapterPluggable pub/sub adapter for horizontal multi-node scaling.
loggerpropertylogger?: LoggerCustom structured logger. Defaults to `console`-based output.
perMessageDeflatepropertyperMessageDeflate?: booleanEnable WebSocket `permessage-deflate` compression.
maxPayloadLengthpropertymaxPayloadLength?: numberMaximum allowed WebSocket message payload size in bytes.
sseOptionspropertysseOptions?: SSEOptionsDefault SSE transport options applied to all channels. See .
onErrorpropertyonError?: (error: unknown, context?: string) => voidGlobal error handler for adapter, WebSocket publish, and other internal errors.
- error
- The caught error value.
- context
- A short string identifying the error site (e.g. `"adapter:publish"`).
RealtimeEnvelope
interfaceStandard realtime message envelope wrapping all events sent over WebSocket or SSE channels.
interface RealtimeEnvelope<T = unknown>5 members
idpropertyid: stringUnique event message identifier (UUID).
eventpropertyevent: stringEvent type name (e.g. `"chat.message"`, `"job:status"`).
topicpropertytopic?: stringOptional pub/sub topic name this event was routed on.
datapropertydata: TStrongly-typed event payload data.
timestamppropertytimestamp: numberUnix millisecond timestamp when the event was dispatched.
RealtimeHandlers
interfaceLifecycle and event callbacks for the WebSocket transport layer.
interface RealtimeHandlers<TData = Record<string, unknown>>createRealtime({ handlers: { open(client) { console.log("Client connected:", client.id); }, message(client, event, data) { console.log(event, data); }, close(client, code, reason) { console.log("Disconnected:", code, reason); }, },});9 members
openpropertyopen?: (client: SocketClient<TData>) => void | Promise<void>Called when a new WebSocket client connects.
messagepropertymessage?: (client: SocketClient<TData>, event: string, data: any) => void | Promise<void>Called for each inbound WebSocket message, after rate limiting, validation, and authorization.
closepropertyclose?: (client: SocketClient<TData>, code: number, reason: string) => void | Promise<void>Called when a client disconnects. `code` and `reason` follow the WebSocket close handshake.
drainpropertydrain?: (client: SocketClient<TData>) => void | Promise<void>Called when the send buffer drains after backpressure.
errorpropertyerror?: (client: SocketClient<TData>, error: unknown) => void | Promise<void>Called on low-level WebSocket transport errors.
pingpropertyping?: (client: SocketClient<TData>, data: Buffer) => void | Promise<void>Called when a WebSocket ping frame is received.
pongpropertypong?: (client: SocketClient<TData>, data: Buffer) => void | Promise<void>Called when a WebSocket pong frame is received.
subscribepropertysubscribe?: (client: SocketClient<TData>, topic: string) => void | Promise<void>Called after a client successfully joins a pub/sub topic.
unsubscribepropertyunsubscribe?: (client: SocketClient<TData>, topic: string) => void | Promise<void>Called after a client leaves a pub/sub topic.
RealtimeRegister
interfaceGlobal interface augmented by application code for typed event names and payloads.
interface RealtimeRegisterdeclare module "./realtime" { interface RealtimeRegister { events: { "chat.message": { user: string; text: string }; "order.status": { orderId: string; status: string }; }; }}RealtimeStats
interfaceSnapshot of realtime engine operational statistics at a given moment.
interface RealtimeStats8 members
websocketConnectionspropertywebsocketConnections: numberNumber of active WebSocket client connections.
sseConnectionspropertysseConnections: numberNumber of active SSE (Server-Sent Events) client connections.
topicspropertytopics: numberNumber of active named SSE channels.
messagesSentpropertymessagesSent: numberCumulative count of messages dispatched outward.
messagesReceivedpropertymessagesReceived: numberCumulative count of messages received from clients.
bytesSentpropertybytesSent: numberCumulative bytes sent across all transports.
bytesReceivedpropertybytesReceived: numberCumulative bytes received from all clients.
activeJobspropertyactiveJobs: numberCount of currently active in-progress job trackers.
SocketClient
interfacePublic interface representing a connected WebSocket client.
interface SocketClient<TData = Record<string, unknown>>14 members
idpropertyid: stringUnique UUID assigned to this connection.
datapropertydata: TDataMutable per-connection session data (auth context, user info, etc.).
rawpropertyraw: ServerWebSocket<TData>The raw underlying `ServerWebSocket` instance (Bun-native).
sendsend<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K]): voidSend a typed registered event.
sendsend(event: string, data: unknown): voidSend an arbitrary named event with unknown payload.
sendJSONsendJSON(event: string, data: unknown): voidSerialize an event as a JSON string and send it.
sendTextsendText(text: string): voidSend a raw UTF-8 text frame.
sendBinarysendBinary(data: ArrayBufferView | ArrayBuffer): voidSend a raw binary frame.
joinjoin(topic: string): Promise<boolean>Subscribe to a pub/sub topic, pending authorization check. Returns `true` on success.
leaveleave(topic: string): voidUnsubscribe from a pub/sub topic.
isSubscribedisSubscribed(topic: string): booleanReturns whether this socket is subscribed to a given topic.
toto(topic: string): RoomBroadcasterReturns a broadcaster that sends to all sockets on `topic` (including this one).
broadcastbroadcast(topic: string): RoomBroadcasterReturns a broadcaster that sends to all sockets on `topic` (excluding this one).
closeclose(code?: number, reason?: string): voidClose the WebSocket with an optional code and reason.
SSEOptions
interfaceConfiguration options for the Server-Sent Events (SSE) transport layer.
interface SSEOptions6 members
headerspropertyheaders?: Record<string, string>Additional HTTP response headers to merge into the SSE response.
heartbeatIntervalpropertyheartbeatInterval?: numberInterval in milliseconds between keep-alive comment pings (`: keep-alive ping`).
retrypropertyretry?: numberClient-side reconnect delay hint in milliseconds, sent as the SSE `retry:` field.
maxBufferedEventspropertymaxBufferedEvents?: numberMaximum number of events buffered for a slow consumer before backpressure is triggered.
historySizepropertyhistorySize?: numberNumber of recent events to retain for `Last-Event-ID` reconnection replay.
onBackpressurepropertyonBackpressure?: (client: SSEClient) => voidCustom backpressure handler invoked when a slow client's buffer is full.
AIStreamChunk
typeAI/LLM streaming chunk frame sent during active token streaming.
type AIStreamChunk = { type: "chunk"; streamId: string; text: string; index: number; }AIStreamDone
typeAI/LLM streaming completion frame emitted when the stream finishes.
type AIStreamDone = { type: "done"; streamId: string; reason: "completed" | "aborted" | "cancelled" | "failed"; usage?: { inputTokens?: number; outputTokens?: number; totalTokens?: number; }; }AIStreamEnvelope
typeDiscriminated union of AI/LLM stream frame types transmitted over SSE.
type AIStreamEnvelope = AIStreamChunk | AIStreamDone | AIStreamErrorAIStreamError
typeAI/LLM streaming error frame emitted when the stream encounters an irrecoverable error.
type AIStreamError = { type: "error"; streamId: string; error: { code: string; message: string; }; }EventPayload
typeResolves the payload type for a given event name `K`.
type EventPayload<K extends string> = IsRegistryConfigured extends true ? K extends keyof RegisteredEvents ? RegisteredEvents[K] : unknown : unknownIsRegistryConfigured
typeType-level boolean indicating if a custom event registry has been configured via module augmentation.
type IsRegistryConfigured = RealtimeRegister extends { events: any; } ? true : falseJobStatus
typeAll possible lifecycle states for a tracked background job.
type JobStatus = "pending" | "running" | "paused" | "completed" | "failed" | "cancelled"RegisteredEvents
typeExtracts the registered events map from , or falls back to a generic event map.
type RegisteredEvents = RealtimeRegister extends { events: infer E extends Record<string, unknown>; } ? E : Record<string, unknown>AuthConfig or JobPayload — declaring your schema once is enough for the rest to follow. See Typed keys.