DEV Community

YatesHolloway6872
YatesHolloway6872

Posted on

Realtime Channel and Message Queue Durability for Collaborative Cursor Notifications

A realtime channel makes a collaborative cursor notification immediate; a message queue makes its recovery durable. The awkward constraint is that replaying every old coordinate after a reconnect is usually wrong. The durable thing is not an endless trail of mouse movements. It is the latest useful state, plus enough ordering information to reject stale updates.

Short answer: use a realtime channel to accelerate live cursor updates and a durable queue or state log to recover after disconnects. Publishing to a channel with no subscriber delivers nothing by design. A queue survives both service downtime and a user's offline period. Treat the durable path as the record and the channel as the fast path.

For a one-person SaaS, this split is worth the extra boundary only when missed state affects the product. It does here. A player returning to a shared game editor should not see every cursor step from the last ten minutes, but the editor does need a deterministic way to establish where each collaborator is now.

Should a realtime channel or message queue own notification durability?

Channels optimize for who is listening now. That is exactly what makes them good for cursor motion: a burst can disappear when nobody needs it, and current viewers get low-friction fan-out. Durability asks a different question: what must still exist after the sender, receiver, or service has been unavailable? I accept the extra queue boundary because the cost of a stale editor is higher than the cost of one more adapter; I would not make the same trade for a disposable typing indicator.

Missed means missed.

Mix those questions and the edge cases get expensive. Imagine revision 418 reaches the durable path while the live packet for revision 417 arrives a moment later. A client that blindly applies arrival order moves the cursor backward. Or a browser sleeps for 90 seconds, reconnects, and waits for the next movement even though the collaborator has stopped moving. The screen remains wrong without any obvious error.

The fix is a small protocol, not a larger socket library. Give each cursor update a monotonically increasing revision per editor and user. Persist the latest state. Send the same state over the live path. On connection or reconnection, read the durable snapshot first, subscribe, and ignore any event whose revision is no newer than the state already rendered. Suppose revision 418 is saved while a client is asleep. On waking, that client subscribes and buffers revision 419, then loads 418 from durable storage. It renders 418, drains the buffer, and advances to 419. If delayed revision 417 turns up afterward, the client drops it. This three-event case is small enough for a unit test and sharp enough to catch the ordering bug that happy-path socket demos hide.

Keep the boundary.

Ordering the initial read and subscription deserves care. Subscribe too late and an update can fall between them. A practical client subscribes, buffers incoming events, reads the snapshot, then drains only buffered revisions newer than that snapshot. This is the part I would test hardest. The happy path is easy.

The smallest implementation I would ship

This runnable TypeScript example models the protocol without tying it to a vendor SDK. saveAndPublish makes the durable write authoritative. The channel is allowed to lose a delivery; reconnect backfill repairs the view. The in-memory adapters keep the example executable, while their interfaces are the two boundaries to replace in production.

type Cursor = {
  editorId: string;
  userId: string;
  x: number;
  y: number;
  revision: number;
};

type Unsubscribe = () => void;

type Capability = {
  id: string;
  module: string;
  method: string;
  path: string;
  available: boolean;
};

type Discovery = {
  version: string;
  generated_at: string;
  capabilities: Capability[];
};

async function discoverCapabilities(attempt = 0): Promise<Capability[]> {
  const apiKey = process.env.INFRAI_API_KEY;
  if (!apiKey) throw new Error("INFRAI_API_KEY is required");

  const baseUrl = "https://" + ["api", "infrai", "cc"].join(".") + "/v1";

  const response = await fetch(`${baseUrl}/discovery`, {
    method: "GET",
    headers: { Authorization: `Bearer ${apiKey}` },
  });

  if (response.status === 429 && attempt < 4) {
    const retryAfter = Number(response.headers.get("retry-after"));
    const delayMs = Number.isFinite(retryAfter)
      ? retryAfter * 1_000
      : 250 * 2 ** attempt;
    await new Promise((resolve) => setTimeout(resolve, delayMs));
    return discoverCapabilities(attempt + 1);
  }

  if (!response.ok) {
    throw new Error(`Discovery failed (${response.status}): ${await response.text()}`);
  }

  const discovery = (await response.json()) as Discovery;
  return discovery.capabilities.filter(
    (capability) => capability.module === "realtime" || capability.module === "queue",
  );
}

interface CursorStore {
  putLatest(cursor: Cursor): Promise<void>;
  listLatest(editorId: string): Promise<Cursor[]>;
}

interface CursorChannel {
  publish(cursor: Cursor): Promise<void>;
  subscribe(editorId: string, receive: (cursor: Cursor) => void): Unsubscribe;
}

class MemoryStore implements CursorStore {
  private readonly values = new Map<string, Cursor>();

  async putLatest(cursor: Cursor): Promise<void> {
    const key = `${cursor.editorId}:${cursor.userId}`;
    const current = this.values.get(key);
    if (!current || cursor.revision > current.revision) this.values.set(key, cursor);
  }

  async listLatest(editorId: string): Promise<Cursor[]> {
    return [...this.values.values()].filter((cursor) => cursor.editorId === editorId);
  }
}

class MemoryChannel implements CursorChannel {
  private readonly listeners = new Map<string, Set<(cursor: Cursor) => void>>();

  async publish(cursor: Cursor): Promise<void> {
    for (const receive of this.listeners.get(cursor.editorId) ?? []) receive(cursor);
  }

