API reference

yatta/realtime

WebSockets, SSE, rooms, pub/sub and AI streaming.

35 exported symbols and 187 members, read from src/types/realtime.ts.

Construct

createRealtime

function

Factory 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.

tsts
import { createRealtime } from "./realtime"; export const realtime = createRealtime({  rateLimit: { messages: 100, windowMs: 10_000 },});

createRealtimeClient

function

Factory helper for initializing a .

createRealtimeClient(options: ClientOptions | string): RealtimeClient
options
Target endpoint string URL or client configuration object.

Returns Configured instance.

tsts
const client = createRealtimeClient("http://localhost:3000/realtime");client.on("open", () => console.log("Connected!"));

getRealtime

function

Retrieves the global default instance, initializing one if not already created.

getRealtime(): RealtimeServer<any>

Returns Global singleton.

sse

function

Creates 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`.

tsts
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

class

Default single-process in-memory pub/sub adapter.

3 members

  • handlersproperty
    handlers:
  • publish
    publish(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.
  • subscribe
    subscribe(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

class

Tracks and broadcasts the lifecycle of a long-running background job over SSE.

tsts
const 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

  • abortControllerproperty
    abortController:
  • cancelCallbacksproperty
    cancelCallbacks: Array<(reason?: string) => void>
  • stateproperty
    state: JobState<TResult>

    Current serializable job state snapshot.

  • signalgetter
    signal: AbortSignal

    `AbortSignal` that is triggered when the job is cancelled. Use to cooperatively stop work.

  • dispatch
    dispatch()
  • start
    start(message?: ): this

    Transitions the job to `running` and broadcasts the update.

    message
    Status message displayed in client UIs.

    Returns `this` for chaining.

  • progress
    progress(percent: number, message?: string, extra?: Record<string, unknown>): this

    Updates 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.

  • pause
    pause(message?: ): this

    Transitions the job to `paused` and broadcasts the update.

    message
    Optional pause message.

    Returns `this` for chaining.

  • resume
    resume(message?: ): this

    Transitions the job from `paused` back to `running`.

    message
    Optional resume message.

    Returns `this` for chaining.

  • done
    done(result?: TResult, message?: ): void

    Marks the job as `completed`, setting progress to 100 and broadcasting the result.

    result
    Optional result payload delivered to subscribed clients.
    message
    Completion message.
  • fail
    fail(error: string | Error): void

    Marks the job as `failed` and broadcasts the error message.

    error
    Error instance or error message string.
  • cancel
    cancel(reason?: ): void

    Cancels the job, triggers the abort signal, and invokes all registered cancel callbacks.

    reason
    Human-readable cancellation reason.
  • onCancel
    onCancel(callback: (reason?: string) => void): this

    Registers a callback to be invoked when the job is cancelled.

    callback
    Function called with the cancellation reason.

    Returns `this` for chaining.

MetricsCollector

class

Internal metrics accumulator tracking connection counts, message throughput, and byte transfer totals.

16 members

  • wsCountproperty
    wsCount:
  • sseCountproperty
    sseCount:
  • msgSentproperty
    msgSent:
  • msgReceivedproperty
    msgReceived:
  • bytesOutproperty
    bytesOut:
  • bytesInproperty
    bytesIn:
  • activeJobsCountproperty
    activeJobsCount:
  • incWS
    incWS()

    Increments the active WebSocket connection counter by 1.

  • decWS
    decWS()

    Decrements the active WebSocket connection counter by 1 (floored at 0).

  • incSSE
    incSSE()

    Increments the active SSE connection counter by 1.

  • decSSE
    decSSE()

    Decrements the active SSE connection counter by 1 (floored at 0).

  • recordSent
    recordSent(bytes: number)

    Records an outbound message and increments cumulative byte totals.

    bytes
    Byte length of the dispatched frame.
  • recordReceived
    recordReceived(bytes: number)

    Records an inbound message and increments cumulative byte totals.

    bytes
    Byte length of the received frame.
  • incJob
    incJob()

    Increments the count of active background job trackers.

  • decJob
    decJob()

    Decrements the count of active background job trackers (floored at 0).

  • snapshot
    snapshot(topicCount: number): RealtimeStats

    Returns a point-in-time snapshot copy of current operational metrics.

    topicCount
    Current number of active named SSE channels.

    Returns Immutable snapshot.

RealtimeBroadcaster

class

Publishes events to a designated topic across distributed pub/sub adapters,

3 members

  • send
    send<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K], id?: string): void

    Publishes 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.
  • send
    send(event: string, data: unknown, id?: string): void

    Publishes 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.
  • send
    send(event: string, data: unknown, id?: string): void

RealtimeClient

class

Isomorphic Realtime Client SDK supporting native WebSockets and Server-Sent Events (SSE).

tsts
const 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

  • wsproperty
    ws: WebSocket | null
  • esproperty
    es: EventSource | null
  • listenersproperty
    listeners:
  • subscribedTopicsproperty
    subscribedTopics:
  • reconnectAttemptsproperty
    reconnectAttempts:
  • isExplicitCloseproperty
    isExplicitClose:
  • connect
    connect(): void
  • initWebSocket
    initWebSocket(): void
  • initSSE
    initSSE(): void
  • scheduleReconnect
    scheduleReconnect(): void
  • on
    on<K extends keyof RegisteredEvents>(event: K, handler: (data: RegisteredEvents[K], env?: RealtimeEnvelope) => void): () => void

    Registers 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.

  • on
    on(event: string, handler: (data: any, env?: RealtimeEnvelope) => void): () => void

    Registers 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.

  • on
    on(event: string, handler: (data: any, env?: RealtimeEnvelope) => void): () => void
  • subscribe
    subscribe(topic: string): () => void

    Subscribes 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.

  • send
    send(event: string, data: unknown): void

    Sends an outbound event to the server. Requires an active WebSocket transport.

    event
    Event name.
    data
    Payload to serialize and dispatch.
  • dispatch
    dispatch(event: string, data: any, envelope?: RealtimeEnvelope)
  • close
    close(): void

    Closes active WebSocket and EventSource connections and disables automatic reconnect attempts.

RealtimeError

class

Base exception thrown by the Yatta Realtime engine for protocol, connection, and stream errors.

class RealtimeError extends Error
tsts
try {  client.send("chat.message", undefined);} catch (error) {  if (error instanceof RealtimeError) {    console.error(error.message, error.code);  }}

RealtimeServer

class

Core realtime engine orchestrating WebSocket and SSE transports, pub/sub routing,

Typically created via rather than instantiated directly.

tsts
import { 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

  • sseChannelsproperty
    sseChannels:
  • jobsproperty
    jobs:
  • rateLimitBucketsproperty
    rateLimitBuckets:
  • messageValidatorsproperty
    messageValidators:
  • metricsproperty
    metrics:

    Live connection and throughput metrics. Query via `.stats()`.

  • adapterproperty
    adapter: PubSubAdapter

    Active pub/sub adapter (in-memory by default, swappable for Redis/NATS).

  • loggerproperty
    logger: Logger

    Structured logger instance used throughout the engine.

  • bindServer
    bindServer(server: Server<unknown>): void

    Binds the live `Bun.Server` instance to this engine, enabling WebSocket pub/sub broadcasts.

    server
    The `Bun.Server` returned by `Bun.serve(...)`.
  • servergetter
    server: Server<unknown>

    Returns the bound `Bun.Server` instance.

  • registerValidator
    registerValidator<K extends keyof RegisteredEvents>(event: K, validator: (data: RegisteredEvents[K]) => boolean): void

    Registers 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);
  • checkAuthorizeJoin
    checkAuthorizeJoin(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.
  • checkAuthorizeSend
    checkAuthorizeSend(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.
  • reportError
    reportError(error: unknown, context?: string): void

    Routes 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"`).
  • isRateLimited
    isRateLimited(clientId: string): boolean

    Checks whether a client has exceeded the configured message rate limit within the current window.

  • upgrade
    upgrade(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.

  • connect
    connect(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,});
  • websocketgetter
    websocket:

    WebSocket handler object to pass directly to `Bun.serve({ websocket })`.

    tsts
    Bun.serve({  fetch: (req, server) => realtime.connect(req, server),  websocket: realtime.websocket,});
  • getChannel
    getChannel(name: string): SSEChannel | undefined

    Returns an existing named SSE channel without creating one.

    name
    Channel/topic name.

    Returns The if it exists, or `undefined`.

  • channel
    channel(name: string): SSEChannel

    Returns 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" });
  • namespace
    namespace(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" });
  • to
    to(topic: string): RealtimeBroadcaster

    Returns 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" });
  • job
    job<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" });
  • stats
    stats(): RealtimeStats

    Captures a real-time operational metrics snapshot of active connections,

    Returns Current telemetry snapshot.

RoomBroadcaster

class

Sends a realtime event to all WebSocket clients subscribed to a pub/sub topic.

1 member

  • send
    send(event: string, data: unknown): void

    Broadcast a named event with payload to all topic subscribers.

    event
    Event name.
    data
    Event payload (JSON-serialized).

SSEChannel

class

Named SSE broadcast channel — a group of connections subscribed to the same topic.

tsts
const ch = realtime.channel("notifications");ch.broadcast("alert", { level: "info", message: "Deploy complete" });console.log(`${ch.size} clients listening`);

9 members

  • subscribersproperty
    subscribers:
  • messageHistoryproperty
    messageHistory: Array<{ id: string; event: string; data: unknown; }>
  • sizegetter
    size: number

    Returns the number of active subscribers currently connected to this channel.

  • subscribe
    subscribe(client: SSEClient<any>): void

    Adds a client to this channel and wires up an abort listener to auto-remove on disconnect.

    client
    SSE client to subscribe.
  • remove
    remove(client: SSEClient<any>): void

    Removes a client from this channel's subscriber set.

    client
    SSE client to remove.
  • replayHistory
    replayHistory(client: SSEClient<any>, lastEventId: string): void

    Replays 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.
  • broadcast
    broadcast<K extends keyof TEvents>(event: K, data: TEvents[K], id?: string): void

    Broadcasts 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" });
  • broadcast
    broadcast(event: string, data: unknown, id?: string): void
  • broadcast
    broadcast(event: string, data: unknown, id?: string): void

SSEClient

class

Represents a single SSE connection to a browser or HTTP client.

class SSEClient
tsts
sse(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

  • isClosedproperty
    isClosed:
  • heartbeatTimerproperty
    heartbeatTimer?: ReturnType<typeof setInterval>
  • bufferedEventCountproperty
    bufferedEventCount:
  • signalproperty
    signal: AbortSignal

    `AbortSignal` that fires when the client disconnects. Use to clean up resources.

  • idproperty
    id: string

    Unique UUID assigned to this SSE connection.

  • dataproperty
    data: TData

    Mutable per-connection session data. Useful for storing auth state or metadata.

  • write
    write(payload: string): boolean
  • send
    send<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K], id?: string): this

    Send 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 name
  • send
    send(event: string, data: unknown, id?: string): this
  • send
    send(data: unknown): this
  • send
    send(arg1: unknown, arg2?: unknown, id?: string): this
  • comment
    comment(text: string): this

    Send an SSE comment line (begins with `:`). Useful for keep-alive pings or debug annotations.

    text
    Comment text.

    Returns `this` for chaining.

  • streamText
    streamText(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 }),});
  • close
    close(): void

    Closes the SSE connection, stopping the heartbeat and terminating the readable stream.

