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/sseis 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 aBroadcastChannelfrom a singlecommittedlifecycle listener and mounts the generatedGET /api/sse/Calculator?stream=<id>endpoint viahono(app, { sse }).packages/client/src/useSse.tsโ a smallEventSourcehook that appliesevent: patchframes withapplyPatchMessageand reconnects onbehindto re-seed from the server's cachedevent: stateframe.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(),
});
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 exactlycachedV + 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 fromcachedV = -1, so the genesis patch (version0, 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(fromoverlay()) and its version equalscachedV. Merge it on top of the cached state, leaving_vunchanged. 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) <= cachedVfor 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โ
_vissnap.event.versionโ the event store's stream version is the single source of truth. Never invent a version.- 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. - Broadcast from snapshots, not projections โ projections are eventually consistent and may lag. Broadcast from the snapshots returned by
app.do(). - 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.