Skip to main content

store

CRUD for the runs record (v1).

Records the lifecycle of each flow execution and provides the key queries used by BatchDiscovery and trigger flows.

Module​

Functions​

get_last_refreshed​

def get_last_refreshed(    cache: CacheProtocol, task_hash: str, run_type: str,) ‑> datetime.datetime | None:

Return the completed_at of the most recent successful run of run_type.

Scoped by type, not by task_hash alone: every run of a datasource shares one task_hash — the indexing runtimes, the background DAG and each of its nodes — so an unscoped query answers "has anything at all completed for this datasource", which is never the question a caller means.

Arguments

  • cache: Cache instance.
  • task_hash: Pipeline task hash.
  • run_type: The type column value to match, e.g. "file_metadata" or a DAG or node name. See runs.v1.schema.

Returns The timestamp, or None if no run of that type has ever completed.

get_run_by_run_id​

def get_run_by_run_id(    cache: CacheProtocol, run_id: str,) ‑> RunRecord | None:

Return a single run record by its UUID, or None if not found.

Arguments

  • cache: Cache instance.
  • run_id: UUID of the run to retrieve.

Returns The matching RunRecord, or None if no run with that ID exists.

get_runs_by_task​

def get_runs_by_task(    cache: CacheProtocol, task_hash: str,) ‑> CollectedRunQueryResult:

Return all run records for task_hash, most recent first.

Arguments

  • cache: Cache instance.
  • task_hash: Pipeline task hash.

Returns A CollectedRunQueryResult containing all run records ordered by started_at descending. The result is falsy when no runs exist.

mark_run_complete​

def mark_run_complete(cache: CacheProtocol, run_id: str, files_processed: int) ‑> None:

Record successful completion of a running flow.

Arguments

  • cache: Cache instance.
  • run_id: UUID of the run to mark complete (returned by mark_run_started).
  • files_processed: Number of files processed.

mark_run_failed​

def mark_run_failed(cache: CacheProtocol, run_id: str, error: str | None = None) ‑> None:

Record a failed flow.

Arguments

  • cache: Cache instance.
  • run_id: UUID of the run to mark failed (returned by mark_run_started).
  • error: Optional error message.

mark_run_started​

def mark_run_started(    cache: CacheProtocol,    task_hash: str,    run_type: str,    stale_threshold_hours: float = 1.0,    tags: dict[str, Any] | None = None,    parent_run_id: str | None = None,    project_id: str | None = None,    dedup: bool = True,) ‑> str | None:

Record that a flow has started for task_hash, or return None.

Uses a single atomic INSERT ... SELECT [WHERE NOT EXISTS] statement so that the "is there already a live run?" check and the "insert this run" step are never separated by a window where another writer can sneak in. This eliminates the TOCTOU race that would exist with a plain SELECT-then-INSERT approach.

A "running" row is considered stale (i.e. orphaned by a crash) when it is older than stale_threshold_hours. Stale rows are ignored by the guard — they should be reaped by _is_run_active in the refresh path.

The guard is scoped by (task_hash, run_type, project_id) so that a running DAG-level run (run_type=dag.name) does not block concurrent node-level runs (run_type=node.name), and a run for one project does not block a run of the same (task_hash, run_type) for a different project linked to the same datasource.

When dedup is False the guard is bypassed and a bookkeeping row is always inserted (the function then never returns None). This is the mode the file_metadata / scan_metadata runtimes use: mutual exclusion now lives at the trigger layer backed by the heartbeat

  • zombie automation, so the runs row is pure bookkeeping. Keeping the row's own guard for those runtimes would re-introduce the crash-lockout it was demoted to fix — a crashed run's stale running row would block a fresh run for up to stale_threshold_hours. All other callers (DAG / node / interactive runs) keep the default dedup=True.

Arguments

  • cache: Cache instance.
  • task_hash: Pipeline task hash.
  • run_type: Required, free-form run-kind discriminator (e.g. "file_metadata", "interactive", a DAG name, or a node name). The dedup guard is scoped by (task_hash, run_type), so this must be specific enough that two runs sharing it are genuinely mutually-exclusive duplicates. There is deliberately no default — a coarse default previously broke the contract.
  • stale_threshold_hours: How old (in hours) a running row must be before it is considered stale and no longer blocks a new run. Defaults to 1 hour.
  • tags: Optional flat metadata dict to attach to this run record.
  • parent_run_id: Optional UUID of the parent run that spawned this run (e.g. the DAG-level run for a node-level run).
  • project_id: Project that owns this run, folded into the dedup scope. None for project-independent runs (e.g. file-metadata indexing). Two projects sharing a (task_hash, run_type) get independent runs because their project_id differs.
  • dedup: When True (default) the INSERT … WHERE NOT EXISTS guard is applied, so a second concurrent start returns None. When False the guard is bypassed and a bookkeeping row is always inserted (the function never returns None); the metadata runtimes use this because mutual exclusion lives at the trigger layer.

Returns The UUID string assigned to this run as run_id, or None when a non-stale "running" row already exists for the same (task_hash, run_type, project_id) (only possible when dedup=True).

Classes​

CollectedRunQueryResult​

class CollectedRunQueryResult(rows: Iterable[RunRecord]):

Materialised run query results.

Eagerly collects all run rows for a task hash into memory. Rows are ordered by started_at descending so records[0] is always the most recent run.

Variables​

  • records : list[RunRecord] - Return all run records, most recent first.

Methods​


as_df​

def as_df(self) ‑> pandas.core.frame.DataFrame:

Return all run records as a flat DataFrame.