Skip to content

StreamingLive data

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; TypeScript watchPrices). Rain is the current lane and requires both a marketId and its on-chain marketAddress.

  • Data Feeds streaming — auxiliary Binance/Chainlink reference tickers (subscribeFeedTicker), the streaming analogue of the REST fetchTicker verb. 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_SUPPORTED error without closing the socket — never a fake stream.

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 closes 4001.

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 codeMeaning
1001Server shutting down (going away).
1008Policy violation — pre-auth message budget (count/bytes) exceeded.
1009Inbound frame exceeded the payload size cap.
4001UNAUTHORIZED — missing, unknown, or revoked API key.
4002INSUFFICIENT_CREDITS — balance cannot cover the next connection-minute.
4003PLATFORM_UNAVAILABLE — metering store unreachable. Fail-closed: never an unmetered stream.
4004CONNECTION_LIMIT — global, per-IP, or per-account concurrent-connection cap reached.
4008RATE_LIMITED — too many failed auth attempts from this IP.

Idle peers that stop answering heartbeat pings are terminated.

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:

  • snapshot is the first frame after (re)subscribe and after backpressure coalescing; update marks 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.
  • subscribeAll requires a venue-wide firehose upstream; venues without one answer an honest NOT_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.

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:

Terminal window
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.

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 carry sourceMetadata.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 discloses sourceMetadata.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).

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();

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. Omitting filters is today’s behaviour, unchanged.
  • pending — the mempool lifecycle only: a pending frame, then confirmed or dropped.
  • all — today’s native confirmed frames plus the lane’s pending and dropped frames. Match a pending frame to its native confirmation by data.transactionHash; the two frames carry different data.id values 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 adds blockNumber.
  • dropped — it did not. The frame adds dropReason: timeout when no receipt arrived in the drop window, or reverted when 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:

Terminal window
predictefy watch trades <marketId> --venue polymarket --status pending

Plan 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.

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.

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.