Skip to content

WebSockets

The client streams the channels documented under WebSockets: what each channel sends, its payloads and its limits are covered there. This page shows the SDK calls. WebSockets are async only, on AsyncSTX.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
ticker = await ws.ticker()
print(ticker.topic, ticker.reply)
asyncio.run(main())

client.websocket() uses the client’s host and key, signs the handshake with that key and sends a User-Agent header, which the exchange requires. The market data channels also accept an unsigned socket; the account channels need the signed one, so they need a key. Connecting joins nothing; each join method returns a Channel. STXWebSocket() also works on its own and takes the same settings as AsyncSTX.

Each message is a ChannelMessage with topic, event and payload. Take them with a callback, one at a time, or with async for:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
ids = [m.market_id for m in await client.markets(status="open", limit=5)]
async with client.websocket() as ws:
# 1. A callback, sync or async.
await ws.trades(market_ids=ids, on_message=lambda msg: print("trade", msg.payload))
book = await ws.orderbook(ids)
# 2. One at a time, with a timeout.
try:
msg = await book.next(timeout=5)
print(msg.event, msg.payload["market_id"])
except asyncio.TimeoutError:
print("no book change in 5 s")
# 3. As an async iterator (runs until the channel closes).
async def consume():
async for msg in book:
print("book", msg.payload["market_id"], msg.payload["bids"][:1])
task = asyncio.create_task(consume())
await asyncio.sleep(3)
task.cancel()
asyncio.run(main())

A slow consumer never blocks the socket: each channel buffers up to queue_size messages (10,000 by default) and drops the oldest beyond that.

Account channels send your current state right after the join. wait_snapshot() returns it, keyed by event name:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
positions = await ws.positions()
snap = await positions.wait_snapshot(timeout=10)
for p in snap["all_positions"]["positions"]:
print(p["market_id"], p["position"], p["open_risk"])
asyncio.run(main())

market_stats carries its history in the join reply instead (channel.reply["markets"]). settlements, orderbook, ticker, trades, markets and market_updates send no snapshot: take the starting state from markets(), orders() and the other client methods.

The aggregated book for the markets you name; market_ids is required. Each book push is the full book for one market: replace what you hold, do not merge.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
ids = [m.market_id for m in await client.markets(status="open", limit=3)]
async with client.websocket() as ws:
book = await ws.orderbook(ids)
print("selected", book.reply["selected_market_ids"])
await book.select_market_ids(ids[:1]) # change markets without rejoining
try:
msg = await book.next(timeout=10)
print(msg.payload["market_id"], "bids", msg.payload["bids"][:2], "offers", msg.payload["offers"][:2])
except asyncio.TimeoutError:
print("quiet book")
asyncio.run(main())

Each level has price, quantity, liquidity (that level), and total_quantity / total_liquidity (cumulative), best first.

A ticker push whenever a market’s last price, top of book, volume or open interest moves. Filter by sports and competitions:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
ticker = await ws.ticker(sports=["Baseball"])
print(ticker.reply)
await ticker.select_filters(sports=None, competitions=["MLB"])
try:
msg = await ticker.next(timeout=10)
p = msg.payload
print(p["market_symbol"], p["last_traded_price"], p["best_bid"], p["best_offer"])
except asyncio.TimeoutError:
print("no ticker change in 10 s")
asyncio.run(main())

Every execution on the exchange, anonymised. action is the taker’s side. This is not your fills: see fills.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
trades = await ws.trades()
print(trades.reply)
try:
msg = await trades.next(timeout=10)
print(msg.payload["market_symbol"], msg.payload["action"], msg.payload["quantity"], "@", msg.payload["price"])
except asyncio.TimeoutError:
print("no trade in 10 s")
asyncio.run(main())

Filter with market_ids= and event_ids=; change them later with select_filters(...).

