Skip to main content

Listen + Pusher Pipeline — Sequence Diagrams

Last updated: 2026-07-26 (bounded Pusher delivery, validated frames, and durable finalization ownership) These diagrams document the real behavior observed during E2E testing with live services (backend, pusher, STT providers, embedding API). Update when the pipeline changes.
The backend-to-Pusher binary protocol accepts only opcodes 100 through 105. Audio frames (101) require the four-byte opcode, an eight-byte timestamp, and PCM payload; JSON opcodes require an object payload, with transcript, finalization, and speaker-sample fields validated before they enter background queues. Malformed or unknown frames close the session with a protocol error instead of reaching application processing. Transcript, speaker-sample, audio-webhook, and private-cloud queues are bounded. Buffered audio also shares a per-session byte budget, so variable-size payloads cannot bypass item-count limits. A failed or cancelled backend-to-Pusher send retains the pending item for reconnect replay; buffers are cleared only after the send completes. Finalization opcode 104 carries the durable job id and dispatch generation, and Pusher claims that Firestore-owned job before processing.

1. Connection + Streaming + Transcription

The STT provider is selected at session start via STT_SERVICE_MODELS env var (e.g., modulate-velma-2,parakeet). The first provider whose deployed model supports the requested mode wins: Parakeet’s live RNNT model is English-only, so multilingual sessions and non-English single-language sessions use Modulate. Listen startup loads transcription preferences once and uses its multi-language setting to request Modulate auto-detection, avoiding a second user-preference read on WebSocket startup. The separately deployed Parakeet batch model retains its multilingual capability. A client Parakeet preference can prioritize Parakeet only when that live-model capability matches; it never overrides the selected service/language/model tuple. All providers implement the STTSocket ABC and are wrapped by the universal GatedSTTSocket VAD gate. Segments buffered in realtime_segment_buffers are sorted by start time before processing — this corrects non-deterministic WebSocket arrival order from providers like Modulate whose internal parallelism can deliver shorter utterances before longer ones. Every /v4/listen stream receives a durable, user-scoped recording_session resource before it creates or resumes a conversation. A client-generated client_conversation_id remains that resource’s ID when present; legacy clients receive a server-generated UUID. The resource atomically binds the recording to exactly one conversation, so reconnects return the original canonical mapping rather than falling back to a user-global Redis pointer. The backend emits a conversation_session event carrying that recording identity:
conversation_session, memory_processing_started, and memory_created carry the same versioned envelope. Clients must accept an envelope only when the recording and conversation IDs match their local session and the sequence is higher than the last accepted sequence. The phase is part of the contract: in_progress, then processing, then a terminal completed/failed/ discarded. Older events without the envelope retain the legacy identity and timestamp compatibility path only. RECORDING_SESSION_MODE=dual_write is the safe default: it persists and compares the durable binding while preserving the legacy proposed route on a mismatch. shadow has the same routing behavior for validation-only rollout; enforce promotes the canonical durable binding and rejects durable-store failures. Mapping conflicts, stale-event discards, and legacy fallback paths use shared fallback telemetry for the cutover gate. The resource contains only IDs, phase, sequence, schema version, and timestamps; it never stores transcript content, credentials, or raw WebSocket payloads. The same recording ID is always a retry of its original generation: after a terminal result the backend returns that canonical binding and its terminal envelope instead of rebinding it. A silence or status rollover on a live WebSocket receives a fresh server-generated recording ID and can therefore bind a new conversation without mutating the terminal generation. If an empty-generation cleanup wins the Firestore transaction race against buffered STT segments or photos, the late write is fenced rather than recreating the terminal parent; the listen handler opens one fresh generation and replays that buffer once. The first-segment started_at value commits in that same content transaction, so this recovery cannot leave a timestamp-only write behind.

Capture-device provenance

Every capture client with WebSocket-upgrade header support sends X-App-Platform and X-Device-Id-Hash on /v4/listen. The backend records that provenance on the conversation stub, allowing extracted canonical memories to participate in the This device filter. Browsers cannot attach custom upgrade headers, so /v4/web/listen carries the hash in its required first auth message; the backend fixes its platform to web before creating the conversation. Missing provenance remains unknown rather than being assigned to a different device.

Live STT observability

For backend-provider /v4/listen sessions, omi_live_stt_accepted_total increments when the listener receives its first nontrivial audio frame. A successful terminal is emitted only after a nonempty transcript has been delivered to the client WebSocket. Provider failures finish the same attempt through terminate_live_stt_session, preserving their bounded failure phase; other attempts that end before delivery are classified during teardown as failure or cancelled. These metrics use only the closed provider, client platform, deployment environment, outcome, and phase labels. They do not include transcript content, user IDs, model revisions, or other high-cardinality request data. Custom-STT sessions are excluded because the backend does not own their provider socket or provider attempt lifecycle.

Dev-only parity-pack capture

