API reference
The public surface, grouped by pipeline stage. Each group is generated from one source file, so everything carrying a docstring appears here.
MarketTickStreamer — Module
MarketTickStreamerProvider-agnostic tick-by-tick market data acquisition for scientific time series research: live WebSocket streaming, historical REST backfill, append-only raw persistence, compaction to analysis-ready files, and paced replay of recorded sessions.
Everything is driven by config/config.toml; credentials come from .env. Entry points: run_stream, run_backfill, replay_source, compact_raw.
Configuration and credentials
MarketTickStreamer.Config — Type
ConfigValidated, typed view of config/config.toml. Constructed via load_config. Field groups mirror the TOML tables; see config/config.toml for semantics.
MarketTickStreamer._resolve_config_date — Method
_resolve_config_date(spec, key, reference) -> DateResolve a configured calendar date. Accepts an ISO date ("2026-08-12") or a sentinel relative to reference: "today", or "today-<N>d" with N a non-negative integer number of calendar days. Sentinels keep a committed configuration from going stale; they count calendar days, not trading days, so a window may resolve onto a weekend or holiday and yield no prints. Throws ArgumentError naming key on any other value.
MarketTickStreamer.load_config — Function
load_config(path = joinpath(PROJECT_ROOT, "config", "config.toml")) -> ConfigRead and validate the TOML configuration. Relative storage/log paths are resolved against the project root. Throws ArgumentError with a specific message on any invalid value.
MarketTickStreamer.load_credentials! — Method
load_credentials!(; env_path = joinpath(PROJECT_ROOT, ".env")) -> (key, secret)Load API credentials from .env (if present) into the environment and return (key_id, secret_key). Looks for ALPACA_API_KEY_ID / ALPACA_SECRET_KEY. Throws an error listing what is missing if either is absent.
Schema and timestamps
Timestamps are Int64 nanoseconds since the UNIX epoch everywhere — the only representation that survives the full range of exchange precision without rounding.
MarketTickStreamer.Bar — Type
BarOne normalized aggregate bar. time_ns is the bar's opening instant, not its close, so a bar and the prints inside it share a time origin.
Like Quote, bars are parsed but not subscribed to by default. They are a convenience the venue computes; anything a bar reports can be derived from the prints this package records, and the derivation is reproducible whereas the venue's aggregation rules are not fully observable.
MarketTickStreamer.Quote — Type
QuoteOne normalized top-of-book quote: the best bid and offer as a venue reported them, with both clocks kept as for Trade — time_ns from the exchange, recv_ns from local receipt.
Quotes are parsed but not subscribed to by default, and not persisted. A quote stream runs an order of magnitude above the trade stream in message count, which is a storage decision rather than a parsing one; reach them through the on_quote callback of live_source.
Sizes are round lots as the tape reports them, not shares.
MarketTickStreamer.Trade — Type
TradeA single normalized trade (tick).
Fields
symbol: ticker symboltime_ns: exchange/participant timestamp, ns since UNIX epoch (UTC)recv_ns: local wall-clock receive timestamp, ns since UNIX epoch (UTC). Thetime_ns → recv_nsgap measures end-to-end latency and is0for records that never crossed the wire (e.g. backfilled).price: trade pricesize: trade size (Float64 to accommodate fractional shares/crypto)exchange: exchange code (provider-specific, e.g. "V" = IEX)conditions: trade condition codes (provider-specific)tape: tape identifier ("A"/"B"/"C" for US equities; may be empty)id: provider trade id (0 when absent)
MarketTickStreamer.now_ns — Method
now_ns() -> Int64Current wall-clock time as Int64 nanoseconds since the UNIX epoch (UTC), built from microsecond-resolution system time.
MarketTickStreamer.ns_to_datetime — Method
ns_to_datetime(ns) -> DateTimeTruncate a nanosecond epoch timestamp to a millisecond-resolution UTC DateTime (for display and coarse grouping only — not for storage).
MarketTickStreamer.ns_to_rfc3339 — Method
ns_to_rfc3339(ns) -> StringRender Int64 nanoseconds since the UNIX epoch as an RFC 3339 UTC timestamp with full nanosecond precision (lossless round-trip with rfc3339_to_ns).
MarketTickStreamer.rfc3339_to_ns — Method
rfc3339_to_ns(s) -> Int64Parse an RFC 3339 timestamp with up to nanosecond fractional precision into Int64 nanoseconds since the UNIX epoch (UTC). Supports Z and ±hh:mm offsets. Throws ArgumentError on malformed input.
MarketTickStreamer.trading_date — Method
trading_date(ns; tz) -> DateExchange-local calendar date of a nanosecond epoch timestamp; used to bucket ticks into per-day files. tz defaults to America/New_York.
Provider interface
A provider adapter implements this interface; the acquisition layer is written against it and never against a vendor. Two adapters ship, and they disagree on every axis the interface abstracts: Alpaca is credentialed, files its tape under the New York calendar, and closes overnight and at weekends; Binance is public, files under UTC, and never closes. A configuration is validated against the selected provider's own ProviderSpec rather than against one vendor's vocabulary.
MarketTickStreamer.KNOWN_PROVIDERS — Constant
Provider names this build knows how to construct.
MarketTickStreamer.AbstractProvider — Type
AbstractProviderSupertype for market-data provider adapters.
A concrete adapter supplies provider_spec and make_provider for the name it answers to, exchange_tz, stream_protocol!, market_clock, historical_trades and its own message parsing. It overrides always_open, session_days and feed_delay_ns where the venue departs from the defaults.
Two adapters ship: providers/alpaca.jl (US equities, New York calendar, credentialed, closes overnight and at weekends) and providers/binance.jl (crypto spot, UTC calendar, public, never closes). They differ on every one of those axes, which is what keeps the interface honest.
MarketTickStreamer.FatalStreamError — Type
FatalStreamError(msg)A stream error that must NOT trigger reconnection (bad credentials, connection-limit exceeded, invalid subscription). Aborts the session.
MarketTickStreamer.LiveSession — Type
LiveSessionHandle for a running live stream: the tick channel plus control state. Obtain via live_source; request shutdown with stop!.
MarketTickStreamer.ProviderSpec — Type
ProviderSpec(; feeds, backfill_feeds, tz, needs_credentials)What the configuration layer must know about a provider before any provider object exists: which feed names it accepts, the calendar its tape is filed under, and whether it needs credentials at all.
Without this the config layer has to hard-code one vendor's vocabulary, which is how provider.feed came to be validated against Alpaca's feed names for every provider, and how a 24-hour venue would have had its days cut at midnight in New York.
MarketTickStreamer.always_open — Method
always_open(p::AbstractProvider) -> BoolWhether the venue never closes. true suppresses the market-hours railings: there is no open to wait for and no close to stop at, and a 24-hour venue asked for its next_close can only answer with a fiction. The session is then bounded by limits.max_session_hours alone.
MarketTickStreamer.exchange_day_start_ns — Method
exchange_day_start_ns(p::AbstractProvider, d::Date) -> Int64First instant of exchange date d on this provider's calendar, in ns since the epoch.
A request window and the key its rows are bucketed by must come from the same clock. Building the window on UTC days while filing rows on exchange dates agrees only while the exchange sits at UTC-4, and silently loses an hour a day otherwise — the defect fixed in 0.2.0. ZonedDateTime resolves the offset per date, including across the 23- and 25-hour transition days.
MarketTickStreamer.exchange_tz — Function
exchange_tz(p::AbstractProvider) -> TimeZoneThe calendar the provider's tape is filed under — the clock that decides which date a print belongs to. Defaults to the provider's ProviderSpec.
MarketTickStreamer.feed_delay_ns — Method
feed_delay_ns(p::AbstractProvider) -> Int64Intrinsic delay of the provider's configured feed (ns): the wall-clock lag between an exchange event and its earliest possible arrival on the wire. Zero for real-time feeds. Market-hours railings shift by this amount so a delayed session waits out the silent post-open window and captures the delayed tape tail after the close.
MarketTickStreamer.live_source — Method
live_source(p::AbstractProvider, cfg::Config;
on_quote = nothing, on_bar = nothing) -> LiveSessionStart the live producer task. Streams ticks into session.channel until the session deadline (limits.max_session_hours), a stop! call, a fatal protocol error, or reconnection exhaustion — whichever comes first. The channel is closed on exit so downstream consumers terminate cleanly.
on_quote and on_bar are handed normalized Quote and Bar values when those channels are subscribed to (stream.channels); neither is by default, and neither is persisted — a quote stream carries an order of magnitude more messages than the trade stream, which is a storage decision taken separately.
Reconnects with jittered exponential backoff (stream.reconnect_base_delay_s * 2^attempt, capped at stream.reconnect_max_delay_s); the attempt counter resets after any connection that actually delivered data.
MarketTickStreamer.make_provider — Method
make_provider(::Val{name}, cfg, key, secret) -> AbstractProviderConstruct the provider called name from configuration and credentials. key/secret are empty strings when the spec says none are needed. Implemented per adapter.
The fallback throws rather than returning anything: a fallback with a value would widen every caller's inferred provider type to include it, and then each provider method downstream would carry a branch with no matching method — which is how JET found this the first time it was written otherwise.
MarketTickStreamer.provider_spec — Method
provider_spec(::Val{name}) -> ProviderSpecStatic description of the provider called name. Implemented per adapter.
MarketTickStreamer.schedule_close_stop! — Method
schedule_close_stop!(s::LiveSession, close_ns; grace_s = 5.0) -> TaskSpawn a guard task that gracefully stop!s the session once the wall clock passes close_ns (ns since epoch, e.g. the market's next_close) plus grace_s. Without this, a streamer left unattended sits on a silent overnight connection until the session deadline.
MarketTickStreamer.session_days — Method
session_days(p::AbstractProvider, start_date, end_date) -> Vector{Date}The dates in [start_date, end_date] on which the venue trades, in order.
Defaults to every calendar date. A venue that rests must opt out, never the reverse: requesting a day the venue was shut costs one empty response, while skipping a day it traded loses that day's tape silently — which is what a hard-coded weekday filter did to a 24-hour venue before this was dispatched.
MarketTickStreamer.stop! — Method
stop!(s::LiveSession)Signal graceful shutdown: sets the stop flag and closes the underlying WebSocket, which unblocks the producer loop; the tick channel then closes, letting sinks drain and finish.
MarketTickStreamer.stream_protocol! — Function
stream_protocol!(ch, p::AbstractProvider, cfg, s::LiveSession;
on_quote = nothing, on_bar = nothing)Provider interface: run one connect→auth→subscribe→stream cycle, pushing normalized Trades into ch. Must return on orderly close, throw FatalStreamError on non-retryable protocol errors, and any other exception on retryable transport failures. Implemented per provider (see providers/alpaca.jl).
on_quote and on_bar, when given, receive normalized Quote and Bar values for the corresponding frames. Neither channel is subscribed to by default, so neither callback fires unless stream.channels asks for it.
Alpaca adapter
MarketTickStreamer.AlpacaProvider — Type
AlpacaProviderConnection descriptor for Alpaca Markets: credentials, feed selection and endpoint roots (overridable via the [alpaca] config table, which lets the test suite point the client at a local mock server).
MarketTickStreamer.condition_map — Method
condition_map(p; ticktype = "trade", tape = "A") -> Dict{String,String}Fetch the provider's own sale-condition decoder from /v2/stocks/meta/conditions/{ticktype}: a map of condition code to description for one tape ("A", "B", "C"). ticktype is "trade" or "quote".
The same character carries different meanings on different tapes, and every vendor normalizes the raw CTA and UTP codes differently, so the glossary is fetched from the provider that produced the data rather than transcribed into this package. What is held here is the much smaller NON_PRICE_CONDITIONS judgement about which codes disqualify a print from a price path.
Requires network access; call it once and cache the result.
MarketTickStreamer.historical_trade_count — Method
historical_trade_count(p, symbol, start_str, end_str;
feed = "sip", page_limit = 10_000,
rate_sleep_s = 0.35) -> IntCount trades on the historical tape for symbol over the inclusive RFC 3339 window [start_str, end_str] without retaining them — the reference side of live-capture coverage checks.
MarketTickStreamer.historical_trades — Method
historical_trades(p, symbol, start_date, end_date;
feed = "sip", page_limit = 10_000, rate_sleep_s = 0.35,
on_page = nothing, each_page = nothing)Download all trades for symbol in [start_date, end_date] (inclusive, exchange dates) from /v2/stocks/{symbol}/trades, following pagination tokens until exhausted. rate_sleep_s throttles between pages (free tier: 200 requests/min). Backfilled records get recv_ns = 0 — they never crossed the wire. on_page(n_page, n_total) is called per page for progress.
By default all trades are accumulated and returned as a Vector{Trade}. With each_page set, each page's trades are handed to each_page(::Vector{Trade}) instead and only the total row count is returned — memory stays bounded by one page regardless of the range.
MarketTickStreamer.market_clock — Method
market_clock(p) -> (; is_open, next_open, next_close)Query /v2/clock on the trading API. Timestamps are returned as the raw RFC 3339 strings Alpaca sends (display only — nothing downstream computes with them).
MarketTickStreamer.subscribe_payload — Method
subscribe_payload(p, symbols, channels) -> StringBuild the JSON subscription message for the requested channels ("trades"/"quotes"/"bars"), e.g. {"action":"subscribe","trades":["AAPL","MSFT"]}.
Binance adapter
Crypto spot, public market data, no credentials. The venue never closes, so the market-hours railings do not apply and the calendar is UTC.
MarketTickStreamer.BinanceProvider — Type
BinanceProviderConnection descriptor for Binance spot market data. Endpoint roots are overridable through the [binance] config table, which is how the test suite points the client at a local mock server.
feed selects the trade granularity: "trade" is every execution, "aggTrade" aggregates executions that filled at one price from a single taker order. They are different populations — see historical_trades.
MarketTickStreamer.historical_trades — Method
historical_trades(p::BinanceProvider, symbol, start_date, end_date;
feed = "aggTrade", page_limit = 1000, rate_sleep_s = 0.1,
on_page = nothing, each_page = nothing)Download aggregate trades for symbol over [start_date, end_date] (inclusive UTC dates) from /api/v3/aggTrades. Signature and semantics match the Alpaca adapter's, so the backfill pipeline is provider-agnostic: each_page streams pages and returns only the row count, otherwise every trade is accumulated.
Pagination is by trade id, not by time window. Binance rejects a startTime/endTime pair spanning an hour or more, so the range is seeded with a single startTime request and then walked forward with fromId = last id + 1 until a row's timestamp passes the end of the range. Walking hour-wide windows instead would issue 24 requests a day and still truncate any hour holding more than page_limit trades, silently.
aggTrade is not the tape. Binance aggregates executions that filled at one price from one taker order into a single row, so an aggregate trade is a taker order's fill, not an execution. Counts and inter-arrival times are therefore not comparable with a raw trade stream, and the choice belongs to the analysis — as with odd lots on an equity tape. Only aggTrade has a time-seekable public endpoint, which is why it is the backfill feed.
MarketTickStreamer.market_clock — Method
market_clock(p::BinanceProvider) -> (; is_open, next_open, next_close)Always open. next_open and next_close are empty because the venue has neither; always_open is true for this provider, so the session orchestration never asks for them.
MarketTickStreamer.ws_url — Method
ws_url(p::BinanceProvider, symbols) -> StringCombined-stream URL for symbols. Binance subscribes through the URL rather than through a post-connection message, so there is no auth or subscribe handshake: the connection is the subscription.
Persistence and compaction
MarketTickStreamer.SPILL_HEADROOM — Constant
SPILL_HEADROOMMultiple of the input size that the spill filesystem must have free before compact_raw will stream to it.
The spill pass trades memory for scratch space, writing every input line back out once. tempdir() is the wrong home for that: on systemd distributions /tmp is a tmpfs sized at half of RAM, so spilling there writes the copy into RAM and, on an input larger than that, takes the machine down — the failure this guard exists to prevent. The default scratch location is therefore the directory holding out_dir, and the space is verified before the first line is read.
MarketTickStreamer.RawSink — Type
RawSinkAppend-only NDJSON writer with size-based file rolling. Create with open_raw_sink; feed via write_batch!; always close_sink! (final flush) on shutdown.
MarketTickStreamer.close_sink! — Method
close_sink!(s::RawSink)Flush and close the sink's file handle.
MarketTickStreamer.compact_raw — Method
compact_raw(raw_paths, out_dir; format = "csv", dedup = true,
mem_fraction = 0.5, tz = tz"America/New_York",
max_live_heap_mb = nothing, scratch_dir = nothing) -> Vector{String}Compact raw NDJSON files into per-symbol, per-trading-day analysis files (out_dir/SYMBOL/YYYY-MM-DD.csv|.arrow), sorted by time_ns. Existing outputs are never overwritten — a #N suffixed sibling is written instead (safesave semantics). Returns the list of files written.
dedup = true drops exact duplicate prints (reconnection double-delivery, overlapping backfill/live captures) via deduplicate_trades, logging the count. Small inputs are compacted in memory; when the estimated footprint exceeds mem_fraction of currently free RAM the input is instead streamed line-by-line into per-(symbol, day) spill files and each group is compacted independently — memory stays bounded by the largest single group. max_live_heap_mb additionally enforces the configured heap ceiling per group (check_live_heap).
The spill pass writes a verbatim copy of the input, so it needs scratch space of the input's own size. scratch_dir places it; the default is the directory containing out_dir, not tempdir() — see SPILL_HEADROOM. The free space is checked up front and compaction refuses to start without it.
MarketTickStreamer.json_to_trade — Method
json_to_trade(line) -> TradeParse one NDJSON line written by trade_to_json back into a Trade.
MarketTickStreamer.open_raw_sink — Method
open_raw_sink(dir, prefix; max_mb = 1024) -> RawSinkOpen a fresh raw NDJSON file dir/prefix_partNNN.jsonl. Existing files are never reopened or overwritten — a new part number is chosen past any that already exist.
MarketTickStreamer.read_raw — Method
read_raw(paths) -> Vector{Trade}Load one or more raw NDJSON files into memory, skipping (and counting via a @warn) corrupt lines — e.g. a partial last line after a hard kill.
MarketTickStreamer.run_sink! — Method
run_sink!(ch::Channel{Trade}, sink::RawSink;
flush_interval_s = 30.0, flush_max_ticks = 5000,
on_flush = nothing) -> IntConsumer loop: drain ch into an in-memory batch and persist whenever the batch reaches flush_max_ticks or flush_interval_s elapses with pending data. Returns the total number of trades written. Terminates (after a final flush) once ch is closed and fully drained — closing the channel is the shutdown signal, so no locks or flags are needed.
on_flush(n_batch, n_total) is called after each disk write (for logging).
MarketTickStreamer.trade_to_json — Method
trade_to_json(t::Trade) -> StringSerialize a Trade to a single NDJSON line (no trailing newline). Timestamps stay Int64 nanoseconds — the round-trip through json_to_trade is lossless.
MarketTickStreamer.write_batch! — Method
write_batch!(s::RawSink, trades) -> RawSinkAppend a batch of trades as NDJSON lines and flush the stream once (buffered writes reach the OS per batch, not per line). Rolls to a new part file when the size limit is exceeded.
Data quality and resource guards
MarketTickStreamer.NON_PRICE_CONDITIONS — Constant
NON_PRICE_CONDITIONSSale-condition codes that never update a bar's open or close price, keyed by tape: "A" and "B" are CTA-processed, "C" is UTP-processed, "O" is the OTC tape. The same character means different things on different tapes, which is why the lists are not shared.
Taken from Alpaca's published lists for its own normalization of the tapes (verified 2026-09-13), not transcribed from the CTA and UTP plan specifications: the plans define raw codes, and every vendor normalizes them. The authoritative decoder for the codes themselves is the provider's own condition_map.
Two things this is deliberately not. It is not the plans' "last-sale eligible" flag, and it is not a per-field eligibility table — the plans track high/low, open/close, volume and last-sale eligibility separately, and a print may update some and not others. It answers one question, the one a price path needs answered: may this print set a price.
Note "9" (corrected consolidated close): excluded here because at tick resolution it is a correction message rather than an execution, though for a daily bar it is precisely the official close.
MarketTickStreamer.check_live_heap — Method
check_live_heap(limit_mb; context = "") -> NothingConfig-gated RAM ceiling: if the live heap exceeds limit_mb, force a garbage collection; if it still exceeds the ceiling, fail loudly with a message naming context — a graceful stop beats an OOM kill. Call from long accumulation loops (backfill pages, compaction groups).
MarketTickStreamer.coverage_report — Method
coverage_report(paths, provider; feed = "sip", page_limit = 10_000,
rate_sleep_s = 0.35) -> DataFrameCompare a live capture against the historical tape: for every symbol in the raw NDJSON paths, count captured trades (after exact-duplicate removal) and query the provider's historical trade count over the same inclusive exchange-time window. coverage = captured / reference; values below 1 quantify feed coverage and stream drops, values above 1 indicate duplicate or spurious prints that dedup did not catch.
MarketTickStreamer.deduplicate_trades — Method
deduplicate_trades(trades) -> Vector{Trade}Remove exact duplicate prints (same symbol, exchange timestamp, id, price, size, exchange), keeping first occurrence and preserving order. Guards against reconnection double-delivery and overlapping backfill/live captures.
MarketTickStreamer.filter_price_forming — Method
filter_price_forming(trades; non_price = NON_PRICE_CONDITIONS) -> Vector{Trade}Keep only the prints for which price_forming holds, preserving order. Use on a loaded capture, never on the way in.
MarketTickStreamer.free_disk_gb — Method
free_disk_gb(path) -> Float64Free disk space (GiB) on the filesystem containing path.
MarketTickStreamer.price_forming — Method
price_forming(t::Trade; non_price = NON_PRICE_CONDITIONS) -> BoolWhether t may set a price, i.e. whether none of its sale conditions is listed for its tape in non_price. A print carrying several conditions is disqualified by any one of them, which is the precedence rule both tape plans state: a single "does not update" overrides every permissive condition beside it.
A print whose tape has no entry in non_price is kept. An unknown tape means no basis on which to exclude, and silently discarding prints on that ground would be the worse error; session_report counts what survives, so the effect stays visible.
This is an analysis-time predicate. Capture is never filtered — the raw layer records every print the venue reported, and which subset constitutes "a trade" is a decision each analysis makes for itself. It is also not one decision: a price path wants price-forming prints only, whereas an arrival process or a waiting-time distribution counts every execution, odd lots and contingent trades included. Filtering the capture would foreclose the second question to answer the first.
MarketTickStreamer.session_report — Method
session_report(paths; gap_threshold_s = 60.0) -> DataFrameAudit raw session files and return one row per symbol:
n_trades,n_duplicates(exact duplicate prints)n_price_forming(prints that may set a price, seeprice_forming; the difference fromn_tradesis odd lots, contingent and derivatively priced trades, corrections and the like — present in the tape, and in the arrival process, but not in a price path)first_time,last_time(UTC, ms precision, from exchange timestamps)n_out_of_order(exchange timestamps decreasing in arrival order)max_gap_s/n_gaps(largest / count of exchange-time gaps exceedinggap_threshold_s— judge against the instrument's typical activity; quiet symbols gap naturally)median_latency_ms/n_negative_latency(receive minus exchange time, live-captured rows only; negative values indicate clock skew.NaNwhen the file is pure backfill)
Inspect this before trusting any captured session.
Resampling onto activity clocks
MarketTickStreamer.dollar_bars — Method
dollar_bars(trades, value; keep_partial = false) -> Vector{Bar}Resample trades onto a traded-value clock: one Bar per value of price times size.
This is the clock that survives changes of scale. It is invariant to splits, and roughly invariant to price-level drift, so a threshold chosen on one sample stays meaningful on another — the property neither tick_bars nor volume_bars has, and the reason dollar bars are the usual default for cross-sample work.
value is in the price's own units; no currency conversion is performed.
MarketTickStreamer.tick_bars — Method
tick_bars(trades, n; keep_partial = false) -> Vector{Bar}Resample trades onto a transaction clock: one Bar per n prints.
The count is of prints as recorded, so whether odd lots and other non-price-forming records count towards it is decided before the call — pass filter_price_forming(trades) to exclude them. On a high-priced name that choice changes the bar count by a factor of three, so it is not incidental.
Sampling by transaction count is the oldest of these clocks: Mandelbrot and Taylor (1967) proposed price changes as a subordinated process running on transaction time, and Ané and Geman (2000) showed returns sampled this way are close to normal where calendar-time returns are heavy-tailed.
Prints are ordered by exchange timestamp, not arrival. The print that completes a bar belongs to it. A trailing incomplete bar is dropped unless keep_partial, since its threshold — and therefore its comparability to the others — is not met.
tick_bars(trades, 500) # one bar per 500 prints
tick_bars(filter_price_forming(trades), 500) # price-forming prints onlyMarketTickStreamer.volume_bars — Method
volume_bars(trades, volume; keep_partial = false) -> Vector{Bar}Resample trades onto a volume clock: one Bar per volume shares.
Clark (1973) introduced volume as the directing process for speculative prices, giving a finite-variance alternative to the stable-Paretian account of heavy tails: the unconditional return distribution is heavy-tailed because it mixes over a random volume clock, not because the underlying increments are.
The caveat is mechanical: a share is not a fixed unit of economic activity across time or across instruments. Share prices drift, splits reset the scale, and a fixed share threshold therefore samples a different amount of value in January than in December. dollar_bars is the usual remedy.
Replay
MarketTickStreamer.replay_source — Method
replay_source(paths; pace = "recorded", speed = 1.0,
capacity = 100_000) -> Channel{Trade}Create a channel that replays the trades recorded in raw NDJSON paths (sorted by receive timestamp).
pace = "recorded": honor the original inter-arrival times ofrecv_ns, compressed byspeed(e.g.speed = 60.0replays an hour in a minute).pace = "max": emit as fast as the consumer takes them (throughput mode).
The channel closes when the recording is exhausted, which cleanly terminates any consumer written against run_sink!-style loops.
Session orchestration
MarketTickStreamer.acquire_session_lock — Method
acquire_session_lock(data_dir) -> lockOne session per data tree: takes a PID-file lock at <data_dir>/.session.lock and throws with a precise message if another live process already holds it (observed failure mode: two backfills interleaving on one raw directory). Stale locks from dead processes are broken automatically; release with close(lock).
MarketTickStreamer.finalize_session_meta — Method
finalize_session_meta(path, status; ticks, raw_files) -> StringRewrite the running sidecar with the final status ("completed" / "interrupted"), tick count, file list, and finish timestamp. The sidecar is a mutable status record by design — the safesave rule protects data products, not status metadata.
MarketTickStreamer.reconcile_sessions! — Method
reconcile_sessions!(raw_dir) -> IntCrash-only startup reconciliation: sidecars still marked running whose recorded process no longer exists on this host are relabeled aborted (their raw NDJSON remains valid to the last flushed line); zero-byte raw stubs are reported. Never deletes anything. Returns the number of sidecars relabeled.
MarketTickStreamer.run_backfill — Method
run_backfill(cfg::Config; provider = nothing) -> Vector{String}Download historical trades for every configured symbol over [backfill.start_date, backfill.end_date] (weekends skipped), persist them through the same raw-NDJSON + compaction path as live data (marked by recv_ns = 0), and return the processed file paths.
Work is done one (symbol, day) at a time: its own raw part set, compacted as soon as that day finishes. Pages are flushed to disk as they arrive, so RAM stays bounded by one page and limits.max_live_heap_mb.
That granularity is what makes backfill.resume mean anything on a long download. The skip test asks whether a day is already present under processed/, so a day must land there the moment it is complete and not before: a run killed after eight of twelve hours resumes at hour eight. Compacting only at the end — as this did until 2026-09-13 — left a killed run with raw data that resume could not see, and the rerun started from nothing.
An interrupted day is not compacted and is therefore downloaded again; its partial raw file is left on disk rather than deleted, and compaction deduplicates if it is ever folded in. Ctrl-C finalizes the provenance sidecar as interrupted and returns the days that did complete.
MarketTickStreamer.run_entrypoint — Method
run_entrypoint(main)Run a command-line entry point's main with clean interrupt handling: SIGINT arrives as an InterruptException and exits without a stack trace.
Delivery is best-effort under threads — the crash-only session lifecycle, not this wrapper, is what actually guarantees that an interrupted run leaves recoverable state. Shared by the scripts under scripts/ and the worked examples so that every entry point behaves the same way at the terminal.
Example
run_entrypoint(() -> main(ARGS))MarketTickStreamer.run_stream — Method
run_stream(cfg::Config; provider = nothing) -> NamedTupleRun one complete live capture session:
- optional market-clock gate (
stream.require_market_open), - live WebSocket source → bounded channel,
- batched raw NDJSON persistence,
- graceful shutdown on Ctrl-C, session deadline, or stream termination.
Returns (; ticks, raw_files). Blocks until the session ends.
MarketTickStreamer.session_id — Method
session_id(cfg::Config; feed = cfg.feed) -> StringUnique identifier for a capture session, <provider>_<feed>_<yyyymmdd-HHMMSS> (UTC start time). Prefixes every raw part file, provenance sidecar, and log file the session produces.
MarketTickStreamer.setup_logging — Method
setup_logging(cfg; session_id) -> AbstractLoggerConsole logger at cfg.log_level, optionally teed to cfg.log_dir/<session_id>.log (plain formatting, no ANSI). The caller installs it with global_logger.
MarketTickStreamer.start_session_meta — Method
start_session_meta(cfg, sid; started_utc) -> StringWrite the sidecar at session START with status = "running" — an abruptly killed process still leaves its provenance on disk.
MarketTickStreamer.tee — Method
tee(src::Channel{Trade}, n; capacity = 10_000, lossy = falses(n),
on_drop = nothing) -> Vector{Channel{Trade}}Fan one tick stream out to n independent consumers (e.g. persistence plus real-time analysis taps). Outputs close when src closes.
Per-output overflow policy: a non-lossy output blocks the fan-out when full (backpressure — data is never silently dropped; use for persistence, which must always win). A lossy output drops the incoming tick instead when its buffer is full (use for analysis taps that must never stall the capture); drops are counted, reported via on_drop(output_index, n_dropped) when given, and warned on first occurrence.
MarketTickStreamer.write_session_meta — Method
write_session_meta(cfg, sid; status, ticks, raw_files, started_utc,
finished_utc = nothing) -> StringPersist a provenance sidecar <raw_dir>/<sid>.meta.toml next to the session's raw files: session summary (id, status, pid, span, tick count, file list), provenance (git commit + dirty flag, Julia and package versions, hostname), a hardware fingerprint (CPU model and logical core count, total memory, Julia and BLAS thread counts, full versioninfo output), and the full effective configuration snapshot. The crash-only lifecycle is start_session_meta → finalize_session_meta, reconciled at startup by reconcile_sessions!.
Live monitoring
MarketTickStreamer.MonitorState — Type
MonitorStateIncremental tail state over one session's raw part files: per-file byte offsets and partial-line carries, cumulative per-symbol counters, a tick-rate history window, and the latest exchange timestamp seen (the tape head). Counters cover data ingested since attach (see from_start).
MarketTickStreamer.latest_session_prefix — Method
latest_session_prefix(dir) -> StringPrefix (session id) of the most recently modified raw part file in dir. Throws if none exist.
MarketTickStreamer.monitor_raw — Method
monitor_raw(dir; session = nothing, refresh_s = 2.0, top_symbols = 10,
rate_window_s = 300.0, from_start = false,
iterations = nothing, io = stdout) -> MonitorStateAttach a live in-terminal dashboard to the session whose raw NDJSON files live under dir (most recently active session unless session gives an id prefix). Tails the part files read-only — safe to run beside a live capture or backfill. Renders every refresh_s: tick totals and rate, tape-head timestamp and lag, a rate history sparkline, and per-symbol counts.
from_start = true ingests the whole existing file first (session totals; costs one full read of the raw data); the default starts at the current end of file (activity since attach). iterations bounds the number of refreshes (nothing = run until interrupted). Returns the final MonitorState.
Visualization
MarketTickStreamer.overview_figure — Method
overview_figure(symbol, days; tz = tz"America/New_York") -> FigureMulti-day diagnostic figure for one symbol from per-day processed tables (days is a vector of (date, DataFrame) pairs, sorted internally by date):
- price path on a concatenated trading-time axis (overnight gaps removed, session boundaries dashed, days labeled at their centers),
- day × session-minute activity heatmap,
- pooled intra-session inter-arrival CCDF (overnight gaps excluded by construction — they would contaminate the waiting-time tail),
- pooled trade-size CCDF with fitted tail exponent.
MarketTickStreamer.save_overview_figures — Function
save_overview_figures(processed_dir, out_dir = joinpath(PROJECT_ROOT, "plots");
symbols = nothing, from = nothing, to = nothing,
formats = ("pdf", "png"), tz = tz"America/New_York",
min_days = 2) -> Vector{String}Render one multi-day overview_figure per symbol from the processed tree (processed_dir/SYMBOL/YYYY-MM-DD.csv|.arrow; safesave #N siblings are ignored — the base file per day is authoritative). symbols, from, and to restrict the sweep. Output: out_dir/SYMBOL_<from>_<to>.pdf|png (safesave). Days are loaded one symbol at a time, so memory stays bounded by one symbol's span.
MarketTickStreamer.save_session_figures — Function
save_session_figures(raw_paths, out_dir = joinpath(PROJECT_ROOT, "plots");
formats = ("pdf", "png"), tz = tz"America/New_York",
min_trades = 10) -> Vector{String}Render one diagnostic figure per (symbol, trading day) found in the raw NDJSON raw_paths, saved as out_dir/SYMBOL_YYYY-MM-DD.pdf|png (safesave — existing files are never overwritten). Groups with fewer than min_trades ticks are skipped with an @info (distribution panels are meaningless). PNG output is written at px_per_unit = 4. Returns the files written.
MarketTickStreamer.session_figure — Method
session_figure(trades; tz = tz"America/New_York") -> FigureBuild the 2×2 diagnostic figure for one symbol's single-day ticks:
- price path over exchange-local time (
HH:MMaxis, min–max decimated), - activity (trades per minute),
- inter-arrival time survival function
P(Δt > x)(log-log, decade ticks), - trade-size survival function
P(S > s)with a fitted tail exponent.
trades must be non-empty and single-symbol (as produced by the grouping in save_session_figures); they are sorted internally by exchange time.
MarketTickStreamer.tick_theme — Method
tick_theme() -> ThemePublication defaults: Computer Modern fonts, boxed axes, inward ticks, no minor ticks, faint dashed grey grid.