market_created and market_updated for every market. The payload maps market id to a market object; market_updated carries only market_id, the timestamps and the fields that changed. Prices arrive in cents on the wire and the SDK converts them to dollar strings, the same format as everywhere else.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
markets = await ws.markets(rule_filters=None, message_types=["market_updated"])
print("rules you can filter on:", len(markets.reply["available_rules"]))
try:
msg = await markets.next(timeout=15)
for market_id, change in msg.payload.items():
print(msg.event, market_id, change)
except asyncio.TimeoutError:
print("no market change in 15 s")
asyncio.run(main())

select_rule_filters([...]) and select_message_types([...]) change the filters without rejoining.

A price series per market, for charts. The history is in the join reply; market_stats pushes changed buckets (upsert by timestamp_us), and market_stats_snapshot replaces a series. price_percent is a percent of max_price, not money.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
ids = [m.market_id for m in await client.markets(status="open", limit=2)]
async with client.websocket() as ws:
stats = await ws.market_stats(ids, range="week")
for series in stats.reply["markets"]:
print(series["market_id"], len(series["points"]), "points")
day = await stats.request_series(ids[:1], range="day")
print("day range:", day["range"])
asyncio.run(main())

created and updated for the markets you watch, and nothing until you do. Watches are sent again automatically after a reconnect. Prices are converted from cents to dollar strings.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
ids = [m.market_id for m in await client.markets(status="open", limit=3)]
async with client.websocket() as ws:
updates = await ws.market_updates(watch=ids)
print("watching", updates.watches)
try:
msg = await updates.next(timeout=15)
print(msg.event, msg.payload["market_id"], msg.payload)
except asyncio.TimeoutError:
print("no update in 15 s")
asyncio.run(main())

Scoped to you. The SDK builds the topic (orders:<user_id>), fetching your user id from GET /api/v1/me the first time it needs it.

all_orders on join, then new_open_order (the whole order) each time one is accepted or filled.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
market = None
async for m in client.iter_markets(status="open", trading=True, sort_by="event_start", sort_direction="desc"):
if m.event_status == "scheduled":
market = m
break
async with client.websocket() as ws:
orders = await ws.orders(market_ids=[market.market_id])
await orders.wait_snapshot()
order = await client.place_order(market.market_id, "buy", "limit", price="0.01", quantity="1")
while True:
msg = await orders.next(timeout=15)
if msg.event == "new_open_order" and msg.payload["id"] == order.id:
print("pushed", msg.payload["status"], msg.payload["price"])
break
await client.cancel_order(order.id)
asyncio.run(main())

market_ids= filters the snapshot and every push; select_market_ids(None) clears it. For cancel-on-disconnect, join with cancel_on_disconnect=True and optionally ping_timeout in milliseconds; the SDK pings at 60% of the granted timeout. See Cancel-on-disconnect.

all_trades on join, then one trade per execution, including trades that later settle or are cancelled.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
fills = await ws.fills(on_message=lambda m: print(m.event))
snap = await fills.wait_snapshot()
trades = snap["all_trades"]["trades"]
print(len(trades), "fills in the snapshot")
for t in trades[:5]:
print(t["id"], t["action"], t["filled"], "@", t["price"], "fee", t["total_fee"])
asyncio.run(main())

all_positions on join, then updated_positions carrying only the positions that changed.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
positions = await ws.positions(on_message=lambda m: print(m.event, len(m.payload["positions"])))
await positions.wait_snapshot()
asyncio.run(main())

new_settlements whenever a market you hold settles. No snapshot; history is client.settlements().

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
settlements = await ws.settlements()
print(settlements.reply)
try:
msg = await settlements.next(timeout=5)
print(msg.payload["settlements"])
except asyncio.TimeoutError:
print("nothing settled in 5 s")
asyncio.run(main())

balances on join, then update when your balance changes and payment_update for payments. Pass account_id= to watch another account you own.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
balances = await ws.balances()
b = (await balances.wait_snapshot())["balances"]
print("available", b["available_balance"], "cash", b["account_balance"])
asyncio.run(main())

Balance pushes follow your activity (orders, fills, settlements, payments), not price moves.

