This is an automated email from the ASF dual-hosted git repository.
pierrejeambrun pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/v3-3-test by this push:
new 48ec0fa9140 Clear all task instances in a task group, not just the
first page (#73547) (#73651)
48ec0fa9140 is described below
commit 48ec0fa91408b234a0cb24fa44a07e793093730d
Author: Pierre Jeambrun <[email protected]>
AuthorDate: Fri Sep 25 11:48:22 2026 +0200
Clear all task instances in a task group, not just the first page (#73547)
(#73651)
Clearing a task group enumerated its task instances through the paginated
task-instance listing and only ever saw the first page, so a group with more
tasks than the page size had the rest silently skipped from both the preview
and the clear. The clear endpoint now takes a task_group_id and resolves the
group's tasks server-side from the dag structure, the same way marking a
task
group already does, so every task is cleared no matter how many there are.
(cherry picked from commit 3eec18d5bc3c8172c5e832b643592154b674456e)
---
.../core_api/datamodels/task_instances.py | 8 ++
.../core_api/openapi/v2-rest-api-generated.yaml | 8 ++
.../core_api/routes/public/task_instances.py | 9 ++
.../core_api/services/public/task_instances.py | 18 ++--
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 12 +++
.../airflow/ui/openapi-gen/requests/types.gen.ts | 4 +
.../TaskInstance/ClearGroupTaskInstanceDialog.tsx | 28 ++-----
.../core_api/routes/public/test_task_instances.py | 98 +++++++++++++++++++++-
.../src/airflowctl/api/datamodels/generated.py | 7 ++
9 files changed, 161 insertions(+), 31 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
index 98244326163..07422e168bd 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/task_instances.py
@@ -216,6 +216,12 @@ class ClearTaskInstancesBody(StrictBaseModel):
description="A list of `task_id` or [`task_id`, `map_index`]. "
"If only the `task_id` is provided for a mapped task, all of its map
indices will be targeted.",
)
+ task_group_id: str | None = Field(
+ default=None,
+ description="Clear every task in this task group. Mutually exclusive
with `task_ids`. "
+ "The group's tasks are resolved on the server from the dag structure,
so all of them are "
+ "targeted regardless of how many there are.",
+ )
dag_run_id: str | None = None
include_upstream: bool = False
include_downstream: bool = False
@@ -249,6 +255,8 @@ class ClearTaskInstancesBody(StrictBaseModel):
raise ValueError("Exactly one of dag_run_id or end_date must be
provided")
if isinstance(data.get("task_ids"), list) and
len(data.get("task_ids")) < 1:
raise ValueError("task_ids list should have at least 1 element.")
+ if data.get("task_ids") and data.get("task_group_id"):
+ raise ValueError("Only one of task_ids or task_group_id may be
provided")
return data
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index ebc0dad8576..f7e902b1793 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -12167,6 +12167,14 @@ components:
description: A list of `task_id` or [`task_id`, `map_index`]. If
only the
`task_id` is provided for a mapped task, all of its map indices
will be
targeted.
+ task_group_id:
+ anyOf:
+ - type: string
+ - type: 'null'
+ title: Task Group Id
+ description: Clear every task in this task group. Mutually exclusive
with
+ `task_ids`. The group's tasks are resolved on the server from the
dag
+ structure, so all of them are targeted regardless of how many
there are.
dag_run_id:
anyOf:
- type: string
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
index 7a4bdfc37a6..7001f67036d 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py
@@ -104,6 +104,7 @@ from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_
from airflow.api_fastapi.core_api.security import GetUserDep,
ReadableTIFilterDep, requires_access_dag
from airflow.api_fastapi.core_api.services.public.task_instances import (
BulkTaskInstanceService,
+ _get_task_group_task_ids,
_get_task_group_task_instances,
_patch_task_group_state,
_patch_task_instance_note,
@@ -898,6 +899,14 @@ def post_clear_task_instances(
if future:
body.end_date = None
+ # A task group has no per-task list at the call site; resolve every task
in it from the dag
+ # structure so all are cleared, not just the first page the UI could
enumerate.
+ if body.task_group_id is not None:
+ body.task_ids = cast(
+ "list[str | tuple[str, int]]",
+ _get_task_group_task_ids(dag_id, body.task_group_id, dag),
+ )
+
if (task_markers_to_clear := body.task_ids) is not None:
mapped_tasks_tuples = {t for t in task_markers_to_clear if
isinstance(t, tuple)}
# Unmapped tasks are expressed in their task_ids (without map_indexes)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
index b00ba4effdb..476ce7bb7f2 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
+++
b/airflow-core/src/airflow/api_fastapi/core_api/services/public/task_instances.py
@@ -180,6 +180,16 @@ def _patch_ti_validate_request(
return dag, list(tis), data
+def _get_task_group_task_ids(dag_id: str, task_group_id: str, dag:
SerializedDAG) -> list[str]:
+ """Return the ids of every task in a task group, resolved from the dag
structure."""
+ task_group = dag.task_group_dict.get(task_group_id)
+ if not task_group:
+ raise HTTPException(
+ status.HTTP_404_NOT_FOUND, f"Task group '{task_group_id}' not
found in DAG '{dag_id}'"
+ )
+ return [task.task_id for task in task_group.iter_tasks()]
+
+
def _get_task_group_task_instances(
dag_id: str,
dag_run_id: str,
@@ -188,13 +198,7 @@ def _get_task_group_task_instances(
session: Session,
) -> list[TI]:
"""Get all task instances in a task group for a specific DAG run."""
- task_group = dag.task_group_dict.get(task_group_id)
- if not task_group:
- raise HTTPException(
- status.HTTP_404_NOT_FOUND, f"Task group '{task_group_id}' not
found in DAG '{dag_id}'"
- )
-
- task_ids = [task.task_id for task in task_group.iter_tasks()]
+ task_ids = _get_task_group_task_ids(dag_id, task_group_id, dag)
query = (
select(TI)
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 62f3c613dea..d120405977d 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -1940,6 +1940,18 @@ export const $ClearTaskInstancesBody = {
title: 'Task Ids',
description: 'A list of `task_id` or [`task_id`, `map_index`]. If
only the `task_id` is provided for a mapped task, all of its map indices will
be targeted.'
},
+ task_group_id: {
+ anyOf: [
+ {
+ type: 'string'
+ },
+ {
+ type: 'null'
+ }
+ ],
+ title: 'Task Group Id',
+ description: "Clear every task in this task group. Mutually
exclusive with `task_ids`. The group's tasks are resolved on the server from
the dag structure, so all of them are targeted regardless of how many there
are."
+ },
dag_run_id: {
anyOf: [
{
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 3b25feb93cf..2fee8a825a2 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -615,6 +615,10 @@ export type ClearTaskInstancesBody = {
string,
number
])> | null;
+ /**
+ * Clear every task in this task group. Mutually exclusive with
`task_ids`. The group's tasks are resolved on the server from the dag
structure, so all of them are targeted regardless of how many there are.
+ */
+ task_group_id?: string | null;
dag_run_id?: string | null;
include_upstream?: boolean;
include_downstream?: boolean;
diff --git
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
index a215486ebe0..47ab623a68a 100644
---
a/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
+++
b/airflow-core/src/airflow/ui/src/components/Clear/TaskInstance/ClearGroupTaskInstanceDialog.tsx
@@ -22,11 +22,7 @@ import { useTranslation } from "react-i18next";
import { CgRedo } from "react-icons/cg";
import { useParams } from "react-router-dom";
-import {
- useDagRunServiceGetDagRun,
- useDagServiceGetDagDetails,
- useTaskInstanceServiceGetTaskInstances,
-} from "openapi/queries";
+import { useDagRunServiceGetDagRun, useDagServiceGetDagDetails } from
"openapi/queries";
import type { LightGridTaskInstanceSummary, TaskInstanceResponse } from
"openapi/requests/types.gen";
import { ActionAccordion } from "src/components/ActionAccordion";
import { useRerunWithLatestVersion } from
"src/components/Clear/useRerunWithLatestVersion";
@@ -68,20 +64,6 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
dagId,
});
- const { data: groupTaskInstances } = useTaskInstanceServiceGetTaskInstances(
- {
- dagId,
- dagRunId: runId,
- taskGroupId: groupId,
- },
- undefined,
- {
- enabled: open,
- },
- );
-
- const groupTaskIds = groupTaskInstances?.task_instances.map((ti) =>
ti.task_id) ?? [];
-
const { data: dagRun } = useDagRunServiceGetDagRun({ dagId, dagRunId: runId
}, undefined, {
enabled: open,
});
@@ -108,7 +90,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
const { data } = useClearTaskInstancesDryRun({
dagId,
options: {
- enabled: open && groupTaskIds.length > 0,
+ enabled: open,
refetchInterval: (query) =>
query.state.data?.task_instances.some((ti: TaskInstanceResponse) =>
isStatePending(ti.state))
? refetchInterval
@@ -123,7 +105,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
include_upstream: upstream,
only_failed: onlyFailed,
run_on_latest_version: runOnLatestVersion,
- task_ids: groupTaskIds,
+ task_group_id: groupId,
},
});
@@ -203,7 +185,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
</Checkbox>
) : undefined}
<Button
- disabled={affectedTasks.total_entries === 0 ||
groupTaskIds.length === 0}
+ disabled={affectedTasks.total_entries === 0}
loading={isPending}
onClick={() => {
mutate({
@@ -218,7 +200,7 @@ export const ClearGroupTaskInstanceDialog = ({ onClose,
open, taskInstance }: Pr
...(note === null ? {} : { note }),
only_failed: onlyFailed,
run_on_latest_version: runOnLatestVersion,
- task_ids: groupTaskIds,
+ task_group_id: groupId,
},
});
}}
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
index e7995edb747..2ba2cedec2d 100644
---
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
+++
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_task_instances.py
@@ -48,7 +48,7 @@ from airflow.models.taskinstancehistory import
TaskInstanceHistory
from airflow.models.taskmap import TaskMap
from airflow.models.team import Team
from airflow.models.trigger import Trigger
-from airflow.sdk import BaseOperator
+from airflow.sdk import BaseOperator, TaskGroup
from airflow.state.metastore import MetastoreBackend
from airflow.utils.platform import getuser
from airflow.utils.state import DagRunState, State, TaskInstanceState
@@ -4266,6 +4266,102 @@ class
TestPostClearTaskInstances(TestTaskInstanceEndpoint):
ti_id = response_data["task_instances"][0]["id"]
_check_task_instance_note(session, ti_id, {"content":
"placeholder-note", "user_id": None})
+ @pytest.mark.parametrize(
+ ("task_group_id", "expected_task_ids"),
+ [
+ pytest.param(
+ "section_1",
+ ["section_1.task_1", "section_1.task_2", "section_1.task_3"],
+ id="flat group",
+ ),
+ pytest.param(
+ "section_2",
+ [
+ "section_2.task_1",
+ "section_2.inner_section_2.task_2",
+ "section_2.inner_section_2.task_3",
+ "section_2.inner_section_2.task_4",
+ ],
+ id="nested group resolves recursively",
+ ),
+ ],
+ )
+ def test_clear_by_task_group_id_targets_every_task_in_the_group(
+ self, test_client, session, task_group_id, expected_task_ids
+ ):
+ """Clearing by task_group_id resolves the whole group server-side from
the dag structure."""
+ self.create_task_instances(session, dag_id="example_task_group")
+ response = test_client.post(
+ "/dags/example_task_group/clearTaskInstances",
+ json={
+ "dry_run": True,
+ "only_failed": False,
+ "dag_run_id": "TEST_DAG_RUN_ID",
+ "task_group_id": task_group_id,
+ },
+ )
+ assert response.status_code == 200
+ response_data = response.json()
+ assert response_data["total_entries"] == len(expected_task_ids)
+ assert sorted(ti["task_id"] for ti in response_data["task_instances"])
== sorted(expected_task_ids)
+
+ def test_clear_by_task_group_id_not_found(self, test_client, session):
+ """An unknown task_group_id returns 404."""
+ self.create_task_instances(session, dag_id="example_task_group")
+ response = test_client.post(
+ "/dags/example_task_group/clearTaskInstances",
+ json={"dry_run": True, "dag_run_id": "TEST_DAG_RUN_ID",
"task_group_id": "nonexistent_group"},
+ )
+ assert response.status_code == 404
+ assert "nonexistent_group" in response.json()["detail"]
+
+ def test_clear_rejects_both_task_ids_and_task_group_id(self, test_client,
session):
+ """task_ids and task_group_id are mutually exclusive."""
+ self.create_task_instances(session, dag_id="example_task_group")
+ response = test_client.post(
+ "/dags/example_task_group/clearTaskInstances",
+ json={
+ "dry_run": True,
+ "dag_run_id": "TEST_DAG_RUN_ID",
+ "task_ids": ["section_1.task_1"],
+ "task_group_id": "section_1",
+ },
+ )
+ assert response.status_code == 422
+
+ def test_clear_by_task_group_id_handles_group_larger_than_page_size(
+ self, test_client, dag_maker, session
+ ):
+ """A group with more tasks than the API page size is cleared in full
(regression for #59235)."""
+ group_size = 60 # deliberately larger than the default 50-row page
the UI used to cap at
+ dag_id = "large_task_group_dag"
+ with dag_maker(session=session, dag_id=dag_id,
start_date=DEFAULT_DATETIME_1, serialized=True):
+ with TaskGroup("big_group"):
+ for index in range(group_size):
+ BaseOperator(task_id=f"task_{index}")
+ dr = dag_maker.create_dagrun(
+ run_id="run_large_group",
+ logical_date=DEFAULT_DATETIME_1,
+ data_interval=(DEFAULT_DATETIME_1, DEFAULT_DATETIME_2),
+ )
+ DagBundlesManager().sync_bundles_to_db()
+ dagbag = DagBag(os.devnull)
+ dagbag.dags = {dag_id: dag_maker.dag}
+ sync_bag_to_db(dagbag, "dags-folder", None)
+ session.flush()
+
+ response = test_client.post(
+ f"/dags/{dag_id}/clearTaskInstances",
+ json={
+ "dry_run": True,
+ "only_failed": False,
+ "dag_run_id": dr.run_id,
+ "task_group_id": "big_group",
+ },
+ )
+ assert response.status_code == 200
+ assert response.json()["total_entries"] == group_size
+
class TestGetTaskInstanceTries(TestTaskInstanceEndpoint):
def test_should_respond_200(self, test_client, session):
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 9efa7cbf5ce..9ffef6d31a7 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -254,6 +254,13 @@ class ClearTaskInstancesBody(BaseModel):
title="Task Ids",
),
] = None
+ task_group_id: Annotated[
+ str | None,
+ Field(
+ description="Clear every task in this task group. Mutually
exclusive with `task_ids`. The group's tasks are resolved on the server from
the dag structure, so all of them are targeted regardless of how many there
are.",
+ title="Task Group Id",
+ ),
+ ] = None
dag_run_id: Annotated[str | None, Field(title="Dag Run Id")] = None
include_upstream: Annotated[bool | None, Field(title="Include Upstream")]
= False
include_downstream: Annotated[bool | None, Field(title="Include
Downstream")] = False