Repository navigation
feature(Stream): Add a live WebSocket stream command - #26
Conversation
`aperiodic stream <dataset>` subscribes to wss://stream.aperiodic.io/v1/stream and prints one JSON line per row, with reconnects, --count and --duration. CI gains a stream_url dispatch input for the live tests. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
szemyd
left a comment
There was a problem hiding this comment.
Adversarial review: aperiodic stream
Reviewed at e8ec1f8. Unit tests pass (go test -race -count=3 -run Stream -skip Live). Findings 2, 3 and 4 were reproduced with throwaway tests against this PR's own fakeStream; the probe code is under finding 4. The companion review on aperiodic-io/client#38 covers the Python side. Findings 1 and 5 apply to both PRs, so fix them the same way in both.
🔴 = should fix before merge, 🟡 = optional.
🔴 1. A 429 during a reconnect ends the stream for good
stream.go:172 (return handshake.StatusCode < 500) treats every 4xx as final, including on a reconnect. After an abnormal drop (1006), or after the client's own 75 s idle timeout, the server probably still counts the old half-open socket toward maxConnections. The reconnect then gets a 429 and the CLI exits 1. On a plan with one connection, one network blip is then likely to end the stream for good.
Fix: keep 401, 403 and 426 final. Retry 429 with the normal backoff once the session has been subscribed at least once, or for a bounded window longer than the server's dead-socket timeout. The first connect can still fail fast on 429. Add a test: ack, drop, then 429 for the next N handshakes, then accept; it should resume.
🔴 2. A server unsubscribed is ignored, so the CLI hangs
stream.go:210-252 has no case for op:"unsubscribed". When the server removes every channel (e.g. a plan downgrade), heartbeats keep the idle timer alive and the CLI prints nothing forever. With --duration it then exits 0.
Repro: ack, then {"op":"unsubscribed","channels":["ohlcv.binance-futures.1m"],"rejected":[{"channel":"ohlcv.binance-futures.1m","code":"not_entitled","message":"plan downgraded"}]}, then heartbeats. Result: code=0, stdout empty, stderr only Subscribed: ....
Fix: handle unsubscribed:
- Print the removed and rejected channels to stderr.
- Track which granted channels are left.
- When none are left, return
&StreamRejectedError{...}so the CLI exits 1. - Add tests for partial and total removal.
The Python client already does this (_unsubscribed).
🔴 3. A subscribe refused with op:"error" hangs the same way
stream.go:249 treats every op:"error" as a warning. If the error carries the subscribe's id ("s1"), the subscription was refused. The CLI keeps waiting on heartbeats forever, then exits 0 under --duration.
Repro: reply to the subscribe with {"op":"error","id":"s1","code":"invalid_message","message":"bad subscribe"} and then heartbeats. Result: code=0, stderr Warning: server error invalid_message: bad subscribe.
Fix: when frame.Op == "error" && frame.ID == "s1" (send the id from a constant), close and return a final error, so the CLI exits 1 with the message. Python does this in _subscribe.
🔴 4. Unknown close codes reconnect in a tight loop
stream.go:217 makes only 4001 and 1008 final; any other close reconnects, including 1000 and every other 4xxx. stream.go:150 also resets attempt after every ack. So a server that acks and then closes (e.g. a future 4003 subscription revoked) gets a reconnect every 250–500 ms forever. In the repro (test backoff of 10–50 ms) that was 110 connections in 1 s.
Fix:
- (a) Treat any close in 4000–4999 as final (
StreamClosedError). - (b) Reset the backoff only after the connection has stayed up for a while (e.g. ≥ 60 s, or after the first data row), not on the ack alone.
Probe used for 2–4 (drop into the package as a _test.go file):
func heartbeatForever(ctx context.Context, conn *websocket.Conn) {
for {
select {
case <-ctx.Done():
return
case <-time.After(100 * time.Millisecond):
send(ctx, conn, `{"op":"heartbeat"}`)
}
}
}
// 2: ack(ctx, conn, sub, ""); send(... unsubscribed ...); heartbeatForever(ctx, conn)
// 3: send(ctx, conn, `{"op":"error","id":"`+sub.ID+`","code":"invalid_message","message":"bad subscribe"}`); heartbeatForever(ctx, conn)
// 4: ack(ctx, conn, sub, ""); _ = conn.Close(4003, "subscription revoked") -> count f.connections() after runCLI("stream","ohlcv","--duration","1s")🔴 5. The CI key rule differs from client#38
.github/workflows/ci.yml:61 picks the staging key only when base_url is set. client#38 picks it when base_url || stream_url is set. If unravel-router dispatches stream_url without base_url, this repo tests the staging stream with the production key, and the live tests get 401. Agree one rule for both repos. The simplest is for router#937 to always send both URLs together, and for both workflows to key off base_url only.
🟡 6. The first connect never fails fast
If the first dial never gets an HTTP response (DNS failure, TLS failure, a mistyped APERIODIC_STREAM_URL, a proxy refusing CONNECT), the CLI retries forever and prints "Connection lost" for a connection it never had. I hit this here: a sandbox proxy refused stream.aperiodic.io. TestCLI_Stream_Live_InvalidKeyIsRefused retried 7 times, then exited 0 at --duration 30s. Suggest: if the first session never dialled successfully, return the error (exit 1). Python fails fast on the first connect.
🟡 7. Nits
cli_stream.go:117: SIGTERM exits 130. The convention is 143 (128+15). Either check which signal arrived, or document that both exit 130.cli_stream.go:86:json.Marshal(row)re-escapes<,>and&insidedataas<etc., so "passed through exactly as the server sent it" isn't byte-exact. Either write the raw bytes yourself, or use anEncoderwithSetEscapeHTML(false), or soften the wording.
Generated by Claude Code
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
All 7 findings are fixed in 8b65f1e. Every behaviour change has a test, and each new test failed before its fix.
The README and 🤖 Generated with Claude Code |
What
aperiodic stream <dataset>streams the live WebSocket feed (wss://stream.aperiodic.io/v1/stream) as JSON lines.{"channel":"ohlcv.binance-futures.1m","snapshot":false,"data":{...}}.datais passed through exactly as the server sent it. Snapshot rows have"snapshot":true.op:errorwarnings and reconnects. Heartbeats are not printed.--symbolsis a comma-separated list. If you leave it out, the subscribe message has nosymbolsfield, so you get every symbol your plan allows.--exchangedefaults tobinance-futuresand--intervalto1m. The dataset and exchange are not checked locally, so the server'sunknown_channelrejection is what you see.--count Nstops after N live rows. Snapshot rows are printed but not counted.--durationstops after a set time. Both send a normal close (1000) and exit 0. Ctrl-C or SIGTERM also closes cleanly, and exits 130.How
stream.go: aStreamerthat runs one session per connection, plus the reconnect loop and typed terminal errors. It usesgithub.com/coder/websocketv1.8.15, the only new dependency, which has no dependencies of its own.X-API-KEYheader.CF-Access-Client-Id/Secretare sent when bothCF_ACCESS_CLIENT_IDandCF_ACCESS_CLIENT_SECRETare set, the same way the REST client does it.APERIODIC_STREAM_URLoverrides the endpoint. There is no--urlflag, because the CLI sets base URLs only through env vars.cli_stream.go: flags, output and help.parseRawArgsis now a thin wrapper over a genericparseArgs, so flags can sit before or after the dataset, the same asraw.Tests
stream_test.go(written first; it failed to build before the implementation): anhttptestWebSocket fake covers:X-API-KEYheader, and that the key is not in the URL--durationstopping with a normal close-race.stream_live_test.go, gated the same way asraw_live_test.go(useLiveStreamdefaultsAPERIODIC_STREAM_URLto production;requireAPIKeywhere a key is needed):ohlcv.binance-futures.1mforperpetual-BTC-USDT:USDTis granted, and a non-snapshot row arrives within 150 s (two minute boundaries, so one missed minute on staging does not block a release)go vetandgofmtall pass. The bad-key live test passed against production. The two key-gated live tests were not run locally because no key was in the environment, so CI is their first run.CI
There is a new optional
workflow_dispatchinput,stream_url(the live-stream WebSocket URL; empty means production). It is passed to the test step asAPERIODIC_STREAM_URL, next tobase_url. After this merges, unravel-router'sintegration-tests.ymlwill pass it with the staging stream URL.The CI API keys (
APERIODIC_API_KEY, andAPERIODIC_STAGING_API_KEYon dispatch) must be on a plan with liveohlcv.binance-futures.1m, orTestCLI_Stream_Live_OHLCVRowArrivesfails.🤖 Generated with Claude Code