@dxv-systems/realtime
v0.3.0
Published
Realtime change signals for DXV apps over AWS AppSync Events — subscribe token minting, SigV4 publish, and a plain-React subscription hook with polling fallback
Maintainers
Readme
@dxv-systems/realtime
Realtime change signals for DXV apps over AWS AppSync Events: subscribe-token minting, SigV4 publish, and a plain-React subscription hook with a polling fallback.
Built against the Contract in docs/plans/2026-09-15-002-feat-realtime-appsync-events-plan.md. A parallel PR builds the AWS Lambda authorizer to the same Contract — the channel-pattern rule in matchChannel and the authorizer's own logic must never drift apart.
npm i @dxv-systems/realtimePrinciple: signals, not data
Publish a small "what went stale" payload and let the client re-read through its normal, permission-checked path (router.refresh(), invalidateQueries). The realtime layer never needs the data model, and a socket only shortens the delay — polling still bounds it.
Env vars
An app opts in by setting all of these (Terraform/infra side, per app). Every server function reads them lazily, at call time, so apps build fine without them and local dev / unconfigured previews work unchanged (publish no-ops, mintSubscribeToken returns null).
| Var | Where | Meaning |
|---|---|---|
| REALTIME_HTTP_ENDPOINT | server + client config | AppSync Events HTTP endpoint, host only (e.g. abc.appsync-api.eu-west-1.amazonaws.com) |
| REALTIME_WS_ENDPOINT | client config | AppSync Events realtime endpoint, host only |
| REALTIME_REGION | server (publish) | SigV4 region |
| REALTIME_NAMESPACE | server | This app's channel namespace, e.g. veijer, portal |
| REALTIME_SIGNING_SECRET | server | HS256 secret for subscribe tokens (per-app, from SSM) |
| REALTIME_PUBLISH_ACCESS_KEY_ID / REALTIME_PUBLISH_SECRET_ACCESS_KEY | server (publish) | IAM keys scoped to appsync:EventPublish on this app's namespace only |
Server (./server, Node only)
import { mintSubscribeToken, realtimeClientConfig, publish } from "@dxv-systems/realtime/server";mintSubscribeToken({ sub, channel, ttlSeconds? })— HS256 JWT (node:crypto, nojsonwebtokendependency). Claims:sub,ns(fromREALTIME_NAMESPACE),ch(thechannelyou pass — a pattern, e.g./veijer/*),iat,exp(≤iat + 900, default 900). Throws ifchannel's first segment isn'tns— that's a caller bug, not a runtime condition. Returnsnullwhen env is unset.realtimeClientConfig({ sub, channel })— the full{ token, expiresAt, httpEndpoint, wsEndpoint }an app's token route returns to the browser, ornull.expiresAtis epoch milliseconds, not the JWT's second-precisionexp, so the client can compare it toDate.now()directly.publish(channel, payload)— SigV4POST /event, 2s timeout, never throws. No-op{ ok: true, skipped: true }when publish env is unset. Body:{ channel, events: [JSON.stringify(payload)] }.publishMany(channels, payload)— fan-out convenience; returns the first failure, or{ ok: true, skipped: true }only if every channel skipped.
Token route (Next.js App Router)
// app/api/realtime/token/route.ts
import { NextResponse } from "next/server";
import { realtimeClientConfig } from "@dxv-systems/realtime/server";
import { getSession } from "@/lib/session";
export async function GET() {
const session = await getSession();
if (!session) return new NextResponse(null, { status: 401 });
const config = realtimeClientConfig({ sub: session.userId, channel: "/veijer/*" });
return NextResponse.json(config); // null when realtime isn't configured — the client just doesn't connect
}Publish after commit (any Node server)
import { publish } from "@dxv-systems/realtime/server";
await db.transaction(async (tx) => { /* write */ });
const result = await publish("/veijer/changes", { type: "changed", paths: ["/dashboard/planning"] });
if (!result.ok) console.error("realtime publish failed", result.error); // never fails the requestIsomorphic (.)
import { matchChannel } from "@dxv-systems/realtime";matchChannel(pattern, channel) — the same wildcard rule the AWS Lambda authorizer enforces (see the Contract). No Node or DOM dependency; safe in any bundle target.
React (./react, "use client", peer dep react >=18)
import { useRealtimeChannel, useLiveRefresh, defaultIsBusy } from "@dxv-systems/realtime/react";useRealtimeChannel(channel, onEvent, { getConfig, onResync? })
Implements the AppSync Events WebSocket protocol: connects to wss://<wsEndpoint>/event/realtime with the aws-appsync-event-ws + base64url header-... subprotocols, sends connection_init, subscribes on connection_ack, re-arms a keepalive watchdog on every ka (closing and reconnecting if one doesn't arrive within connectionTimeoutMs), and calls onEvent with the parsed payload on data. connected reflects a live subscription (subscribe_success), not just a live socket — an ack without a successful subscribe is not "connected". Reconnects with capped exponential backoff + jitter (1s → 30s). onEvent, getConfig and onResync are held in refs, so none need to be memoized by the caller.
Any gap between a close and the next subscribe_success can drop events silently — the server has no way to redeliver what it published while you were disconnected. onResync fires on every subscribe_success after the first one for that hook instance (i.e. every reconnect), so callers can re-read state the same way onEvent would trigger, closing that gap. Wire it to the same refresh/invalidate as onEvent:
useRealtimeChannel(channel, onEvent, { getConfig, onResync: () => router.refresh() });The socket is capped at one hour (MAX_CONNECTION_AGE_MS) and rotated then — a fresh, planned reconnect with no backoff. This is not a token-refresh mechanism: the AppSync Lambda authorizer only runs on connect/subscribe, never on an already-open subscription, so a token nearing its 15-minute exp does not need the socket cycled (that was tried and reverted — it caused an event-loss gap roughly every 14.5 minutes for no benefit). The one-hour cap exists only to bound how long a subscription outlives a revoked/rotated signing secret. getConfig is re-called fresh on every connect regardless, so the rotated connection always carries a current token.
getConfig returning null means "realtime isn't configured here" — the hook stays paused (connected: false, no socket, no error) and retries getConfig every 30s in case config becomes available later. A subscribe_error or broadcast_error logs once via console.warn (with the error payload) before closing and backing off — a misconfigured channel pattern shows up in the console instead of silently never connecting.
Shared connection: one WebSocket per app, not per hook. Each hook registers a subscriber on its channel; the first subscriber for a channel sends subscribe, the last sends unsubscribe, and the socket closes only when no channel has subscribers left. Two hooks on the same channel share one subscription and both receive its events — the second one is connected immediately and gets no onResync, because it has missed nothing. A subscribe_error still closes the whole socket (the usual cause is an expired token, and only a reconnect mints a new one), so every channel re-subscribes together and every hook that was already subscribed is resynced. Subscription ids are assigned per connection in registration order (s1, s2, …). __resetRealtimeConnection() is exported for tests only — it drops the module-level connection so each test starts from nothing.
connected comes from useRealtimeChannel, but useRealtimeChannel's onEvent/onResync want requestRefresh from useLiveRefresh — which needs connected. Break the cycle with a ref: declare useRealtimeChannel first, route its callbacks through a ref, and populate that ref from useLiveRefresh's return value (both hooks already hold their own callbacks in refs internally, so neither needs requestRefresh to be memoized).
"use client";
import { useRef } from "react";
import { useRouter, usePathname } from "next/navigation";
import { useRealtimeChannel, useLiveRefresh } from "@dxv-systems/realtime/react";
export function LiveRefresh() {
const router = useRouter();
const pathname = usePathname();
const requestRefreshRef = useRef<() => void>(() => {});
function matches(payload: unknown) {
const paths = (payload as { paths?: string[] })?.paths ?? [];
return paths.some((p) => pathname.startsWith(p));
}
const { connected } = useRealtimeChannel(
"/veijer/*",
(e) => {
if (matches(e)) requestRefreshRef.current();
},
{
getConfig: () => fetch("/api/realtime/token").then((r) => (r.ok ? r.json() : null)),
onResync: () => requestRefreshRef.current(), // a reconnect may have missed events — re-read on every resubscribe
},
);
const { requestRefresh } = useLiveRefresh({
refresh: () => router.refresh(),
intervalMs: 2 * 60_000,
connectedIntervalMs: 10 * 60_000,
connected,
// isBusy omitted — the built-in heuristic (ported from Veijer PR 426) is used
});
requestRefreshRef.current = requestRefresh;
return null;
}Vite / TanStack Query
import { useRef } from "react";
import { useQueryClient } from "@tanstack/react-query";
import { useRealtimeChannel, useLiveRefresh } from "@dxv-systems/realtime/react";
function useFeatureRequestEvents(tenantId: string, frId: string) {
const queryClient = useQueryClient();
const invalidate = () => queryClient.invalidateQueries({ queryKey: ["feature-request", frId] });
const requestRefreshRef = useRef<() => void>(() => {});
const { connected } = useRealtimeChannel(
`/portal/t-${tenantId}/fr-${frId}`,
() => requestRefreshRef.current(),
{
getConfig: () => fetch("/api/realtime/token").then((r) => (r.ok ? r.json() : null)),
onResync: () => requestRefreshRef.current(), // a reconnect may have missed events — re-fetch on every resubscribe
},
);
const { requestRefresh } = useLiveRefresh({
refresh: invalidate,
intervalMs: 30_000,
connectedIntervalMs: 2 * 60_000,
connected,
});
requestRefreshRef.current = requestRefresh;
}useLiveRefresh({ refresh, intervalMs, connectedIntervalMs, connected, isBusy? }) → { requestRefresh }
Ported from app-veijer-trappen-b7eb's reviewed LiveRefresh component (PR 426) so the two stay behaviourally identical. Polls refresh() every connected ? connectedIntervalMs : intervalMs, only while the tab is visible. On visibilitychange to visible, refreshes immediately if the last refresh is older than 30s. Skips a tick (single pending retry, not a stack of them — retries 15s later) while busy.
requestRefresh() runs that same gated path (visibility check, busy check + the single pending 15s retry, last-refresh-time update) on demand, so a realtime push event can trigger a refresh without bypassing the busy gate or duplicating its logic. Stable identity across re-renders — safe to pass directly as onEvent/onResync handlers. Calls within 500ms of each other coalesce into a single refresh (trailing), so a server action publishing several events in a row doesn't fan out into several router.refresh() calls. A call made while the tab is hidden is not lost: it's deferred and runs on the next return-to-tab, regardless of the 30s staleness rule above (that rule only governs the tick's own catch-up, not a pending push).
Busy, by default (omit isBusy to get this): a held-down pointer always counts as busy (stale after 5 minutes — a safety net against a stuck flag, not a real drag). Otherwise, a focused input/textarea/select/contenteditable or an open dialog ([role="dialog"] or [aria-modal="true"] — matches Radix, which doesn't use the native <dialog open> attribute) only counts as busy within 60s of the last keydown/input/pointerdown/wheel — an idle wall display with stale focus in a search box, or a dialog left open unattended, keeps refreshing regardless. The pure decision function is exported as shouldSkipRefresh(facts: BusyFacts) for testing; useLiveRefresh owns the interaction-tracking listeners itself, scoped to its own effect (registered on mount, removed on unmount — nothing attaches at module import time).
Pass isBusy to override this entirely with your own () => boolean. defaultIsBusy() is still exported as a stateless fallback for checking "is the user busy" outside useLiveRefresh (e.g. gating some unrelated action) — it has no interaction-window tracking of its own (a plain function can't own listeners), so prefer omitting isBusy from useLiveRefresh over passing defaultIsBusy explicitly.
Development
npm run build --workspace @dxv-systems/realtime
npm run test --workspace @dxv-systems/realtimeOutput is CommonJS, like every other package here (tsconfig.json sets module: CommonJS, no "type" field in package.json — identical to @dxv-systems/turnstile). ./react's "use client" directive survives the build but lands after the "use strict" the CommonJS emit adds: dist/react.js starts with "use strict"; then "use client";. @dxv-systems/turnstile's own dist/react.js has the exact same order and is already consumed successfully as @dxv-systems/turnstile/react by app-veijer-trappen-b7eb (via src/components/auth/turnstile-field.tsx) — a working precedent for this exact CJS build shape, so it's left as-is rather than changed to ESM output for this package alone. Wrap ./react in an app-owned "use client" component regardless; do not rely on the re-exported directive to establish the boundary for you.
Protocol notes and assumptions
datamessageeventfield. The AppSync docs' own protocol reference shows adatamessage'seventfield as a one-element array of a JSON string ("event": ["\"my event content\""]), while the publish docs describe the published payload as a plain array of stringified events. The hook handles both a bare string and an array defensively (Array.isArray(raw) ? raw : [raw], thenJSON.parseeach item) rather than picking one and risking silently dropping every event in production if the doc example was literal.- Publish headers. The Events HTTP publish doc's own example sends only
content-type: application/json(plus the auth header) — noaccept, nocontent-encoding, nocharset. The earlier draft of this package copiedaccept/content-encoding: amz-1.0/charset=UTF-8from a different section of the WebSocket protocol doc (the IAM subprotocol's connect-signing example), which does not apply to the plain HTTP/eventendpoint. Fixed to match the actual Events HTTP example:{ "content-type": "application/json" }, plus whateveraws4fetchadds for SigV4 (host/x-amz-date/Authorization). aws4fetchretries disabled.aws4fetch'sAwsClientretries any5xx/429response up to 10 times with growing backoff by default — that would blow past this package's 2s publish budget.publish()constructs it withretries: 0.- No token-driven reconnect. An earlier draft closed the socket ~30s before the JWT's
expto "refresh" it. Removed: the Lambda authorizer only runs onEVENT_CONNECT/EVENT_SUBSCRIBE, never against an already-open subscription, so that bought nothing but an event-loss gap every ~14.5 minutes. Reconnects now only happen for real reasons (network drop,katimeout,subscribe_error) plus the one-hourMAX_CONNECTION_AGE_MScap described above.
