feat(api): WebSocket /v1/stream live tick fan-out - #58
Merged
Merged
Conversation
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).
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"}`).
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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
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 }tsis Unix epoch milliseconds (integer) — different from the REST endpoints' RFC 3339 string. PRD lines 1054–1056 show the frontend doingtick.ts / 1000forlightweight-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
internal/stream/hub.goPSUBSCRIBE index:*:streamper process; per-tick decode + fan-outinternal/stream/conn.gointernal/stream/iplimit.goX-Forwarded-Foraware client IPcmd/api/main.goapp.Get("/v1/stream", adaptor.HTTPHandler(...))via fiber v3adaptormiddlewarefiber v3 ships no native WS package; the
adaptormiddleware bridges to a plainnet/httphandler that runs the gorilla upgrade.Reliability posture
PingInterval + PongTimeoutis the longest a healthy conn can stay silent before the server closes).IPLimiter. Acquire / Release brackets the lifetime of the goroutine pair.Deps
github.com/gorilla/websocketv1.5.3 (PRD tech_stack pin)Env vars
WS_MAX_CONNS_PER_IP5Acceptance smoke (live, Deribit → ingestion → engine → API → Python WS client)
Two ticks fanned end-to-end on a 65-second window with PRD wire shape exact. Channel names lowercase,
tsms int,confidencefloat,type:"tick"envelope.What's not in scope (deferred)
unsubscribeaction (clients open a new conn for now; cheap to add when frontend asks)volx_api_ws_*Prometheus counters (lands with engine-exporter symmetry PR)Test plan
go build ./...cleango vet ./...cleanFixes #24