Metadata-Version: 2.4
Name: bbo-tap
Version: 0.1.0
Summary: Best bid/ask websockets from Binance, Bybit, OKX, Gate and KuCoin perps, captured verbatim to Redpanda/Kafka
Keywords: crypto,market-data,websocket,bbo,order-book,kafka,redpanda,binance,okx,bybit
Author: Vadym O
Author-email: Vadym O <justdoitpicks@gmail.com>
License-Expression: MIT
License-File: LICENSE
Classifier: Development Status :: 3 - Alpha
Classifier: Environment :: Console
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: Intended Audience :: Financial and Insurance Industry
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Python :: 3.14
Classifier: Topic :: Office/Business :: Financial :: Investment
Classifier: Typing :: Typed
Requires-Dist: confluent-kafka>=2.16.0
Requires-Dist: pyyaml>=6.0.3
Requires-Dist: websockets>=17.2
Requires-Python: >=3.13
Project-URL: Homepage, https://github.com/vudya1000/bbo-tap
Project-URL: Source, https://github.com/vudya1000/bbo-tap
Project-URL: Issues, https://github.com/vudya1000/bbo-tap/issues
Project-URL: Changelog, https://github.com/vudya1000/bbo-tap/releases
Description-Content-Type: text/markdown

# bbo-tap

