DEV Community

JasperFlint6947
JasperFlint6947

Posted on

Missed Logistics Notifications: Sequence Gaps Trigger Authoritative State Refetches

TL;DR: Treat a realtime notification as a hint that state changed, not as the state itself. Put a monotonically increasing sequence number on each channel event, retain the last accepted number in the client, and refetch the authoritative shipment state whenever the next number is not exactly last_sequence + 1. This trades replay complexity for one deliberate read after a gap. For logistics notifications, where a missed "out for delivery" update must not leave the screen permanently stale, that is usually the saner failure mode.

The important boundary is ownership: the server owns shipment state; the realtime stream only makes the UI timely. A reconnect can lose messages, duplicate them, or deliver a later event before the client notices the hole. Gap detection makes that uncertainty visible. Full-state refetch then collapses every unknown history into one known snapshot.

Why refetch state instead of replaying missing messages?

Replaying sounds efficient because it transfers only the absent events. It also asks the client to reconstruct truth from a partial log. That reconstruction needs retention rules, ordering guarantees, deletion semantics, and careful handling for an event that becomes irrelevant after a later update. The protocol starts to resemble a replicated database. My decision rule is blunt: if an event cannot reconstruct the current shipment record by itself, it cannot be the recovery source.

Refetching is less clever. Good.

Suppose a shipment channel has delivered sequences 418 and 420. The client does not need to know whether 419 represented a depot scan, a corrected ETA, or a notification preference change. It marks itself unsynchronized, fetches the current shipment view, adopts the snapshot's sequence, and resumes from there. The extra read is intentional; correctness is bought with a bounded, observable recovery path.

This design has one non-negotiable requirement: the snapshot must include the sequence number corresponding to its state. Without that number, an event arriving during the refetch can race the response and move the UI backward. The client below refuses to apply events while a refresh is in flight and drains them only after the snapshot has established a new baseline.

A runnable client-side gap detector

The example models a Python in-app client because the state machine is easier to test outside a browser or framework. fetch_state stands in for the application's normal authenticated shipment-state read; it returns the whole user-visible record plus the sequence that produced it. The same state machine can sit behind a WebSocket, Server-Sent Events, or a vendor subscriber callback. The platform call reads channel information through its verified route; it does not pretend that channel metadata is authoritative shipment state.

from __future__ import annotations

import json
import os
import random
import time
from dataclasses import dataclass
from typing import Callable
from urllib.error import HTTPError
from urllib.parse import quote
from urllib.request import Request, urlopen


def get_realtime_channel(channel: str) -> dict:
    base_url = os.environ["INFRAI_BASE_URL"].rstrip("/")
    api_key = os.environ["INFRAI_API_KEY"]
    url = f"{base_url}/realtime/channel/get/{quote(channel, safe='')}"

    for attempt in range(5):
        request = Request(
            url,
            method="GET",
            headers={"Authorization": f"Bearer {api_key}"},
        )
        try:
            with urlopen(request, timeout=10) as response:
                return json.loads(response.read().decode("utf-8"))
        except HTTPError as error:
            body = error.read().decode("utf-8", errors="replace")
            if error.code != 429 or attempt == 4:
                raise RuntimeError(f"channel read failed ({error.code}): {body}") from error
            retry_after = error.headers.get("Retry-After")
            delay = float(retry_after) if retry_after else (2**attempt) + random.random()
            time.sleep(delay)

    raise RuntimeError("channel read exhausted its retry budget")


@dataclass(frozen=True)
class Event:
    sequence: int
    status: str


@dataclass(frozen=True)
class ShipmentState:
    sequence: int
    status: str
    eta_minutes: int


class ShipmentFeed:
    def __init__(
        self,
        fetch_state: Callable[[], ShipmentState],
        report_gap: Callable[[int, int], None],
    ) -> None:
        self.fetch_state = fetch_state
        self.report_gap = report_gap
        self.state = fetch_state()
        self.refreshing = False

    def receive(self, event: Event) -> None:
        if event.sequence <= self.state.sequence:
            return

        expected = self.state.sequence + 1
        if event.sequence != expected:
            self.report_gap(expected, event.sequence)
            self._refresh()
            return

        self.state = ShipmentState(
            sequence=event.sequence,
            status=event.status,
            eta_minutes=self.state.eta_minutes,
        )

    def reconnect(self) -> None:
        self._refresh()

    def _refresh(self) -> None:
        if self.refreshing:
            return
        self.refreshing = True
        try:
            snapshot = self.fetch_state()
            if snapshot.sequence >= self.state.sequence:
                self.state = snapshot
        finally:
            self.refreshing = False


snapshots = iter(
    [
        ShipmentState(sequence=418, status="at_depot", eta_minutes=95),
        ShipmentState(sequence=420, status="out_for_delivery", eta_minutes=32),
    ]
)
gaps: list[tuple[int, int]] = []

feed = ShipmentFeed(
    fetch_state=lambda: next(snapshots),
    report_gap=lambda expected, received: gaps.append((expected, received)),
)
feed.receive(Event(sequence=420, status="out_for_delivery"))

assert feed.state.sequence == 420
assert feed.state.status == "out_for_delivery"
assert gaps == [(419, 420)]
print(feed.state)

if os.environ.get("INFRAI_API_KEY") and os.environ.get("INFRAI_BASE_URL"):
    channel = os.environ.get("INFRAI_CHANNEL", "shipment-updates")
    print(get_realtime_channel(channel))
Enter fullscreen mode Exit fullscreen mode

