Skip to content
DocsconceptsDocumentation

Real-time

How subscriptions, deltas, and hibernation fit together.

Last updated:

Every Lunora query doubles as a subscription. Calling useQuery(api.x.y, args) in React opens a multiplexed WebSocket to your Worker, registers the query with the resolving Durable Object, and re-renders the component whenever a mutation writes to a table the query reads.

import { api } from "@/lunora/_generated/api";
import { useMutation, useQuery } from "@lunora/react";

export const Chat = ({ channelId }: { channelId: string }) => {
    const messages = useQuery(api.messages.list, { channelId });
    const { mutate: send } = useMutation(api.messages.send);

    if (!messages) return <p>Loading…</p>;
    return (
        <ul>
            {messages.map((m) => (
                <li key={m._id}>{m.text}</li>
            ))}
        </ul>
    );
};

Hibernation

Subscriptions are stored on the WebSocket via state.serializeAttachment(...). When the Durable Object hibernates (no traffic for ~10s) the WebSocket is suspended without losing the subscription registry; the next message wakes it up and resumes pushing updates exactly where they were going.

This is why your DO bill drops to near-zero between bursts: idle subscribers pay for storage, not compute.

How an update reaches a subscriber

A subscription is registered by function path, not by table, and updates arrive by re-running the query.

Each socket carries a SocketAttachment mapping subId → SubscriptionQuery, which survives hibernation. When a mutation writes through ctx.db, the write hook records the tables it touched (recordChangedTable). Once it commits, the shard re-executes every subscription that read one of those tables — under that subscriber's own identity, never the writer's — and pushes the new result as a data frame when it differs from the last one that socket saw.

Two consequences worth knowing:

  • Row-level filtering is the query's job, not the router's. A write to messages refreshes every subscription that reads messages; each one then returns only what its own handler (and its RLS policies) selects for that subscriber. Nothing tries to match a mutation's row against a subscription's args.
  • Re-execution is not free. A query that reads a hot table re-runs on every write to it. The shard's opt-in reactive cache memoizes results by (functionPath, args, identity) so identical subscribers don't each pay for the handler, and a reconnecting client's sinceSeq skips the re-snapshot entirely when nothing it reads has changed.

There is also a lower-level delta frame — ShardDO.broadcastDelta(delta) fanned out to sockets whose subscription registered a matching table string. It predates re-execution and is there for hand-written ShardDO subclasses that don't override executeSubscription. Generated shards don't use it, and a @lunora/client subscription registers its function path as the table, so nothing addressed by a real table name reaches one.

Reconnect

The client (@lunora/client) uses exponential backoff with jitter and resumes by bookmark: it sends the last delta sequence it acknowledged so the server can replay anything that was missed during the disconnect.

When WebSockets are blocked

One-shot calls (query, mutation, action) ride HTTP POST, not the socket. So a corporate proxy or captive portal that refuses the Upgrade handshake does not break them — it breaks reactivity, and only that: every live query stops moving while the app otherwise works.

After three consecutive connect attempts that never reach open, the client stops waiting on the socket for freshness and re-runs each subscribed query over the batch-RPC endpoint on an interval, applying the results through the same path a server data frame takes. connectionStatus() reports "polling", and the moment a socket does open, polling stops and the normal resubscribe takes over.

const client = new LunoraClient({
    url,
    // Defaults shown. `intervalMs: 0` disables the fallback entirely.
    pollingFallback: { afterFailedAttempts: 3, intervalMs: 5_000 },
});

It is a degradation, not a second transport, and the gap matters:

  • Only plain live queries are polled. Shapes (@lunora/db collections), durable streams and whispers are socket-only surfaces with no request/response form to re-run — they stay dark until a socket opens.
  • Freshness is bounded by the interval, not by the write. There is no CDC cursor on this path: every tick is a full snapshot, so a poll costs what the query costs, repeatedly.
  • It does not make an app offline-capable. A poll is an HTTP request; when nothing can reach the origin, the offline queue is what carries you.

Treat "polling" in a status indicator as live, but slower and partial — which is why it is its own status rather than folded into "connected".

The trigger is deliberately "never opened", not "disconnected". A socket that opens and drops is a flaky link, and the reconnect backoff is the right answer to that; falling back to polling there would trade a self-healing channel for a permanently worse one.

Durable streams

A .stream() procedure is ephemeral by default: the producer belongs to the socket that opened it, so closing the tab ends the run and a reconnect starts over. That is right for a progress ticker and wrong for anything expensive: a model's answer must not disappear because someone refreshed.

Declare the stream durable and the run stops belonging to the socket:

export const answer = query.input({ threadId: v.id("threads"), prompt: v.string() }).stream(
    async function* ({ ctx, args }) {
        for await (const token of callTheModel(args.prompt)) {
            yield token;
        }
    },
    { durable: true },
);

