|
| 1 | +""" |
| 2 | +tests/test_airflow_dag.py |
| 3 | +-------------------------- |
| 4 | +Unit and integration tests for Apache Airflow DAG in dags/data_pipeline_dag.py. |
| 5 | +""" |
| 6 | + |
| 7 | +import sys |
| 8 | +from pathlib import Path |
| 9 | + |
| 10 | +import pytest |
| 11 | + |
| 12 | +# Skip module if apache-airflow is not installed in the current environment |
| 13 | +pytest.importorskip("airflow") |
| 14 | + |
| 15 | +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) |
| 16 | + |
| 17 | +from dags.data_pipeline_dag import _get_safe_run_id, dag, default_args |
| 18 | + |
| 19 | + |
| 20 | +class TestAirflowDagDefinition: |
| 21 | + """Test suite for Airflow DAG structure and helper functions.""" |
| 22 | + |
| 23 | + def test_dag_structure_and_metadata(self): |
| 24 | + assert dag is not None |
| 25 | + assert dag.dag_id == "dataprep_pipeline" |
| 26 | + assert dag.schedule_interval == "@daily" |
| 27 | + assert len(dag.tasks) == 5 |
| 28 | + |
| 29 | + def test_dag_task_dependencies(self): |
| 30 | + task_dict = {t.task_id: t for t in dag.tasks} |
| 31 | + assert set(task_dict.keys()) == { |
| 32 | + "ingest_data", |
| 33 | + "validate_data", |
| 34 | + "clean_data", |
| 35 | + "transform_data", |
| 36 | + "load_data", |
| 37 | + } |
| 38 | + |
| 39 | + # Check linear dependency flow: ingest -> validate -> clean -> transform -> load |
| 40 | + ingest_task = task_dict["ingest_data"] |
| 41 | + validate_task = task_dict["validate_data"] |
| 42 | + clean_task = task_dict["clean_data"] |
| 43 | + transform_task = task_dict["transform_data"] |
| 44 | + load_task = task_dict["load_data"] |
| 45 | + |
| 46 | + assert validate_task in ingest_task.downstream_list |
| 47 | + assert clean_task in validate_task.downstream_list |
| 48 | + assert transform_task in clean_task.downstream_list |
| 49 | + assert load_task in transform_task.downstream_list |
| 50 | + |
| 51 | + def test_get_safe_run_id_sanitizes_illegal_chars(self): |
| 52 | + context_with_colon = {"run_id": "scheduled__2026-08-01T12:00:00+00:00"} |
| 53 | + safe_id = _get_safe_run_id(context_with_colon) |
| 54 | + assert ":" not in safe_id |
| 55 | + assert "+" not in safe_id |
| 56 | + assert safe_id == "scheduled__2026-08-01T12_00_00_00_00" |
| 57 | + |
| 58 | + def test_get_safe_run_id_default_fallback(self): |
| 59 | + context_empty = {} |
| 60 | + safe_id = _get_safe_run_id(context_empty) |
| 61 | + assert safe_id == "default_run" |
0 commit comments