DEV Community

FlorianBlake3536
FlorianBlake3536

Posted on

Live Caption Room Delivery: Segment Timing and Reconnect Recovery

TL;DR

For customer-support rooms, publish each caption segment with stable room-local ordering and media-relative timing, acknowledge it only after a durable append, and reconnect with a cursor-based backfill before switching the client to the live stream. The transport carries events; the log defines what can be recovered.

Start with the bill, because caption payloads are deceptively small. A useful first-pass model is retained_bytes = segment_rate * average_segment_bytes * active_seconds * rooms * retention_days, but raw bytes are rarely the only term. Object count, write operations, replication, and index entries can outweigh a few hundred bytes of text when every interim hypothesis becomes its own durable object.

Consider a planning example, not a benchmark: 40 concurrently active support rooms, four interim segments per second, 350 bytes per encoded event, eight active hours per day, and 30 days of retention produce 48.384 GB of raw events and 138,240,000 event records before replicas, indexes, or storage framing. The important number is the record count. Keep the live stream granular, append events to a partitioned log, compact finalized captions into larger minute-scale objects, and retain only final segments plus a short correction window. That change attacks the dominant term without delaying the on-screen captions.

The loss is real.

If every interim recognition hypothesis is discarded after the correction window, an investigation can reconstruct what the customer actually saw only if display revisions were separately audited; it cannot recreate every transient word choice made by the recognizer. Teams with legal replay requirements should retain that audit stream. Everyone else should avoid paying indefinitely for data that has no recovery or compliance use.

How should a room channel publish caption segments with timing after reconnect?

Treat a caption as a domain event, not a string pushed through a socket. Each event needs a stable room ID, a session ID, a source ID, a monotonically increasing source sequence, start and end offsets against one declared media timeline, text, and a final/interim flag. A stable event ID derived from those identity fields makes a retry recognizable. Wall-clock receipt time is still useful for operations, but it must not decide caption order: clocks jump, networks reorder packets, and two speakers can legitimately overlap.

For a Node.js publisher, the runtime-specific part is small: serialize the event, durably append with an idempotency key, then fan it out and resolve the publish promise.

Don't resolve on socket write alone.

A WebSocket frame proves neither durable storage nor application-level consumption; RFC 6455 defines the transport framing and connection behavior, not a replay contract. The same distinction applies when WebRTC carries the live media. WebRTC data channels expose ordered and partially reliable delivery choices, while the application still owns room history, cursor meaning, and duplicate handling.

Use one sequence domain per caption source, rather than pretending that overlapping agents and customers have a single natural speech order. The room view can merge sources by media offset, then use source ID and source sequence as deterministic tie-breakers. If the user refreshes after seeing source agent-7 sequence 812 and source customer-3 sequence 641, the reconnect request sends both cursors. The server returns later events for each source, including corrections, and the client applies them idempotently.

Timing needs an equally strict contract. start_ms and end_ms should be non-negative offsets from the room session's media epoch, with end_ms >= start_ms. If capture restarts and establishes a new media epoch, issue a new session ID instead of silently resetting offsets inside the old session. Otherwise a perfectly delivered segment can appear minutes earlier in the transcript, and no amount of socket reliability will repair the semantic error.

Short segments are fine. Ambiguous timelines aren't.

The storage model behind reliable backfill

A recoverable design has three related views with different jobs. The append log is authoritative for reconnect and preserves ordered events. The live channel is an acceleration path that reduces display latency. Compacted transcript objects serve long-range reads, exports, and retention policies without forcing an object store to hold millions of tiny records. These views may share infrastructure, but their contracts should remain separate so an optimization in one does not quietly weaken another.

Concern Contract Failure mode named up front Practical control
Identity One stable event ID for one logical segment revision A retry creates duplicate words Conditional append or unique key
Ordering Monotonic sequence per room session and source Arrival order scrambles overlapping speech Reject regressions; merge deterministically
Timing Offsets use one declared media epoch A capture restart moves captions backward Rotate the session ID
Durability Acknowledgement follows durable append The UI shows text that reconnect cannot recover Append before fan-out acknowledgement
Backfill Cursor means last applied sequence Inclusive replay repeats the boundary event Query strictly after the cursor
Retention Final transcript and audit needs are explicit Compaction deletes evidence a policy requires Separate transcript and audit retention

The write boundary deserves suspicion. A publisher retry after an ambiguous timeout can submit the same event twice, so append_if_absent must return the existing event when the identity and payload match, and reject an identity collision when they do not. A consumer must also be idempotent because a reconnect may overlap with an in-flight live delivery. Exactly-once language tends to hide these two ordinary checks; stable identity plus repeatable application is easier to inspect.

Partitioning by room session keeps local ordering cheap and bounds a reconnect scan. It also creates a hot-partition risk during large support events, so capacity tests should include one unusually busy room rather than only an even distribution across thousands of quiet rooms. I'm not sure where a particular system's hot-room ceiling lies without its payload distribution, storage engine, replication policy, and latency target. A load test with recorded segment sizes and reconnect bursts resolves that uncertainty.

Compaction must preserve correction semantics. If interim sequence 811 says reset password and final sequence 812 replaces it with reset the password, the compacted transcript stores the final text plus enough source and timing metadata to resume after its high-water mark. It should not manufacture a new ordering scheme. Keep a manifest that records the compacted-through sequence for every source, publish the object, and only then make older log records eligible for deletion under policy. Deleting first creates a backfill hole if object publication never becomes visible to readers.

A transport-neutral recovery loop

