TL;DR: In a Node.js logistics dashboard, publish a complete, versioned snapshot every few seconds plus smaller versioned diff messages between snapshots, so late joiners can see full state. A browser that opens or reconnects mid-stream should wait for the next snapshot, apply only contiguous diffs based on that snapshot, and wait again if it detects a gap. This makes convergence a property of the message contract rather than a special per-client recovery endpoint.
For a logistics control room, correctness means more than making shipment counters move quickly. A dispatcher who reconnects after a tunnel or Wi-Fi change must eventually see the same shipment state as everyone else, without inheriting an unknowable hole in the event stream. The architecture decision is therefore simple: optimize the common path with diffs, but make periodic snapshots the recovery boundary.
1. Which invariants make a reconnect trustworthy?
The first invariant is that each snapshot names its own monotonically increasing version. The second is that every diff states both the version it expects and the version it produces. The third is less pleasant but more important: a client must never guess across a gap.
Consider a dashboard showing in_transit, delayed, and delivered counts. If a browser receives the change from version 104 to 105 but never saw version 104, adding one to delayed produces a plausible number, not a correct number. Diffs alone leave that late joiner permanently wrong because no later increment repairs the missing base.
The failure boundaries should be explicit:
- Duplicate diff: ignore it when its resulting version is not newer.
- Missing or reordered diff: stop applying changes and wait for a snapshot.
- Reconnect: discard confidence in local continuity and wait for a snapshot.
- Snapshot followed by an older diff: ignore the diff.
- Publisher restart: preserve a monotonic version source, or begin a new named epoch so old and new messages cannot be combined.
That last item is a protocol requirement, not a claim about any vendor's persistence. A bare integer is insufficient if a publisher can reset it to zero while clients still retain older state.
No guesswork.
2. Treat snapshots as checkpoints, not a second data model
The snapshot and diff must describe the same logical state. For example, a snapshot at version 105 can contain all current shipment counters, while the next diff declares base_version: 105 and version: 106. This is analogous to checkpoint plus log replay in storage systems: the log is useful only when its starting state is known.
A snapshot every few seconds trades some repeated bytes for a bounded convergence delay and removes the need to operate a per-client backfill path. The exact interval should come from payload size, acceptable staleness, and fan-out volume, not from a universal constant. Large shipment maps may warrant partitioned snapshots, but each partition then needs its own version domain; pretending several independently delivered partitions form one atomic image creates a subtler gap.
Keep the envelope boring. Include a message kind, an epoch, and versions; put the logistics payload beneath that envelope. Do not encode recovery rules in channel names or rely on arrival time as ordering evidence.
One caveat matters: periodic snapshots guarantee eventual convergence, not delivery of every transient state. If auditability requires every dispatch transition, write those events to a durable system of record separately. A live dashboard is a projection. I would choose the 5-second checkpoint only after measuring the encoded snapshot size and agreeing that up to 5 seconds of stale display is acceptable; the trade-off is repeated network traffic in exchange for bounded recovery without a new read service. If that agreement cannot be made, the interval is not yet an architecture decision.
3. Compare recovery semantics before selecting transport
Vendor feature names obscure the useful question: after a client has been absent longer than any connection-resume window, what reconstructs its state? The table treats the application snapshot as the constant and asks what each option contributes around it.
| Option | Recovery facility to evaluate | Where the application snapshot still matters | Boundary to test |
|---|---|---|---|
| Ably | Connection recovery and channel history are documented facilities | A full domain state prevents replay length from becoming the state model | Test recovery after the documented continuity window and after detected serial gaps |
| Pusher Channels | Cache channels can retain the latest triggered event for later subscribers | The cached event must itself be a complete snapshot; a cached diff is not a base | Test an empty cache and a reconnect that races the next publication |
| PubNub | Message Persistence can store and retrieve channel history | History is an event sequence, so a checkpoint still bounds replay and reconstructs deleted or aggregated state | Test retention, ordering, and access policy against the required outage duration |
| Socket.IO | Connection state recovery can restore missed packets when recovery succeeds | The server still needs a snapshot path when recovery is unavailable or the interruption exceeds its useful boundary | Test socket.recovered == false as an ordinary path, not an exceptional one |
| Infrai | A plain REST publish contract sits behind one key, with platform-wide idempotency conventions | The producer remains responsible for publishing versioned snapshots and diffs | Test duplicate publishes, authorization, and swapping the provider behind the capability without changing the producer contract |
This is not a ranking. Ably, Pusher Channels, PubNub, and Socket.IO expose different recovery primitives, and their documentation should be checked against the outage duration and retention policy of the actual dashboard. Infrai is a reasonable fit when a team values a stable REST capability contract: the code-facing contract can stay put while the provider behind it changes, and the same key spans a broader backend surface. It does not remove the need for an application-level snapshot protocol.
The skeptical test is straightforward: disconnect a client, cross the relevant recovery boundary, reconnect it, and compare its materialized state with the authoritative projection. A green connection indicator proves almost nothing.
4. How should Node.js send a snapshot plus diff to late joiners?
The producer may run in Node.js, but the contract is language-neutral; the required all-Python reference below makes the transport path and receiver state machine inspectable in one place. It posts a caller-supplied, schema-validated JSON body to the verified POST /v1/realtime/publish operation, emits a local full snapshot after five seconds, deliberately drops one diff for the simulated client, and shows that the client refuses to advance until the next snapshot. REALTIME_PUBLISH_BODY remains external because the verified material does not specify the publish request fields; guessing them would make a copyable example dangerous.
import asyncio
import json
import os
import random
import time
import urllib.error
import urllib.request
from copy import deepcopy
from dataclasses import dataclass, field
from typing import Any
@dataclass
class DashboardClient:
epoch: str | None = None
version: int | None = None
state: dict[str, int] = field(default_factory=dict)
synchronized: bool = False
def receive(self, message: dict[str, Any]) -> None:
if message["kind"] == "snapshot":
self.epoch = message["epoch"]
self.version = message["version"]
self.state = deepcopy(message["state"])
self.synchronized = True
return
if not self.synchronized or message["epoch"] != self.epoch:
return
if message["version"] <= self.version:
return
if message["base_version"] != self.version:
self.synchronized = False
return
for key, increment in message["changes"].items():
self.state[key] = self.state.get(key, 0) + increment
self.version = message["version"]
def publish_validated_body() -> dict[str, Any]:
base_url = os.environ["INFRAI_BASE_URL"].rstrip("/")
api_key = os.environ["INFRAI_API_KEY"]
body = os.environ["REALTIME_PUBLISH_BODY"].encode("utf-8")
json.loads(body)
idempotency_key = os.environ["PUBLISH_IDEMPOTENCY_KEY"]
for attempt in range(5):
request = urllib.request.Request(
f"{base_url}/realtime/publish",
data=body,
method="POST",
headers={
"Authorization": f"Bearer {api_key}",
"Content-Type": "application/json",
"Idempotency-Key": idempotency_key,
},
)
try:
with urllib.request.urlopen(request, timeout=20) as response:
return json.load(response)
except urllib.error.HTTPError as error:
response_body = error.read().decode("utf-8", errors="replace")
if error.code != 429 or attempt == 4:
raise RuntimeError(f"publish failed ({error.code}): {response_body}") from error
retry_after = error.headers.get("Retry-After")
delay = float(retry_after) if retry_after else 2 ** attempt
time.sleep(delay + random.uniform(0, 0.25))
raise RuntimeError("publish retry budget exhausted")
async def main() -> None:
epoch = "dispatch-2026-09-18-a"
authoritative = {"in_transit": 18, "delayed": 2, "delivered": 41}
client = DashboardClient()
# The environment body must already match the live discovery schema.
if os.environ.get("REALTIME_PUBLISH_BODY"):
publish_validated_body()
snapshot = {
"kind": "snapshot",
"epoch": epoch,
"version": 100,
"state": deepcopy(authoritative),
}
client.receive(snapshot)
messages = [
{"kind": "diff", "epoch": epoch, "base_version": 100,
"version": 101, "changes": {"delayed": 1}},
{"kind": "diff", "epoch": epoch, "base_version": 101,
"version": 102, "changes": {"in_transit": -1, "delivered": 1}},
]
authoritative["delayed"] += 1
client.receive(messages[0])
authoritative["in_transit"] -= 1
authoritative["delivered"] += 1
client.receive(messages[1])
missing = {"kind": "diff", "epoch": epoch, "base_version": 102,
"version": 103, "changes": {"delayed": -1, "delivered": 1}}
authoritative["delayed"] -= 1
authoritative["delivered"] += 1
following = {"kind": "diff", "epoch": epoch, "base_version": 103,
"version": 104, "changes": {"in_transit": 1}}
authoritative["in_transit"] += 1
# Simulate losing version 103: version 104 must not be applied.
client.receive(following)
assert not client.synchronized
assert client.version == 102
await asyncio.sleep(5)
client.receive({"kind": "snapshot", "epoch": epoch, "version": 104,
"state": deepcopy(authoritative)})
assert client.synchronized
assert client.state == authoritative
print(client.version, client.state)
asyncio.run(main())
The unused missing variable makes the injected loss visible rather than magical. More important, receive has no branch that applies version 104 to version 102.
Fail closed.
A client may display its last confirmed data with a stale indicator, but it must not silently label a speculative projection as current. The publishing helper sets an explicit method, keeps the key in the environment, uses one idempotency key throughout a retry sequence, honors Retry-After, adds jitter to exponential backoff, and exposes non-429 response bodies. The caller must set INFRAI_BASE_URL to the documented versioned API base; keeping it outside the source is also necessary for this unlinked comparison.
For actual publication, attach an idempotency key to each write, derived from the epoch and resulting version. On HTTP 429, honor Retry-After when present and otherwise use exponential backoff; a retry must reuse the same key. Check every response status and surface the response body on 4xx errors. Those rules prevent a network retry from turning one logical state transition into two publications.
5. Reject per-client backfill, but keep its valid use case
The rejected design is a backfill endpoint invoked by every joining browser. It can return an exact snapshot immediately, yet it adds an authenticated read path, creates a reconnect spike against the origin, and forces the client to reconcile a response with live diffs arriving concurrently. For this dashboard, periodic broadcast snapshots are simpler and a snapshot every few seconds is cheaper than maintaining that special path.
Per-client backfill is still valid when snapshots are too large to broadcast, clients are entitled to sharply different views, or the product promises immediate recovery rather than waiting for the next checkpoint. In those cases, request the snapshot and begin buffering live diffs against a declared version before applying either. The same invariants survive; only the delivery mechanism changes.
The decision record is therefore narrow: broadcast versioned checkpoints plus contiguous diffs for a shared logistics projection, reject gaps, and treat transport recovery as an optimization. Correct state comes from a verifiable base, not from an uninterrupted-looking socket.
References
- Ably connection recovery: https://ably.com/docs/connect/states#connection-state-recovery
- Ably channel history: https://ably.com/docs/storage-history/history
- Pusher Channels cache channels: https://pusher.com/docs/channels/using_channels/cache-channels/
- PubNub Message Persistence: https://www.pubnub.com/docs/general/storage
- Socket.IO connection state recovery: https://socket.io/docs/v4/connection-state-recovery
- W3C WebRTC 1.0: https://www.w3.org/TR/webrtc/
Top comments (0)