---
title: "Surface Data Plane — Architecture"
description: "How every surface data request becomes a queued job, how the surface's data worker executes it, and how the result streams back to the browser — with timings and timeout handling."
created: 2026-10-01
updated: 2026-10-01
authors: ThinkingCap
topics: [Research Development]
status: published
canonical: https://console.thinkingcap.com/guest/Research-Development/CapGPT/Surface-Data-Plane
summary: Surface data requests flow through a queue-backed architecture where browsers never access databases directly. Each request becomes a row in a Postgres queue, a dedicated worker claims and executes it, and results stream back via server-sent events. This design isolates the data layer, enables request deduplication across users, and allows independent scaling of each surface's query load.
audio: https://thinkingcap.blob.core.windows.net/guest-home/summaries/2601cb2727dad42501aa0b90f3f0da7b30f8ac35c7f173b28989a372258535e0.mp3
audio_full: https://thinkingcap.blob.core.windows.net/guest-home/summaries/full-721fa1a09d0ea3033200907fab03328ec0ee05366b5b6c8fc01614a4a81fb4f9.mp3
---

# Surface Data Plane

Every Learner-View surface and every console Work Surface serves its data through a
**queue-backed data plane**: the browser never queries a database or calls the LMS
directly. A request becomes a row in `maestro_queues`, a per-surface **data worker**
claims and executes it, and the answer streams back to the browser over SSE.

The goals: keep the data layer out of the edge, make every request idempotent and
coalescible, and let the fleet scale each surface's query load independently.

One shared implementation — `server/lib/dataPlane/` in this repo — is the canonical
source, synced to every surface repo (11 `lv-surface-*` + 13 `tc-surface-*`).

## Architecture at a glance

```mermaid
flowchart LR
    subgraph Browser
        A[Surface component<br/>dataClient.query]
    end

    subgraph Edge["Surface BFF (web service)"]
        B["POST /api/&lt;surface&gt;/q<br/>verify JWT · build envelope"]
        C[Reply pump<br/>LISTEN + SSE]
    end

    subgraph QueueDB["MaestroQueues PG (own database)"]
        D[(maestro_queues)]
        E[(surface_result_blobs<br/>bytea · 24h TTL)]
    end

    subgraph Worker["Data worker (1 container per surface, RO cred only)"]
        F[claimOne<br/>SKIP LOCKED + lease]
        G[runOperation<br/>gates re-enforced]
        H[completeJob<br/>UPDATE + pg_notify in one tx]
    end

    I[(Client DB<br/>read-only)]
    J[capcom<br/>/worker/job-finished]

    A -->|POST operation+params| B
    B -->|INSERT · dedupe_key coalesces| D
    B -->|202 jobId + streamUrl| A
    A -->|SSE /q/:jobId/stream| C
    F -->|poll 100ms→2s| D
    F --> G -->|SQL (RO)| I
    G --> H
    H -->|result row| D
    H -->|pg_notify| C
    H -.->|big payloads| E
    C -->|SSE done| A
    A -.->|GET /q/:jobId/blob| C
    H -.->|wall · cpu · RSS · queue wait| J
```

The queue lives in its **own Postgres database** (MaestroQueues) so queue load never
touches application databases. The worker holds **only the read-only client
credential** — write-ish paths stay in the BFF; the two never share a process.

## The lifecycle, step by step

```mermaid
sequenceDiagram
    autonumber
    participant UI as Browser (dataClient)
    participant BFF as Surface BFF
    participant Q as maestro_queues
    participant W as Data worker
    participant DB as Client DB (RO)

    UI->>BFF: POST /api/<surface>/q {operation, params}
    BFF->>BFF: verify JWT · envelope from token claims<br/>dedupeKey = sha256(surface|op|ver|client|scopeKey|params)
    BFF->>Q: INSERT ... ON CONFLICT (dedupe_key) DO NOTHING
    alt identical job already live
        Q-->>BFF: existing jobId (coalesced)
    else new job
        Q-->>BFF: new jobId
    end
    BFF-->>UI: 202 {jobId, streamUrl}  (~ms)

    UI->>BFF: GET /q/:jobId/stream (SSE)
    BFF->>Q: immediate row check (done-before-LISTEN race)
    BFF-->>UI: event: queued

    loop poll: 100ms idle backoff → 2s max
        W->>Q: claimOne() — UPDATE ... FOR UPDATE SKIP LOCKED
    end
    Q-->>W: job (status=in_progress, lease=120s, dequeue_count+1)

    alt TTL cache hit (reference lists / search only)
        W-->>Q: completeJob(cache: 'hit')  (~instant)
    else execute
        W->>DB: runOperation (single-flight, gates re-enforced)
        DB-->>W: rows
        opt PDF / big payload
            W->>Q: putBlob → surface_result_blobs (24h)
        end
        W->>Q: completeJob — UPDATE done + pg_notify, ONE transaction
    end

    Q-->>BFF: pg_notify 'surface_<id>_done' (jobId)
    BFF->>Q: SELECT result row
    BFF-->>UI: event: done {result}
    UI->>UI: render
    W-->>J: telemetry (fire-and-forget) — wall · cpu · RSS · queue wait
```

