Engines

Realtime

Clients connect over WebSocket or plain HTTP streaming — the same handlers serve both, chosen automatically from the request.

Setup

yatta/func/realtime.tsts
import { createRealtime } from "yatta/realtime"; export const realtime = createRealtime({  handlers: {    open(client) {      console.log(`connected: ${client.id}`);    },     message(client, event, data) {      client.send(event, data);   // echo    },     close(client) {      console.log(`disconnected: ${client.id}`);    },  },});

Wiring it up

yatta/main.tsts
import { realtime } from "./func/realtime";import routers from "./func/routerHelper"; const server = Bun.serve({  port: 4000,  fetch: (req, srv) => routers(req, srv),  websocket: realtime.websocket,   // required for upgrades});

routerHelper routes /realtime to the engine, which picks WebSocket or SSE based on the request.

Broadcasting

Publishts
// Everyone subscribed to the topicrealtime.to("orders").publish("order.updated", { id: 1, total: 42 }); // A named roomconst room = realtime.room("support");room.publish("agent.joined", { name: "Ada" }); // An SSE channel specificallyrealtime.channel("notifications").broadcast("ping", { at: Date.now() });

Authentication

The hook runs on every connection. Return the session data to allow it, or null to reject with a 401.

authenticatets
export const realtime = createRealtime({  async authenticate(req) {    const token = req.headers.get("Authorization")?.replace("Bearer ", "");    if (!token) return null;     const session = await auth.getSession(token);    if (!session) return null;     return { userId: session.user.id, roles: session.user.roles };  },   handlers: {    open: (c) => console.log(c.data.userId),    message: (c, e, d) => c.send(e, d),  },});

Client data is typed, so client.data.userId is a string rather than unknown.

Typed sessionts
export const realtime = createRealtime<{ userId: string }>({  /* … */});

Authorization

Guardsts
export const realtime = createRealtime({  authorize: {    // Before a client joins a topic    join: (client, topic) => client.data.roles.includes("admin"),     // Before an inbound message reaches the handler    send: (client, event, data) => event !== "admin.delete",  },   validate: (event, data) => {    if (event === "chat.message" && typeof data !== "object") {      return { valid: false, error: "payload must be an object" };    }    return true;  },});

Topics and rooms

Client sidets
// Joining is permission-checked by authorize.joinclient.join("orders"); // Rooms scope a broadcasterconst room = realtime.room("support");room.broadcast("agent.joined", { name: "Ada" }); // Namespace every topic under a prefix — handy for multi-tenant appsconst tenant = realtime.namespace("tenant:acme");tenant.to("orders").publish("updated", { id: 1 });

Streaming LLM output

streamText emits a chunk event per token and always finishes with a terminal done or error frame.

yatta/backend/ai.tsts
import { API, createAPI } from "yatta/api";import { SSEClient } from "yatta/realtime"; const api = createAPI(); api.get(async (ctx) => {  const client = new SSEClient(ctx.req);   await client.streamText(llm.stream("Tell me a joke"), {    eventName: "ai:stream",    streamId: "session-123",    usage: () => ({ inputTokens: 10, outputTokens: 25 }),  });   return new Response(client.readable, {    headers: { "Content-Type": "text/event-stream" },  });}); export default api;
Note
A disconnect mid-stream aborts the iterator and emits a done frame with reason: "aborted", so clients never hang waiting for a stream that stopped.

Job progress

Tracking a jobts
const tracker = realtime.job(jobId); tracker.progress(25, "Uploading");tracker.progress(100, "Done");tracker.done({ url });tracker.fail(new Error("Upload failed"));

Client SDK

Browserts
import { Realtime } from "yatta/realtime"; const rt = new Realtime("/realtime"); rt.on("order.updated", (data) => {  console.log("order", data.id);}); // Falls back to SSE, reconnects with exponential backoffrt.connect();

Backpressure

A slow consumer is capped rather than allowed to grow the heap without bound. On overflow you are notified and can close the connection.

Limitsts
export const realtime = createRealtime({  sseOptions: {    maxBufferedEvents: 1024,    heartbeatInterval: 15_000,    onBackpressure: (client) => {      console.warn("client too slow, closing");      client.close();    },  },});