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

Reply via email to