API reference: acquisition and storage
The public surface, grouped by pipeline stage. Each group is generated from one source file, so everything carrying a docstring appears here. The analysis side — data quality, activity clocks, replay and the diagnostic figures — is on the next page.
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 a TOML configuration (config/config.toml in a clone); credentials come from the environment file it names. Entry points: run_stream, run_backfill, replay_source, compact_raw.
Plotting is optional. The diagnostic figures (session_figure, overview_figure and their save_* drivers) are implemented by a package extension that loads with CairoMakie, and the terminal plots of monitor_raw by one that loads with UnicodePlots; neither is a dependency of the package.
Configuration and credentials
MarketTickStreamer.DEFAULT_CONFIG — Constant
DEFAULT_CONFIGPath of the configuration shipped with the package, config/config.toml. In a clone of the repository it is the working configuration; in an installed copy it is a template to copy from, since the package directory is read-only there.
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._cfg_value — Method
_cfg_value(tbl, table, key, T, default) -> TRead tbl[key], falling back to default, and require it to be of type T (one of Bool, Int, Float64, String). Throws ArgumentError naming table.key otherwise, so a mistyped value cannot reach the pipeline.
MarketTickStreamer._reject_unknown_keys — Method
_reject_unknown_keys(raw)Check the parsed TOML against CONFIG_KEYS and the known provider tables. Throws ArgumentError naming the offending table or key, so a misspelt entry cannot be silently ignored in favour of a default.
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 = DEFAULT_CONFIG) -> ConfigRead and validate the TOML configuration. Relative storage.data_dir, logging.log_dir and credentials.env_file are resolved against the directory of path. Every key is checked for type and documented bounds, and unknown tables or keys are rejected, so a misspelt key cannot fall back to its default unnoticed. Throws ArgumentError naming the offending key.
The no-argument form reads the configuration shipped with the repository and suits a clone of it. In an installed copy of the package it throws an ArgumentError: there DEFAULT_CONFIG is a read-only template to copy into a project of your own.
MarketTickStreamer.load_credentials! — Method
load_credentials!(cfg::Config) -> (key, secret)
load_credentials!(; env_path) -> (key, secret)Load API credentials from the environment file (credentials.env_file of the configuration, or env_path), if it exists, into the process 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 opening instant and close_ns the closing one, both in nanoseconds since the Unix epoch on the exchange clock, so a bar and the prints inside it share a time origin.
For an activity bar (tick_bars, volume_bars, dollar_bars) time_ns is the timestamp of its first print and close_ns that of the print that completed it. The duration close_ns - time_ns is then a random variable — the time the market needed to transact a fixed amount of activity — and is the observable a waiting-time study reads off an activity clock. For a venue-computed bar close_ns is the end of the venue's aggregation window.
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!.
stop is atomic: it is written by whoever ends the session and read by the producer, the watchdog and the guard tasks, which run on other threads. The logging layer reads the same flag to mute transport-teardown noise, so the flag belongs to one session and not to the process.
MarketTickStreamer.ProviderSpec — Type
ProviderSpec(feeds, backfill_feeds, tz, needs_credentials, max_page_limit,
regular_session, price_unit, size_unit)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, whether it needs credentials at all, the largest REST page the venue serves, and the venue's regular session.
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.
max_page_limit exists because a venue may clamp an over-limit request instead of rejecting it, and a clamped page is indistinguishable from the last page of the tape.
regular_session is (open, close) in exchange-local hours, (0.0, 24.0) for a venue that never closes. The multi-day figures lay their days out on it; nothing in acquisition depends on it, since session railings come from the venue's own clock.
price_unit and size_unit are what the venue's prices and sizes are denominated in, as the figure axes should name them: dollars and shares on a US equity tape, quote asset and base asset on a crypto pair.
MarketTickStreamer.RestPolicy — Type
RestPolicy(request_timeout_s, connect_timeout_s, max_retries)
RestPolicy() # 30 s, 10 s, 5 retries
RestPolicy(cfg::Config) # the `[rest]` tableHow a provider's REST requests are bounded: the read timeout and the connection timeout, both in whole seconds, and the number of retries on a transient failure. Without a read timeout a request against a server that accepts the connection and then goes silent blocks its task for good, and with it a backfill.
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,
stop = Threads.Atomic{Bool}(false)) -> 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.
stop lets the caller own the session's stop flag, which run_stream shares with its logger.
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). rest bounds its REST requests (RestPolicy).
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 a vendor is free to re-map the raw CTA and UTP codes — this one delivers them unchanged — 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 table of the codes that 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. rest bounds its REST requests (RestPolicy).
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. page_limit above the venue maximum of 1000 is rejected here, because the server would clamp it silently and the short-page termination rule would then end the walk after the first page.
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_HEADROOMDefault multiple of the input size that the spill filesystem must have free before compact_raw will stream to it (spill_headroom).
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._safesave — Method
_safesave(write_fn, path) -> StringWrite a file through write_fn(tmp_path) so that path always holds the newest complete result and nothing already there is lost.
The content goes to <base>.partial<ext> in the same directory first (the extension is kept last because writers such as FileIO.save infer the format from it). If path already exists it is then preserved as the first free <base>_#N<ext>, N from 1 — a hard link where the filesystem allows one, a copy otherwise — and finally the partial file is renamed onto path. The rename replaces atomically on one filesystem, so at every instant path is either the previous complete file or the new complete one, never absent and never half-written. Backups are numbered in order of displacement: _#1 is the oldest.
If write_fn throws, the partial file is removed, path is left untouched and the exception propagates.
These are the semantics of DrWatson's safesave — newest data canonical, prior results kept as numbered backups — without the dependency.
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,
footprint_factor = 4.0, spill_headroom = SPILL_HEADROOM,
max_open_spill_files = 256, compression = 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. An existing output is never destroyed: the new file takes the canonical name and the one it displaces is kept as a numbered backup YYYY-MM-DD_#N.csv|.arrow (safesave semantics, see _safesave). 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.
footprint_factor is the estimated in-memory size of the parsed input per input byte, spill_headroom the scratch space required as a multiple of the input size, and max_open_spill_files the number of spill files held open at once: further groups are served by closing the least recently written file and reopening it in append mode, so the pass works within the process's file-descriptor limit however many symbol-days the input spans. The pipeline takes all of them from [limits].
Arrow output stores the symbol, exchange, conditions and tape columns dictionary-encoded, which costs nothing to read and takes a liquid US equity day from 65 to 44 bytes per print. compression (:zstd or :lz4, Arrow only) compresses the record batches as well — 8 bytes per print with :zstd on the same day — at the price of memory-mapped reads: an uncompressed file is mapped and touched lazily, a compressed one is decompressed on load (64 ms against under 1 ms for that day). Uncompressed is the default.
MarketTickStreamer.compact_raw — Method
compact_raw(cfg::Config, raw_paths; dedup = true, scratch_dir = nothing)
-> Vector{String}Compact into cfg.processed_dir with everything else taken from the configuration: format, the provider's calendar, the heap ceiling and the spill thresholds of [limits]. This is the form the pipeline and the scripts use, so a capture is never filed under a calendar other than its provider's.
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.processed_files — Method
processed_files(processed_dir, symbol; from = nothing, to = nothing)
-> Vector{String}The per-day files of symbol in a processed tree (processed_dir/SYMBOL/YYYY-MM-DD.csv|.arrow), in date order, optionally restricted to [from, to]. Safesave backups (_#N) and partial files are not listed: the base file of a day is the newest and the only one a reader should open. Throws an ArgumentError when the symbol has no directory.
replay_source(processed_files(cfg.processed_dir, "AAPL"; from = Date(2026, 5, 1)))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.
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" or "failed"), 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. Any other error finalizes it as failed and propagates; the days compacted before the error remain valid and a rerun resumes after them.
Logging is scoped to the call: the session runs under its own logger and the caller's global logger is untouched.
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. Logging is scoped to the call: the session runs under its own logger and the caller's global logger is untouched.
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, shutting_down = Threads.Atomic{Bool}(false))
-> (logger, io)Console logger at cfg.log_level, optionally teed to cfg.log_dir/<session_id>.log (plain formatting, no ANSI, flushed per record). Returns the logger and the open log file, or nothing without file logging; the caller runs the session under with_logger(logger) and closes io when it ends. Nothing is installed globally, so a session leaves the caller's logger as it found it and two sessions in one process keep separate logs.
While shutting_down[] is true, records from HTTP.jl internals are dropped: during a deliberate shutdown an EOFError from a socket we closed ourselves is expected, not an incident.
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.
A consumer that closes its output is dropped from the fan-out with a warning and the other outputs keep receiving, so an analysis tap that dies cannot end the persistence branch. If every output is closed the fan-out stops with an error record. A failure of the fan-out task itself is logged (errormonitor) rather than lost.
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, and the base name of a copy of the resolved Manifest.toml written alongside as <sid>.manifest.toml), 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, and per-symbol counts. With UnicodePlots loaded (using UnicodePlots) the counts become a bar chart and a rate-history plot is added; without it the dashboard is plain text.
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.