Skip to content

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]) -> Finding

MonitoringLedger

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) -> bool

Close 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) -> bool

Close 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) -> bool

Attach 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) -> str

Upsert 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) -> str

Idempotently open this run and return its id.

start_stage

start_stage(stage_name: str, attempt: int = 1, started_at: dt.datetime | None = None) -> str

Idempotently open one named attempt and return its stable id.