StreamingLive data
Streaming
Predictefy streams four things over one WebSocket, one auth handshake, and one metered connection:
-
Prediction-market streaming — live venue order books and trades. Multiple clients watching the same market share one upstream venue subscription, so you never burn a venue’s rate limits by scaling out consumers.
-
Venue option-price streaming — live price frames from capability-qualified venue lanes (
subscribePrice; TypeScriptwatchPrices). Rain is the current lane and requires both amarketIdand its on-chainmarketAddress. -
Data Feeds streaming — auxiliary Binance/Chainlink reference tickers (
subscribeFeedTicker), the streaming analogue of the RESTfetchTickerverb. See Data Feeds streaming below. -
Cross-venue arbitrage streaming — one shared executable-arbitrage surface (
subscribeArbitrage), recomputed server-side on a disclosed cadence rather than tick by tick, and gated on the same plan feature as the REST verb. See Cross-venue arbitrage streaming below. -
Endpoint path:
/v1/stream(WebSocket upgrade). -
Trade subscriptions use a native venue stream when available, or a disclosed chain-scan tape for Rain, XO, and PRED. A venue with neither answers an honest
NOT_SUPPORTEDerror without closing the socket — never a fake stream.
Auth handshake
Section titled “Auth handshake”Connect with the same pk_live_… API key as the REST API:
Authorization: Bearer pk_live_…header (preferred), or- a first-frame auth message for browser clients (which cannot set WebSocket
headers):
{ "op": "auth", "apiKey": "pk_live_…" }. It must be the first frame, within 10 seconds, or the socket closes4001.
Keys are never read from the URL — ?apiKey= is deliberately unsupported (raw keys
in URLs leak into proxy and observability logs).
Auth is checked before any venue subscription is created. Failures close the socket:
| Close code | Meaning |
|---|---|
1001 | Server shutting down (going away). |
1008 | Policy violation — pre-auth message budget (count/bytes) exceeded. |
1009 | Inbound frame exceeded the payload size cap. |
4001 | UNAUTHORIZED — missing, unknown, or revoked API key. |
4002 | INSUFFICIENT_CREDITS — balance cannot cover the next connection-minute. |
4003 | PLATFORM_UNAVAILABLE — metering store unreachable. Fail-closed: never an unmetered stream. |
4004 | CONNECTION_LIMIT — global, per-IP, or per-account concurrent-connection cap reached. |
4008 | RATE_LIMITED — too many failed auth attempts from this IP. |
Idle peers that stop answering heartbeat pings are terminated.
Wire protocol (JSON text frames)
Section titled “Wire protocol (JSON text frames)”Client → server:
// `marketId` is the VENUE-NATIVE book id (see "Which id goes in marketId?"// below) — for polymarket the CLOB asset/token id = the outcome's `outcomeId`.{ "op": "subscribe", "channel": "orderbook", "venue": "polymarket", "marketId": "<asset_id>" }{ "op": "unsubscribe", "channel": "orderbook", "venue": "polymarket", "marketId": "<asset_id>" }{ "op": "subscribe", "channel": "trades", "venue": "<venue>", "marketId": "<market>" }{ "op": "subscribeAll", "channel": "orderbook", "venue": "<venue>" }{ "op": "unsubscribeAll", "channel": "orderbook", "venue": "<venue>" }// Venue option-price frames — marketAddress is the on-chain market contract:{ "op": "subscribePrice", "venue": "rain", "marketId": "<market>", "marketAddress": "0x…" }{ "op": "unsubscribePrice", "venue": "rain", "marketId": "<market>", "marketAddress": "0x…" }// Data Feeds reference tickers — feed + symbol, no venue/marketId:{ "op": "subscribeFeedTicker", "feed": "binance", "symbol": "BTC/USDT" }{ "op": "unsubscribeFeedTicker", "feed": "chainlink", "symbol": "BTC/USD" }// Cross-venue executable arbitrage — optional filters (executableOnly, venues, minEdge):{ "op": "subscribeArbitrage", "executableOnly": true, "venues": ["polymarket", "kalshi"], "minEdge": 0.02 }{ "op": "subscribeArbitrage" }{ "op": "unsubscribeArbitrage" }Server → client:
{ "type": "subscribed", "channel": "orderbook", "venue": "polymarket", "marketId": "…" }{ "type": "unsubscribed", "channel": "orderbook", "venue": "polymarket", "marketId": "…" }// books: full-state frames; prices are probabilities in [0,1]{ "type": "snapshot", "venue": "polymarket", "marketId": "…", "data": { "bids": [{ "price": 0.4, "size": 10 }], "asks": [], "timestamp": 1780000000000 }, "ts": 1780000000123 }{ "type": "update", /* same shape as snapshot */ }{ "type": "trade", "venue": "…", "marketId": "…", "data": { /* trade */ }, "ts": 1780000000123 }// Venue option-price ack + data:{ "type": "subscribed", "channel": "price", "venue": "rain", "marketId": "…" }{ "type": "price", "venue": "rain", "marketId": "…", "marketAddress": "0x…", "data": { /* venue option prices */ }, "ts": 1780000000123 }// Data Feeds ticker ack + data (feed/symbol, not venue/marketId):{ "type": "subscribed", "channel": "feedTicker", "feed": "binance", "symbol": "BTC/USDT" }{ "type": "feedTicker", "feed": "binance", "symbol": "BTC/USDT", "data": { "symbol": "BTC/USDT", "last": 61714.63, "asOf": "…", "provenance": { "source": "binance-ws" }, "sourceMetadata": { "transport": "websocket" } }, "ts": 1780000000123 }{ "type": "error", "code": "NOT_SUPPORTED", "message": "…", "venue": "…", "marketId": "…" }Notes:
snapshotis the first frame after (re)subscribe and after backpressure coalescing;updatemarks live ticks. Both carry the full book, never deltas.- Protocol-level problems (
BAD_MESSAGE,NOT_SUPPORTED,NOT_SUBSCRIBED,SUBSCRIPTION_LIMIT) are non-fatal: the socket stays open. subscribeAllrequires a venue-wide firehose upstream; venues without one answer an honestNOT_SUPPORTED.- Trades stream from a native fills channel or the disclosed chain-scan tape served for Rain,
XO, and PRED. Other venues without either capability return
NOT_SUPPORTED. - Active logical subscriptions are capped by plan: 2 Free, 20 Builder, and
100 Pro. A lower service safety cap can also apply. Exceeding either returns
a non-fatal
SUBSCRIPTION_LIMIT.
Which id goes in marketId?
Section titled “Which id goes in marketId?”marketId is the venue-native id for the exact book you want — it is passed
straight through to the venue and is not always the unified marketId from
fetchMarkets. For polymarket it is the CLOB asset/token id, which the unified
API returns as the outcome’s outcomeId (fetchMarket →
outcomes[].outcomeId) — the same id you pass to fetchOrderBook. Other venues
use their own native id (for example hyperliquid uses the coin symbol like
BTC).
Subscribing with the wrong id (e.g. a polymarket market-level marketId) returns a
subscribed ack and then zero data frames — the venue never recognizes it as a
book, so it looks identical to a quiet market. When in doubt, use the value you’d
pass to fetchOrderBook.
Going the other way: venue id to canonical identity
Section titled “Going the other way: venue id to canonical identity”Streaming frames carry venue-native ids, so correlating them back to Predictefy’s catalog needs a
lookup. POST /v1/mappings does it in bulk:
curl -s "$PREDICTEFY_API_URL/v1/mappings" \ -H "Authorization: Bearer pk_live_YOUR_KEY" \ -H "content-type: application/json" \ -d '{"pairs":[{"venue":"polymarket","marketId":"0xabc…"}]}'It accepts 1 to 200 {venue, marketId} pairs and returns one result per pair in input order,
so you can zip the response straight onto your request array. A known market returns its canonical
marketPk, its current status, and its clusterId when it has been matched to markets on other
venues.
Unknown pairs and sandbox-venue pairs are not dropped: they stay in the response with
marketPk, clusterId and status set to null. Positional correlation therefore holds even when
some ids are unrecognised — you never have to re-align two arrays of different lengths.
The route is gated by READS_ENABLE_MAPPINGS and answers 404 where it is not enabled.
Data Feeds streaming
Section titled “Data Feeds streaming”subscribeFeedTicker streams auxiliary reference tickers — the WebSocket analogue of
the REST GET /api/feeds/{feed}/fetchTicker verb. Here
{feed} means one of the curated reference feeds (binance or chainlink), not a venue id.
This is reference data, separate from prediction-market venues (it never touches markets,
clusters, or history) and rides the same auth and per-connection-minute credits as venue
subscriptions.
binance— spot reference prices over Binance’s public market-data WebSocket; frames carrysourceMetadata.transport = "websocket". Symbols:BTC/USDT,ETH/USDT,SOL/USDT,XRP/USDT.chainlink— on-chain oracle prices, poll-backed (there is no Chainlink push stream): the service polls the on-chain feed on a bounded interval and emits only when the round advances. Every frame disclosessourceMetadata.transport = "poll"so you always know it is poll-derived, not a push.
The official SDK wraps the handshake:
import { Predictefy } from '@predictefy/sdk';
const client = new Predictefy({ apiKey: process.env.PREDICTEFY_API_KEY });
const close = client.watchFeedTicker({ feed: 'binance', symbol: 'BTC/USDT' }, (ticker) => console.log(ticker.last, ticker.sourceMetadata?.transport),);// later: close();The key is sent over the safe first-frame handshake — never a ?apiKey= URL. An
unknown feed or unsupported symbol answers a non-fatal NOT_SUPPORTED (the socket stays
open).
Cross-venue arbitrage streaming
Section titled “Cross-venue arbitrage streaming”subscribeArbitrage streams the executable-arbitrage surface — the WebSocket analogue
of the REST fetchArbitrage verb. It is cross-venue, so it takes no
venue and no marketId.
Optional filters. subscribeArbitrage accepts optional filters:
{ executableOnly?: boolean, venues?: string[], minEdge?: number }. Filters apply
server-side; omitting them preserves the default unfiltered stream.
Pair rows. Each base cluster emits every ordered cross-venue pair, with 90 rows as the
pathology guard. The publisher makes one batched live-book read per venue for the selected outcome
books and reuses those results across the pairs. Each row’s clusterId is the composite
${clusterId}:${venueA}:${venueB}. Because the base cluster id can contain :, strip the last
two colon-delimited segments to recover it. The REST-only pairsPerCluster query parameter can
cap these rows for fetchArbitrage; it is not a subscribeArbitrage filter.
What the subscription delivers. There is one client operation: subscribeArbitrage. It sends
the complete selected surface — the full surface when no filters are set — as snapshots and
sequence-ordered deltas. A complete snapshot is due every 30 seconds; on that recompute pass it is
sent immediately before the pass’s delta. Every successful recompute sends a delta, including an
empty upserts/removes delta when no row changed. Either kind can span messages no larger than
200KiB, tagged with part: { i, n }.
Every arbitrage data message carries type: "arbitrage", the frame in data, and a socket-write
ts. Both data variants carry exchange, seq, computedAt, publishedAt, intervalMs,
heartbeatMs, contracts, limit, and part. New and reconnecting subscribers receive the
current coherent snapshot after the acknowledgement; if a client falls behind, the relay repairs
it with a current snapshot before resuming deltas. The full frame shape is in the
WebSocket API reference.
The cadence is honest, not tick-by-tick. This is one shared server-side recompute:
intervalMs is the real cadence (3000 ms by default). Every successful pass publishes a delta;
when the priced surface did not change, that delta has empty upserts and removes. A pass with a
due snapshot publishes the snapshot first and then its delta under the next seq, so the timestamps
keep publisher liveness observable without pretending the market moved.
Read heartbeatMs off the frame; it is a conservative liveness bound, not a fixed constant. It
is the first recompute tick at or after the server’s 30000 ms target: 30000 ms at the default 3000
ms interval, but 40000 ms if the interval were 20000 ms. Successful per-pass deltas normally arrive
more often, at intervalMs; an empty delta is liveness, not a market change. Sustained silence
beyond the heartbeatMs you were actually sent means a publisher outage or an entitlement
teardown (below), not a quiet market.
You can measure the lane yourself. Each frame carries computedAt (the pass finished),
publishedAt (the server handed it to its relay), and the envelope’s ts (written to your
socket). On a live delivery, those gaps measure compute-to-publish and relay-to-socket latency;
nothing on that path deliberately buffers, batches, or waits for a timer. On retained replay,
ts − publishedAt instead measures the last-known frame’s age. The recompute interval is the live
lane’s only deliberate delay, and it is there to bound upstream venue API cost.
The label discipline is the REST verb’s. A row is labeled arbitrage only with positive
net edge, both legs depth-executable at the requested size against live asks, and resolution
equivalence verified. Everything else is served as an indicative price discrepancy with the
per-leg reasons codes explaining why. The surface is never trimmed down to the winners — the
indicative rows are part of it, with their evidence.
A new subscriber gets a snapshot of the current surface right after its subscribed ack. The replay
keeps the frame’s original publishedAt; only the outer ts records the new socket write. Compare
that age with heartbeatMs: a retained frame older than the advertised heartbeat is stale and does
not claim that the publisher is still live.
Staleness contract. The retained frame’s publishedAt is the publisher-liveness signal: a
healthy publisher refreshes it on every successful recompute, including an empty delta. If
publishedAt trails the envelope’s ts by more than roughly 90 seconds (the current SDK default),
treat the publisher as stale and fall back to REST. A publisher that goes permanently dark after
publishing once therefore continues to yield an aging retained frame instead of reverting to
NOT_SUPPORTED; that is intentional, and client-side age detection is the safeguard. The
TypeScript SDK will surface this condition as a staleness event in the current SDK release train.
Publisher readiness requires at least one valid frame observed by this relay. If none has ever
arrived, subscribeArbitrage gets a non-fatal NOT_SUPPORTED error instead of a success ack, even
when Redis itself is reachable; the client can retry after the publisher is enabled. A deployment
with no Redis relay returns the same code. The channel is separately gated on the same arbitrage
plan feature as the REST verb — a key whose plan does not include it gets
PLAN_UPGRADE_REQUIRED, and the Free plan does not include it.
The entitlement is re-checked while the socket is open, roughly once a minute. If the plan
stops entitling the feature mid-stream, the arbitrage subscription is torn down and you get
the same PLAN_UPGRADE_REQUIRED frame; if the key itself is revoked or rotated you get
UNAUTHORIZED instead, since that needs re-authentication rather than an upgrade. Either way
the socket stays open and every other subscription on it keeps streaming, and a re-check that
cannot complete never interrupts you.
The TypeScript SDK wraps it like the ticker lane. Pass onError for any long-lived watcher:
an entitled socket with nothing to report and a refused one look identical without it.
const close = client.watchArbitrage( ({ frame }) => { for (const row of frame.rows) { console.log(row.label, row.executable, row.reasons); } }, { onError: (e) => console.error(e.code, e.message) },);// later: close();Pending fills (Polymarket)
Section titled “Pending fills (Polymarket)”A pending fill is a Polymarket settlement observed in the Polygon public mempool before it mines. Polymarket settles fills as ordinary Polygon transactions, and those transactions are public for a moment before a block includes them, so decoding one yields its fills slightly early. The measured lead is about 1.7 s median and 2.5 s at p90. Polymarket is the only venue with this lane; every other venue’s trade stream is unchanged.
Subscribing. A trades subscription takes an optional filters object:
{ "op": "subscribe", "channel": "trades", "venue": "polymarket", "marketId": "<asset_id>", "filters": { "status": "pending" }}confirmed— the default. Omittingfiltersis today’s behaviour, unchanged.pending— the mempool lifecycle only: apendingframe, thenconfirmedordropped.all— today’s native confirmed frames plus the lane’spendinganddroppedframes. Match a pending frame to its native confirmation bydata.transactionHash; the two frames carry differentdata.idvalues because the native tape has its own id scheme.
Any other channel or venue, an unknown key inside filters, or an unknown status value answers a
non-fatal BAD_MESSAGE. A trades subscription is keyed by venue and market, so to change
status, unsubscribe and subscribe again. The subscribed ack echoes "status" for pending
and all, and carries a disclosure whose provenance is "mempool".
Lifecycle. A fill arrives first as pending and then resolves exactly once, under the same
data.id and data.transactionHash:
confirmed— it mined. The frame addsblockNumber.dropped— it did not. The frame addsdropReason:timeoutwhen no receipt arrived in the drop window, orrevertedwhen the transaction mined with a failed status, so no fill happened.
Every frame on this lane carries a top-level status, seenAt (the ISO time the fill was first
seen), and "provenance": "mempool". data.id is <txHash>:<fillIndex>, where index 0 is the
taker fill and the makers follow in settlement order. data.wallet is the filling order’s wallet
and data.counterparty is the other side when it is known. The exact frames are in the
WebSocket API reference.
The TypeScript SDK takes the same filter. Pass onError: without it a refused subscription looks
exactly like a quiet market.
const close = client.watchTrades( { venue: 'polymarket', marketId: '<asset_id>', status: 'pending' }, ({ trade, status, seenAt, dropReason }) => console.log(status, seenAt, trade.price, trade.shares, dropReason ?? ''), { onError: (e) => console.error(e.code, e.message) },);// later: close();The CLI exposes it as --status:
predictefy watch trades <marketId> --venue polymarket --status pendingPlan and credits. Pending fills require the Pro plan or higher. A key below that gets a
non-fatal PLAN_UPGRADE_REQUIRED frame and no subscription — the socket stays open:
{ "type": "error", "code": "PLAN_UPGRADE_REQUIRED", "message": "the \"pending fills\" feature requires the Pro plan or higher (current plan: \"builder\")", "venue": "polymarket", "marketId": "<asset_id>", "channel": "trades"}While a connection holds a pending or all trades subscription, that connection-minute costs
800 credits instead of 2 — the pending minute replaces the base minute rather than adding to
it. The first pending minute is charged when the subscription is accepted, on top of the
connect-time minute you already prepaid. A balance that cannot cover it gets an
INSUFFICIENT_CREDITS error frame and the subscription is refused; the socket stays open and
every other subscription on it keeps streaming. See Credits & billing.
Availability. Pending frames pause while the mempool watcher is unavailable. Subscriptions stay open and no error frame is sent, so treat a gap as missing observations rather than as a quiet market. The feature can also be switched off centrally.
Backpressure
Section titled “Backpressure”A slow consumer never causes unbounded buffering:
- Order books are coalesced — intermediate updates are dropped and only the latest
book per subscribed market is parked; when your socket drains, that latest book
arrives as a single
snapshot. - Trades are dropped (not parked) while backpressured — the trade stream is lossy under backpressure by design.
- Price frames are dropped (not parked) while backpressured — the next venue price frame supersedes the last.
- Feed tickers are dropped (not parked) while backpressured — same policy as trades. Reference tickers are low-frequency and the next frame supersedes the last.
- Arbitrage frames — clients that fall behind under backpressure receive a fresh snapshot automatically through hub-side stale-client recovery once the connection catches up.
Credits
Section titled “Credits”Streaming costs 2 credits per connection-minute, prepaid: minute #1 is charged at
connect, each subsequent minute on the minute. A minute in which the connection holds a
pending fills subscription costs 800 credits instead of the 2,
charged the same way. When your balance cannot cover the next
minute you get an INSUFFICIENT_CREDITS error frame, then close 4002. See
Credits & billing.