### 1 · Enqueue (browser → BFF → queue)

- Component calls `dataClient.query(operation, params)` → `POST /api/<surface>/q`.
- BFF builds the **envelope from verified token claims** — never the request body:
  operation, resolved client, actorEmail, staff flag, signed branch scope.
- `dedupeKey` = sha256 over surface · operation · version · client · **scopeKey** ·
  canonical(params). scopeKey hashes the operator's signed scope (plus the staff
  flag), so two admins with different branch scopes never share a job or a cached
  result — the cache-safety rule.
- `INSERT ... ON CONFLICT (dedupe_key) WHERE status IN ('pending','in_progress') DO
  NOTHING`: an identical live job **absorbs the request** and the caller attaches to
  the same row. Ten admins opening the same panel = one query executed.
- Response: `202 {jobId, streamUrl}` — one INSERT, milliseconds.

### 2 · Claim (worker)

- One worker container per surface type, claiming `type='surface-<id>'` rows in a
  poll loop (100ms idle backoff doubling to 2s).
- The claim is a single atomic statement:
  `UPDATE ... WHERE id = (SELECT ... FOR UPDATE SKIP LOCKED LIMIT 1)` — oldest
  claimable row (`pending`, **or `in_progress` with a lapsed lease**), setting
  `locked_by`, `lease_expires_at = now() + 120s`, `dequeue_count + 1`.

### 3 · Execute

- **TTL cache check** (reference lists and search only — operations with
  `ttlSeconds > 0`): a hit completes the job instantly with `cache: 'hit'`.
- **Single-flight**: concurrent executions of the same dedupeKey inside the process
  share one promise — the second caller is reported `coalesced`.
- `runOperation(operation, params, authCtx)` re-enforces **every gate** from the
  envelope's authContext worker-side, then runs the handler's SQL over the RO
  credential.
- PDFs and other large payloads go to `surface_result_blobs` (bytea, 24h expiry);
  the result carries a `blobRef` instead of the payload.

### 4 · Complete

- `completeJob` commits `status='done' + result` **and** `pg_notify(channel, jobId)`
  **in one transaction** — NOTIFY fires on commit, so a waiter can never be woken
  before the row is readable.
- Telemetry (wall clock, process CPU, RSS at finish, queue wait) is posted to
  capcom's `/worker/job-finished` fire-and-forget — it must never break the job path.

### 5 · Reply (BFF → browser)

- A dedicated LISTEN connection (never the pool) subscribes to the surface's done
  channel. On notify it matches the jobId against this instance's waiter map,
  SELECTs the row, and sends `event: done` over SSE.
- `GET /q/:jobId` is a plain poll fallback; `GET /q/:jobId/blob` serves claim-check
  downloads. The result lives on the row, so a dropped SSE reconnects without loss.

## Coalescing — three independent layers

```mermaid
flowchart TD
    R[Incoming identical requests] --> L1
    L1{1 · Transport<br/>dedupe_key unique<br/>on live rows}
    L1 -->|conflict| ATTACH[Attach to the in-flight job<br/>coalesced: true]
    L1 -->|no conflict| L2{2 · Single-flight<br/>in-process promise<br/>per dedupeKey}
    L2 -->|in flight| SHARE[Share the running promise]
    L2 -->|idle| L3{3 · TTL cache<br/>ttlSeconds &gt; 0 ops only}
    L3 -->|hit| HIT[Complete instantly<br/>cache: 'hit']
    L3 -->|miss| EXEC[runOperation → complete<br/>then store under TTL]
```

1. **Transport** — `dedupe_key` partial unique index on live rows: cross-process,
   cross-instance coalescing at enqueue time.
2. **Single-flight** — one in-flight promise per dedupeKey inside the worker:
   covers the lease-lapse double-claim window.
