Skip to main content
Version: Current

Real-Time with act-http/sse

@rotorsoft/act-http/sse broadcasts incremental state updates over Server-Sent Events. Each event-sourced commit emits a domain patch (a partial state update); the broadcast layer forwards those patches keyed by event version, so subscribers can apply them to their cached state without ever refetching.

SSE is one of two HTTP-shaped integration paths. For the outbound side โ€” webhooks, downstream services, message buses โ€” see External integration patterns. The two are independent; an app can run both.

npm install @rotorsoft/act-http

Migrating from @rotorsoft/act-sse? That standalone package is deprecated on npm and no longer published from this repo โ€” @rotorsoft/act-http/sse is the only home. The surface is identical; swap the import specifier and you're done. Versions already installed keep working. See the 1.x migration guide.

Clearing a fieldโ€‹

A reducer clears a field by patching it away, and @rotorsoft/act-patch accepts either spelling:

.patch({ Cleared: () => ({ result: 0, left: undefined, operator: undefined }) })

Both undefined and null mean delete this key. On the wire only null survives โ€” JSON.stringify drops undefined-valued keys, and every SSE transport serializes frames as JSON. The broadcast layer normalizes undefined to null when it builds a frame (#1471), so both spellings behave identically for live subscribers and you can keep writing whichever reads better.

The normalization applies to the frame only. Server-side cached state is produced by apply_patch, which handles undefined natively, so a reconnecting client's reseed has the key absent rather than set to null.

If you build a transport of your own over these patches, do the same normalization โ€” a patch that crosses a JSON boundary must encode deletes as null.

See it runningโ€‹

The multi-transport calculator demo wires SSE end-to-end next to tRPC, Hono REST, and OpenAPI:

  • packages/server/src/server.ts โ€” publishes every commit to a BroadcastChannel from a single committed lifecycle listener and mounts the generated GET /api/sse/Calculator?stream=<id> endpoint via hono(app, { sse }).
  • packages/client/src/useSse.ts โ€” a small EventSource hook that applies event: patch frames with applyPatchMessage and reconnects on behind to re-seed from the server's cached event: state frame.
  • packages/client/src/Calculator.tsx โ€” renders the SSE-fed state live under the keypad.

Run pnpm dev:http from the repo root and press keys at http://localhost:3000 โ€” the live panel updates on every commit without refetching, from either transport or another browser tab.

Architectureโ€‹

app.do() โ†’ snapshots (each carries its domain patch)
โ”‚
โ–ผ
deriveState(snap) โ† app-specific (overlay presence, etc.)
state._v = snap.event.version โ† event store version = single source of truth
โ”‚
โ–ผ
broadcast.publish(streamId, state, patches)
โ”‚
โ”œโ”€โ”€ version-key each patch: { [baseV+1]: patch1, [baseV+2]: patch2, ... }
โ””โ”€โ”€ push to all SSE subscribers
โ”‚
โ–ผ
Client: applyPatchMessage(msg, cached)
โ”‚
โ”œโ”€โ”€ contiguous โ†’ deep-merge patches in version order
โ”œโ”€โ”€ stale โ†’ skip (client already ahead)
โ””โ”€โ”€ behind โ†’ resync (client missed versions)

The version contractโ€‹

_v on every state object is always snap.event.version โ€” the event store's monotonic stream version. There is no separate counter, no clock-based ordering. The event store is the single source of truth for ordering; the broadcast layer is just a fan-out.

Server-sideโ€‹

BroadcastChannelโ€‹

Manages per-stream subscriber sets and an LRU state cache for reconnects:

import { BroadcastChannel } from "@rotorsoft/act-http/sse";
import type { BroadcastState, PatchMessage } from "@rotorsoft/act-http/sse";

type AppState = BroadcastState & {
// your domain state fields
name: string;
status: string;
};

const broadcast = new BroadcastChannel<AppState>({
cacheSize: 50, // LRU entries; default 50
// Called when a subscriber callback throws. The frame still reaches every
// other subscriber and the publish still returns โ€” a bad consumer must not
// break the publisher. Defaults to the framework's log() port.
onSubscriberError: (error, streamId) => metrics.sseSubscriberError.inc(),
// Called when overlay() finds no cached baseline. Live subscribers get a
// `_resync` frame so they refetch, but a steady stream of these means
// cacheSize is too small for the working set. Defaults to a no-op.
onOverlayMiss: (streamId) => metrics.sseOverlayMiss.inc(),
});
note

The snake_case member names these classes originally shipped with (publish_overlay, get_state, get_subscriber_count, cache_size, get_online, is_online) still work but are deprecated aliases โ€” scheduled for removal in the next major. Use the short names shown here (overlay, state, subscriberCount, cacheSize, online, isOnline).

Broadcasting a commitโ€‹

After every app.do(), forward each emitted snapshot's domain patch:

const snaps = await app.do("CreateItem", target, input);
const last = snaps.at(-1)!;

// 1. Derive the broadcast view (typically snap.state plus overlays)
const state: AppState = {
...last.state,
_v: last.event!.version, // MUST come from event.version
};

// 2. Collect each emitted snapshot's patch, in commit order
const patches = snaps
.map((s) => s.patch)
.filter(Boolean) as Partial<AppState>[];

// 3. Publish โ€” sends a version-keyed PatchMessage to all subscribers
broadcast.publish(streamId, state, patches);
// Reactions drain automatically if you've wired
// app.on("committed", () => app.settle()) at bootstrap.

publish() writes the new state to the LRU cache (so reconnects can read it) and pushes a PatchMessage<AppState> to subscribers. The keys of the message are absolute event versions (baseV + 1, baseV + 2, โ€ฆ), so subscribers can apply them directly to their cached state without computing offsets.

Overlays (non-event state changes)โ€‹

Some state changes don't have a corresponding event โ€” typically presence ("alice is online") or computed-field refreshes. Use overlay():

broadcast.overlay(streamId, {
onlineUsers: presence.online(streamId),
});

This applies the overlay to the cached state, leaves _v unchanged, and emits a single-key patch message at the cached version, tagged with an _overlay: true marker. The marker is what lets a live, caught-up client apply it: without it, a same-version message is indistinguishable from a stale patch the client already has, so applyPatchMessage would drop it. With it, applyPatchMessage merges the overlay on top of the client's current state (keeping _v), so presence reaches already-connected viewers, not just reconnecting ones.

Presenceโ€‹

Overlay data survives the next commit: publish() carries overlay-contributed keys onto the new cached state unless the domain state speaks to them (#1473), so a reconnecting client reseeds with the same presence a live client is holding. A domain state that sets or drops one of those keys still wins โ€” the store is authoritative for its own fields.

That carry lives in the LRU cache, so it lasts as long as the entry does. If the stream ages out of the cache the overlay data is genuinely gone, and the server says so rather than letting it vanish quietly: evicting an entry that carried overlay keys fires onOverlayMiss and sends the stream's live subscribers a resync, so they refetch instead of showing presence the server has forgotten (#1648). A steady trickle of those means cacheSize is too small for the working set.

online() returns a Set, which JSON.stringify encodes as {} โ€” so the broadcast layer converts a Set to an array before it reaches the cache or the wire (#1472, completed for publish() in #1646). Clients receive ["alice", "bob"], and a reconnecting client's reseed matches what a live one holds. The two paths run the same normalization; the one thing the cache deliberately does not adopt is the wire's delete encoding, since a reseed must show a cleared key as absent rather than null. A Map is left alone: unlike a Set it has no unambiguous JSON encoding, so pass one already shaped the way you want it sent.

PresenceTracker is a ref-counted online-status tracker designed for multi-tab clients (each tab opens its own SSE; add / remove maintain a per-identity counter):

import { PresenceTracker } from "@rotorsoft/act-http/sse";

const presence = new PresenceTracker();

// On SSE connect
presence.add(streamId, identityId);

// On SSE disconnect
presence.remove(streamId, identityId);

// Query
presence.online(streamId); // Set<string>
presence.isOnline(streamId, identityId); // boolean

If you use the generated transportsโ€‹

You may not need to write a subscription handler at all. Both trpc(app, { sse }) and hono(app, { sse }) from @rotorsoft/act-http accept an sse: { channel: broadcast } option that walks the registry and emits one subscription โ€” or one streaming GET /api/sse/<stateName>?stream=<id> endpoint โ€” per registered state, all reading from your BroadcastChannel. You keep owning publication (broadcast.publish(...) after commits); the generator owns subscription, accounting, cleanup, and the wire format. See Auto-generated API surfaces ยง Real-time subscriptions, and the runnable demo in packages/server + packages/client (pnpm dev:http from the repo root).

The section below is the custom-server path โ€” the same loop the generator writes for you, hand-rolled for hosts the generators don't cover.

tRPC subscription (custom server)โ€‹

act-http/sse doesn't dictate the wire format โ€” your tRPC handler decides. A typical pattern yields the cached state on connect, then forwards each patch message. Wrap the two shapes in a small app-level envelope so the client can tell them apart:

import type { PatchMessage } from "@rotorsoft/act-http/sse";

type Envelope<S> =
| { kind: "snap"; state: S }
| { kind: "patch"; msg: PatchMessage<S> };

export const onStateChange = publicProcedure
.input(z.object({ streamId: z.string(), identityId: z.string().optional() }))
.subscription(async function* ({ input, signal }) {
const { streamId, identityId } = input;
let resolve: (() => void) | null = null;
let pending: PatchMessage<AppState> | null = null;

const cleanup = broadcast.subscribe(streamId, (msg) => {
pending = msg;
resolve?.();
resolve = null;
});

if (identityId) presence.add(streamId, identityId);

try {
// Initial snapshot for first paint
const cached = broadcast.state(streamId);
if (cached) yield { kind: "snap", state: cached } satisfies Envelope<AppState>;

while (!signal?.aborted) {
if (!pending) {
await new Promise<void>((r) => {
resolve = r;
signal?.addEventListener("abort", () => r(), { once: true });
});
}
if (signal?.aborted) break;
if (pending) {
const msg = pending;
pending = null;
yield { kind: "patch", msg } satisfies Envelope<AppState>;
}
}
} finally {
cleanup();
if (identityId) presence.remove(streamId, identityId);
}
});

Client-sideโ€‹

applyPatchMessageโ€‹

:::note Safe to import from a browser The sse subpath's client half โ€” applyPatchMessage, patch, and the wire types โ€” depends only on @rotorsoft/act-patch, which has no dependencies of its own. The server half (BroadcastChannel) reaches the framework through a dynamic import, so a bundler never pulls Node APIs such as AsyncLocalStorage into a client bundle. CI enforces this with pnpm check:browser-safe, which walks the built subpath's static import graph and fails on any Node builtin. :::

import { applyPatchMessage } from "@rotorsoft/act-http/sse";

onData: (env) => {
if (env.kind === "snap") {
utils.getState.setData({ streamId }, env.state);
return;
}
const cached = utils.getState.getData({ streamId });
const result = applyPatchMessage(env.msg, cached);

if (result.ok) {
utils.getState.setData({ streamId }, result.state);
} else if (result.reason === "behind") {
utils.getState.invalidate({ streamId }); // missed versions โ€” refetch
}
// "stale" โ†’ no-op; the cache is already past these versions
};

applyPatchMessage(msg, cached) returns { ok: true, state } | { ok: false, reason: "stale" | "behind" }:

  • Contiguous โ€” min(msg.keys) is exactly cachedV + 1. Apply patches in version order via the deep-merge from @rotorsoft/act-patch; final _v = max(msg.keys). A fresh client with no baseline resumes from cachedV = -1, so the genesis patch (version 0, the first event of any stream) is contiguous and folds onto init state โ€” a first patch at version โ‰ฅ 1 is instead behind, since it can't be built from init alone.
  • Overlay โ€” the message carries _overlay: true (from overlay()) and its version equals cachedV. Merge it on top of the cached state, leaving _v unchanged. This is how presence and computed-field refreshes reach a caught-up client; an ordinary same-version message (no marker) stays stale.
  • Stale โ€” max(msg.keys) <= cachedV for a client that has a baseline (and not a current-version overlay). The client is already ahead (e.g., a mutation response landed before the SSE patch arrived). No-op. A fresh client is never stale.
  • Behind โ€” min(msg.keys) > cachedV + 1. The client missed versions and must resync via a full refetch.

Key rulesโ€‹

  1. _v is snap.event.version โ€” the event store's stream version is the single source of truth. Never invent a version.
  2. One broadcast function โ€” every code path that calls app.do() should funnel through the same publish helper. Multiple publish sites with different state shapes is how double-apply bugs start.
  3. Broadcast from snapshots, not projections โ€” projections are eventually consistent and may lag. Broadcast from the snapshots returned by app.do().
  4. Presence is an overlay, not an event โ€” use overlay() so connect/disconnect doesn't pollute the event log.

The double-apply bugโ€‹

If a projection falls back to the broadcast cache on a miss, it reads state that already has event patches applied. Re-applying those same patches corrupts counters and indices.

// BUG โ€” broadcast cache holds post-event snapshots
let state = projCache.get(id) ?? broadcast.state(id); // โ† already patched!
mutator(state); // patches applied a second time

// FIX โ€” fall back to durable storage only
let state = projCache.get(id) ?? (await db.select(id)) ?? defaultState();
mutator(state);

The broadcast cache exists for reconnect seeding and for overlay()'s read-modify-write. Everything else should go through the database.