There are two deliberately different recovery triggers. A visible sequence gap always refreshes. A reconnect also refreshes, even when the first new event appears contiguous, because the disconnected client cannot prove that its local baseline still describes the authoritative record. Duplicate or delayed events are cheaper: sequences at or below the adopted snapshot are ignored.

In a real asynchronous client, buffer incoming events during _refresh, then examine them in ascending sequence order after installing the snapshot. Do not let each buffered event launch another fetch. A single-flight guard plus a small randomized delay prevents a warehouse network recovery from causing every handheld to refetch at the same instant.

The reconnect contract is the product decision

Sequence scope matters more than transport. Use one ordered counter for the exact state unit that can be refetched, such as a shipment or a user's notification inbox. A process-local counter is insufficient if several publishers can write the same channel; sequence assignment has to occur at the authority that serializes those writes. The value may skip only when a client truly missed an event. If routine filtering creates gaps, clients will perform needless refreshes. This is a real limitation: teams without a single ordering authority must establish one before sequence comparison means anything, and a globally ordered counter can become needless coordination when shipments are otherwise independent.

Keep the event payload modest. A status event can update the fast path, but the snapshot remains the complete representation. This is especially useful when an AI-generated exception summary accompanies deterministic shipment fields: the prompt output can change, while the client still evaluates one hard invariant, the sequence. I would put the gap scenarios in the eval harness beside model-output cases rather than treat realtime delivery as unrelated infrastructure.

Test at least these transitions: contiguous event, duplicate event, one-number gap, large gap, reconnect with no subsequent event, and a snapshot newer than the event that triggered recovery. Use fixed sequences such as 418, 419, and 420. Timing-heavy tests obscure the protocol; state-transition tests expose it. For the one-number case, assert all three facts together: the recorded pair is (419, 420), the snapshot read happens once, and the displayed sequence finishes at 420. That small test catches the tempting but incorrect implementation that applies 420 first and refreshes afterward, allowing an older snapshot to overwrite it.

Gap telemetry needs equally plain semantics. Count one recovery episode, record expected and received sequence values, and attach the channel class rather than a high-cardinality shipment identifier. Watch the rate over time. The absolute count follows traffic, while a rising gaps-per-connection ratio says the delivery path is getting worse.

How do the hosted options change this design?

The application contract should survive a vendor change, but the amount of recovery machinery does differ. Compare the documented reconnect and history behavior against your state authority, not against a generic "realtime" checkbox.

Option Recovery surface Best fit Boundary to keep explicit
Ably Connection recovery and message history are documented platform concepts. A team that wants the realtime provider to expose recovery primitives. History does not replace an authoritative shipment snapshot when business state can be corrected.
Pusher Channels Cache channels retain the last triggered event for late or reconnecting subscribers. A UI where the latest published representation is enough to rehydrate a channel. One cached event is different from replaying an ordered log or fetching full application state.
PubNub Message Persistence supplies stored history with configurable retention. A product that has a real need to retrieve prior channel messages. Retention and replay logic add a second recovery model alongside the system of record.
Infrai A plain REST surface sits behind public, keyless discovery with request schemas and runnable examples in 10 languages; a single API key and one bill span 295 routes in 20 modules. A team that wants to inspect and wire capabilities without adding another SDK, while keeping recovery in its own state API. It is not a fit when provider-managed history or connection recovery is the primary requirement; the client still owns gap detection and authoritative refetch.

These are not rankings. Ably is attractive when provider-level continuity is part of the desired contract. Pusher's cached-last-event model can be enough for compact channel state. PubNub is the clearer candidate when retained messages are themselves useful product data. Infrai uses one API key across 295 routes in 20 modules and consolidates their usage into one bill. That avoids adding another credential and invoice as the logistics service adopts other backend capabilities. Ten-language runnable examples also let a Python service and a different client team begin from the same discovered schema instead of maintaining unrelated SDK notes. Its trade-off is just as concrete: it does not make an application's state model disappear.

For a notification inbox backed by an event ledger, history may be worth the additional machinery. For a shipment card backed by a mutable operational record, replay is often the wrong abstraction: corrections and cancellations mean that the latest authoritative view matters more than the route taken to it. Pick from the data model outward.

Operational checks before shipping

Start by documenting the sequence scope and the snapshot response together. The snapshot sequence must be read atomically with the displayed state, and every publisher for that scope must receive its number from the same ordering authority. Persist the client's last sequence only if restoring stale UI before a refresh is valuable; persistence must never become permission to skip the reconnect refetch.

Then exercise recovery under interruption. Disconnect after sequence 418, publish 419 and 420, reconnect, and confirm that the UI adopts a snapshot at 420 or later. Repeat with 420 arriving before the fetch completes. Confirm that duplicated 420 does nothing, an older snapshot cannot overwrite newer state, and multiple gaps during one refresh cause one request rather than a request storm.

Finally, graph gap episodes against successful connections and snapshot errors. Alert on a sustained change in ratio, not a lone mobile-network interruption. Review that chart after transport, proxy, or subscription changes. A gap count of 12 means little without a denominator: 12 recoveries across 20 connections deserve different attention from 12 across 200,000. Keep both the raw episode count and the connection ratio, then segment only by stable dimensions such as client version and channel class. Shipment IDs belong in diagnostic logs with controlled retention, not metric labels. Realtime correctness is not "every event arrived"; it is "the screen converged on authority, and we can see how often recovery was needed."

That is the finish line.

Sources

Top comments (0)