[![CI](https://github.com/vudya1000/bbo-tap/actions/workflows/ci.yml/badge.svg)](https://github.com/vudya1000/bbo-tap/actions/workflows/ci.yml)
[![PyPI](https://img.shields.io/pypi/v/bbo-tap)](https://pypi.org/project/bbo-tap/)
[![Python](https://img.shields.io/pypi/pyversions/bbo-tap)](https://pypi.org/project/bbo-tap/)
[![License: MIT](https://img.shields.io/badge/license-MIT-blue)](https://github.com/vudya1000/bbo-tap/blob/main/LICENSE)

Best bid/ask websockets from **Binance, Bybit, OKX, Gate and KuCoin** USDT perpetuals, captured
**verbatim** to a Kafka-API topic (Redpanda, Kafka). One small asyncio process; no parsing, no
normalization: every received frame is one record, so whatever you build downstream (SQL,
stream processing) works on exactly what the exchanges sent.

- Explicit symbol table: one unified name (`BTC-USDT`) mapped to each venue's own symbol
- One topic partition per exchange, each in receive order
- Reconnects with backoff, proactive 12 h reconnect, idle detection, fast clean shutdown
- Fails loudly: exits non-zero when the broker is unreachable, never drops frames silently

## Install

```sh
pip install bbo-tap        # or: uv tool install bbo-tap / uvx bbo-tap
```

Or the container: `docker pull ghcr.io/vudya1000/bbo-tap`.

## Quick start

```sh
bbo-tap --example-config > config.yaml     # edit brokers and symbols
bbo-tap --config config.yaml
```

The broker needs the topics described in [Output](#output) first, e.g. with Redpanda's `rpk`:

```sh
rpk topic create raw.ws -p 5
rpk topic create ref.symbols -p 1 -c cleanup.policy=compact
```

## Config

```yaml
redpanda:
  brokers: [localhost:9092]
  topic: raw.ws               # optional, default raw.ws
  symbols_topic: ref.symbols  # optional, default ref.symbols
exchanges:                    # every venue listed keeps its partition, enabled or not
  binance: {enabled: true, partition: 0}
  bybit:   {enabled: true, partition: 1}
  okx:     {enabled: true, partition: 2}
  gate:    {enabled: true, partition: 3}
  kucoin:  {enabled: true, partition: 4}
symbols:                      # unified name -> each venue's own symbol; leave a venue out if not listed there
  BTC-USDT: {binance: BTCUSDT, bybit: BTCUSDT, okx: BTC-USDT-SWAP, gate: BTC_USDT, kucoin: XBTUSDTM}
  ETH-USDT: {binance: ETHUSDT, bybit: ETHUSDT, okx: ETH-USDT-SWAP, gate: ETH_USDT, kucoin: ETHUSDTM}
```

Symbols are mapped explicitly, never derived by a naming rule: KuCoin calls BTC `XBT`, and
contracts like `1000PEPE` differ between venues. The file is checked at startup and rejected on an
unknown venue, a row naming a venue not listed under `exchanges`, two venues on one partition, a
venue symbol used twice, an enabled venue with no symbols, or duplicate YAML keys.

`BBO_TAP_BROKERS` (comma-separated) overrides `redpanda.brokers`, so one symbol list serves every
environment.

## Output

Every received frame (data, acks, errors, pongs, notices) becomes one record on `raw.ws`:

```json
{"exchange":"okx","conn_id":"okx-0-1759820000131","seq":918273,"ts_recv_ns":1759820000131402000,"payload":"<exact frame text>"}
```

| Field | Meaning |
| --- | --- |
| `exchange` | Venue name; also the record key |
| `conn_id` | `{exchange}-{shard}-{connect_unix_ms}`, new on every connect |
| `seq` | 0, 1, 2 … per `conn_id`: a gap is a lost frame, a repeat is a harmless retry |
| `ts_recv_ns` | Wall clock right after the socket read (run NTP); also the Kafka timestamp, in ms |
| `payload` | The frame text, unmodified (KuCoin's binary frames decoded as UTF-8) |

- **`raw.ws`**: one partition per exchange, set from config (not hashed), so each partition is one
  venue in receive order. Recommended retention: 7 days.
- **`ref.symbols`** (compacted): written on every start, one record per unified symbol, key
  `BTC-USDT`, value `{"binance":"BTCUSDT","bybit":"BTCUSDT",…}`. Join on it to map venue symbols
  to unified ones; frames themselves are never tagged.
- **Delivery is at least once.** Every data frame on all five venues is a full top-of-book
  snapshot, not a delta, so the latest frame per (exchange, symbol) is the state at that time:
  duplicates don't matter, take the latest `ts_recv_ns`. Producer: `acks=1`, no idempotence,
  zstd, `linger.ms=50`.
- **Loss is never silent.** Frames queue locally (up to 512 MB) while the broker is unreachable;
  when the queue is full, reading pauses instead of dropping. After 30 s without delivery the
  process exits with code 1. At startup it checks that both topics exist and `raw.ws` has every
  configured partition.

## Exchanges

| Venue | Stream | Subscribe | App ping |
| --- | --- | --- | --- |
| Binance USDⓈ-M | `<symbol>@bookTicker` on `/public/stream` | One batch | None (server pings) |
| Bybit V5 linear | `orderbook.1.<symbol>` | One request per symbol | `{"op":"ping"}` / 20 s |
| OKX V5 SWAP | `bbo-tbt` | One batch | `ping` / 20 s |
| Gate USDT futures | `futures.book_ticker`, decimal sizes header | One request per symbol | `futures.ping` / 10 s |
| KuCoin futures (Pro, no token) | `obu` depth 1, binary frames | One request per symbol, 5/s | `{"op":"ping"}` / 18 s |

Bybit and Gate fail a whole request when one symbol in it is unknown, hence one request per
symbol. Binance and KuCoin accept unknown symbols silently, which is why the live check exists
(see [Development](#development)).

**Connections**: up to 25 symbols each; reconnect on socket error, server close, failed protocol
ping, or 30 s without any frame; backoff 1 s doubling to 60 s with full jitter, at most one new
connection per second per venue; proactive reconnect every 12 h (the short gap is accepted).
Frames/s per connection is logged every 60 s.

## Running in production

```sh
docker run -d --name bbo-tap \
  -e BBO_TAP_BROKERS=redpanda:9092 \
  -v ./config.yaml:/etc/bbo-tap/config.yaml:ro \
  --memory 1g --stop-timeout 20 --restart unless-stopped \
  ghcr.io/vudya1000/bbo-tap:latest
```

- **Exit codes**: 0 stopped by SIGTERM/SIGINT (sockets dropped at once, broker flushed), 1 broker
  failure, 2 bad config. Run it under something that restarts on non-zero.
- **Health**: `--heartbeat FILE` touches the file every 10 s only while records flow; the image's
  healthcheck fails when it is older than 60 s (alive but stuck, e.g. every connection down).
- **Resources** (measured, 25 symbols × 5 venues, ~1.8k frames/s): ~17 % of one core, ~60 MB RSS,
  ~1 TB/month inbound traffic, ~8 GB/day in the broker after compression. Allow 1 GB of memory for
  the outage queue.

## Development

```sh
uv sync
uv run ruff format . && uv run ruff check . && uv run mypy
uv run pytest                     # offline, plus broker tests if Redpanda is on localhost:9092 (else skipped)
uv run pytest --live -s           # 30 s against the real exchanges
uv run pytest --live --live-config config.yaml --live-duration 120   # your symbol list; quiet coins need longer
```

The live check verifies, per venue: acks, at least one data frame per symbol, pongs, no error
frames, no early close, and mids agreeing across venues (catches a ticker meaning different coins
on different venues). It needs network access to the exchanges, several of which block US IPs,
so it is not part of CI.

`tests/formats/<exchange>/` holds captured frames of every kind, with notes; the classifier the
live check uses is tested against them.

### Releasing

1. Run the live check from a machine the exchanges accept (not CI):
   `uv run pytest --live --live-config <your config> --live-duration 120`.
2. `uv version --bump patch` (or `minor`), commit, then
   `git tag v$(uv version --short) && git push origin main --tags`.
3. The release workflow checks the tag against the version, runs the checks, builds, publishes to
   TestPyPI, waits for approval in the `pypi` environment, publishes to PyPI, pushes
   `ghcr.io/vudya1000/bbo-tap:<version>` and `:latest` (amd64 + arm64), and creates a GitHub release.

## License

MIT
