Skip to content

feature(Stream): Add a live WebSocket stream command - #26

Merged
szemyd merged 3 commits into
mainfrom
feat/stream
Oct 7, 2026
Merged

szemyd merged 3 commits into
mainfrom
feat/stream

Conversation

@szemyd

@szemyd szemyd commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

What

aperiodic stream <dataset> streams the live WebSocket feed (wss://stream.aperiodic.io/v1/stream) as JSON lines.

aperiodic stream ohlcv --exchange binance-futures --interval 1m \
  --symbols perpetual-BTC-USDT:USDT,perpetual-ETH-USDT:USDT [--count N] [--duration 90s]
  • stdout: one line per row, {"channel":"ohlcv.binance-futures.1m","snapshot":false,"data":{...}}. data is passed through exactly as the server sent it. Snapshot rows have "snapshot":true.
  • stderr: granted and rejected channels, op:error warnings and reconnects. Heartbeats are not printed.
  • --symbols is a comma-separated list. If you leave it out, the subscribe message has no symbols field, so you get every symbol your plan allows. --exchange defaults to binance-futures and --interval to 1m. The dataset and exchange are not checked locally, so the server's unknown_channel rejection is what you see.
  • --count N stops after N live rows. Snapshot rows are printed but not counted. --duration stops after a set time. Both send a normal close (1000) and exit 0. Ctrl-C or SIGTERM also closes cleanly, and exits 130.
  • Exits 1 with the server's message on handshake 401, 403 or 429 (and any other 4xx), when every channel is rejected, or on close 4001 (plan lapsed or key rotated) or 1008 (rate-limit abuse). None of these is retried.
  • Reconnects after network drops, abnormal or other closes, handshake 5xx, or 75 s with no message (heartbeats arrive about every 30 s). It uses capped exponential backoff with jitter (0.5 s doubling to 30 s, then a random wait between half and all of that) and sends the subscription again. The backoff resets after each successful ack.

How

  • stream.go: a Streamer that runs one session per connection, plus the reconnect loop and typed terminal errors. It uses github.com/coder/websocket v1.8.15, the only new dependency, which has no dependencies of its own.
  • The key goes only in the X-API-KEY header. CF-Access-Client-Id/Secret are sent when both CF_ACCESS_CLIENT_ID and CF_ACCESS_CLIENT_SECRET are set, the same way the REST client does it. APERIODIC_STREAM_URL overrides the endpoint. There is no --url flag, because the CLI sets base URLs only through env vars.
  • cli_stream.go: flags, output and help. parseRawArgs is now a thin wrapper over a generic parseArgs, so flags can sit before or after the dataset, the same as raw.

Tests

  • stream_test.go (written first; it failed to build before the implementation): an httptest WebSocket fake covers:
    • the X-API-KEY header, and that the key is not in the URL
    • the subscribe payload, with and without symbols
    • ack, snapshot, heartbeat, error and data frames
    • handshake 401, 403 and 429: exit 1 with no retry
    • every channel rejected: exit 1
    • partial rejection: a warning, then rows still stream
    • a drop, then reconnect and resubscribe
    • no reconnect after a 4001 or 1008 close
    • --duration stopping with a normal close
    • a missing key or dataset, and help
    • It passes 5 times in a row with -race.
  • stream_live_test.go, gated the same way as raw_live_test.go (useLiveStream defaults APERIODIC_STREAM_URL to production; requireAPIKey where a key is needed):
    • a bad key gets a 401
    • an unknown dataset is rejected
    • ohlcv.binance-futures.1m for perpetual-BTC-USDT:USDT is granted, and a non-snapshot row arrives within 150 s (two minute boundaries, so one missed minute on staging does not block a release)
  • Run locally: the unit tests, go vet and gofmt all 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_dispatch input, stream_url (the live-stream WebSocket URL; empty means production). It is passed to the test step as APERIODIC_STREAM_URL, next to base_url. After this merges, unravel-router's integration-tests.yml will pass it with the staging stream URL.

The CI API keys (APERIODIC_API_KEY, and APERIODIC_STAGING_API_KEY on dispatch) must be on a plan with live ohlcv.binance-futures.1m, or TestCLI_Stream_Live_OHLCVRowArrives fails.

🤖 Generated with Claude Code

szemyd and others added 2 commits October 7, 2026 16:25
`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 szemyd left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 & inside data as < etc., so "passed through exactly as the server sent it" isn't byte-exact. Either write the raw bytes yourself, or use an Encoder with SetEscapeHTML(false), or soften the wording.

Generated by Claude Code

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@szemyd

szemyd commented Oct 7, 2026

Copy link
Copy Markdown
Contributor Author

All 7 findings are fixed in 8b65f1e. Every behaviour change has a test, and each new test failed before its fix. go vet, gofmt and the unit tests pass, and the unit tests passed 8 times in a row with -race. The same rules are being implemented in the Python client.

# Finding Fix Test
1 A 429 on a reconnect is final Handshake 401, 403 and 426 (and any other 4xx) stay final. A 429 is final only before any session has been subscribed; after that it is retried with the normal backoff. 5xx is retried. TestStream_Retries429OnceSubscribed (ack, drop, three 429s, then it resumes), TestStream_HandshakeRefusalsOnTheFirstConnectExitWithoutRetrying (now covers 426)
2 unsubscribed is ignored The removed and rejected channels are printed to stderr, and the client tracks which granted channels remain. When none remain it returns StreamRejectedError and exits 1. TestStream_PartialUnsubscribeKeepsStreaming, TestStream_TotalUnsubscribeExits
3 A subscribe refused with op:"error" hangs An op:"error" whose id equals streamSubscribeID closes the stream and returns StreamSubscribeError, which exits 1 with the code and message. Any other op:"error" is still a warning. TestStream_SubscribeRefusedByAnErrorExits
4a Unknown close codes reconnect Any close code from 4000 to 4999, and 1008, is final. A server close with 1000, 1001, 1011, 1012 or 1013, a network error and the idle timeout all reconnect. TestStream_DoesNotReconnectAfterTerminalCloses (adds 4003 and 4999), TestStream_ReconnectsAfterServerCloses
4b The backoff resets on the ack The backoff now resets only after the connection delivers a data row or stays up for 60 s or more. TestStream_BacksOffWhenAckedSessionsCloseAtOnce (ack then close: the attempt count keeps growing)
5 The CI key rule No change in behaviour. Only base_url picks the staging key, and stream_url never does. A comment in ci.yml now says so, so the router must send both URLs together. n/a
6 The first connect never fails fast If the first connect gets no HTTP response at all (DNS, TLS, proxy or a bad URL), the CLI exits 1. "Connection lost" now appears only when a connection existed; otherwise the message is "Connection failed". TestStream_FirstConnectWithoutAResponseFailsFast
7 Nits The CLI now checks which signal arrived: SIGINT exits 130 and SIGTERM exits 143. Rows are written with an Encoder that has SetEscapeHTML(false). The README now says data is unchanged apart from whitespace. TestStream_SignalExitCodes, TestStream_SendsKeyHeaderAndPrintsRows (asserts <a&b> is not escaped)

The README and stream help now list these exit and reconnect rules.

🤖 Generated with Claude Code

@szemyd
szemyd merged commit 2422e3c into main Oct 7, 2026
9 checks passed
@szemyd
szemyd deleted the feat/stream branch October 7, 2026 15:46
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant