betting_combat.store.monitoring
The run ledger in monitoring.betting_combat: run, stage, target_receipt, finding_summary.
This repository’s own port of common-dagster’s MonitoringRepo (the calls a step and the
common_dagster wrapper make: start_run, finish_run, start_stage,
finish_stage, record_target_receipt, record_check_result,
record_finding_summaries), written through this repository’s own client. The SQL, the stable
ids (sha256 of the same parts) and the rules are the same, so the rows common-dagster writes
itself (its run-status sensors and check results, through common-db) and the rows the steps write
here line up (tests/store/test_sql.py):
- arguments are validated first (a bad call raises
ValueError); - a write that fails is logged and returns
False— observability never fails the workload (the monitoring database has a 5 s connect budget and a 10 s statement budget,conf.settings’postgres_monitoring_statement_timeout_ms); - ids are deterministic, so every write is idempotent.
Schema: betting_combat (common-dagster’s rule: the schema is the container name, so the
rows common_dagster writes itself land in the same tables). Difference from it: the
findings are Finding values (finding_type, severity, finding_count) or
plain mappings with those keys, validated the same way.
Classes
Finding
One grouped finding: a kind, its severity and how many times it occurred.
of
of(value: Finding | Mapping[str, object]) -> FindingMonitoringLedger
The four ledger tables of one run, written through db.monitoring.
finish_run
finish_run(status: str, error: str | None = None, finished_at: dt.datetime | None = None) -> boolClose the run once (and any stage left open under it).
finish_stage
finish_stage(stage_id: str, status: str, message: str | None = None, finished_at: dt.datetime | None = None) -> boolClose one exact attempt without erasing prior failed attempts.
record_check_result
record_check_result(receipt_id: str, checks_total: int, checks_failed: int, message: str | None = None) -> boolAttach observed check counts without changing target completion.
record_finding_summaries
record_finding_summaries(stage_id: str, findings: Sequence[Finding | Mapping[str, object]], receipt_id: str | None = None, recorded_at: dt.datetime | None = None) -> tuple[str, ...]Persist grouped findings (one row per kind and severity, never one per item).
record_target_receipt
record_target_receipt(stage_id: str, target_type: str, target_name: str, target_key: str | None = None, status: str = 'complete', rows_written: int | None = None, rows_rejected: int | None = None, rows_total_after: int | None = None, checks_total: int | None = None, checks_failed: int | None = None, message: str | None = None, recorded_at: dt.datetime | None = None) -> strUpsert one idempotent receipt for the exact target slice; returns its id.
start_run
start_run(flow_name: str, job_name: str, job_type: str, run_date: dt.date | None = None, started_at: dt.datetime | None = None) -> strIdempotently open this run and return its id.
start_stage
start_stage(stage_name: str, attempt: int = 1, started_at: dt.datetime | None = None) -> strIdempotently open one named attempt and return its stable id.