This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new c2f600f6ec6 Add airflowctl tasks state command (#71206)
c2f600f6ec6 is described below
commit c2f600f6ec64717af551fde127ed19f34eed7a04
Author: Haseeb Malik <[email protected]>
AuthorDate: Sun Oct 4 17:59:25 2026 -0400
Add airflowctl tasks state command (#71206)
* Add airflowctl tasks state command
* Document what airflowctl tasks state prints for an unset state
* Point airflowctl at --map-index when a mapped task instance is not found
For a mapped task queried without --map-index the API returns a 404
explaining the task is mapped, but tasks state and tasks failed-deps reported
only a bare 'not found', leaving users guessing.
Generated-by: Claude Opus 5
---------
Co-authored-by: Jarek Potiuk <[email protected]>
---
.../airflowctl_tests/test_airflowctl_commands.py | 2 +
airflow-ctl/docs/images/command_hashes.txt | 2 +-
airflow-ctl/docs/images/output_tasks.svg | 78 +++---
airflow-ctl/src/airflowctl/ctl/cli_config.py | 17 ++
.../src/airflowctl/ctl/commands/task_command.py | 53 +++-
.../airflow_ctl/ctl/commands/test_task_command.py | 275 ++++++++++++++++++++-
6 files changed, 380 insertions(+), 47 deletions(-)
diff --git
a/airflow-ctl-tests/tests/airflowctl_tests/test_airflowctl_commands.py
b/airflow-ctl-tests/tests/airflowctl_tests/test_airflowctl_commands.py
index 7e5c94014f3..377a5d8a87d 100644
--- a/airflow-ctl-tests/tests/airflowctl_tests/test_airflowctl_commands.py
+++ b/airflow-ctl-tests/tests/airflowctl_tests/test_airflowctl_commands.py
@@ -107,6 +107,8 @@ TEST_COMMANDS = [
'tasks failed-deps example_bash_operator runme_0 --logical-date
"{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 state example_bash_operator runme_0 "manual__{date_param}"',
+ 'tasks state example_bash_operator runme_0 --logical-date "{date_param}"',
'tasks clear example_bash_operator --dag-run-id "manual__{date_param}"
--task-ids runme_0 -o json',
# Task Instances commands
'taskinstances get example_bash_operator "manual__{date_param}" runme_0',
diff --git a/airflow-ctl/docs/images/command_hashes.txt
b/airflow-ctl/docs/images/command_hashes.txt
index bd58cb3c8cd..4c9906ab2b2 100644
--- a/airflow-ctl/docs/images/command_hashes.txt
+++ b/airflow-ctl/docs/images/command_hashes.txt
@@ -10,7 +10,7 @@ jobs:d4af478f28dae48ee18d43b998edc345
pools:19efe105b9515ab1926ebcaf0e028d71
providers:34502fe09dc0b8b0a13e7e46efdffda6
taskinstances:bea84117114c2438eb7e7026f6bb7042
-tasks:ea587dc805cadbce81cd640d357f2661
+tasks:eb70701dfe1b9baeda39d5eea07fde5f
variables:f8fc76d3d398b2780f4e97f7cd816646
version:31f4efdf8de0dbaaa4fac71ff7efecc3
plugins:4864fd8f356704bd2b3cd1aec3567e35
diff --git a/airflow-ctl/docs/images/output_tasks.svg
b/airflow-ctl/docs/images/output_tasks.svg
index 9f30d658f38..08d22e642ed 100644
--- a/airflow-ctl/docs/images/output_tasks.svg
+++ b/airflow-ctl/docs/images/output_tasks.svg
@@ -1,4 +1,4 @@
-<svg class="rich-terminal" viewBox="0 0 933 367.2"
xmlns="http://www.w3.org/2000/svg">
+<svg class="rich-terminal" viewBox="0 0 933 391.59999999999997"
xmlns="http://www.w3.org/2000/svg">
<!-- Generated with Rich https://www.textualize.io -->
<style>
@@ -19,90 +19,94 @@
font-weight: 700;
}
- .terminal-2578352511-matrix {
+ .terminal-761188829-matrix {
font-family: Fira Code, monospace;
font-size: 20px;
line-height: 24.4px;
font-variant-east-asian: full-width;
}
- .terminal-2578352511-title {
+ .terminal-761188829-title {
font-size: 18px;
font-weight: bold;
font-family: arial;
}
- .terminal-2578352511-r1 { fill: #ff8700 }
-.terminal-2578352511-r2 { fill: #c5c8c6 }
-.terminal-2578352511-r3 { fill: #808080 }
-.terminal-2578352511-r4 { fill: #68a0b3 }
+ .terminal-761188829-r1 { fill: #ff8700 }
+.terminal-761188829-r2 { fill: #c5c8c6 }
+.terminal-761188829-r3 { fill: #808080 }
+.terminal-761188829-r4 { fill: #68a0b3 }
</style>
<defs>
- <clipPath id="terminal-2578352511-clip-terminal">
- <rect x="0" y="0" width="914.0" height="316.2" />
+ <clipPath id="terminal-761188829-clip-terminal">
+ <rect x="0" y="0" width="914.0" height="340.59999999999997" />
</clipPath>
- <clipPath id="terminal-2578352511-line-0">
+ <clipPath id="terminal-761188829-line-0">
<rect x="0" y="1.5" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-1">
+<clipPath id="terminal-761188829-line-1">
<rect x="0" y="25.9" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-2">
+<clipPath id="terminal-761188829-line-2">
<rect x="0" y="50.3" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-3">
+<clipPath id="terminal-761188829-line-3">
<rect x="0" y="74.7" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-4">
+<clipPath id="terminal-761188829-line-4">
<rect x="0" y="99.1" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-5">
+<clipPath id="terminal-761188829-line-5">
<rect x="0" y="123.5" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-6">
+<clipPath id="terminal-761188829-line-6">
<rect x="0" y="147.9" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-7">
+<clipPath id="terminal-761188829-line-7">
<rect x="0" y="172.3" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-8">
+<clipPath id="terminal-761188829-line-8">
<rect x="0" y="196.7" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-9">
+<clipPath id="terminal-761188829-line-9">
<rect x="0" y="221.1" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-10">
+<clipPath id="terminal-761188829-line-10">
<rect x="0" y="245.5" width="915" height="24.65"/>
</clipPath>
-<clipPath id="terminal-2578352511-line-11">
+<clipPath id="terminal-761188829-line-11">
<rect x="0" y="269.9" width="915" height="24.65"/>
</clipPath>
+<clipPath id="terminal-761188829-line-12">
+ <rect x="0" y="294.3" width="915" height="24.65"/>
+ </clipPath>
</defs>
- <rect fill="#292929" stroke="rgba(255,255,255,0.35)" stroke-width="1"
x="1" y="1" width="931" height="365.2" rx="8"/>
+ <rect fill="#292929" stroke="rgba(255,255,255,0.35)" stroke-width="1"
x="1" y="1" width="931" height="389.6" rx="8"/>
<g transform="translate(26,22)">
<circle cx="0" cy="0" r="7" fill="#ff5f57"/>
<circle cx="22" cy="0" r="7" fill="#febc2e"/>
<circle cx="44" cy="0" r="7" fill="#28c840"/>
</g>
- <g transform="translate(9, 41)"
clip-path="url(#terminal-2578352511-clip-terminal)">
+ <g transform="translate(9, 41)"
clip-path="url(#terminal-761188829-clip-terminal)">
- <g class="terminal-2578352511-matrix">
- <text class="terminal-2578352511-r1" x="0" y="20" textLength="73.2"
clip-path="url(#terminal-2578352511-line-0)">Usage:</text><text
class="terminal-2578352511-r3" x="85.4" y="20" textLength="195.2"
clip-path="url(#terminal-2578352511-line-0)">airflowctl tasks</text><text
class="terminal-2578352511-r2" x="280.6" y="20" textLength="24.4"
clip-path="url(#terminal-2578352511-line-0)"> [</text><text
class="terminal-2578352511-r4" x="305" y="20" textLength="24.4"
clip-path="url(# [...]
-</text><text class="terminal-2578352511-r2" x="915" y="44.4" textLength="12.2"
clip-path="url(#terminal-2578352511-line-1)">
-</text><text class="terminal-2578352511-r2" x="0" y="68.8" textLength="292.8"
clip-path="url(#terminal-2578352511-line-2)">Perform Tasks operations</text><text
class="terminal-2578352511-r2" x="915" y="68.8" textLength="12.2"
clip-path="url(#terminal-2578352511-line-2)">
-</text><text class="terminal-2578352511-r2" x="915" y="93.2" textLength="12.2"
clip-path="url(#terminal-2578352511-line-3)">
-</text><text class="terminal-2578352511-r1" x="0" y="117.6" textLength="256.2"
clip-path="url(#terminal-2578352511-line-4)">Positional Arguments:</text><text
class="terminal-2578352511-r2" x="915" y="117.6" textLength="12.2"
clip-path="url(#terminal-2578352511-line-4)">
-</text><text class="terminal-2578352511-r4" x="24.4" y="142" textLength="85.4"
clip-path="url(#terminal-2578352511-line-5)">COMMAND</text><text
class="terminal-2578352511-r2" x="915" y="142" textLength="12.2"
clip-path="url(#terminal-2578352511-line-5)">
-</text><text class="terminal-2578352511-r4" x="48.8" y="166.4" textLength="61"
clip-path="url(#terminal-2578352511-line-6)">clear</text><text
class="terminal-2578352511-r2" x="268.4" y="166.4" textLength="475.8"
clip-path="url(#terminal-2578352511-line-6)">Clear task instances of a Dag by its ID</text><text
class="terminal-2578352511-r2" x="915" y="166.4" textLength="12.2"
clip-path="url(#terminal-2578352511-line-6)">
-</text><text class="terminal-2578352511-r4" x="48.8" y="190.8"
textLength="134.2"
clip-path="url(#terminal-2578352511-line-7)">failed-deps</text><text
class="terminal-2578352511-r2" x="268.4" y="190.8" textLength="610"
clip-path="url(#terminal-2578352511-line-7)">Returns the unmet dependencies for a task instance</text><text
class="terminal-2578352511-r2" x="915" y="190.8" textLength="12.2"
clip-path="url(#terminal-2578352511-line-7)">
-</text><text class="terminal-2578352511-r4" x="48.8" y="215.2"
textLength="219.6"
clip-path="url(#terminal-2578352511-line-8)">states-for-dag-run</text><text
class="terminal-2578352511-r2" x="915" y="215.2" textLength="12.2"
clip-path="url(#terminal-2578352511-line-8)">
-</text><text class="terminal-2578352511-r2" x="268.4" y="239.6"
textLength="597.8"
clip-path="url(#terminal-2578352511-line-9)">Get the status of all task instances in a Dag run</text><text
class="terminal-2578352511-r2" x="915" y="239.6" textLength="12.2"
clip-path="url(#terminal-2578352511-line-9)">
-</text><text class="terminal-2578352511-r2" x="915" y="264" textLength="12.2"
clip-path="url(#terminal-2578352511-line-10)">
-</text><text class="terminal-2578352511-r1" x="0" y="288.4" textLength="97.6"
clip-path="url(#terminal-2578352511-line-11)">Options:</text><text
class="terminal-2578352511-r2" x="915" y="288.4" textLength="12.2"
clip-path="url(#terminal-2578352511-line-11)">
-</text><text class="terminal-2578352511-r4" x="24.4" y="312.8"
textLength="24.4" clip-path="url(#terminal-2578352511-line-12)">-h</text><text
class="terminal-2578352511-r2" x="48.8" y="312.8" textLength="24.4"
clip-path="url(#terminal-2578352511-line-12)">, </text><text
class="terminal-2578352511-r4" x="73.2" y="312.8" textLength="73.2"
clip-path="url(#terminal-2578352511-line-12)">--help</text><text
class="terminal-2578352511-r2" x="268.4" y="312.8" textLength="378.2"
clip-path="ur [...]
+ <g class="terminal-761188829-matrix">
+ <text class="terminal-761188829-r1" x="0" y="20" textLength="73.2"
clip-path="url(#terminal-761188829-line-0)">Usage:</text><text
class="terminal-761188829-r3" x="85.4" y="20" textLength="195.2"
clip-path="url(#terminal-761188829-line-0)">airflowctl tasks</text><text
class="terminal-761188829-r2" x="280.6" y="20" textLength="24.4"
clip-path="url(#terminal-761188829-line-0)"> [</text><text
class="terminal-761188829-r4" x="305" y="20" textLength="24.4"
clip-path="url(#termina [...]
+</text><text class="terminal-761188829-r2" x="915" y="44.4" textLength="12.2"
clip-path="url(#terminal-761188829-line-1)">
+</text><text class="terminal-761188829-r2" x="0" y="68.8" textLength="292.8"
clip-path="url(#terminal-761188829-line-2)">Perform Tasks operations</text><text
class="terminal-761188829-r2" x="915" y="68.8" textLength="12.2"
clip-path="url(#terminal-761188829-line-2)">
+</text><text class="terminal-761188829-r2" x="915" y="93.2" textLength="12.2"
clip-path="url(#terminal-761188829-line-3)">
+</text><text class="terminal-761188829-r1" x="0" y="117.6" textLength="256.2"
clip-path="url(#terminal-761188829-line-4)">Positional Arguments:</text><text
class="terminal-761188829-r2" x="915" y="117.6" textLength="12.2"
clip-path="url(#terminal-761188829-line-4)">
+</text><text class="terminal-761188829-r4" x="24.4" y="142" textLength="85.4"
clip-path="url(#terminal-761188829-line-5)">COMMAND</text><text
class="terminal-761188829-r2" x="915" y="142" textLength="12.2"
clip-path="url(#terminal-761188829-line-5)">
+</text><text class="terminal-761188829-r4" x="48.8" y="166.4" textLength="61"
clip-path="url(#terminal-761188829-line-6)">clear</text><text
class="terminal-761188829-r2" x="268.4" y="166.4" textLength="475.8"
clip-path="url(#terminal-761188829-line-6)">Clear task instances of a Dag by its ID</text><text
class="terminal-761188829-r2" x="915" y="166.4" textLength="12.2"
clip-path="url(#terminal-761188829-line-6)">
+</text><text class="terminal-761188829-r4" x="48.8" y="190.8"
textLength="134.2"
clip-path="url(#terminal-761188829-line-7)">failed-deps</text><text
class="terminal-761188829-r2" x="268.4" y="190.8" textLength="610"
clip-path="url(#terminal-761188829-line-7)">Returns the unmet dependencies for a task instance</text><text
class="terminal-761188829-r2" x="915" y="190.8" textLength="12.2"
clip-path="url(#terminal-761188829-line-7)">
+</text><text class="terminal-761188829-r4" x="48.8" y="215.2" textLength="61"
clip-path="url(#terminal-761188829-line-8)">state</text><text
class="terminal-761188829-r2" x="268.4" y="215.2" textLength="390.4"
clip-path="url(#terminal-761188829-line-8)">Get the state of a task instance</text><text
class="terminal-761188829-r2" x="915" y="215.2" textLength="12.2"
clip-path="url(#terminal-761188829-line-8)">
+</text><text class="terminal-761188829-r4" x="48.8" y="239.6"
textLength="219.6"
clip-path="url(#terminal-761188829-line-9)">states-for-dag-run</text><text
class="terminal-761188829-r2" x="915" y="239.6" textLength="12.2"
clip-path="url(#terminal-761188829-line-9)">
+</text><text class="terminal-761188829-r2" x="268.4" y="264"
textLength="597.8"
clip-path="url(#terminal-761188829-line-10)">Get the status of all task instances in a Dag run</text><text
class="terminal-761188829-r2" x="915" y="264" textLength="12.2"
clip-path="url(#terminal-761188829-line-10)">
+</text><text class="terminal-761188829-r2" x="915" y="288.4" textLength="12.2"
clip-path="url(#terminal-761188829-line-11)">
+</text><text class="terminal-761188829-r1" x="0" y="312.8" textLength="97.6"
clip-path="url(#terminal-761188829-line-12)">Options:</text><text
class="terminal-761188829-r2" x="915" y="312.8" textLength="12.2"
clip-path="url(#terminal-761188829-line-12)">
+</text><text class="terminal-761188829-r4" x="24.4" y="337.2"
textLength="24.4" clip-path="url(#terminal-761188829-line-13)">-h</text><text
class="terminal-761188829-r2" x="48.8" y="337.2" textLength="24.4"
clip-path="url(#terminal-761188829-line-13)">, </text><text
class="terminal-761188829-r4" x="73.2" y="337.2" textLength="73.2"
clip-path="url(#terminal-761188829-line-13)">--help</text><text
class="terminal-761188829-r2" x="268.4" y="337.2" textLength="378.2"
clip-path="url(#term [...]
</text>
</g>
</g>
diff --git a/airflow-ctl/src/airflowctl/ctl/cli_config.py
b/airflow-ctl/src/airflowctl/ctl/cli_config.py
index a4d6afd35e4..4515723fd01 100755
--- a/airflow-ctl/src/airflowctl/ctl/cli_config.py
+++ b/airflow-ctl/src/airflowctl/ctl/cli_config.py
@@ -1228,6 +1228,23 @@ TASK_COMMANDS = (
ARG_MAP_INDEX,
),
),
+ ActionCommand(
+ name="state",
+ help="Get the state of a task instance",
+ description=(
+ "Get the state of a task instance. "
+ "Select the run with either run_id or --logical-date (pass exactly
one). "
+ "Prints the state value, or None when the task instance has no
state yet."
+ ),
+ func=lazy_load_command("airflowctl.ctl.commands.task_command.state"),
+ args=(
+ ARG_DAG_ID,
+ ARG_TASK_ID,
+ ARG_RUN_ID,
+ ARG_LOGICAL_DATE,
+ ARG_MAP_INDEX,
+ ),
+ ),
ActionCommand(
name="states-for-dag-run",
help="Get the status of all task instances in a Dag run",
diff --git a/airflow-ctl/src/airflowctl/ctl/commands/task_command.py
b/airflow-ctl/src/airflowctl/ctl/commands/task_command.py
index 1c6a7004903..f03a83e5103 100644
--- a/airflow-ctl/src/airflowctl/ctl/commands/task_command.py
+++ b/airflow-ctl/src/airflowctl/ctl/commands/task_command.py
@@ -31,6 +31,29 @@ if TYPE_CHECKING:
from airflowctl.api.datamodels.generated import TaskInstanceResponse
+def _is_mapped_task_error(error: ServerResponseError) -> bool:
+ """Whether the API rejected the lookup because the task is mapped and no
map index was given."""
+ try:
+ detail = error.response.json().get("detail")
+ except ValueError:
+ return False
+ return isinstance(detail, str) and "is mapped" in detail
+
+
+def _task_instance_not_found_message(
+ dag_id: str, run_id: str, task_id: str, map_index: int, error:
ServerResponseError
+) -> str:
+ """Build the message shown when a task instance is not found."""
+ map_index_part = f" with map index {map_index}" if map_index >= 0 else ""
+ message = (
+ f"Task instance for task {task_id!r}{map_index_part} in Dag run "
+ f"{run_id!r} of Dag {dag_id!r} not found"
+ )
+ if map_index < 0 and _is_mapped_task_error(error):
+ message += ". The task is mapped; pass --map-index to select one of
its task instances"
+ return message
+
+
def _format_task_instance(ti: TaskInstanceResponse, has_mapped_instances:
bool) -> dict[str, str]:
data = {
"dag_id": ti.dag_id,
@@ -75,10 +98,8 @@ def failed_deps(args, api_client=NEW_API_CLIENT) -> None:
)
except ServerResponseError as e:
if e.response.status_code == 404:
- map_index_part = f" with map index {args.map_index}" if
args.map_index >= 0 else ""
rich.print(
- f"[red]Task instance for task {args.task_id!r}{map_index_part}
in Dag run "
- f"{run_id!r} of Dag {args.dag_id!r} not found[/red]"
+ f"[red]{_task_instance_not_found_message(args.dag_id, run_id,
args.task_id, args.map_index, e)}[/red]"
)
sys.exit(1)
raise
@@ -111,3 +132,29 @@ def states_for_dag_run(args, api_client=NEW_API_CLIENT) ->
None:
data=[_format_task_instance(ti, has_mapped_instances) for ti in
task_instances],
output=args.output,
)
+
+
+@provide_api_client(kind=ClientKind.CLI)
+def state(args, api_client=NEW_API_CLIENT) -> None:
+ """Get the state of a task instance."""
+ run_id = resolve_dag_run_id(api_client, args)
+
+ try:
+ task_instance = api_client.task_instances.get(
+ dag_id=args.dag_id,
+ dag_run_id=run_id,
+ task_id=args.task_id,
+ map_index=args.map_index,
+ suppress_error_log=True,
+ )
+ except ServerResponseError as e:
+ if e.response.status_code == 404:
+ rich.print(
+ f"[red]{_task_instance_not_found_message(args.dag_id, run_id,
args.task_id, args.map_index, e)}[/red]"
+ )
+ sys.exit(1)
+ raise
+
+ # Unset states print as None to stay drop-in compatible with the deprecated
+ # `airflow tasks state`, which prints its nullable `ti.state` column
directly.
+ print(task_instance.state.value if task_instance.state else None)
diff --git a/airflow-ctl/tests/airflow_ctl/ctl/commands/test_task_command.py
b/airflow-ctl/tests/airflow_ctl/ctl/commands/test_task_command.py
index 9b15695209c..5839f87a65f 100644
--- a/airflow-ctl/tests/airflow_ctl/ctl/commands/test_task_command.py
+++ b/airflow-ctl/tests/airflow_ctl/ctl/commands/test_task_command.py
@@ -34,10 +34,12 @@ from airflowctl.api.operations import ServerResponseError
from airflowctl.ctl import cli_parser
from airflowctl.ctl.commands import task_command
+MAPPED_TASK_DETAIL = "Task instance is mapped, add the map_index value to the
URL"
-def _make_server_error(status_code: int) -> ServerResponseError:
+
+def _make_server_error(status_code: int, detail: str = "boom") ->
ServerResponseError:
request = httpx.Request("GET",
"http://testserver/api/v2/dags/test_dag/dagRuns/test_run")
- response = httpx.Response(status_code, request=request, json={"detail":
"boom"})
+ response = httpx.Response(status_code, request=request, json={"detail":
detail})
return ServerResponseError(message="boom", request=request,
response=response)
@@ -309,19 +311,30 @@ class TestFailedDeps:
)
@pytest.mark.parametrize(
- ("extra_args", "expected_message"),
+ ("extra_args", "detail", "expected_message"),
[
- ([], "Task instance for task 'test_task' in Dag run 'test_run' of
Dag 'test_dag' not found"),
+ (
+ [],
+ "boom",
+ "Task instance for task 'test_task' in Dag run 'test_run' of
Dag 'test_dag' not found",
+ ),
(
["--map-index", "3"],
+ "boom",
"Task instance for task 'test_task' with map index 3 in Dag
run 'test_run' "
"of Dag 'test_dag' not found",
),
+ (
+ [],
+ MAPPED_TASK_DETAIL,
+ "Task instance for task 'test_task' in Dag run 'test_run' of
Dag 'test_dag' not found. "
+ "The task is mapped; pass --map-index to select one of its
task instances",
+ ),
],
)
- def test_failed_deps_task_instance_not_found(self, extra_args,
expected_message, capsys):
+ def test_failed_deps_task_instance_not_found(self, extra_args, detail,
expected_message, capsys):
api_client = self._make_api_client()
- api_client.task_instances.get_dependencies.side_effect =
_make_server_error(404)
+ api_client.task_instances.get_dependencies.side_effect =
_make_server_error(404, detail)
with pytest.raises(SystemExit, match="1"):
task_command.failed_deps(
@@ -643,3 +656,253 @@ class TestStatesForDagRun:
task_command.states_for_dag_run(self.parser.parse_args(argv),
api_client=api_client)
assert ctx.value is error
+
+
+class TestState:
+ parser = cli_parser.get_parser()
+ dag_id = "test_dag"
+ run_id = "test_run"
+ task_id = "test_task"
+ logical_date = datetime.datetime(2025, 1, 1, tzinfo=datetime.timezone.utc)
+
+ def _make_task_instance(self, state: TaskInstanceState | None) ->
TaskInstanceResponse:
+ return TaskInstanceResponse(
+ id=uuid.uuid4(),
+ task_id=self.task_id,
+ dag_id=self.dag_id,
+ dag_run_id=self.run_id,
+ map_index=-1,
+ logical_date=self.logical_date,
+ run_after=self.logical_date,
+ state=state,
+ start_date=None,
+ end_date=None,
+ try_number=1,
+ max_tries=0,
+ task_display_name=self.task_id,
+ dag_display_name=self.dag_id,
+ pool="default_pool",
+ pool_slots=1,
+ executor_config="{}",
+ duration=None,
+ hostname=None,
+ unixname=None,
+ queue=None,
+ priority_weight=None,
+ operator=None,
+ operator_name=None,
+ queued_when=None,
+ scheduled_when=None,
+ pid=None,
+ executor=None,
+ note=None,
+ rendered_map_index=None,
+ trigger=None,
+ triggerer_job=None,
+ dag_version=None,
+ )
+
+ def _make_api_client(self, state: TaskInstanceState | None =
TaskInstanceState.SUCCESS) -> mock.MagicMock:
+ api_client = mock.MagicMock()
+ api_client.dag_runs.list.return_value.dag_runs =
[mock.MagicMock(dag_run_id=self.run_id)]
+ api_client.task_instances.get.return_value =
self._make_task_instance(state=state)
+ return api_client
+
+ def test_state_by_run_id(self, capsys):
+ api_client = self._make_api_client(state=TaskInstanceState.SUCCESS)
+
+ task_command.state(
+ self.parser.parse_args(["tasks", "state", self.dag_id,
self.task_id, self.run_id]),
+ api_client=api_client,
+ )
+
+ api_client.dag_runs.list.assert_not_called()
+ api_client.task_instances.get.assert_called_once_with(
+ dag_id=self.dag_id,
+ dag_run_id=self.run_id,
+ task_id=self.task_id,
+ map_index=-1,
+ suppress_error_log=True,
+ )
+ assert capsys.readouterr().out == "success\n"
+
+ def test_state_by_logical_date(self, capsys):
+ api_client = self._make_api_client(state=TaskInstanceState.RUNNING)
+
+ task_command.state(
+ self.parser.parse_args(
+ [
+ "tasks",
+ "state",
+ self.dag_id,
+ self.task_id,
+ "--logical-date",
+ self.logical_date.isoformat(),
+ ]
+ ),
+ api_client=api_client,
+ )
+
+ api_client.dag_runs.list.assert_called_once_with(
+ dag_id=self.dag_id,
+ logical_date_gte=self.logical_date,
+ logical_date_lte=self.logical_date,
+ order_by="-id",
+ limit=1,
+ suppress_error_log=True,
+ )
+ api_client.task_instances.get.assert_called_once_with(
+ dag_id=self.dag_id,
+ dag_run_id=self.run_id,
+ task_id=self.task_id,
+ map_index=-1,
+ suppress_error_log=True,
+ )
+ assert capsys.readouterr().out == "running\n"
+
+ def test_state_with_map_index(self, capsys):
+ api_client = self._make_api_client(state=TaskInstanceState.SUCCESS)
+
+ task_command.state(
+ self.parser.parse_args(
+ ["tasks", "state", self.dag_id, self.task_id, self.run_id,
"--map-index", "3"]
+ ),
+ api_client=api_client,
+ )
+
+ api_client.task_instances.get.assert_called_once_with(
+ dag_id=self.dag_id,
+ dag_run_id=self.run_id,
+ task_id=self.task_id,
+ map_index=3,
+ suppress_error_log=True,
+ )
+ assert capsys.readouterr().out == "success\n"
+
+ def test_state_prints_none_when_task_instance_has_no_state(self, capsys):
+ api_client = self._make_api_client(state=None)
+
+ task_command.state(
+ self.parser.parse_args(["tasks", "state", self.dag_id,
self.task_id, self.run_id]),
+ api_client=api_client,
+ )
+
+ assert capsys.readouterr().out == "None\n"
+
+ @pytest.mark.parametrize(
+ "extra_args",
+ [
+ [],
+ ["test_run", "--logical-date", "2025-01-01T00:00:00+00:00"],
+ ],
+ ids=["neither", "both"],
+ )
+ def test_state_requires_exactly_one_of_run_id_and_logical_date(self,
extra_args, capsys):
+ api_client = self._make_api_client()
+
+ with pytest.raises(SystemExit, match="1"):
+ task_command.state(
+ self.parser.parse_args(["tasks", "state", self.dag_id,
self.task_id, *extra_args]),
+ api_client=api_client,
+ )
+
+ api_client.task_instances.get.assert_not_called()
+ assert _normalize_rich_output(capsys.readouterr().out) == (
+ "Provide either run_id or --logical-date, but not both"
+ )
+
+ @pytest.mark.parametrize(
+ ("logical_date", "expected_message"),
+ [
+ ("not-a-date", "Invalid --logical-date: 'not-a-date'"),
+ ("2025-01-01T00:00:00", "--logical-date must include a timezone
offset"),
+ ],
+ ids=["unparsable", "naive"],
+ )
+ def test_state_rejects_bad_logical_date(self, logical_date,
expected_message, capsys):
+ api_client = self._make_api_client()
+
+ with pytest.raises(SystemExit, match="1"):
+ task_command.state(
+ self.parser.parse_args(
+ ["tasks", "state", self.dag_id, self.task_id,
"--logical-date", logical_date]
+ ),
+ api_client=api_client,
+ )
+
+ api_client.dag_runs.list.assert_not_called()
+ assert _normalize_rich_output(capsys.readouterr().out) ==
expected_message
+
+ @pytest.mark.parametrize("list_failure", ["no_matching_run",
"dag_not_found_404"])
+ def test_state_dag_run_not_found_by_logical_date(self, list_failure,
capsys):
+ api_client = self._make_api_client()
+ if list_failure == "no_matching_run":
+ api_client.dag_runs.list.return_value.dag_runs = []
+ else:
+ api_client.dag_runs.list.side_effect = _make_server_error(404)
+
+ with pytest.raises(SystemExit, match="1"):
+ task_command.state(
+ self.parser.parse_args(
+ [
+ "tasks",
+ "state",
+ self.dag_id,
+ self.task_id,
+ "--logical-date",
+ self.logical_date.isoformat(),
+ ]
+ ),
+ api_client=api_client,
+ )
+
+ api_client.task_instances.get.assert_not_called()
+ assert _normalize_rich_output(capsys.readouterr().out) == (
+ "Dag run for test_dag with logical date
'2025-01-01T00:00:00+00:00' not found"
+ )
+
+ @pytest.mark.parametrize(
+ ("extra_args", "detail", "expected_message"),
+ [
+ (
+ [],
+ "boom",
+ "Task instance for task 'test_task' in Dag run 'test_run' of
Dag 'test_dag' not found",
+ ),
+ (
+ ["--map-index", "3"],
+ "boom",
+ "Task instance for task 'test_task' with map index 3 in Dag
run 'test_run' "
+ "of Dag 'test_dag' not found",
+ ),
+ (
+ [],
+ MAPPED_TASK_DETAIL,
+ "Task instance for task 'test_task' in Dag run 'test_run' of
Dag 'test_dag' not found. "
+ "The task is mapped; pass --map-index to select one of its
task instances",
+ ),
+ ],
+ )
+ def test_state_task_instance_not_found(self, extra_args, detail,
expected_message, capsys):
+ api_client = self._make_api_client()
+ api_client.task_instances.get.side_effect = _make_server_error(404,
detail)
+
+ with pytest.raises(SystemExit, match="1"):
+ task_command.state(
+ self.parser.parse_args(
+ ["tasks", "state", self.dag_id, self.task_id, self.run_id,
*extra_args]
+ ),
+ api_client=api_client,
+ )
+
+ assert _normalize_rich_output(capsys.readouterr().out) ==
expected_message
+
+ def test_state_reraises_non_404_error(self):
+ api_client = self._make_api_client()
+ api_client.task_instances.get.side_effect = _make_server_error(500)
+
+ with pytest.raises(ServerResponseError):
+ task_command.state(
+ self.parser.parse_args(["tasks", "state", self.dag_id,
self.task_id, self.run_id]),
+ api_client=api_client,
+ )