3. **TTL cache** — only for operations the manifest marks cacheable. Drill-down and
   learner-PII operations are `ttlSeconds = 0`: **always fresh**, never stored
   (Douglas's rule: lists may be stale, user drill-down always fresh).

Cache entries carry tags (`c:<client>`, `u:<userId>`) for O(1) invalidation, capped
at 500 entries with oldest-first eviction.

## Timing — measured live

From `worker_runs` telemetry, last 6 hours on 2026-10-01:

| Step | Latency |
|---|---|
| Enqueue → 202 | milliseconds (one INSERT) |
| Queue wait (enqueue → claim) | **avg 0.5 – 1.4s** (poll backoff + claim) |
| Execution (claim → done) | **avg 0.08 – 0.6s**, max observed 1.2s |
| Notify → SSE event | milliseconds (push; 5s sweep worst case) |
| **End-to-end, created → done** | **~1–2s typical** — accounting 0.9s, users 1.2s, learnerviews 1.2s; skills 11.9s avg (heavier queries) |
| Cache hit | near-instant (skips execution) |

Traffic is sparse (tens of jobs/hour across the fleet), so queue wait is dominated
by the worker's idle poll backoff, not contention.

## Job states

```mermaid
stateDiagram-v2
    [*] --> pending : enqueue
    pending --> in_progress : claimOne\n(SKIP LOCKED, lease 120s)
    in_progress --> done : completeJob\n(result + pg_notify, one tx)
    in_progress --> pending : failJob\n(30s backoff, dequeue_count+1)
    in_progress --> pending : lease lapses\n(worker died — that IS the retry)
    pending --> poison : dequeue_count ≥ max_dequeues
    in_progress --> done : HttpError 400/403/404\n(completes with error payload —\nan answer, not a failure)
    done --> [*]
    poison --> [*] : surfaced on capcom cards
```

## Timeouts and failure handling

```mermaid
flowchart TD
    subgraph WorkerSide["Worker side"]
        CRASH[Worker dies — Spot churn] --> LEASE[Lease stops renewing<br/>heartbeat every 60s]
        LEASE --> LAPSE[lease_expires_at passes<br/>≤ 120s]
        LAPSE --> RECLAIM[Row claimable again<br/>re-driven by any worker]
        ERR[Unexpected error] --> BACKOFF[failJob → pending<br/>30s backoff]
        BACKOFF --> COUNT{dequeue_count<br/>≥ max_dequeues?}
        COUNT -->|no| RECLAIM
        COUNT -->|yes| POISON[poison — parked, surfaced]
        EXPECTED[HttpError 400/403/404] --> ANSWER[completeJob with error<br/>never retried]
    end

    subgraph BrowserSide["Browser side"]
        T60[SSE waiter deadline 60s] --> SWEEP[5s sweep re-reads rows<br/>covers missed NOTIFY]
        SWEEP -->|past deadline| TIMEOUT[event: error TIMEOUT<br/>job keeps running]
        TIMEOUT --> ROW[result stays on the row<br/>poll / reconnect can still fetch]
        DISC[browser disconnect] --> DROP[waiter dropped · job unaffected]
        LDROP[LISTEN connection drops] --> RECON[reconnect 2s / 5s<br/>sweep fills the gap]
    end
```

| Constant | Value | Where |
|---|---|---|
| Lease | 120s (`MAESTRO_LEASE_SECONDS`) | claimOne |
| Heartbeat | 60s (`MAESTRO_RENEW_SECONDS`) | renewHeld loop |
| Failure backoff | 30s (`available_at`) | failJob |
| Poison cutoff | `max_dequeues` (per queue row) | claim filter |
| SSE deadline | 60s per waiter | registerWaiter |
| Sweep cadence | 5s | replyPump |
| Poll backoff | 100ms → 2s max (`DATAPLANE_IDLE_MS`) | worker loop |
| Blob retention | 24h, reaped hourly at max idle | surface_result_blobs |
| Cache cap | 500 entries, oldest-first | cache.ts |

Key properties:

- **Crash = retry, not loss.** Stop heartbeating and the lease lapses; the row
  becomes claimable and is re-driven. No dead-letter machinery is needed for
  worker death — the lease *is* the retry.
- **Errors split two ways.** Expected client errors (400/403/404 `HttpError`)
  *complete* the job with an error payload — an answer, not a failure, and never a
  path to poison. Unexpected errors release the row and climb `dequeue_count`
  toward the poison cutoff.
- **A browser timeout never kills work.** The 60s SSE deadline only synthesizes a
  `TIMEOUT` event for the waiting component; the job finishes regardless and the
  result remains fetchable.
- **Telemetry is fire-and-forget.** A capcom outage degrades charts, never jobs.

## Fleet shape

- One `worker_types` pair (live + dev) per surface: queue types `surface-<id>` /
  `surface-<id>-dev`, command `node dist/dataWorker.js`, Maestro-managed VMSS
  placement (Spot — lease semantics are what make that safe), typically one
  container per type.
- Dev placements drain the stable MaestroQueues DB; `SURFACE_QUEUE_TYPE` and
  `MAESTRO_QUEUES_URL` select tier at runtime — the image is identical.
- Required env: `MAESTRO_QUEUES_URL` (unset → BFF 503s the plane and the worker
  refuses to boot), `SURFACE_QUEUE_TYPE`, app/patch DB URLs, JWT secret, the RO
  client-DB credential — **never** a writer credential.
- Cards and charts on the Kebb wallboard read `maestro_queues` directly plus the
  per-job telemetry above (queue wait, wall, CPU, RSS).

## Code map

| Concern | File (per surface repo) |
|---|---|
| Envelope, dedupeKey, scopeKey | `server/lib/dataPlane/envelope.ts` (LV) · `worker/src/lib/dataPlane/envelope.ts` (console) |
| Producer + consumer queue client | `…/dataPlane/queue.ts` |
| Reply pump (LISTEN + SSE waiters) | `…/dataPlane/replyPump.ts` |
| TTL cache + single-flight | `…/dataPlane/cache.ts` |
| Worker entrypoint | `worker/src/dataWorker.ts` (boot via `dataWorkerEntry.ts`) |
| HTTP front door | `worker/src/routes/dataPlane.ts` |
| Telemetry | `…/dataPlane/telemetry.ts` |
| Queue schema | capcom `apps/api/db/maestroqueues-schema.sql` |