The following Python sketch makes the concurrency boundary explicit. The interfaces are deliberately generic. subscribe must register the listener before it returns, high_watermarks must read a consistent set of durable cursors, and list_after must return records no later than those watermarks in source-sequence order. A production implementation also authenticates room membership before either path runs.

from dataclasses import dataclass
from typing import AsyncIterator, Mapping


@dataclass(frozen=True)
class CaptionSegment:
    room_id: str
    session_id: str
    source_id: str
    sequence: int
    start_ms: int
    end_ms: int
    text: str
    is_final: bool

    @property
    def event_id(self) -> str:
        return f"{self.room_id}:{self.session_id}:{self.source_id}:{self.sequence}"


def validate(segment: CaptionSegment) -> None:
    if segment.sequence < 1:
        raise ValueError("sequence must be positive")
    if segment.start_ms < 0 or segment.end_ms < segment.start_ms:
        raise ValueError("invalid media-relative timing")
    if not segment.text:
        raise ValueError("caption text must not be empty")


async def publish(segment: CaptionSegment, log, rooms) -> CaptionSegment:
    validate(segment)
    stored = await log.append_if_absent(
        partition=(segment.room_id, segment.session_id),
        key=segment.event_id,
        value=segment,
    )
    await rooms.publish(segment.room_id, stored)
    return stored


async def reconnect(
    room_id: str,
    session_id: str,
    applied: Mapping[str, int],
    log,
    rooms,
) -> AsyncIterator[CaptionSegment]:
    # Subscribe first. Events after the watermark stay queued for the live phase.
    live = await rooms.subscribe(room_id)
    limits = await log.high_watermarks(room_id, session_id)

    async for segment in log.list_after(
        partition=(room_id, session_id),
        cursors=applied,
        through=limits,
    ):
        yield segment

    async for segment in live:
        if segment.session_id != session_id:
            continue
        if segment.sequence > limits.get(segment.source_id, 0):
            yield segment
Enter fullscreen mode Exit fullscreen mode

Subscribing before reading the watermark closes a subtle race. Walk through a reconnect from cursor 812 rather than trusting the names of the methods. The server first registers the live listener. Segment 813 is durably appended and queued for that listener; the watermark read then returns 813. Backfill asks for records strictly after 812 and no later than 813, so it yields 813. The queued live copy is later ignored because its sequence is not greater than the captured watermark. Now consider segment 814 arriving one instruction later, after the watermark read: it misses the bounded backfill by design, remains in the already active listener queue, and passes the sequence > 813 check. Both boundary events appear once. Reverse the first two operations and 813 can land after the history query but before listener registration, leaving a silent hole that neither side knows to request. Remove the sequence check and the safe ordering delivers 813 twice. The algorithm therefore depends on three contracts working together — listener registration before watermark capture, a backfill upper bound, and idempotent client application — rather than on favorable network timing. This is the sort of reconnect bug that passes a calm demo and fails during a busy handoff, particularly when several agents are producing captions while a supervisor opens the room transcript.

A slow client still needs a policy. Bound its outbound queue by bytes and age, stop incremental delivery when the bound is crossed, and send a resync instruction carrying the last safely applied cursor. Dropping an arbitrary middle segment is worse because the stream then looks healthy while the transcript is incomplete. The client should render only after validating identity and timing, keep the highest contiguous sequence per source, and request a gap fill when it receives 814 after 812. Event 814 may be buffered briefly; it must not advance the durable cursor past missing 813.

The catch is that this design is not suitable when the transcript is intentionally ephemeral and losing captions on refresh is acceptable; a memory-only channel is simpler there. At the other extreme, regulated contact centers may need immutable display audits, legal holds, regional placement, and independently verifiable deletion. In that case, stick with a storage system and operating model that can prove those controls, even if its reconnect path is more involved. The architecture is a baseline, not evidence of compliance.

Test the gaps, not the happy path

Unit tests should cover duplicate publish, identity collision, sequence regression, invalid timing, and replacement of interim text by a final segment. Integration tests should pause delivery after a known cursor, publish more events, reconnect, and assert that the resulting event-ID set is complete and unique. Then introduce the hard interleaving: establish the live subscription, append one event immediately before the watermark read, append another immediately after it, and verify both appear once.

Use property-based tests for longer histories. Generate multiple sources, retries, overlapping timing, disconnect points, and bounded reordering; after recovery, compare the client state with a fold over the authoritative log. This catches assumptions that a dozen hand-written examples miss, especially around empty histories and sources that join after a reconnect cursor was created.

Observe four quantities separately: durable append latency, fan-out latency, reconnect backfill latency, and cursor lag. A single end-to-end latency percentile cannot tell storage pressure from a congested subscriber. Add counters for deduplicated publishes, sequence gaps, forced resyncs, and compaction watermarks. Log event IDs and cursors, not full caption text by default, because support transcripts can contain customer data and broad diagnostic logging creates a second retention system by accident.

Deployment deserves one ugly test: terminate a publisher after the durable append but before it receives the acknowledgement. Its retry must return the same logical event and produce one transcript result. Then terminate a reconnect worker between backfill and live delivery. The client should reconnect with the last contiguous cursor and converge again.

No special recovery token is needed.

Cost verification belongs in the same test plan. Measure stored bytes per finalized minute, records written per spoken minute, index bytes, write operations, compaction reads and writes, and reconnect egress at several retention windows. Your mileage may vary because speech activity, interim frequency, language, framing overhead, and replica count all move the result. The decision rule is still crisp: retain data because it supports replay, audit, or an explicit product requirement; compact it when granularity no longer serves those purposes; delete it only after the compacted high-water mark is readable and policy permits deletion.

References

Further reading

Top comments (0)