Rust SDK
fathom-charts handles authentication, the initial snapshot, resuming after a disconnect, duplicates, history pagination, exports and rate-limit retries for you, on Tokio. This page describes every option and method; the behaviour the four SDKs share is in SDKs.
Install
cargo add fathom-charts
cargo add tokio --features full
cargo add futures serde serde_json --features serde/derive- Rust 1.88 or later, on Tokio: a runtime with its time and I/O drivers, as
#[tokio::main]gives. Building a stream outside a runtime returnsCONFIGURATION. - TLS is rustls: nothing native to install.
- Subscriptions and paginated results are
futures::Streams: read them withStreamExt::next.
Indicator types
cargo install fathom-charts --features cli
fathom-charts-types https://api.fathomcharts.com/v1/indicators src/fathom_charts_types.rsThe command takes the catalog’s URL or a file saved from GET /v1/indicators; declare the file as a module (mod fathom_charts_types;). It needs serde (with derive) and serde_json. Per indicator: <Name>Params and <Name>Data (big-trades → BigTradesParams, BigTradesData), with the catalog’s descriptions as doc comments; strings with a fixed set of values become enums, optional fields Option<_>, and a tuple with optional elements, such as a footprint level [price, bid, ask, unknown?], a Vec<i64> whose shape its doc comment gives. INDICATOR_IDS lists the ids, and each indicator implements the Indicator trait (ID, Params, Data). Generate the file again when the catalog changes.
Real time: FathomChartsStream
use fathom_charts::{FathomChartsStream, Mode, StreamEvent, SubscribeOptions};
use futures::StreamExt;
use serde_json::json;
mod fathom_charts_types;
use fathom_charts_types::BigTradesData;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let stream = FathomChartsStream::builder("wss://stream.fathomcharts.com/v1")
.api_key(std::env::var("FATHOM_CHARTS_API_KEY")?)
.on_state_change(|state| println!("connection {state:?}"))
.build()?;
let mut big_trades = stream.subscribe::<BigTradesData>(
SubscribeOptions::new("NQ.front", "big-trades")
.params(json!({ "minimum": 30 }))
.mode(Mode::Confirmed),
)?;
while let Some(event) = big_trades.next().await {
if let StreamEvent::Upsert(upsert) = event? {
println!("{} {:?} {}", upsert.cursor, upsert.data.side, upsert.data.volume);
}
}
Ok(())
}Options
Methods of FathomChartsStream::builder(url), durations as Duration, then build()?:
| Option | Purpose |
|---|---|
url | Required. The stream address: wss://stream.fathomcharts.com/v1. |
api_key(key) | API key, sent in the Authorization: Bearer header. |
stream_token(f) | Instead of api_key, in an application distributed to end users, which must not hold the key: an async function returning a fresh single-use token that your backend gets from POST /v1/stream-tokens (FathomChartsRest::stream_token), as Result<String, _>. It is called on every connection, reconnections included. Give exactly one of api_key and stream_token, otherwise build() returns CONFIGURATION. |
web_socket(f) | Optional: a function that opens the connection and returns a Socket (by default Socket::connect, with the Authorization header). A socket whose pausable is false takes the max_buffered path below. |
reconnect(Reconnect { … }) | Optional: initial_delay (first reconnection delay, 250 ms), max_delay (backoff ceiling, 30 s), quota_retry (how long a 4029 close code is retried, 120 s), quota_delay (minimum delay between two attempts after a 4029, 5 s). Reconnect::default() gives these values. |
heartbeat_timeout(d) | Silence after which the connection is considered lost, then reopened. 15 s by default: the server sends an hb every 5 seconds. |
high_water_mark(n) | Number of events waiting to be read, all subscriptions together, above which the SDK stops reading the socket and keeps only the latest version of each object in progress. While it waits, it sends ping messages so that the server keeps the connection. 10,000 by default. |
max_buffered(n) | For a socket that cannot be paused: number of waiting events above which the SDK closes the connection (close code 4000, reason client buffer full), then reopens it from the last cursor received once half of high_water_mark has been read. Nothing waiting is lost. 4 × high_water_mark by default. |
random(f) | Optional, for tests: the source of the reconnection jitter, a function returning a number in [0, 1). |
on_state_change(f) | Called on every change of connection_state(). |
on_notice(f) | Called with each Notice from the server (kind, id, effective_at, ends_at, message): a restart, a maintenance window announced or cancelled. See Updates and maintenance. |
stream.close() closes every subscription and the connection. stream.connection_state() is Idle, Connecting, Open, Reconnecting or Closed; stream.maintenance() lists the maintenance windows announced on the connection and not over, by start.
subscribe()
stream.subscribe::<T>(options)? opens a subscription and returns a Subscription<T>. SubscribeOptions::new(instrument, indicator) holds the fields of the subscribe message, with .params(json!({…})) (or serde_json::to_value of a generated <Name>Params), .mode(Mode::Confirmed) (Mode::Live by default) and .from(…) (From::Live by default, From::Cursor { cursor, engine } or From::Time(time)). The SDK picks the sub name itself. T is the indicator’s data type (BigTradesData…), serde_json::Value by default.
A Subscription<T> is a futures::Stream of Result<StreamEvent<T>, Error>. StreamEvent is Subscribed, Snapshot, Upsert, Remove, Status, Reset or Notice; their fields are those of Events, in snake case (tick_size_nanos, lag_ms; is_final for final). An error on the subscription is yielded once as an Err, then the stream ends. Data that does not fit T is yielded as an UPSTREAM_ERROR, and the subscription ends.
Subscription
| Member | Purpose |
|---|---|
cursor() | The cursor of the last event handed to you: your resume point. None until a cursor has arrived. It also moves forward with hb messages when no event is waiting to be read. |
topic() | The topic of the latest Subscribed. |
engine() | The engine of the latest Subscribed: the computation your state comes from. Save it with cursor() to resume later. |
id() | The sub name the SDK chose for this subscription (s1, s2…). |
request() | The options of subscribe(), completed with the default values. |
close() | Stops the subscription: the SDK sends unsubscribe, and the stream ends. |
Dropping the subscription closes it. Reconnection, resuming and close codes: see Connection lifecycle; to resume after your process restarts, see Resume after a restart.
History: FathomChartsRest
use fathom_charts::FathomChartsRest;
let rest = FathomChartsRest::new("https://api.fathomcharts.com", std::env::var("FATHOM_CHARTS_API_KEY")?)?;FathomChartsRest::builder(base_url) sets the options:
| Option | Purpose |
|---|---|
base_url | Required. The API address: https://api.fathomcharts.com. |
api_key(key) | Required. API key, sent in the Authorization: Bearer header. |
rate_limit_retries(n) | Number of retries of a request the server did not process, or of a read that failed (see Rate limit and retries). 3 by default. |
client(c) | Optional: your own reqwest::Client (proxy, timeouts…). |
Methods
Queries are PageQuery::new(instrument, indicator, from, to) and, for exports, RangeQuery::new(…), with .params(…), .mode(…) and the other fields of the History page (limit, cursor, until_cursor, snapshot), serialized as the API expects.
| Method | Returns |
|---|---|
page::<T>(&query).await | One page: instrument, tick_size_nanos, snapshot, items, next, compute_units. |
pages::<T>(query) | A stream of the pages, following next, as you read them. |
history::<T>(query) | A stream of the results across pages: HistoryItem::Mutation (t is Upsert or Remove, each with its cursor), preceded, with snapshot, by HistoryItem::Snapshot, the state at from. |
estimate(&query).await | The cost estimate, charging nothing. |
export(&query, resume).await | One export response: the zstd-compressed file as is, with its resume token. |
export_ndjson(query, options) | A stream of the decompressed export, in chunks (Bytes) of whole lines, with automatic resume. ExportOptions: max_resumes (5), resume_delay (1 s). See Exports. |
stream_token().await | A stream token: token, expires_at. |
usage().await | Your usage this month and your limits. |
instruments().await | Every instrument served. |
status().await | The state of the computation, of each instrument, and the maintenance windows announced. |
offers().await | The plans and the limits of each combination. |
indicators().await | The indicator catalog. |
Cancellation is Rust’s: dropping the future of a call, or a stream of pages, history or export_ndjson, stops it, including while it waits to retry.
Errors
Every error the SDK returns, or ends a stream with, is a fathom_charts::Error. Its codes are listed in Errors: code() is an ErrorCode (InvalidParameters, Forbidden, Configuration, Http(status)…), code().as_str() gives it as on the wire, and ERROR_CODES lists them; CloseCode lists the WebSocket close codes.
| Method | Content |
|---|---|
code() | The stable error code. |
message() | A human-readable explanation. |
status() | The HTTP status, over REST. |
close_code() | The WebSocket close code, when the connection was closed. |
retry_after() | The delay of the Retry-After header, when the response has one. |
problem(), problem_errors() | The full problem+json document of the response, and its invalid fields; None otherwise. |
source() | The original error, when there is one (of the HTTP client, of the WebSocket…). |
Cursors
use fathom_charts::compare_cursors;
compare_cursors("20720.111416.0", "20720.111417.0")?; // Ordering::Less