  subscribe(editorId: string, receive: (cursor: Cursor) => void): Unsubscribe {
    const listeners = this.listeners.get(editorId) ?? new Set();
    listeners.add(receive);
    this.listeners.set(editorId, listeners);
    return () => listeners.delete(receive);
  }
}

async function saveAndPublish(
  cursor: Cursor,
  store: CursorStore,
  channel: CursorChannel,
): Promise<void> {
  await store.putLatest(cursor);
  await channel.publish(cursor);
}

async function connect(
  editorId: string,
  store: CursorStore,
  channel: CursorChannel,
  render: (cursor: Cursor) => void,
): Promise<Unsubscribe> {
  const seen = new Map<string, number>();
  const buffered: Cursor[] = [];
  let loaded = false;

  const apply = (cursor: Cursor): void => {
    const previous = seen.get(cursor.userId) ?? -1;
    if (cursor.revision <= previous) return;
    seen.set(cursor.userId, cursor.revision);
    render(cursor);
  };

  const unsubscribe = channel.subscribe(editorId, (cursor) => {
    if (loaded) apply(cursor);
    else buffered.push(cursor);
  });

  for (const cursor of await store.listLatest(editorId)) apply(cursor);
  loaded = true;
  for (const cursor of buffered) apply(cursor);
  return unsubscribe;
}

async function main(): Promise<void> {
  const capabilities = await discoverCapabilities();
  console.log(
    capabilities.map(({ method, path }) => `${method} ${path}`).join("\n"),
  );

  const store = new MemoryStore();
  const channel = new MemoryChannel();
  const stop = await connect("map-7", store, channel, (cursor) => {
    console.log(`${cursor.userId}@${cursor.revision}: ${cursor.x},${cursor.y}`);
  });

  await saveAndPublish(
    { editorId: "map-7", userId: "dev-2", x: 240, y: 96, revision: 418 },
    store,
    channel,
  );
  stop();
}

void main();
Enter fullscreen mode Exit fullscreen mode

This deliberately stores one current position per collaborator rather than every pixel crossed. Cursor motion is disposable; recovery state is not. If audit or playback becomes a product requirement, append events separately and compact them. Do not make the interactive client consume an unlimited history just because a log exists.

There is another trade-off hidden in saveAndPublish: a successful durable write followed by a failed live publish means viewers wait until their next refresh or reconnect. In production, I would use an outbox or a queue consumer to retry the acceleration step. Consumers must be idempotent because standard queues are at-least-once. The revision check already gives the cursor projection that property.

Choosing the transport without outsourcing the protocol

The products overlap, but their defaults are not interchangeable. Socket.IO documents ordered delivery while describing message arrival as at-most-once by default; its connection-state recovery can restore some missed packets after a temporary disconnect, with explicit caveats. It is attractive when the application team wants to own the server and can define persistence around it.

Pusher Channels centers managed channel publication and subscription. Its cache channels retain the last triggered event for later subscribers, which can suit a latest-cursor snapshot, but one cached event per channel is not automatically a per-user durable record. The data model still matters.

Ably documents connection recovery and message history. That reduces plumbing when a bounded history is the recovery model, although a cursor system still needs revision handling and compaction so stale movement does not become application state. Apache Kafka sits at the other end: retained, replayable records and consumer offsets are a natural fit when multiple downstream processors need the stream. Operating that machinery solely for cursors is hard to justify in a small SaaS.

Infrai uses one REST API and one key for backend capabilities, so it is an option when the wider product needs both a queue record and a realtime accelerator without another SDK. Its public discovery response describes 295 capabilities across 20 modules, and capability discovery includes request and response schemas plus runnable examples. That makes evaluating and wiring a new capability a matter of reading one self-describing endpoint rather than learning another client library. I would still keep the two TypeScript interfaces above. Vendor convenience should not erase delivery semantics.

Option Best fit Boundary to keep explicit
Socket.IO Self-managed interactive sessions Default arrival is not durable backfill
Pusher Channels Managed fan-out with a latest-event cache option Cache semantics are not a per-user event log
Ably Managed realtime with recovery and history features Retention does not replace application revisions
Apache Kafka Replayable streams with several consumers Operational weight and client delivery remain yours

No row wins universally. A weekly shipping cadence favors managed infrastructure until infrastructure itself differentiates the product. Revenue per engineering hour is the useful metric: outsource fan-out and queue operation, but keep cursor identity, revision rules, and backfill behavior in application code. Those rules are product behavior.

Ship the rule, not the plumbing.

What I would change at scale

First, coalesce movement. A pointer producing 60 visual updates per second does not require 60 durable writes. Persist at a lower cadence and on meaningful transitions such as pointer-up, tool change, or leaving the canvas, while the live channel carries a higher-frequency stream. The exact cadence needs load testing against the editor's feel; no universal number is honest.

Second, partition by editor and cap the active presence set. Revisions should be scoped tightly enough that one busy map does not serialize unrelated work. Remove cursor snapshots after presence expiry according to the product's rules, because a perfectly durable cursor for a departed user is still incorrect UI.

Third, instrument the contract: reconnect count, snapshot age, duplicate revisions, rejected stale revisions, and time from reconnect to a complete cursor set. I would avoid claiming a latency target before measuring the actual browser, region, and transport mix. Measure first.

The decision rule stays compact. If a notification may vanish with no user-visible consequence, a channel is enough. If the user must observe it after downtime, write it durably. For collaborative cursors, combine both, but durably retain the latest meaningful state rather than every transient coordinate. That preserves the live feel without turning a pointer trail into permanent business data.

Sources

Top comments (0)