Everything the five channels above carry, on one join: four snapshots, then the same events under the same names. Do not also join a per-type channel (you would get each message twice), and use orders instead if you need cancel-on-disconnect.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
account = await ws.account()
snap = await account.wait_snapshot()
print(sorted(snap))
print("available", snap["balances"]["available_balance"])
asyncio.run(main())

user_updated right after the join, then whenever your profile changes.

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async with client.websocket() as ws:
info = await ws.user_info()
profile = (await info.wait_snapshot())["user_updated"]
print(profile["userStatus"], profile["firstName"])
asyncio.run(main())

join(topic, payload) joins any documented topic, and push(event, payload) sends it a control event and returns the reply. For example, event volume and the order slip:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
markets = await client.markets(status="open", trading=True, limit=3)
async with client.websocket() as ws:
events = await ws.join("events", {"event_ids": [m.event_id for m in markets]})
print(events.reply["events"])
slip = await ws.join(f"order_slip:{await ws.user_id()}")
m = markets[0]
reply = await slip.push(
"add_order",
{"market_id": m.market_id, "qty": 1, "side": "buy", "limit_price": 0.5,
"max_price": float(m.max_price)},
)
print("ref", reply["ref"])
msg = await slip.next(timeout=10)
print(msg.event, msg.payload["updates"][0]["risk_with_fee"])
asyncio.run(main())

The order slip takes numbers, not strings, as its channel page describes; it costs an order and places nothing.

The SDK:

  • sends a heartbeat on the phoenix topic every 25 s (the server closes a socket silent for 60 s), and reconnects if one goes unanswered;
  • pings each joined channel every 30 s (channel_ping_interval), which also keeps the session behind orders and account alive;
  • on orders with cancel-on-disconnect, pings at 60% of the granted ping_timeout, a separate and much shorter deadline;
  • after a drop, reconnects with backoff, signs the handshake again, and rejoins every channel with its current filters and watches. Snapshots arrive again after the rejoin.

What it cannot do is know what you missed while disconnected. Pass on_reconnect and call orders() (and anything else you show) again there:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
async def resync():
open_orders = await client.orders(status=["open", "delayed"])
print("reconnected; open orders:", len(open_orders))
async with client.websocket(on_reconnect=resync) as ws:
await ws.orders()
await ws.fills()
print("connected:", ws.connected, "reconnects so far:", ws.reconnects)
asyncio.run(main())

Channels are not ordered relative to each other: a fill can arrive before the order update that explains it. Key off ids.

The intervals, the reconnect backoff and the join timeout are keyword arguments to client.websocket() and STXWebSocket:

import asyncio
from stx import AsyncSTX, ReconnectPolicy
async def main():
async with AsyncSTX() as client:
ws = client.websocket(
heartbeat_interval=25.0, # seconds between socket heartbeats
channel_ping_interval=30.0, # ping on every joined channel; None turns it off
reconnect_policy=ReconnectPolicy(initial_backoff=0.5, max_backoff=30),
join_timeout=10.0,
)
async with ws:
await ws.ticker()
print("connected:", ws.connected)
asyncio.run(main())

A refused join raises STXChannelException with the server’s reason in .reply, for example {"reason": "market_ids_required"} or {"reason": "unauthorized"}. So does a refused control event such as select_market_ids([]) on orderbook.

run_forever() blocks until close() is called or reconnects give up:

import asyncio
from stx import AsyncSTX
async def main():
async with AsyncSTX() as client:
ids = [m.market_id for m in await client.markets(status="open", limit=5)]
async with client.websocket() as ws:
await ws.orderbook(ids, on_message=lambda m: print("book", m.payload["market_id"]))
await ws.fills(on_message=lambda m: print("fills", m.event))
asyncio.get_running_loop().call_later(10, lambda: asyncio.ensure_future(ws.close()))
await ws.run_forever()
asyncio.run(main())
v1.5.9Changelogllms.txtllms-full.txt