Skip to main content

deployment

The served Prefect deployment that runs a background v9 DAG.

The pod must stay responsive while a DAG runs: its heartbeat, message handlers and liveness tickers share one event loop, and a step that blocks that loop takes the pod's datasets offline on the Hub. A separate process makes that impossible by construction, and a served deployment is how the rest of this system already gets one — file_metadata, scan_metadata and file_metadata_refresh all run this way.

Every parameter here is JSON, because Prefect persists deployment parameters in its database and hands them to the subprocess as JSON. That is also why neither a credential nor ehr_config is among them: the first must not be written to that database in clear, and the second carries a private key. Both are fetched at run time from the pod control server instead, against the nonce the pod left in the cache for this lineage (flows.dag.child_auth).

The ordering in the flow body is therefore fixed: pod_key_path is a parameter because opening the cache needs it for the encryption binding, and the nonce lives in the cache — anything else is circular.

State that cannot cross a process boundary is rebuilt here: the Hub session, the cache handle, the datasource, the pod's hooks and the telemetry buffers. The DAG itself is rebuilt from the FlowSpec the triggering link persisted, so the child and the pod agree by construction rather than by a pickle.

Module​

Functions​

background_dag_runtime​

async def background_dag_runtime(    pod_name: str,    pod_identifier: str,    hub_url: str,    task_hash: str,    project_id: str,    datasource_name: str,    monitor_task_id: str,    run_id: str,    trigger: str,    is_rerun: bool,    username: str | None = None,    pod_key_path: str | None = None,    concurrency_slot: str | None = None,    slot_timeout_seconds: float | None = None,    pod_control_url: str | None = None,) ‑> dict[str, Any] | State:

Run one background v9 DAG, in the subprocess Prefect launched for it.

This flow run is the run — the duplicate guard, the concurrency slot and the context assembly all happen in its body (see flows.dag.executor.execute_background_dag), so a setup failure is a failed flow run rather than a dropped trigger, and a duplicate is a terminal skipped run rather than a silent return.

The body is delegated to rather than nested for a reason: this run already carries the lineage's tags, and a subflow carrying them too would be read by the duplicate guard as a second live run of the same lineage.

Arguments

  • pod_name: Pod whose background cache this run writes to.
  • pod_identifier: username/pod_name, for Hub reporting and telemetry.
  • hub_url: Hub base URL.
  • task_hash: The (pod, datasource) hash keying this run.
  • project_id: Project the datasource is linked to.
  • datasource_name: Datasource this run processes.
  • monitor_task_id: Hub monitor task id to report under.
  • run_id: The id this run's DAG-level Run row is opened under. Minted by the pod, not here: a run that dies cannot close its own rows, and the id is what lets the pod close this run's rows without touching a sibling's on the same task_hash.
  • trigger: Why this run was started, as a BackgroundTrigger value.
  • is_rerun: Whether to apply rerun semantics — hold the concurrency slot and stamp the lineage.
  • username: Hub username, used to rebuild the session and tag telemetry.
  • pod_key_path: Path to pod_rsa.pem. Bound into the cache-encryption layer, and the source of the pod keys the datasource manager needs.
  • concurrency_slot: Prefect concurrency limit to occupy, if any.
  • slot_timeout_seconds: How long to wait for that slot.
  • pod_control_url: Base URL of the pod control server that launched this run, where its credentials and EHR configuration come from.

Returns The DAG's step results, or a terminal State when nothing ran.

Raises

  • ConfigError: If no flow spec is recorded for this lineage, or the run was given no pod_key_path to build its datasource from.