Skip to content

feat(api): WebSocket /v1/stream live tick fan-out - #58

Merged
obchain merged 2 commits into
mainfrom
feat/24-ws-stream
May 26, 2026
Merged

obchain merged 2 commits into
mainfrom
feat/24-ws-stream

Conversation

@obchain

@obchain obchain commented May 26, 2026

Copy link
Copy Markdown
Owner

Summary

Last API-tier endpoint per PRD §6 / issue #24. Browsers + dashboards get every BVOL/EVOL tick the engine publishes (60s cadence) within sub-second of the Redis PUBLISH.

After this lands, M1 backend surface is feature-complete. Remaining M1 work is frontend (#25–#27), Python parity (#21), CI (#28).

Pipeline

Engine → Redis PUBLISH index:{id}:stream ─┐
                                          ▼
                                       Hub.Run  (one PSUBSCRIBE for all indices)
                                          │
                                          ▼
                            fan-out → registered Conn.send channels
                                          │
                                          ▼
                              gorilla WS frames → browser

One JSON decode per pubsub message rather than once per recipient.

Wire format (PRD §6, line 988)

Client → server:

{ "action": "subscribe", "channels": ["bvol", "evol"] }

Server → client (per tick):

{ "type": "tick", "channel": "bvol",
  "value": 37.37, "ts": 1779782592444, "confidence": 1.0 }

ts is Unix epoch milliseconds (integer) — different from the REST endpoints' RFC 3339 string. PRD lines 1054–1056 show the frontend doing tick.ts / 1000 for lightweight-charts' seconds-since-epoch input. The hub converts engine RFC 3339 → ms int once per pubsub message.

Error shape (sent to client on bad input; connection stays open):

{ "type": "error", "code": "bad_request", "message": "..." }

Module additions

File Responsibility
internal/stream/hub.go one PSUBSCRIBE index:*:stream per process; per-tick decode + fan-out
internal/stream/conn.go per-conn read/write goroutines, gorilla upgrade, bounded send queue, ping/pong keepalive
internal/stream/iplimit.go per-IP active-connection cap (PRD anon: 5) + X-Forwarded-For aware client IP
cmd/api/main.go app.Get("/v1/stream", adaptor.HTTPHandler(...)) via fiber v3 adaptor middleware

fiber v3 ships no native WS package; the adaptor middleware bridges to a plain net/http handler that runs the gorilla upgrade.

Reliability posture

  • Slow consumer = drop-newest. Send queue is 32 frames per conn; on a full queue the newest frame is dropped and the connection stays open. Slow client cannot create head-of-line blocking across other subscribers.
  • Keepalive. 30 s ping, 10 s pong timeout (PingInterval + PongTimeout is the longest a healthy conn can stay silent before the server closes).
  • Per-IP cap. 5 anon active conns per IP via IPLimiter. Acquire / Release brackets the lifetime of the goroutine pair.
  • Origin check permissive at v1. Read-only API, no cookie auth, no CSRF surface. Tightens when paid-tier auth ships in M3.

Deps

  • github.com/gorilla/websocket v1.5.3 (PRD tech_stack pin)

Env vars

Var Default Use
WS_MAX_CONNS_PER_IP 5 anon per-IP cap (PRD §6 anon: 5; paid keys skip this in M3)

Acceptance smoke (live, Deribit → ingestion → engine → API → Python WS client)

subscribed; waiting for ticks...
recv: {"type":"tick","channel":"bvol","value":37.37,"ts":1779782592444,"confidence":1}
recv: {"type":"tick","channel":"evol","value":50.36,"ts":1779782592444,"confidence":1}

Two ticks fanned end-to-end on a 65-second window with PRD wire shape exact. Channel names lowercase, ts ms int, confidence float, type:"tick" envelope.

What's not in scope (deferred)

  • unsubscribe action (clients open a new conn for now; cheap to add when frontend asks)
  • Per-channel ack on successful subscribe (frontend just waits for first tick)
  • volx_api_ws_* Prometheus counters (lands with engine-exporter symmetry PR)
  • Auth / API-key path (M3)

Test plan

  • go build ./... clean
  • go vet ./... clean
  • Live smoke (above) — two ticks received in PRD shape
  • Per-IP cap enforced (verified by opening 6 conns from localhost; 6th gets HTTP 429)

Fixes #24

Last API-tier endpoint per PRD §6 / issue #24. Browsers + dashboards
get every BVOL/EVOL tick the engine publishes (60s cadence) within
sub-second of the engine's Redis PUBLISH.

Pipeline:

  Engine → Redis PUBLISH index:{id}:stream ─┐
                                            ▼
                                         Hub.Run  (one PSUBSCRIBE for all indices)
                                            │
                                            ▼
                              fan-out → registered Conn.send channels
                                            │
                                            ▼
                                gorilla WS frames → browser

Layer additions:

- `internal/stream/hub.go` — `Hub` keeps a single Redis
  PSUBSCRIBE open and dispatches each parsed tick to subscribed
  connections. JSON decode happens once per pubsub message rather
  than once per recipient.

- `internal/stream/conn.go` — per-conn `Conn` with bounded
  send queue (32 frames), gorilla WS read/write goroutines,
  ping/pong keepalive (30s ping, 10s pong timeout). Drop-newest on
  slow consumer so a wedged client cannot create head-of-line
  blocking across other subscribers.

- `internal/stream/iplimit.go` — `IPLimiter` caps per-IP active
  connections; PRD §6 anon limit pinned at 5 via `cfg.WSMaxConnsPerIP`.
  `clientIP` honors `X-Forwarded-For` for the reverse-proxy
  deployments (Caddy / Cloudflare per PRD §13).

- `cmd/api/main.go` — `app.Get("/v1/stream", adaptor.HTTPHandler(...))`.
  fiber v3 has no native WS package; the `adaptor` middleware
  bridges to a plain `net/http` handler that does the gorilla
  upgrade.

Wire format (PRD §6 pinned):

  client → server:
    { "action": "subscribe", "channels": ["bvol", "evol"] }
  server → client (per tick):
    { "type":"tick", "channel":"bvol",
      "value": 37.37, "ts": 1779782592444, "confidence": 1.0 }

`ts` is Unix epoch milliseconds (integer) — *different* from the
REST endpoints' RFC 3339 string form. PRD lines 1054–1056 show the
frontend doing `tick.ts / 1000` for `lightweight-charts`'
seconds-since-epoch input. The hub converts engine's RFC 3339
timestamps to ms ints once per message.

Error shape:

  { "type":"error", "code":"bad_request", "message":"..." }

Returned on malformed JSON or unknown actions; client connection
stays open.

Deps:
- gorilla/websocket v1.5.3

Env vars:
- WS_MAX_CONNS_PER_IP (default 5; matches PRD §6 anon)

Acceptance smoke (Deribit → ingestion → engine → API → Python WS client):

  $ python -c 'import asyncio, websockets, json; ...'
  subscribed; waiting for ticks...
  recv: {"type":"tick","channel":"bvol","value":37.37,"ts":1779782592444,"confidence":1}
  recv: {"type":"tick","channel":"evol","value":50.36,"ts":1779782592444,"confidence":1}

Two ticks fanned end-to-end with PRD wire shape exact.

After this: M1 backend surface is feature-complete. Remaining M1
work: frontend (#25–#27), Python parity (#21), CI (#28).
@obchain obchain added area:api Go fiber service lang:go Go priority:p0 Blocker, must ship this phase type:feature New functionality labels May 26, 2026
MED-1: dropped the dead `!ok` arm in `writeLoop`'s `case frame, ok
:= <-c.send`. The `send` channel is never closed anywhere — the
only channel `close()`d is `c.closed` via `sync.Once`. The
unreachable branch could mislead future work into thinking
`send` closure was a designed exit signal and trip up an
attempt to add a legitimate close (double-close panic risk).
Added an inline note on the real exit paths.

LOW-1: introduced a separate `WriteTimeout = 10s` constant for
data-frame writes; `PongTimeout` now only covers the keepalive
arm where it semantically belongs. The two budgets can now be
tuned independently — e.g., if a future paid-tier wants a
tighter pong cadence for faster client-failure detection, that
no longer drags the data-frame stall budget with it.

Re-verified live: subscribe → bad-action error frame
round-trip works (`{"code":"bad_request","message":"unknown
action (allowed: subscribe)","type":"error"}`).
@obchain
obchain merged commit 48b2cb0 into main May 26, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:api Go fiber service lang:go Go priority:p0 Blocker, must ship this phase type:feature New functionality

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Go API: WS /v1/stream

1 participant