For local conformance debugging only, the listen runtime may write a restricted local cassette. It starts after listen has selected its STT service and accepts only an exact Firebase UID from OMI_PARITY_PACK_ALLOWED_PRINCIPALS. Capture is off unless all of OMI_ENV_STAGE=dev, OMI_PARITY_PACK_CAPTURE=1, and an absolute, outside-repository OMI_PARITY_PACK_ROOT are present. No Helm or production setting enables it. When enabled, the receiver records decoded client PCM (including each multi-channel PCM stream), successfully forwarded STT audio, and provider transcript callbacks. The runtime persists the anonymous cassette during teardown. Events and raw audio retained in memory are bounded (1,000 events or 8 MiB of audio); once either limit is reached, remaining capture events are dropped. Capture initialization, observation, and persistence errors are isolated from the listen path.

Terminal STT provider failure

Provider initialization failure, a latched dead provider socket, or an audio send failure is terminal for the current client WebSocket. The backend first awaits a bounded service_status event (status=stt_failed, semantic outcome, provider, retryability, and one of the documented reason codes), then closes the client with code 1011 and reason transcription_service_unavailable. The send boundary clears its local audio buffer only after the provider socket accepts the chunk. On failure, the bytes remain in that buffer until normal session teardown frees it; they are never cleared while the client connection stays open and appears healthy. The client close is the recovery signal for mobile and desktop reconnect/fallback logic. The backend deliberately does not replace a provider socket in-session because maintaining transcript timestamp continuity across that handoff requires a separate state model.

2. Conversation Lifecycle (Silence Timeout Path)

This is the normal path when the client stays connected but stops speaking. Key design rule (#6061): Listen NEVER processes conversations locally. All conversation processing routes through pusher or the dedicated durable finalizer worker. Before either handoff, listen persists a Firestore finalization job. The legacy live pusher path claims that same job and transfers it to pusher process-scoped ownership, so closing the originating listener after the opcode-104 handoff cannot cancel it. When the durable dispatch flag is enabled, Cloud Tasks wakes the worker using only the opaque job id and generation. pending_conversation_requests remains a low-latency reconnect aid, not the durable source of truth. Before a job can become completed, the finalizer atomically claims a durable fanout boundary with a persistent idempotency key. It passes that key to the integration endpoint and records fanout_status=completed only after the call returns. Lease recovery retries the same key, so a crash after an integration call cannot silently drop fanout or create an unkeyed duplicate.

2.1 Durable finalization ownership

conversation_finalization_jobs/{job_id} is the finalization ledger. It is created atomically with the in_progress → processing transition and is keyed by (uid, conversation_id, finalization_revision). A task is a named, at-least-once wake-up only; duplicate delivery, reconnect replay, and manual reconciliation must first acquire the Firestore lease. The job and task payload never contain transcript text, raw BYOK keys, headers, or raw exception text. blocked_byok is an explicit recovery state: durable offline BYOK is not supported without a separately reviewed credential broker; the worker must never silently use Omi keys. See Listen finalization jobs for configuration, replay, metrics, and the operational runbook.

3. Disconnect Path

What happens when the WS connection closes. Teardown invariant: flush any remaining multi-channel mix via audio_bytes_send before pusher_close() so ListenPusherSession._flush() can deliver tail audio to pusher.

3.1 Pusher Reconnect & Pending Flush

When pusher reconnects after a disconnection, all buffered conversations are replayed.

4. Speaker ID Lifecycle (2-Session Flow)

Speaker identification requires two sessions: one to store the embedding, one to match against it. The user’s own voice is also identified via their speech profile embedding (loaded at session start alongside person embeddings).

5. Private Cloud Sync (Audio Upload)

When private_cloud_sync_enabled is set for the user.

Conversation playback artifact

At conversation completion (process_conversation, and _finalize_sync_audio_files for offline sync), an audio-merge Cloud Task (schema_version 2, name amc-{conversation_id}-{fingerprint}) builds one dense per-conversation MP3playback/{uid}/{conversation_id}/conversation.mp3 — containing only captured audio: intra-part gaps (<90s) stay silence-filled, inter-part gaps (>90s) are collapsed. The handler stamps the conversation doc with conversation_audio: a spans manifest ({file_id, wall_offset, artifact_offset, len} per part, wall offsets relative to started_at — the TranscriptSegment.start basis), both durations (duration = wall clock, captured_duration = audio only), and the audio_files_fingerprint it was built from. Staleness is fingerprint-driven: when late chunks change audio_files (pusher batch flush after completion, offline sync, conversation merge), the stamped fingerprint no longer matches and both the call sites and the /v1/sync/audio/{id}/urls poll re-enqueue a rebuild under a new task name (defeating the named-task tombstone). /urls returns the artifact as a top-level conversation_audio object alongside the per-part audio_files list, which remains for older app versions.

6. Event Wire Protocol

Server → Client (JSON over WS text frames)

Event Types

Client → Server

7. Timing Constants

8. WebSocket Task Supervision

routers/listen/runtime.py and pusher.py use an asyncio.wait(FIRST_COMPLETED) supervisor loop instead of asyncio.gather() to manage background tasks. This prevents ghost connections where a hung background task blocks cleanup forever. Supervisor exits on:
  • Client disconnect (receive task completes)
  • Background task crash (exception)
  • Lifetime task normal completion (e.g. heartbeat inactivity timeout)
Supervisor re-waits on:
  • Finite task normal completion (e.g. process_pending_conversations, speaker_identification_task)
After supervisor exit, remaining tasks drain with BG_DRAIN_TIMEOUT (30s) before force-cancel. The connection gauge (inc/dec) is always paired in try/finally.