Skip to content

Live capture

Run a live source into a directory of segments, and read back the manifest.json that records them. See Capture live data for how segments, rolls, reconnects and restarts work.

run_capture async

run_capture(
    make_source: SourceFactory, config: CaptureConfig
) -> CaptureRun

Capture live data into a directory of segments, for as long as asked.

make_source builds a fresh source for each segment: a source keeps connection and book state for one run and cannot be reused. config gives the pair, the capture directory (out_dir), the total run time (minutes), and the roll limits (roll_minutes, roll_mb).

A new out_dir (or an empty one) starts a new capture. One holding a manifest.json for the same source, pair and level is continued; any other contents are refused with :class:~ob_analytics.exceptions.ConfigError. So is a directory another capture process is writing to: the capture holds a lock on out_dir while it runs.

SIGINT/SIGTERM stop the capture: the running segment is closed with its closing rows and the manifest is finished.

CaptureRun dataclass

CaptureRun(
    out_dir: Path,
    manifest: CaptureManifest,
    error: str | None = None,
)

What :func:run_capture returns.

Attributes:

Name Type Description
out_dir Path

The capture directory.

manifest CaptureManifest

The manifest as it was last written.

error str or None

Why the run stopped early, when it could not start at all (the first snapshot failed). None for a run that ended normally or was stopped by a signal -- gaps and failed segments along the way are in the manifest, not here.

CaptureConfig dataclass

CaptureConfig(
    pair: str,
    out_dir: Path,
    minutes: float = 10.0,
    keep_raw: bool = True,
    roll_minutes: float | None = None,
    roll_mb: float | None = None,
)

User-facing capture-run parameters.

Run-level knobs shared by every venue — what symbol, where to write, for how long. Per-venue settings (e.g. the CCXT exchange id) are not here; they are typed :class:~ob_analytics.config.SourceSettings carried by the source itself (CcxtSource(settings=CcxtSettings(...))), which is what replaced the former untyped extras dict.

read_manifest

read_manifest(path: Path) -> CaptureManifest | None

Read the manifest of the segmented capture at path.

path is the capture directory (or its manifest.json). Returns None when there is no manifest: a single capture made by :func:~ob_analytics.live._runner.run_capturer, or a processed output. Raises :class:ValueError for a manifest this version cannot read.

CaptureManifest dataclass

CaptureManifest(
    source: str,
    pair: str,
    level: str,
    started: Timestamp,
    declarations: dict[str, Any] = dict(),
    roll_minutes: float | None = None,
    roll_mb: float | None = None,
    ended: Timestamp | None = None,
    restarts: int = 0,
    segments: list[Segment] = list(),
    gaps: list[Gap] = list(),
)

What manifest.json records about a segmented capture.

Read one with :func:read_manifest. source, pair and level identify the capture: a restart into the same directory must match them. declarations is what the source declares about its feed (the same fields a segment's meta.json carries). restarts counts the times the capture was started again into this directory.

dropped property

dropped: int

Messages lost across every segment: unusable plus never received.

unfinished property

unfinished: list[Segment]

Segments a dead capture process left open (closed by the restart).

open_segments property

open_segments: list[Segment]

Segments not yet closed: still being captured, or left by a process that died and has not been restarted. Their files have no closing rows yet, and their row counts are not in the manifest.

write

write(root: Path) -> None

Write manifest.json into root, replacing the old one atomically.

The file is written beside the old one and renamed over it, so a reader -- or a restart after a crash -- never sees half a manifest.

note_streaming

note_streaming(segment: Segment) -> None

Record the gap, if any, before segment's first live event.

The segment before it with coverage either is still streaming (a roll: the two overlap, so there is no gap) or stopped before segment started to stream, and the time between is a gap.

note_uncovered_end

note_uncovered_end(ended: Timestamp) -> None

Record a gap if the capture ended while no segment was streaming.

The last segment that streamed is looked at. If it ran to the end (finished/stopped) nothing is missing. Otherwise -- a failure, or a roll whose next segment never streamed -- the time from its last event to the end of the capture is a gap.

segment_dirs

segment_dirs(root: Path) -> list[Path]

The directories under root of closed segments that hold data, in order.

root is the capture directory, or a process output made from it (which keeps the same segment names). A segment that wrote no book rows -- its snapshot failed -- is left out: there is nothing to read. So is an open segment (see :attr:open_segments): it has no closing rows yet, so the checks would read every order still resting as a fault.

checks

checks() -> tuple[QualityCheck, ...]

The capture-level checks ob-analytics audit adds to each segment's.

Each is a :attr:~ob_analytics.analytics.Severity.WARNING: a gap or a dropped message is data the capture does not have, which the segments on either side of it still describe correctly.

render

render() -> str

Return a fixed-width, human-readable report block.

Segment dataclass

Segment(
    name: str,
    started: Timestamp,
    ended: Timestamp | None = None,
    stream_started: Timestamp | None = None,
    stream_ended: Timestamp | None = None,
    heartbeat: Timestamp | None = None,
    end_reason: EndReason | None = None,
    error: str | None = None,
    n_book_events: int = 0,
    n_trade_events: int = 0,
    dropped: int = 0,
    sequence_missing: int = 0,
)

One segment of a capture, as the manifest records it.

Attributes:

Name Type Description
name str

The segment's directory name (seg-0001, ...).

started, ended Timestamp or None

When the segment was opened and closed. ended is None while it runs.

stream_started, stream_ended Timestamp or None

The stretch of time the segment covers: from its first live event to the moment it was asked to stop, or its stream stopped by itself. stream_started is None for a segment that never streamed (a snapshot that failed, say), or that was asked to stop before its first live event: it wrote rows, but covered nothing.

heartbeat Timestamp or None

Last time the running capture said the segment was alive. After a crash it is the latest time the segment is known to cover.

end_reason EndReason or None

Why it ended; None while it runs.

error str or None

The error that ended it, if any.

n_book_events, n_trade_events int

Rows written to the book file and to trades.csv.

dropped int

Messages the source received but could not use (its dropped counter in meta.json).

sequence_missing int

Venue sequence numbers the source never received (its sequence_missing counter), for sources that count them live.

Gap dataclass

Gap(
    after: str,
    before: str | None,
    start: Timestamp,
    end: Timestamp,
    cause: str,
)

A stretch of time no segment covers.

Attributes:

Name Type Description
after str

The segment the gap follows.

before str or None

The segment that ends it; None when the capture ended inside it.

start, end Timestamp

When coverage stopped and when it resumed (or the capture ended).

cause str

How the segment before it ended: "disconnect", "restart" (the capture process died), "roll" (a roll whose next segment was slow to stream, or never did), or "stopped" / "finished" (the capture was stopped and later started again into the same directory).

EndReason

Bases: str, Enum

Why a segment ended.

Attributes:

Name Type Description
FINISHED

The capture reached the end of its run time.

STOPPED

A signal (Ctrl-C, SIGTERM) stopped the capture.

ROLLED_TIME, ROLLED_SIZE

The segment reached its time or size limit. The next segment was streaming before this one stopped, so a roll loses nothing.

FAILED

The source raised: most often a lost connection.

ENDED_EARLY

The source stopped streaming before it was asked to, with no error.

UNFINISHED

The capture process died while the segment was open. The restart closed it: see :func:~ob_analytics.live.run_capture.