Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,8 @@ def date_param():
'tasks states-for-dag-run example_bash_operator "manual__{date_param}"',
'tasks states-for-dag-run example_bash_operator --logical-date "{date_param}"',
'tasks clear example_bash_operator --dag-run-id "manual__{date_param}" --task-ids runme_0 -o json',
"tasks list example_bash_operator",
"tasks list example_bash_operator --order-by=-task_id",
# Task Instances commands
'taskinstances get example_bash_operator "manual__{date_param}" runme_0',
'taskinstances get-dependencies example_bash_operator "manual__{date_param}" runme_0',
Expand Down
2 changes: 1 addition & 1 deletion airflow-ctl/docs/images/command_hashes.txt
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ jobs:a5b644c5da8889443bb40ee10b599270
pools:19efe105b9515ab1926ebcaf0e028d71
providers:34502fe09dc0b8b0a13e7e46efdffda6
taskinstances:bea84117114c2438eb7e7026f6bb7042
tasks:ea587dc805cadbce81cd640d357f2661
tasks:78f8bf0f0fbad664c37aa7489aab478f
variables:f8fc76d3d398b2780f4e97f7cd816646
version:31f4efdf8de0dbaaa4fac71ff7efecc3
plugins:4864fd8f356704bd2b3cd1aec3567e35
Expand Down
78 changes: 41 additions & 37 deletions airflow-ctl/docs/images/output_tasks.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
6 changes: 6 additions & 0 deletions airflow-ctl/src/airflowctl/api/operations.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
ProviderCollectionResponse,
QueuedEventCollectionResponse,
QueuedEventResponse,
TaskCollectionResponse,
TaskDependencyCollectionResponse,
TaskInstanceCollectionResponse,
TaskInstanceResponse,
Expand Down Expand Up @@ -736,6 +737,11 @@ def clear(
)
return TaskInstanceCollectionResponse.model_validate_json(self.response.content)

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


class VariablesOperations(BaseOperations):
"""Variable operations."""
Expand Down
1 change: 1 addition & 0 deletions airflow-ctl/src/airflowctl/ctl/help_texts.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ taskinstances:

tasks:
clear: "Clear task instances of a Dag by its ID"
list: "List all tasks of a Dag"

variables:
get: "Retrieve a variable by its key"
Expand Down
53 changes: 53 additions & 0 deletions airflow-ctl/tests/airflow_ctl/api/test_operations.py
Original file line number Diff line number Diff line change
Expand Up @@ -92,11 +92,13 @@
QueuedEventCollectionResponse,
QueuedEventResponse,
ReprocessBehavior,
TaskCollectionResponse,
TaskDependencyCollectionResponse,
TaskDependencyResponse,
TaskInstanceCollectionResponse,
TaskInstanceResponse,
TaskInstanceState,
TaskResponse,
TriggerDAGRunPostBody,
VariableBody,
VariableCollectionResponse,
Expand Down Expand Up @@ -1829,6 +1831,40 @@ class TestTasksOperations:
task_instances=[task_instance_response],
total_entries=1,
)
task_collection_response = TaskCollectionResponse(
tasks=[
TaskResponse(
task_id="task_1",
task_display_name="task_1",
owner="airflow",
start_date=None,
end_date=None,
trigger_rule=None,
depends_on_past=False,
wait_for_downstream=False,
retries=None,
queue=None,
pool=None,
pool_slots=None,
execution_timeout=None,
retry_delay=None,
retry_exponential_backoff=0,
priority_weight=None,
weight_rule=None,
ui_color=None,
ui_fgcolor=None,
template_fields=None,
downstream_task_ids=None,
doc_md=None,
operator_name=None,
params=None,
class_ref=None,
is_mapped=None,
extra_links=[],
)
],
total_entries=1,
)

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

@pytest.mark.parametrize(
("order_by", "expected_params"),
[
pytest.param(None, {}, id="no-order-by"),
pytest.param("-task_id", {"order_by": "-task_id"}, id="with-order-by"),
],
)
def test_list(self, order_by, expected_params):
def handle_request(request: httpx.Request) -> httpx.Response:
assert request.url.path == f"/api/v2/dags/{self.dag_id}/tasks"
assert dict(request.url.params) == expected_params
return httpx.Response(200, json=json.loads(self.task_collection_response.model_dump_json()))

client = make_api_client(transport=httpx.MockTransport(handle_request))
response = client.tasks.list(self.dag_id, order_by=order_by)
assert response == self.task_collection_response


class TestVariablesOperations:
key = "key"
Expand Down
Loading
Loading