Three things change:

  • Every chunk is persisted before it is sent, under a monotonic seq. A reconnecting client replays the frames it missed out of the shard's SQLite and then continues live; the consumer's for await never sees the interruption.
  • The producer outlives the socket. Close the tab mid-run and the run finishes anyway; the transcript is waiting when you come back.
  • A live run is shared. Its identity is (identity, functionPath, args), so a second tab of the same signed-in user attaches to the run already in flight instead of paying for a second generation. A different identity always gets its own run. A finished run is not shared: a later caller asking the same question gets a fresh answer, because a transcript is the record of one execution, not a cached response.

On the client, opt in per call so the reconnect resumes rather than fails:

const { chunks, status } = useStream(api.chat.answer, { threadId, prompt }, { durable: true });

Transcripts are eligible for trimming after 24 hours by default, per procedure ({ durable: { ttlMs } } to change it) — the sweep runs on an attach and at most once an hour, so expiry is an earliest-time, not a deadline.

A run is capped on two axes so a runaway generator cannot fill the shard's storage: 50,000 chunks and 64 MiB of encoded chunk bytes. Whichever binds first ends the run with STREAM_TOO_LONG (HTTP 507), so a stream yielding ~1 MiB frames (image or audio chunks) stops around chunk 65, nowhere near the chunk count. Yield smaller frames, or drop durable.

If the Durable Object is evicted mid-run, what happens depends on who asks next. A client resuming that transcript gets STREAM_INTERRUPTED: its tail cannot be spliced back on, and re-generating would duplicate what it already has. A client asking fresh simply gets a new run: the dead one has no claim on its key. When you need a producer that survives eviction itself, the run belongs in a workflow.

Whispers (ephemeral peer messages)

A whisper relays a payload between sockets on a shard without touching SQLite, the CDC log, or any query. Nothing is stored, nothing re-runs, and the sender never receives its own message — which is what makes it cheap enough for live cursors, typing indicators and presence pings.

const stop = client.whisperSubscribe("cursors", (data, from) => {
    moveCursor(from, data as { x: number; y: number });
});

client.whisper("cursors", { x, y });

from is the sender's server-stamped user id (absent when anonymous) and is unforgeable. A send is fire-and-forget: it is dropped when the socket is down (whispers are never queued), when the payload exceeds the per-frame size cap, or when the sender is over its rate budget.

Authorizing topics

By default a topic's only boundary is the shard: any client that can open a socket to it can join, read and inject on any topic name. Declare an onWhisper authorizer to change that — the shard runs it before a socket joins a topic and before it broadcasts to one:

// lunora/whisper.ts
import { onWhisper } from "@lunora/server";

export const authorize = onWhisper(async (ctx, event) => {
    // event: { action: "subscribe" | "send", topic, connectionId, shardKey, userId, context? }
    const [kind, roomId] = event.topic.split(":");

    if (kind !== "room") {
        return true; // presence:* and cursors stay open
    }

    return (await ctx.db.roomMembers.findFirst({ where: { roomId, userId: ctx.auth.userId } })) !== null;
});

It is a query, so ctx.db is RLS-scoped to the connecting user and the check cannot write. Only a literal true allows: any other return value — or a throw — denies, and multiple exports are AND-ed. Declaring one governs every topic on the shard, so route on the topic name for the namespaces that should stay open.

A denied frame is dropped silently. Neither whisper_subscribe nor whisper has an ack, so there is nothing to report a denial on, and adding an error frame would let a client probe which topics exist. Gate your UI on a query you can read a verdict from, not on whether whispers arrive. (The per-socket cap below is the one refusal that does answer — it is a resource limit, decided before the authorizer runs, so it discloses nothing about the topic.)

Three ceilings worth knowing. The verdict is memoised per socket, action and topic, so revoking access does not evict a member already joined — they keep receiving until the socket closes or the Durable Object hibernates. That memo is capped at 256 distinct topic/action pairs per socket, because the authorizer is a query and a client naming a fresh topic on every frame would otherwise run one per name: past the cap the socket is refused with TOO_MANY_WHISPER_TOPICS and the authorizer does not run. (That refusal is the one whisper frame that does answer — it names only the cap, never the topic or the verdict — and the memo is in-memory, so a reconnect starts fresh.) And a whisper leaves no durable trace, so it can neither be replayed nor audited. Anything that must stop the instant access is revoked, or that you need a record of, belongs behind a query or mutation with RLS — not on a topic.

Why no MQTT / Pub/Sub

Real-time fan-out in Lunora is entirely Durable-Object-based: hibernated WebSocket subscriptions on ShardDO, type-safe and integrated end-to-end with your queries, with no external broker to run. Cloudflare Pub/Sub (an MQTT broker) is therefore a non-goal. It is in beta with gated onboarding and no Worker binding, and the only capability it adds over the DO path is native-MQTT device ingest (IoT clients speaking MQTT directly), a narrow niche. We'll revisit only if Pub/Sub reaches GA and a concrete need to ingest from native-MQTT devices appears; until then there is nothing to configure.