TypedSocket

class

Concrete implementation of wrapping a raw Bun `ServerWebSocket`.

class TypedSocket implements SocketClient<TData>

13 members

  • idproperty
    id: string

    Unique UUID for this socket connection.

  • datagetter
    data: TData
  • datasetter
    data:
  • send
    send(event: string, data: unknown): void

    Send a named event JSON envelope to this client.

  • sendJSON
    sendJSON(event: string, data: unknown): void

    Serialize and send a named event as a JSON `RealtimeEnvelope` to this specific client.

    event
    Event name.
    data
    Event payload.
  • sendText
    sendText(text: string): void

    Send a raw UTF-8 text WebSocket frame to this client.

    text
    Raw text to send.
  • sendBinary
    sendBinary(data: ArrayBufferView | ArrayBuffer): void

    Send a raw binary WebSocket frame to this client.

    data
    Binary buffer to send.
  • join
    join(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.

  • leave
    leave(topic: string): void

    Unsubscribe from a pub/sub topic.

    topic
    Topic name to leave.
  • isSubscribed
    isSubscribed(topic: string): boolean

    Check whether this socket is subscribed to a given topic.

    topic
    Topic name to check.
  • to
    to(topic: string): RoomBroadcaster

    Returns a broadcaster that sends to all sockets on the topic (including this socket).

    topic
    Topic name.
  • broadcast
    broadcast(topic: string): RoomBroadcaster

    Returns a broadcaster that sends to all sockets on the topic, **excluding** this socket.

    topic
    Topic name.
  • close
    close(code?: number, reason?: string): void

    Close 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

interface

Configuration options for the isomorphic SDK.

interface ClientOptions

4 members

  • urlproperty
    url: string

    Target server endpoint URL (e.g. `"http://localhost:3000/realtime"` or `"ws://localhost:3000/realtime"`).

  • transportproperty
    transport?: "auto" | "websocket" | "sse"

    Transport mode: `"websocket"`, `"sse"`, or `"auto"` (prefers WebSocket if available). Defaults to `"auto"`.

  • protocolsproperty
    protocols?: string | string[]

    Optional sub-protocols for WebSocket handshakes.

  • reconnectproperty
    reconnect?: { enabled?: boolean; minDelay?: number; maxDelay?: number; factor?: number; jitter?: boolean; }

    Exponential backoff and jitter configuration for automatic reconnection.

JobState

interface

Serializable snapshot of a job's current lifecycle state, broadcast to clients via SSE.

interface JobState<TResult = unknown>

8 members

  • jobIdproperty
    jobId: string

    Unique job identifier.

  • statusproperty
    status: JobStatus

    Current lifecycle status.

  • percentproperty
    percent: number

    Progress percentage (0–100).

  • messageproperty
    message: string

    Human-readable status message for display in UIs.

  • resultproperty
    result?: TResult

    Result value emitted on successful completion.

  • errorproperty
    error?: string

    Error message string emitted on failure.

  • updatedAtproperty
    updatedAt: number

    Unix millisecond timestamp of the last state update.

  • extraproperty
    extra?: Record<string, unknown>

    Optional arbitrary extra metadata for custom progress data.

Logger

interface

Structured logger interface for realtime engine diagnostics.

interface Logger

4 members

  • debug
    debug(...args: unknown[]): void

    Emit debug-level trace messages.

  • info
    info(...args: unknown[]): void

    Emit informational operational messages.

  • warn
    warn(...args: unknown[]): void

    Emit warning-level messages.

  • error
    error(...args: unknown[]): void

    Emit error-level messages.

PubSubAdapter

interface

Pluggable publish/subscribe adapter interface enabling horizontal scaling across multiple server nodes.

interface PubSubAdapter
tsts
class 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

  • publish
    publish(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.
  • subscribe
    subscribe(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.

  • close
    close(): Promise<void>?

    Optional cleanup hook for closing client connections.

RateLimitConfig

interface

Configuration for WebSocket inbound message rate limiting.

interface RateLimitConfig
tsts
createRealtime({  rateLimit: {    messages: 50,    windowMs: 10_000,    onRateLimit: (client) => client.send("error:rate_limit", { message: "Slow down!" }),  },});

3 members

  • messagesproperty
    messages: number

    Maximum number of messages allowed per client within `windowMs`.

  • windowMsproperty
    windowMs: number

    Sliding window duration in milliseconds over which `messages` are counted.

  • onRateLimitproperty
    onRateLimit?: (client: SocketClient<any>) => void

    Custom handler invoked when a client exceeds the rate limit.

RealtimeConfig

interface

Top-level configuration for / .

interface RealtimeConfig<TData = Record<string, unknown>>
tsts
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

  • handlersproperty
    handlers?: RealtimeHandlers<TData>

    WebSocket lifecycle and message event handlers. See .

  • authenticateproperty
    authenticate?: (req: Request) => Promise<TData | null> | TData | null

    Authentication hook invoked on every incoming connection (WebSocket upgrade or SSE).

  • authorizeproperty
    authorize?: { join?: (client: SocketClient<TData>, topic: string) => boolean | Promise<boolean>; send?: (client: SocketClient<TData>, event: string, data: unknown) => boolean | Promise<boolean>; }

    Fine-grained authorization guards.

  • validateproperty
    validate?: (event: string, data: unknown) => boolean | { valid: boolean; error?: string; }

    Global inbound message payload validator.

  • rateLimitproperty
    rateLimit?: RateLimitConfig

    WebSocket inbound message rate limiting. See .

  • adapterproperty
    adapter?: PubSubAdapter

    Pluggable pub/sub adapter for horizontal multi-node scaling.

  • loggerproperty
    logger?: Logger

    Custom structured logger. Defaults to `console`-based output.

  • perMessageDeflateproperty
    perMessageDeflate?: boolean

    Enable WebSocket `permessage-deflate` compression.

  • maxPayloadLengthproperty
    maxPayloadLength?: number

    Maximum allowed WebSocket message payload size in bytes.

  • sseOptionsproperty
    sseOptions?: SSEOptions

    Default SSE transport options applied to all channels. See .

  • onErrorproperty
    onError?: (error: unknown, context?: string) => void

    Global 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

interface

Standard realtime message envelope wrapping all events sent over WebSocket or SSE channels.

interface RealtimeEnvelope<T = unknown>

5 members

  • idproperty
    id: string

    Unique event message identifier (UUID).

  • eventproperty
    event: string

    Event type name (e.g. `"chat.message"`, `"job:status"`).

  • topicproperty
    topic?: string

    Optional pub/sub topic name this event was routed on.

  • dataproperty
    data: T

    Strongly-typed event payload data.

  • timestampproperty
    timestamp: number

    Unix millisecond timestamp when the event was dispatched.

RealtimeHandlers

interface

Lifecycle and event callbacks for the WebSocket transport layer.

interface RealtimeHandlers<TData = Record<string, unknown>>
tsts
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

  • openproperty
    open?: (client: SocketClient<TData>) => void | Promise<void>

    Called when a new WebSocket client connects.

  • messageproperty
    message?: (client: SocketClient<TData>, event: string, data: any) => void | Promise<void>

    Called for each inbound WebSocket message, after rate limiting, validation, and authorization.

  • closeproperty
    close?: (client: SocketClient<TData>, code: number, reason: string) => void | Promise<void>

    Called when a client disconnects. `code` and `reason` follow the WebSocket close handshake.

  • drainproperty
    drain?: (client: SocketClient<TData>) => void | Promise<void>

    Called when the send buffer drains after backpressure.

  • errorproperty
    error?: (client: SocketClient<TData>, error: unknown) => void | Promise<void>

    Called on low-level WebSocket transport errors.

  • pingproperty
    ping?: (client: SocketClient<TData>, data: Buffer) => void | Promise<void>

    Called when a WebSocket ping frame is received.

  • pongproperty
    pong?: (client: SocketClient<TData>, data: Buffer) => void | Promise<void>

    Called when a WebSocket pong frame is received.

  • subscribeproperty
    subscribe?: (client: SocketClient<TData>, topic: string) => void | Promise<void>

    Called after a client successfully joins a pub/sub topic.

  • unsubscribeproperty
    unsubscribe?: (client: SocketClient<TData>, topic: string) => void | Promise<void>

    Called after a client leaves a pub/sub topic.

RealtimeRegister

interface

Global interface augmented by application code for typed event names and payloads.

tsts
declare module "./realtime" {  interface RealtimeRegister {    events: {      "chat.message": { user: string; text: string };      "order.status": { orderId: string; status: string };    };  }}

RealtimeStats

interface

Snapshot of realtime engine operational statistics at a given moment.

interface RealtimeStats

8 members

  • websocketConnectionsproperty
    websocketConnections: number

    Number of active WebSocket client connections.

  • sseConnectionsproperty
    sseConnections: number

    Number of active SSE (Server-Sent Events) client connections.

  • topicsproperty
    topics: number

    Number of active named SSE channels.

  • messagesSentproperty
    messagesSent: number

    Cumulative count of messages dispatched outward.

  • messagesReceivedproperty
    messagesReceived: number

    Cumulative count of messages received from clients.

  • bytesSentproperty
    bytesSent: number

    Cumulative bytes sent across all transports.

  • bytesReceivedproperty
    bytesReceived: number

    Cumulative bytes received from all clients.

  • activeJobsproperty
    activeJobs: number

    Count of currently active in-progress job trackers.

SocketClient

interface

Public interface representing a connected WebSocket client.

interface SocketClient<TData = Record<string, unknown>>

14 members

  • idproperty
    id: string

    Unique UUID assigned to this connection.

  • dataproperty
    data: TData

    Mutable per-connection session data (auth context, user info, etc.).

  • rawproperty
    raw: ServerWebSocket<TData>

    The raw underlying `ServerWebSocket` instance (Bun-native).

  • send
    send<K extends keyof RegisteredEvents>(event: K, data: RegisteredEvents[K]): void

    Send a typed registered event.

  • send
    send(event: string, data: unknown): void

    Send an arbitrary named event with unknown payload.

  • sendJSON
    sendJSON(event: string, data: unknown): void

    Serialize an event as a JSON string and send it.

  • sendText
    sendText(text: string): void

    Send a raw UTF-8 text frame.

  • sendBinary
    sendBinary(data: ArrayBufferView | ArrayBuffer): void

    Send a raw binary frame.

  • join
    join(topic: string): Promise<boolean>

    Subscribe to a pub/sub topic, pending authorization check. Returns `true` on success.

  • leave
    leave(topic: string): void

    Unsubscribe from a pub/sub topic.

  • isSubscribed
    isSubscribed(topic: string): boolean

    Returns whether this socket is subscribed to a given topic.

  • to
    to(topic: string): RoomBroadcaster

    Returns a broadcaster that sends to all sockets on `topic` (including this one).

  • broadcast
    broadcast(topic: string): RoomBroadcaster

    Returns a broadcaster that sends to all sockets on `topic` (excluding this one).

  • close
    close(code?: number, reason?: string): void

    Close the WebSocket with an optional code and reason.

SSEOptions

interface

Configuration options for the Server-Sent Events (SSE) transport layer.

interface SSEOptions

6 members

  • headersproperty
    headers?: Record<string, string>

    Additional HTTP response headers to merge into the SSE response.

  • heartbeatIntervalproperty
    heartbeatInterval?: number

    Interval in milliseconds between keep-alive comment pings (`: keep-alive ping`).

  • retryproperty
    retry?: number

    Client-side reconnect delay hint in milliseconds, sent as the SSE `retry:` field.

  • maxBufferedEventsproperty
    maxBufferedEvents?: number

    Maximum number of events buffered for a slow consumer before backpressure is triggered.

  • historySizeproperty
    historySize?: number

    Number of recent events to retain for `Last-Event-ID` reconnection replay.

  • onBackpressureproperty
    onBackpressure?: (client: SSEClient) => void

    Custom backpressure handler invoked when a slow client's buffer is full.

AIStreamChunk

type

AI/LLM streaming chunk frame sent during active token streaming.

type AIStreamChunk = { type: "chunk"; streamId: string; text: string; index: number; }

AIStreamDone

type

AI/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

type

Discriminated union of AI/LLM stream frame types transmitted over SSE.

AIStreamError

type

AI/LLM streaming error frame emitted when the stream encounters an irrecoverable error.

type AIStreamError = { type: "error"; streamId: string; error: { code: string; message: string; }; }

EventPayload

type

Resolves the payload type for a given event name `K`.

type EventPayload<K extends string> = IsRegistryConfigured extends true ? K extends keyof RegisteredEvents ? RegisteredEvents[K] : unknown : unknown

IsRegistryConfigured

type

Type-level boolean indicating if a custom event registry has been configured via module augmentation.

type IsRegistryConfigured = RealtimeRegister extends { events: any; } ? true : false

JobStatus

type

All possible lifecycle states for a tracked background job.

type JobStatus = "pending" | "running" | "paused" | "completed" | "failed" | "cancelled"

RegisteredEvents

type

Extracts 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>
Tip
Most of the types above are inferred. You rarely import AuthConfig or JobPayload — declaring your schema once is enough for the rest to follow. See Typed keys.