Skip to content

Commit 5424ba7

Browse files
committed
Add airflowctl tasks list command
Listing a Dag's tasks currently requires the local `airflow tasks list` command, which parses Dag files on the caller's machine and needs direct metadata database access. Exposing it through airflowctl lets operators list tasks against the Public API under RBAC instead, per AIP-94. related: #68402
1 parent f0d50bd commit 5424ba7

7 files changed

Lines changed: 120 additions & 38 deletions

File tree

airflow-ctl-tests/tests/airflowctl_tests/test_airflowctl_commands.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,8 @@ def date_param():
102102
'tasks states-for-dag-run example_bash_operator "manual__{date_param}"',
103103
'tasks states-for-dag-run example_bash_operator --logical-date "{date_param}"',
104104
'tasks clear example_bash_operator --dag-run-id "manual__{date_param}" --task-ids runme_0 -o json',
105+
"tasks list example_bash_operator",
106+
"tasks list example_bash_operator --order-by=-task_id",
105107
# Task Instances commands
106108
'taskinstances get example_bash_operator "manual__{date_param}" runme_0',
107109
'taskinstances get-dependencies example_bash_operator "manual__{date_param}" runme_0',

airflow-ctl/docs/images/command_hashes.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ jobs:a5b644c5da8889443bb40ee10b599270
1010
pools:19efe105b9515ab1926ebcaf0e028d71
1111
providers:34502fe09dc0b8b0a13e7e46efdffda6
1212
taskinstances:bea84117114c2438eb7e7026f6bb7042
13-
tasks:ea587dc805cadbce81cd640d357f2661
13+
tasks:78f8bf0f0fbad664c37aa7489aab478f
1414
variables:f8fc76d3d398b2780f4e97f7cd816646
1515
version:31f4efdf8de0dbaaa4fac71ff7efecc3
1616
plugins:4864fd8f356704bd2b3cd1aec3567e35

airflow-ctl/docs/images/output_tasks.svg

Lines changed: 41 additions & 37 deletions
Loading

airflow-ctl/src/airflowctl/api/operations.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
ProviderCollectionResponse,
7070
QueuedEventCollectionResponse,
7171
QueuedEventResponse,
72+
TaskCollectionResponse,
7273
TaskDependencyCollectionResponse,
7374
TaskInstanceCollectionResponse,
7475
TaskInstanceResponse,
@@ -736,6 +737,11 @@ def clear(
736737
)
737738
return TaskInstanceCollectionResponse.model_validate_json(self.response.content)
738739

740+
def list(self, dag_id: str, order_by: str | None = None) -> TaskCollectionResponse | ServerResponseError:
741+
"""List tasks of a Dag."""
742+
self.response = self.client.get(f"dags/{dag_id}/tasks", params=_build_query_params(order_by=order_by))
743+
return TaskCollectionResponse.model_validate_json(self.response.content)
744+
739745

740746
class VariablesOperations(BaseOperations):
741747
"""Variable operations."""

airflow-ctl/src/airflowctl/ctl/help_texts.yaml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@ taskinstances:
9393

9494
tasks:
9595
clear: "Clear task instances of a Dag by its ID"
96+
list: "List all tasks of a Dag"
9697

9798
variables:
9899
get: "Retrieve a variable by its key"

airflow-ctl/tests/airflow_ctl/api/test_operations.py

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,11 +92,13 @@
9292
QueuedEventCollectionResponse,
9393
QueuedEventResponse,
9494
ReprocessBehavior,
95+
TaskCollectionResponse,
9596
TaskDependencyCollectionResponse,
9697
TaskDependencyResponse,
9798
TaskInstanceCollectionResponse,
9899
TaskInstanceResponse,
99100
TaskInstanceState,
101+
TaskResponse,
100102
TriggerDAGRunPostBody,
101103
VariableBody,
102104
VariableCollectionResponse,
@@ -1829,6 +1831,40 @@ class TestTasksOperations:
18291831
task_instances=[task_instance_response],
18301832
total_entries=1,
18311833
)
1834+
task_collection_response = TaskCollectionResponse(
1835+
tasks=[
1836+
TaskResponse(
1837+
task_id="task_1",
1838+
task_display_name="task_1",
1839+
owner="airflow",
1840+
start_date=None,
1841+
end_date=None,
1842+
trigger_rule=None,
1843+
depends_on_past=False,
1844+
wait_for_downstream=False,
1845+
retries=None,
1846+
queue=None,
1847+
pool=None,
1848+
pool_slots=None,
1849+
execution_timeout=None,
1850+
retry_delay=None,
1851+
retry_exponential_backoff=0,
1852+
priority_weight=None,
1853+
weight_rule=None,
1854+
ui_color=None,
1855+
ui_fgcolor=None,
1856+
template_fields=None,
1857+
downstream_task_ids=None,
1858+
doc_md=None,
1859+
operator_name=None,
1860+
params=None,
1861+
class_ref=None,
1862+
is_mapped=None,
1863+
extra_links=[],
1864+
)
1865+
],
1866+
total_entries=1,
1867+
)
18321868

18331869
def test_clear(self):
18341870
expected_body = self.clear_task_instances.model_dump(mode="json", exclude_none=True)
@@ -1845,6 +1881,23 @@ def handle_request(request: httpx.Request) -> httpx.Response:
18451881
response = client.tasks.clear(self.dag_id, self.clear_task_instances)
18461882
assert response == self.task_instance_collection_response
18471883

1884+
@pytest.mark.parametrize(
1885+
("order_by", "expected_params"),
1886+
[
1887+
pytest.param(None, {}, id="no-order-by"),
1888+
pytest.param("-task_id", {"order_by": "-task_id"}, id="with-order-by"),
1889+
],
1890+
)
1891+
def test_list(self, order_by, expected_params):
1892+
def handle_request(request: httpx.Request) -> httpx.Response:
1893+
assert request.url.path == f"/api/v2/dags/{self.dag_id}/tasks"
1894+
assert dict(request.url.params) == expected_params
1895+
return httpx.Response(200, json=json.loads(self.task_collection_response.model_dump_json()))
1896+
1897+
client = make_api_client(transport=httpx.MockTransport(handle_request))
1898+
response = client.tasks.list(self.dag_id, order_by=order_by)
1899+
assert response == self.task_collection_response
1900+
18481901

18491902
class TestVariablesOperations:
18501903
key = "key"

0 commit comments

Comments
 (0)