Engines
Realtime
Clients connect over WebSocket or plain HTTP streaming — the same handlers serve both, chosen automatically from the request.
Setup
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
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
// 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.
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.
export const realtime = createRealtime<{ userId: string }>({ /* … */});Authorization
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
// 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.
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
const tracker = realtime.job(jobId); tracker.progress(25, "Uploading");tracker.progress(100, "Done");tracker.done({ url });tracker.fail(new Error("Upload failed"));Client SDK
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.
export const realtime = createRealtime({ sseOptions: { maxBufferedEvents: 1024, heartbeatInterval: 15_000, onBackpressure: (client) => { console.warn("client too slow, closing"); client.close(); }, },});