A data pipeline can finish with a green Airflow DAG run while a business-critical task inside that run has failed. In the specific case examined here, the mechanism is not mysterious: an all_done cleanup task is the only leaf, it succeeds after an earlier failure, and Airflow's documented DAG-run aggregation therefore sees a successful leaf. The Airflow all_done DAG run success is real scheduler state, but it does not necessarily prove successful business processing. Apache Airflow 3.1.5 explicitly warns about this leaf-state behavior.
That creates two different acceptance questions. First, does the graph propagate the failures the team intends to block? Second, does independent evidence show that the required business output actually exists and is correct? An Airflow leaf task failure masking test answers the first. An output-contract test answers the second. Neither is a substitute for the other.
This playbook builds both tests around Airflow 3.1.5 and its Task SDK 1.1.5. It compares the original cleanup-only graph, a deliberately inadequate watcher, and a correctly wired Airflow watcher pattern validation. It then adds an independent manifest gate that catches a task returning success without delivering its required output.
No lab run was executed for this article. Every result below is therefore labeled expected, not observed. The acceptance choices are operationally explicit: ACCEPT a qualified graph, REWIRE defective failure propagation, HOLD incomplete or contradictory evidence, or REPLAY approved business work after repair and reconciliation.
Decide What a Successful Run Must Prove
Start by refusing to use the word “success” without naming the contract it refers to.
Airflow's DAG-run state is an orchestration contract. The business output is an application contract. Cleanup is a resource-handling contract. A release gate needs all three. This is a narrower problem than the general role of Airflow within the data-engineering tool stack: the question here is whether one tested DAG revision reports the states that its owners intend.
For this fixture, define the required business path as:
extract -> transform -> publish -> cleanupextract, transform, and publish are required business tasks. cleanup is required operational work and is deliberately allowed to run after upstream failure. The release policy in this lab also treats a cleanup failure as blocking; another organization could choose a different cleanup policy, but it must be declared rather than inferred afterward.
The business contract is equally concrete:
expected IDs: R-001, R-002, R-003
expected count: 3
expected transformed values: 100, 200, 300
expected manifest: present under this run's output namespace
expected checksum: recomputed from the canonical transformed recordsThe acceptance decisions then become testable:
Decision | Meaning |
ACCEPT | Graph coverage is correct, actual final states conform to the tested case, and required output/cleanup evidence passes. |
REWIRE | The graph can produce a misleading terminal run state because required failures do not reach a blocking leaf. |
HOLD | Evidence is absent, contradictory, incomplete, or the output contract fails despite apparently acceptable scheduler state. |
REPLAY | A corrected revision is already qualified, prior side effects are reconciled, and specifically approved business work may be run again. |
Under the Airflow 3.1.5 DAG-run contract, the run becomes terminal after its tasks are terminal, and its final state is determined from its leaf tasks. A run succeeds when its leaves are success or skipped; it fails when a leaf is failed or upstream_failed. The same documentation specifically warns that a successful all_done leaf can make the run successful even when an earlier task failed.
That documented behavior is the starting condition, not evidence that any particular unexecuted fixture produced the expected result.
Freeze a Disposable Scheduler-Backed Baseline
A trigger-rule qualification is meaningful only when the tested runtime can be reconstructed. Pin the environment before interpreting a color in the UI.
Airflow 3.1.5 is the controlled baseline required for this acceptance exercise, not a claim about the latest release. The 3.1.5 package metadata pins apache-airflow-core==3.1.5 and apache-airflow-task-sdk==1.1.5; the corresponding Task SDK documentation identifies airflow.sdk as the Airflow 3 DAG-authoring interface.
Use this environment manifest as the proposed qualification baseline:
lab_environment:
airflow: "3.1.5"
airflow_core: "3.1.5"
airflow_task_sdk: "1.1.5"
airflow_provider_standard: "1.10.0"
provider_usage_in_dag: "none"
python: "3.12.12"
executor: "LocalExecutor"
metadata_database: "PostgreSQL 16.15"
default_timezone: "UTC"
dag_runtime:
schedule: null
catchup: false
start_date: "2026-09-01T00:00:00Z"
retries: 0
artifact_identity:
dag_revision_a: "all_done-lab/A"
dag_revision_b: "all_done-lab/B"
dag_revision_c: "all_done-lab/C"
deployment_artifact: "record immutable image/repository digest at execution"
storage:
shared_root: "/opt/airflow/lab-shared"
ownership: "disposable acceptance-test resources only"The provider pin above is an environment pin from the Airflow 3.1.5 Python 3.12 constraints; the example DAG itself uses no provider operator. The Python version is a selected lab pin within Airflow 3.1.5's supported Python range, not a general recommendation.
The deployment artifact must become more precise when the lab is actually executed: record the immutable container digest or equivalent build identifier rather than accepting a mutable tag. Until that value is captured, an evidence bundle is incomplete and should be HOLD, not quietly upgraded to “reproducible.”
This scheduler-backed environment is intentionally boring. LocalExecutor is fixed so executor behavior is not another experimental variable. PostgreSQL is fixed so metadata storage is not another variable. Retries are zero so a deliberate failure is a final failure rather than an intermediate up_for_retry.
The shared path matters even with this simple baseline. Airflow's best-practices guidance warns against assuming local files are shared when tasks may run on different machines and recommends common storage for data that must cross task boundaries. The lab therefore uses one explicitly mounted, owned output location accessible to all test tasks. This is the concrete qualification fixture; workflow orchestration in the broader data platform is useful context, but it cannot replace a frozen runtime manifest.
Construct the Four-Task Cleanup Graph
The fixture must be small enough that state propagation can be read directly from the edges. The following DAG file constructs all three revisions from the same business implementation.
Revision A has cleanup as its only leaf. Revision B adds the deliberately insufficient watcher only after cleanup. Revision C gives the watcher direct upstream edges from the entire approved monitored population.
The DAG uses the Airflow 3 Task SDK interfaces for DAG, task, and runtime context, consistent with the Task SDK 1.1.5 public authoring interface.
# dags/all_done_acceptance_lab.py
from future import annotations
import hashlib
import json
import os
from pathlib import Path
from typing import Any
import pendulum
from airflow.exceptions import AirflowException
from airflow.sdk import DAG, get_current_context, task
from airflow.utils.trigger_rule import TriggerRule
SHARED_ROOT = Path(
os.environ.get("AIRFLOW_LAB_SHARED_ROOT", "/opt/airflow/lab-shared")
)
BUSINESS_TASK_IDS = ("extract", "transform", "publish")
MONITORED_TASK_IDS = ("extract", "transform", "publish", "cleanup")
ALLOWED_FAULTS = {"none", "transform", "publish", "cleanup", "missing_output"}
INPUT_RECORDS = [
{"id": "R-001", "value": 10},
{"id": "R-002", "value": 20},
{"id": "R-003", "value": 30},
]
def canonical_bytes(value: Any) -> bytes:
return json.dumps(
value, sort_keys=True, separators=(",", ":")
).encode("utf-8")
def checksum(value: Any) -> str:
return hashlib.sha256(canonical_bytes(value)).hexdigest()
def run_directory() -> Path:
ctx = get_current_context()
run_id = str(ctx["run_id"])
dag_id = str(ctx["dag"].dag_id)
# Never use arbitrary run_id text as a filesystem path component.
run_key = hashlib.sha256(run_id.encode("utf-8")).hexdigest()[:24]
return SHARED_ROOT / dag_id / run_key
def configuration() -> dict[str, Any]:
ctx = get_current_context()
dag_run = ctx["dag_run"]
conf = dict(dag_run.conf or {})
fault = str(conf.get("fault", "none"))
if fault not in ALLOWED_FAULTS:
raise AirflowException(f"Unsupported lab fault: {fault}")
return {"fault": fault}
def write_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_suffix(path.suffix + ".tmp")
temporary.write_text(
json.dumps(payload, indent=2, sort_keys=True) + "\n",
encoding="utf-8",
)
temporary.replace(path)
def build_lab_dag(revision: str) -> DAG:
if revision not in {"A", "B", "C"}:
raise ValueError(f"Unsupported revision: {revision}")
dag_id = f"all_done_acceptance_lab_{revision.lower()}"
with DAG(
dag_id=dag_id,
schedule=None,
catchup=False,
start_date=pendulum.datetime(2026, 9, 1, tz="UTC"),
default_args={"retries": 0},
tags=["acceptance-lab", f"revision-{revision.lower()}"],
) as dag:
@task(task_id="extract", retries=0)
def extract() -> list[dict[str, Any]]:
ctx = get_current_context()
root = run_directory()
records = [dict(record) for record in INPUT_RECORDS]
write_json(
root / "audit" / "input_manifest.json",
{
"dag_id": ctx["dag"].dag_id,
"run_id": ctx["run_id"],
"dag_revision": revision,
"ids": [r["id"] for r in records],
"count": len(records),
"input_checksum": checksum(records),
},
)
tmp = root / "temporary" / "staging.marker"
tmp.parent.mkdir(parents=True, exist_ok=True)
tmp.write_text("owned disposable resource\n", encoding="utf-8")
return records
@task(task_id="transform", retries=0)
def transform(records: list[dict[str, Any]]) -> list[dict[str, Any]]:
if configuration()["fault"] == "transform":
raise AirflowException("LAB_INJECTED: transform failure")
return [
{
"id": record["id"],
"normalized_value": record["value"] * 10,
}
for record in records
]
@task(task_id="publish", retries=0)
def publish(records: list[dict[str, Any]]) -> dict[str, Any]:
ctx = get_current_context()
root = run_directory()
fault = configuration()["fault"]
if fault == "publish":
raise AirflowException("LAB_INJECTED: publish failure")
# Negative correctness control: task returns normally but does
# not create the required output manifest.
if fault == "missing_output":
write_json(
root / "audit" / "publish_negative_control.json",
{
"run_id": ctx["run_id"],
"dag_revision": revision,
"required_output_intentionally_omitted": True,
},
)
return {"published": False, "negative_control": True}
output = {
"dag_id": ctx["dag"].dag_id,
"run_id": ctx["run_id"],
"dag_revision": revision,
"ids": [r["id"] for r in records],
"count": len(records),
"records": records,
"output_checksum": checksum(records),
}
# The business manifest appears only after successful processing.
write_json(root / "business" / "output_manifest.json", output)
return {"published": True, "output_checksum": output["output_checksum"]}
@task(
task_id="cleanup",
retries=0,
trigger_rule=TriggerRule.ALL_DONE,
)
def cleanup() -> None:
ctx = get_current_context()
root = run_directory()
tmp = root / "temporary" / "staging.marker"
existed = tmp.exists()
if existed:
tmp.unlink()
write_json(
root / "audit" / "cleanup.json",
{
"run_id": ctx["run_id"],
"dag_revision": revision,
"temporary_resource_existed": existed,
"temporary_resource_handled": not tmp.exists(),
"business_audit_artifacts_deleted": False,
},
)
if configuration()["fault"] == "cleanup":
raise AirflowException("LAB_INJECTED: cleanup failure")
extracted = extract()
transformed = transform(extracted)
published = publish(transformed)
cleaned = cleanup()
published >> cleaned
if revision in {"B", "C"}:
@task(
task_id="watcher",
retries=0,
trigger_rule=TriggerRule.ONE_FAILED,
)
def watcher() -> None:
# Running is itself evidence that one direct monitored
# parent met ONE_FAILED. The watcher must then remain red.
raise AirflowException(
"Monitored task failure propagated by watcher"
)
watcher_arg = watcher()
watcher_task = dag.get_task("watcher")
if revision == "B":
# Intentionally insufficient negative control.
cleaned >> watcher_arg
if revision == "C":
# Explicit direct edges; watcher is never in this list.
for task_id in MONITORED_TASK_IDS:
dag.get_task(task_id) >> watcher_task
return dag
dag_a = build_lab_dag("A")
dag_b = build_lab_dag("B")
dag_c = build_lab_dag("C")The failure injections are conspicuous and limited to owned inert resources. Nothing catches AirflowException, manually changes a task state, reads the metadata database from task code, or manufactures a successful result.
Keep Audit Output Outside the Cleanup Target
The temporary/ directory and the acceptance artifacts deliberately have different ownership semantics.
cleanup may remove temporary/staging.marker. It must not remove audit/input_manifest.json, audit/cleanup.json, or business/output_manifest.json. That separation prevents a successful cleanup from destroying the very artifacts needed to determine whether publish happened.
Run scoping also matters. The filesystem directory is derived from a hash of run_id; the original run_id remains inside the manifests. One run therefore cannot satisfy another run's acceptance check merely because a previous manifest happened to exist at a shared global path.
For a real object store, database schema, or warehouse, use the same principle: an immutable or versioned evidence namespace and a separately owned temporary namespace. This lab does not claim that a filesystem is appropriate production publication architecture; it is inert shared storage chosen to make the acceptance contract observable.
Reproduce a Green Run With a Failed Business Task
First establish the positive control. Trigger Revision A with fault=none. The manually derived expectation is:
Task | Expected final state | Output effect |
extract | success | Input audit manifest present |
transform | success | Deterministic transformed records returned |
publish | success | Business output manifest present |
cleanup | success | Temporary marker removed; cleanup record present |
DAG run | success | Sole leaf cleanup succeeded |
Now inject the transform failure without changing anything else:
airflow dags trigger all_done_acceptance_lab_a \
--run-id "lab-a-transform-fail-001" \
--conf '{"fault":"transform"}'This command is a runnable instruction, not a claim that it was executed for this article.
The expected states are:
Node | Expected state |
extract | success |
transform | failed |
publish | upstream_failed |
cleanup | success |
DAG run | success |
Required business manifest | absent |
Airflow distinguishes failed, upstream_failed, skipped, queued, scheduled, and retry-related states; those terms should not be collapsed into “didn't run.” The Airflow 3.1.5 task-state documentation defines upstream_failed specifically for a task prevented from running because an upstream task failed or became upstream_failed. It separately defines scheduled, queued, running, success, failed, skipped, and up_for_retry.
The fixture's expected sequence is therefore straightforward. transform raises. publish, whose normal trigger rule requires its upstream to succeed, cannot execute and becomes upstream_failed. cleanup is allowed to execute once its upstream is done regardless of outcome, and the cleanup callable itself succeeds.
The important finding is not merely “a task failed.” It is the expected contradiction:
Orchestration summary: DAG run success
Business evidence: required output absent
Task evidence: transform failedThat is a REWIRE result for Revision A's failure-reporting design, even though the injected-failure run was expected to demonstrate precisely this defect.
Read the Leaves Without Ignoring the Rest of the Graph
The expected green DagRun is derived independently from Airflow's documented leaf rule, not from assumptions about a UI color.
In Revision A:
extract -> transform -> publish -> cleanup
[leaf]Only cleanup has no downstream task. The Airflow 3.1.5 DAG-run documentation says the DAG run's terminal state is based on those leaf nodes and calls out the exact risk of an all_done leaf succeeding after a mid-DAG failure.
That does not mean every DAG containing all_done cleanup is defective. The defect is the combination of graph shape and release intent: the only terminal evidence Airflow aggregates does not encode a failure the owners intended the run to block.
A green run is therefore one evidence field. It is not permission to ignore the task ledger or output contract.
Make a Plausible Watcher Repair Fail First
A common repair intuition is: “Add a watcher after cleanup. Give it one_failed. Now failures will make the run red.”
Revision B deliberately demonstrates why that is incomplete:
extract -> transform -> publish -> cleanup -> watcher
success ONE_FAILEDThe watcher has exactly one direct parent: cleanup.
Suppose transform fails. publish becomes upstream_failed, while cleanup runs under all_done and succeeds. From the watcher's perspective, its only parent succeeded. There is no failed direct upstream parent satisfying one_failed, so the watcher is expected to be skipped.
Airflow's 3.1.5 Best Practices watcher guidance emphasizes precisely this dependency rule: trigger-rule evaluation concerns direct upstream tasks, and a watcher must be connected appropriately to the tasks it is intended to monitor. The documentation notes that a failure in a non-direct ancestor is not enough just because the watcher is transitively downstream.
The expected negative-control result is:
Revision B, transform failure | Expected |
transform | failed |
publish | upstream_failed |
cleanup | success |
watcher | skipped |
Leaf | watcher |
DAG run | success |
Output manifest | absent |
Qualification decision | REWIRE |
This is intentionally not the recommended topology. It exists because a good acceptance suite should disprove a plausible wrong solution before approving the right one.
The distinction is easy to obscure in a beginner example where the focus is simply constructing dependencies; that introductory Airflow project context serves a different purpose. Here, the acceptance question is not whether the arrows form a syntactically valid DAG. It is whether every intended blocking state has a direct route into the terminal failure-propagation mechanism.
Wire the Watcher to Its Full Monitored Population
Revision C retains cleanup(trigger_rule="all_done") and changes only the watcher coverage:
extract ------\
| \
transform ------\
| \
publish ----------> watcher [ONE_FAILED, always raises]
| /
cleanup --------/The business chain still remains:
extract -> transform -> publish -> cleanupThe additional arrows are failure-observation dependencies. They do not replace the business edges.
MONITORED_TASK_IDS is explicit:
MONITORED_TASK_IDS = (
"extract",
"transform",
"publish",
"cleanup",
)For each ID, the DAG builder adds a direct edge to watcher. The watcher itself is intentionally absent from that list, avoiding a self-dependency or cycle.
This follows the Airflow 3.1.5 watcher pattern: a ONE_FAILED watcher is downstream of the tasks it monitors and deliberately fails when triggered. The Airflow trigger-rule documentation defines one_failed against upstream failures and, importantly, says it does not wait for every upstream task to finish once its condition is met.
That timing fact changes how the graph should be described. Do not say “the watcher waits for everything and then decides the run.” It may become runnable as soon as a direct parent supplies the triggering failure condition. Nor does the watcher's failure cancel cleanup. cleanup remains independently governed by its own upstream dependencies and all_done trigger rule.
Likewise, do not call the run terminal just because the watcher has failed. The run reaches its terminal state only when the run-level completion criteria are satisfied. A watcher is failure propagation, not a transactional abort switch.
Direct Parents Are the Coverage Boundary
Do not trust a visual inspection of arrows as the sole coverage test. Assert the graph mechanically when the DAG is parsed or tested.
Put the following independent test in the repository. Notice that the approved list is repeated in test policy rather than imported from MONITORED_TASK_IDS; importing the implementation's own list and comparing it with itself would prove little.
# tests/test_all_done_graph.py
from dags.all_done_acceptance_lab import dag_a, dag_b, dag_c
APPROVED_C_TASKS = {
"extract",
"transform",
"publish",
"cleanup",
}
WATCHER = "watcher"
def test_revision_a_has_cleanup_as_only_leaf() -> None:
leaves = {task.task_id for task in dag_a.leaves}
assert leaves == {"cleanup"}
def test_revision_b_is_the_intended_negative_control() -> None:
watcher = dag_b.get_task(WATCHER)
assert watcher.upstream_task_ids == {"cleanup"}
def test_revision_c_watcher_directly_covers_all_approved_tasks() -> None:
watcher = dag_c.get_task(WATCHER)
actual_non_watcher = set(dag_c.task_ids) - {WATCHER}
assert actual_non_watcher == APPROVED_C_TASKS
assert watcher.upstream_task_ids == APPROVED_C_TASKS
assert WATCHER not in watcher.upstream_task_ids
# For this controlled revision, watcher must be the only leaf.
leaves = {task.task_id for task in dag_c.leaves}
assert leaves == {WATCHER}This test performs two distinct checks.
First, it verifies that the expected business/cleanup population has not silently changed. If somebody adds charge_customer, write_ledger, or even an innocently named postprocess task, actual_non_watcher no longer equals the approved set. The test fails until ownership decides whether that task is monitored.
Second, it verifies direct watcher parentage. A newly added task cannot be considered covered merely because there happens to be a transitive route through cleanup.
That turns watcher dependency coverage into a release artifact rather than a convention remembered by the original DAG author.
Test the All-Success and Cleanup-Failure Controls
A watcher that always fails must also demonstrate that it does not run on a clean path.
For Revision C with fault=none, the expected states are:
Task | Expected final state |
extract | success |
transform | success |
publish | success |
cleanup | success |
watcher | skipped |
DAG run | success |
Output contract | pass |
A skipped watcher is correct here. None of its direct parents failed, so ONE_FAILED is unsatisfied. Airflow's Best Practices watcher example similarly describes the watcher being skipped when the monitored tasks succeed.
Next test fault=cleanup.
Expected behavior:
extract success
transform success
publish success
cleanup failed
watcher failed
DagRun failed
manifest presentpublish has already completed, so a business output may exist. That does not make the run acceptable under this lab's declared policy because cleanup itself belongs to MONITORED_TASK_IDS. A cleanup failure is therefore release-blocking.
This distinction should be written down instead of hidden in trigger-rule mechanics. Some cleanup failures may represent leaked temporary resources that absolutely must block promotion. Another workflow might decide cleanup is non-blocking and alert separately. Either can be engineered; what is unacceptable is allowing graph topology to make the policy accidentally.
The minimum expected fault matrix across the three revisions is:
Case | Revision A | Revision B | Revision C | Output contract |
All business tasks succeed | run success | run success | run success | pass |
Transform fails | run success | run success | run failed | fail/absent |
Publish fails | run success | run success | run failed | fail/absent |
Cleanup fails | run failed | watcher fails; run failed | watcher fails; run failed | may pass |
Optional task skipped | extension only | extension only | expected run success if skip allowed | pass |
Publish returns success but required output absent | run success | run success | run success | fail |
Every cell is an expected outcome until executed. No entry above is an observed scheduler result.
Revision C earns qualification only if actual scheduler-backed evidence later matches these expectations.
Keep Skips and Retry States Out of the Wrong Bucket
skipped, failed, and upstream_failed carry different meanings. Airflow documents them as different task-instance states; a qualification ledger should preserve that vocabulary rather than reduce everything to red/green.
Keep the optional-branch test isolated from the four-task baseline. For example, a separate extension can add a controlled branch whose optional work is deliberately not selected:
# Isolated extension only; do not add this to baseline A/B/C.
from airflow.sdk import task
@task.branch(task_id="choose_optional_path", retries=0)
def choose_optional_path() -> str:
return "optional_bypass"
@task(task_id="optional_work", retries=0)
def optional_work() -> None:
raise RuntimeError("Should be skipped in this control")
@task(task_id="optional_bypass", retries=0)
def optional_bypass() -> None:
pass
choice = choose_optional_path()
work = optional_work()
bypass = optional_bypass()
choice >> [work, bypass]In a fully qualified extension, both branch structure and watcher coverage must be declared explicitly. The expected optional_work=skipped state is not a failed task merely because its callable did not execute.
Airflow's trigger-rule documentation also describes how skipped states interact with downstream trigger rules. That matters because a legitimate optional skip and a required business task that unexpectedly becomes skipped are operationally different even if both have the same Airflow state.
A release contract therefore needs a per-task policy, not just a global set of “allowed terminal states.”
For the core fixture, retries remain zero. Airflow defines up_for_retry as an intermediate state after failure when another attempt remains. A later retry-enabled qualification must judge the final task-instance outcome and preserve attempt history. It must not interpret attempt one failing as a permanent business failure if attempt two is validly scheduled and succeeds.
Conversely, it must not hide failed attempts from the audit trail. “Final state succeeded after retry” and “never failed” are different operational histories even though both may satisfy the release contract.
A Skipped Watcher Is Not a Failed Watcher
Interpret final states against the role of each task:
Final state | Acceptance interpretation |
Required business task success | Necessary but not sufficient; validate its output contract. |
Required business task failed | Blocking failure. |
Required business task upstream_failed | Required work did not complete; blocking unless policy explicitly says otherwise. |
Required business task skipped | Usually blocking unless that exact skip is an approved condition. |
Optional task skipped | Acceptable when the branch policy expected it. |
Watcher failed | Failure was propagated to the run-level blocking surface. |
Watcher skipped | No direct monitored parent satisfied ONE_FAILED; not proof of output correctness. |
Task scheduled or queued | Not evidence that its callable executed. |
Task up_for_retry | Non-final in a retry-enabled qualification. |
That last distinction prevents a subtle evidence error. A row showing queued proves that Airflow reached a scheduling stage; it does not prove business code ran. The acceptance ledger must record actual final state and attempt evidence rather than equating “task instance existed” with “task executed.”
Validate Business Outputs Independently of Task Success
Revision C fixes state propagation for failures represented as failed monitored task states. It does not make success synonymous with correct data.
The missing_output fault proves the boundary. In that mode, the publish callable deliberately returns normally but does not create business/output_manifest.json.
Expected states:
extract success
transform success
publish success
cleanup success
watcher skipped
DagRun successThat is correct Airflow behavior for the code that actually ran. There is no failed task for ONE_FAILED to observe.
The external gate must still reject the run.
A minimal independent validator can be kept outside the DAG:
# tools/validate_business_output.py
from future import annotations
import hashlib
import json
import sys
from pathlib import Path
EXPECTED_IDS = ["R-001", "R-002", "R-003"]
EXPECTED_RECORDS = [
{"id": "R-001", "normalized_value": 100},
{"id": "R-002", "normalized_value": 200},
{"id": "R-003", "normalized_value": 300},
]
def canonical_bytes(value: object) -> bytes:
return json.dumps(
value, sort_keys=True, separators=(",", ":")
).encode("utf-8")
def expected_checksum() -> str:
return hashlib.sha256(canonical_bytes(EXPECTED_RECORDS)).hexdigest()
def validate(path: Path, expected_run_id: str) -> None:
if not path.is_file():
raise SystemExit(f"HOLD: required manifest absent: {path}")
manifest = json.loads(path.read_text(encoding="utf-8"))
checks = {
"run_id": manifest.get("run_id") == expected_run_id,
"ids": manifest.get("ids") == EXPECTED_IDS,
"count": manifest.get("count") == 3,
"records": manifest.get("records") == EXPECTED_RECORDS,
"checksum": manifest.get("output_checksum") == expected_checksum(),
}
failed = [name for name, ok in checks.items() if not ok]
if failed:
raise SystemExit(
"HOLD: output contract mismatch: " + ", ".join(failed)
)
print("OUTPUT_CONTRACT_PASS")
if name == "__main__":
if len(sys.argv) != 3:
raise SystemExit(
"usage: validate_business_output.py MANIFEST RUN_ID"
)
validate(Path(sys.argv[1]), sys.argv[2])This validator does not query task state to decide whether output is correct. It recomputes the deterministic records and checksum from an independently declared business expectation. That separation is important: checking a manifest's checksum against a checksum copied from the same untrusted manifest would be circular.
The same idea applies beyond this inert fixture. A warehouse publication might validate expected partition identity, row cardinality constraints, business keys, or reconciled totals. Those are application contracts, not trigger rules. Cloud-native pipeline design foundations provide broader context, but this qualification deliberately remains provider-neutral and does not turn the lab into a managed-service tutorial.
The Airflow Best Practices guide similarly recommends self-checks that verify produced data rather than assuming task completion proves correctness.
The watcher is therefore a state-propagation mechanism, not a correctness oracle.
Capture a Run-Level Evidence Bundle
A credible qualification result should be reconstructible without a screenshot.
For every task in every run, capture at least:
dag_id
run_id
DAG revision
task_id
try number
trigger rule
direct parents
final state
log pointerAlso capture a distinct run-level record containing the actual terminal DagRun state, runtime/deployment manifest, business-output validation result, and cleanup evidence.
For Revision C, the static dependency ledger is:
Task | Trigger rule | Direct parents |
extract | all_success | none |
transform | all_success | extract |
publish | all_success | transform |
cleanup | all_done | publish |
watcher | one_failed | extract, transform, publish, cleanup |
Airflow 3.1.5 exposes public API v2 endpoints for getting a DagRun, retrieving its task instances, and fetching task logs; these are appropriate external evidence surfaces and avoid having task code query the metadata database.
A compact external collector can poll the already-triggered run:
# tools/collect_run_evidence.py
from future import annotations
import json
import os
import time
import urllib.parse
import urllib.request
from pathlib import Path
API = os.environ["AIRFLOW_API_BASE"].rstrip("/")
TOKEN = os.environ["AIRFLOW_API_TOKEN"]
DAG_ID = os.environ["LAB_DAG_ID"]
RUN_ID = os.environ["LAB_RUN_ID"]
REVISION = os.environ["LAB_DAG_REVISION"]
OUT = Path(os.environ.get("LAB_EVIDENCE_DIR", "./evidence"))
GRAPH = {
"extract": ("all_success", []),
"transform": ("all_success", ["extract"]),
"publish": ("all_success", ["transform"]),
"cleanup": ("all_done", ["publish"]),
"watcher": (
"one_failed",
["extract", "transform", "publish", "cleanup"],
),
}
def get_json(path: str) -> dict:
request = urllib.request.Request(
API + path,
headers={"Authorization": f"Bearer {TOKEN}"},
)
with urllib.request.urlopen(request, timeout=30) as response:
return json.load(response)
dag = urllib.parse.quote(DAG_ID, safe="")
run = urllib.parse.quote(RUN_ID, safe="")
while True:
run_record = get_json(f"/api/v2/dags/{dag}/dagRuns/{run}")
state = run_record["state"]
if state in {"success", "failed"}:
break
time.sleep(2)
task_response = get_json(
f"/api/v2/dags/{dag}/dagRuns/{run}/taskInstances"
)
instances = task_response["task_instances"]
OUT.mkdir(parents=True, exist_ok=True)
with (OUT / "task_evidence.jsonl").open("w", encoding="utf-8") as handle:
for ti in sorted(instances, key=lambda item: item["task_id"]):
task_id = ti["task_id"]
trigger_rule, parents = GRAPH[task_id]
try_number = ti.get("try_number")
# A log URL is evidence only where an attempt actually exists.
log_pointer = None
if isinstance(try_number, int) and try_number > 0:
encoded_task = urllib.parse.quote(task_id, safe="")
log_pointer = (
f"{API}/api/v2/dags/{dag}/dagRuns/{run}/"
f"taskInstances/{encoded_task}/logs/{try_number}"
)
row = {
"dag_id": DAG_ID,
"run_id": RUN_ID,
"dag_revision": REVISION,
"task_id": task_id,
"try_number": try_number,
"trigger_rule": trigger_rule,
"direct_parents": parents,
"final_state": ti["state"],
"log_pointer": log_pointer,
}
handle.write(json.dumps(row, sort_keys=True) + "\n")
(OUT / "dag_run.json").write_text(
json.dumps(
{
"dag_id": DAG_ID,
"run_id": RUN_ID,
"dag_revision": REVISION,
"terminal_state": state,
"api_record": run_record,
},
indent=2,
sort_keys=True,
)
+ "\n",
encoding="utf-8",
)Do not put AIRFLOW_API_TOKEN in the evidence output. Record non-secret runtime configuration separately, including the immutable deployment digest once known.
For task instances that never entered a worker attempt (such as an upstream_failed or trigger-rule-skipped task), do not invent “attempt 1” or a fake log link. Preserve the actual API field or null representation.
An acceptance bundle might finally contain:
environment.json
deployment_identity.json
graph_coverage_test.txt
task_evidence.jsonl
dag_run.json
audit/input_manifest.json
audit/cleanup.json
business/output_manifest.json # only when publication occurred
business_validation.txt
decision.jsonMissing task rows, missing runtime identity, mismatched revisions, or missing business validation produce HOLD.
Do Not Use Debugging Shortcuts as Execution Evidence
dag.test() is useful before the scheduler-backed qualification, but its scope must remain explicit.
The Airflow 3.1.5 debugging documentation describes dag.test() as a way to execute a DAG in a serialized local process for debugging; it also documents executor-related options and mark_success_pattern, which can mark matching tasks successful instead of executing them.
Use the local check for fast feedback:
# Local developer check only.
from dags.all_done_acceptance_lab import dag_c
dag_c.test(
run_conf={"fault": "transform"},
)Then run the same revision under the declared scheduler and LocalExecutor. The local test does not replace that qualification.
Never use mark_success_pattern on monitored tasks in an acceptance result. A task marked successful without its callable running is exactly the wrong evidence for a test whose purpose is to connect execution, task state, and output.
Likewise, do not manually convert a failed transform to success, catch and suppress its injected exception, or manipulate metadata from task code. Such actions may be legitimate debugging tools in other contexts; here they invalidate the experiment.
Assign ACCEPT, REWIRE, HOLD, or REPLAY
Once evidence exists, the decision should require very little interpretation.
Use this gate:
Evidence condition | Decision |
Revision A/B allows an intended blocking failure while DagRun becomes success | REWIRE |
Revision C direct-parent coverage assertion fails | REWIRE |
Revision C expected fault produces watcher failure and failed DagRun, with complete evidence | ACCEPT that failure-propagation test |
All-success Revision C is green and output/cleanup contracts pass | ACCEPT |
Required output absent or incorrect despite green tasks | HOLD |
Required task is unexpectedly skipped | HOLD |
Runtime/package/revision identity missing | HOLD |
Actual state contradicts expected matrix | HOLD, investigate rather than rationalize |
Corrected revision qualified, prior outputs reconciled, replay scope approved | REPLAY |
An injected transform failure is not itself a failed qualification. For Revision C, the expected outcome of that negative test is a failed task, a failed watcher, and a failed DagRun. Matching that expectation demonstrates that the graph reports the tested failure.
By contrast, Revision A's green transform-failure run is an expected experimental result and a failed graph design. That distinction matters. “The test behaved as predicted” is not the same as “the tested revision is acceptable.”
The output-negative case has a different interpretation again. A successful DagRun with a missing required manifest does not establish a watcher defect, because no monitored task entered a failed state. It establishes that orchestration state alone is insufficient as the business release gate. The correct run decision is HOLD.
No screenshot can resolve those distinctions by itself.
Repair the Graph Without Rewriting the Old Run
Once Revision A or B is shown to be defective, do not “repair history.”
Preserve:
original DAG revision
original deployment identity
original run ID
original task states
original logs
original output/audit evidence
original terminal DagRun stateThen introduce the corrected topology as a versioned deployment and qualify it separately.
Changing the DAG file does not retroactively change what the old run meant when it executed. The historical run's leaf set, task instances, output, and deployment revision remain evidence about that old execution.
For replay, generate a new run ID and therefore a new output namespace. Do not overwrite the old manifest.
An approval record can be as small as:
{
"source_run_id": "lab-a-transform-fail-001",
"source_revision": "A",
"replay_revision": "C",
"approved_tasks": ["extract", "transform", "publish", "cleanup"],
"output_reconciliation": "complete",
"non_idempotent_work_reviewed": true,
"decision": "REPLAY"
}The real approved task set may be smaller. Do not automatically clear every task in a failed production workflow merely because Airflow makes clearing convenient. Some business actions (payments, notifications, ledger inserts, third-party API submissions) can be non-idempotent.
The replay sequence is therefore:
Preserve the original evidence bundle.
Identify every external side effect that may have occurred.
Reconcile partial business output against the source run ID.
Qualify the corrected graph revision.
Approve the smallest safe replay scope.
Trigger a new run with a new run-scoped output identity.
Apply the same task-state and business-output acceptance gate to that new run.
Only after those checks does REPLAY become an action rather than a synonym for “try again.”
Cleanup Is Not Business Rollback
Suppose publish writes three inert rows to a run-scoped lab table and fails while writing its manifest. cleanup can still remove the staging marker and succeed.
That tells you the temporary resource was handled. It tells you nothing about whether the three rows disappeared.
The same principle holds for real external effects. Removing a temporary file does not roll back a committed database transaction. Deleting a scratch directory does not retract a message already sent to another system. Closing a connection does not reverse a completed API call.
Before replay, reconcile publication independently. In the inert fixture that can mean inspecting the run-scoped output namespace and either proving it empty or recording exactly which deterministic IDs exist. In a production workflow, the owning team must apply the target system's actual idempotency, transaction, or compensation model.
Do not make cleanup unconditionally fail to force the run red; that defeats its resource-handling role. Keep cleanup's all_done behavior and repair state propagation separately.
Keep Coverage Current as the Pipeline Grows
Revision C's acceptance is deliberately narrow. It qualifies the identified graph revision, runtime, monitored population, trigger rules, and output contract tested in the evidence bundle. It does not establish future correctness after the DAG changes.
Treat at least these changes as gate-reopening events:
a new business task;
a new leaf;
a different direct dependency;
a changed trigger rule;
a change from required to optional semantics, or vice versa;
retries being enabled;
a changed output schema, count rule, identity rule, or checksum calculation;
a new external side effect;
a different DAG revision or deployment artifact.
The DAG owner should own graph-state semantics. The business-output owner should own the independent output contract. Release automation can combine their evidence, but it should not erase that ownership distinction.
A graph-diff test should ask not only “did task IDs change?” but also “did watcher parentage change?” The parse-time assertion already makes a newly added unmonitored task fail qualification. Keep that fail-closed behavior.
This is one concrete facet of scalable data-pipeline responsibilities: operational growth should not quietly expand the unreviewed failure surface.
The important discipline is modest: every newly meaningful state must appear somewhere in the release contract. A new successful leaf, a newly allowed skip, or a changed cleanup policy can invalidate a previous acceptance even when the Python file still parses and every unit test unrelated to state propagation remains green.
Build the Data Engineering Skills Behind Trustworthy Runs
The reliable answer to an all_done cleanup problem is not “trust green,” nor is it “make cleanup fail.” It is to separate evidence surfaces.
Airflow 3.1.5 documents why a successful all_done leaf can yield a successful DagRun despite an earlier failure. The watcher pattern repairs that reporting only when every intended task is a direct monitored parent. The watcher still cannot detect a callable that returns successfully while producing wrong or absent business data, so the release gate also needs independent output evidence.
A trustworthy acceptance bundle therefore reconciles the exact DAG revision and deployment, direct graph coverage, one final state per task, actual terminal DagRun state, attempt/log evidence, run-scoped output identity, deterministic business validation, and cleanup evidence. Revision A and the negative-control Revision B should be REWIRE; incomplete or contradictory evidence should HOLD; qualified Revision C can be ACCEPT; and repaired business work should REPLAY only after side-effect reconciliation.
For engineers building the underlying skills, Refonte Learning's Data Engineering Program includes foundations in data warehousing and ETL, big-data technologies, and data-pipeline design; those are relevant foundations for reasoning about evidence-driven pipeline operations. The program description does not establish that this exact Airflow 3.1.5 watcher lab is part of its curriculum.
