kaxil commented on code in PR #74350:
URL: https://github.com/apache/airflow/pull/74350#discussion_r4198251299
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py:
##########
@@ -1491,5 +1579,10 @@ def delete_task_instance(
)
TI.delete_attempts(
- dag_id=dag_id, run_id=dag_run_id, task_id=task_id,
map_index=map_index, session=session
+ dag_id=dag_id,
+ run_id=dag_run_id,
+ task_id=task_id,
+ map_index=map_index,
Review Comment:
The existence check above uses `scope.region_index`, but this passes the raw
`map_index` query param, which defaults to -1. For a loop pass addressed as
`?region_id=R®ion_index=2`, the lookup finds the row and then
`delete_attempts` deletes `region_index = -1` in R, which matches nothing, so
the endpoint reports success and the pass stays. A mapped slot addressed by
region coordinates hits the same thing. Should this be
`map_index=scope.region_index`? A delete test with region coordinates would
catch it.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py:
##########
@@ -393,7 +427,8 @@ def get_task_instance_tries(
TI.dag_id == dag_id,
TI.run_id == dag_run_id,
TI.task_id == task_id,
- TI.map_index == map_index,
+ TI.region_index == scope.region_index,
+ TI.region_id == scope.region_id,
Review Comment:
Before this, `/tries` listed every attempt at the map index whatever its
region. Now it only sees the region the live producer resolves to. For a run
from 3.3, clearing a whole mapped task archives the sentinel-region rows and
re-expands into a new region, so the pre-upgrade attempts and their logs drop
out of `/tries` unless the client already knows to pass
`region_id=00000000-0000-0000-0000-000000000000`, and nothing in the response
points there. Is that intended? If so, maybe list archived attempts from the
sentinel region as well, or at least document it on the route.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/log.py:
##########
@@ -196,7 +206,9 @@ def get_external_log_url(
TaskInstance.task_id == task_id,
TaskInstance.dag_id == dag_id,
TaskInstance.run_id == dag_run_id,
- TaskInstance.map_index == map_index,
+ TaskInstance.region_index == scope.region_index,
+ TaskInstance.region_id == scope.region_id,
+ TaskInstance.try_number == try_number,
Review Comment:
This now filters on `try_number`, but unlike `get_log` above it doesn't set
`.execution_options(include_all_attempts=True)`, so the `working_set IS TRUE`
criteria added by `_restrict_to_current_attempts` hides every archived try.
After a retry, `externalLogUrl/1` returns 404 where it used to return a link
(the old query ignored the try and used the live row). As far as I can trace,
`test_execution_lists_live_work_and_selects_archived_tries` can't pass on this,
and the `history` lookup in `test_external_link_uses_the_selected_archived_try`
goes through the same filter, so it returns None after the setup's
`dag.clear()` and needs the option too.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/task_instances.py:
##########
@@ -668,14 +734,18 @@ def get_task_instances(
)
load_legacy_rendered_fields(task_instances, session=session)
return TaskInstanceCollectionResponse(
- task_instances=task_instances,
+ task_instances=[
+ task_coordinate_response(TaskInstanceResponse, ti, resolver)
for ti in task_instances
+ ],
total_entries=total_entries,
total_entries_limit=total_entries_limit,
next_cursor=(
- encode_cursor(task_instances[-1], order_by) if has_next and
task_instances else None
+ encode_cursor(TaskCoordinateView(task_instances[-1],
resolver), order_by)
Review Comment:
The cursor is encoded through `TaskCoordinateView`, so it uses
`resolver.public_map_index` (the pinned definition's `get_needs_expansion()`).
The ordering and the `map_index` filter use `public_map_index_expression`
(region `node_id == task_id`), and `get_execution` has a third inline copy of
the region rule. They agree today, but a mapped slot cleared with run-on-latest
after `.expand()` was removed keeps its region and index while being re-pinned
to an unmapped definition. SQL then sorts and filters it as 3 while the
response and cursor say -1, and the next page can repeat rows. Could the
response side use the stored-region rule as well (`_is_mapped_region` already
does that), so the two agree by construction?
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/execution.py:
##########
@@ -0,0 +1,116 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+from typing import Annotated
+from uuid import UUID
+
+from fastapi import Depends, HTTPException, Query, status
+from sqlalchemy import select
+
+from airflow.api_fastapi.common.db.common import SessionDep
+from airflow.api_fastapi.common.parameters import QueryLimit, QueryOffset
+from airflow.api_fastapi.common.router import AirflowRouter
+from airflow.api_fastapi.core_api.datamodels.execution import (
+ ExecutionCollectionResponse,
+ ExecutionRegionResponse,
+ ExecutionTaskResponse,
+)
+from airflow.api_fastapi.core_api.openapi.exceptions import
create_openapi_http_exception_doc
+from airflow.api_fastapi.core_api.security import DagAccessEntity,
requires_access_dag
+from airflow.api_fastapi.core_api.services.public.execution import
get_execution_members
+from airflow.models.dagrun import DagRun
+from airflow.models.dynamic_region import SENTINEL_REGION_ID,
load_region_ancestry
+
+execution_router = AirflowRouter(
+ tags=["DagRun"],
+ prefix="/dags/{dag_id}/dagRuns/{dag_run_id}",
+)
+
+
+@execution_router.get(
+ "/execution",
Review Comment:
Does this need to be a public v2 endpoint? The UI is its only consumer, and
`ExecutionRegionResponse` commits `forked_from_region_id`, `resumes_from_index`
and `parent_region_index` to the stable API while the loop-clear model is still
moving. The TI listing already accepts and returns `region_id`/`region_index`,
so the new data here is mainly the ancestry. Serving that from the private UI
API would let the region model change before the release ships without a v2
deprecation.
##########
airflow-core/src/airflow/api_fastapi/core_api/services/public/task_coordinates.py:
##########
@@ -111,3 +113,34 @@ def resolve_task_scope(
map_index=tasks[0].region_index,
region_id=tasks[0].region_id,
)
+
+
+def _coordinate_resolver(dag_bag: DagBagDep, session: SessionDep) ->
TaskCoordinateResolver:
+ return TaskCoordinateResolver(dag_bag, session)
+
+
+CoordinateResolverDep = Annotated[TaskCoordinateResolver,
Depends(_coordinate_resolver)]
+
+
+def _task_scope(
+ dag_id: str,
+ dag_run_id: str,
+ task_id: str,
+ resolver: CoordinateResolverDep,
+ map_index: int = -1,
+ region_id: Annotated[UUID | None, Query()] = None,
+ region_index: Annotated[int | None, Query(ge=-1)] = None,
+) -> TaskScope:
+ return resolve_task_scope(
+ dag_id=dag_id,
+ run_id=dag_run_id,
+ task_id=task_id,
+ session=resolver.session,
+ dag_bag=resolver.dag_bag,
+ map_index=map_index,
Review Comment:
Two things about the addressing rules now that this dependency sits on the
`/{map_index}` routes. When both `region_id` and `region_index` are sent,
`resolve_task_scope` returns before looking at `map_index`, so `GET
.../taskInstances/m/7?region_id=M®ion_index=2` quietly returns slot 2 (the
PATCH path rejects conflicting body and query coordinates with a 400). And for
a loop task with no pass selected, as far as I can trace the client gets 400
without `region_id` but 409 "Multiple live producers" with `region_id` alone,
on routes where 409 otherwise means a state conflict (HITL's "already been
updated", for example). Maybe reject `map_index not in (-1, region_index)` with
a 400, and return 400 for the missing-index case in both shapes?
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/ui/gantt.py:
##########
@@ -26,6 +27,10 @@
class GanttTaskInstance(BaseModel):
"""Task instance data for Gantt chart."""
+ id: UUID
+ region_id: UUID
+ region_index: int
+ map_index: int
Review Comment:
These four required fields break the exact-dict assertion in
`test_gantt.py::test_should_response_200`, and through the regenerated
`GanttTaskInstance` type, the uncast `allTries` literals in
`layouts/Details/Gantt/utils.test.ts` (tsc runs in `pnpm lint`). Both are fixed
later in #74352, so this commit is red on its own. Separately, `map_index` is
always -1 here since the query filters on `public_map_index_expression(...) ==
-1`, and the UI only reads `id` and the region fields. Could it be dropped?
##########
airflow-core/src/airflow/models/xcom.py:
##########
@@ -267,6 +267,7 @@ def set(
task_id: str,
run_id: str,
map_index: int = -1,
+ region_id: UUID | None = SENTINEL_REGION_ID,
Review Comment:
Is `None` meant to be allowed here? With `region_id=None`,
`select_producers` drops the region filter and `session.scalar` takes whichever
live owner comes first among regions that share the map index. The only caller
passes a concrete `scope.region_id`, so `region_id: UUID = SENTINEL_REGION_ID`
would close that off. The docstring's param list is also missing `region_id`.
##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/extra_links.py:
##########
@@ -65,19 +69,16 @@ def get_extra_links(
task_id: str,
session: SessionDep,
dag_bag: DagBagDep,
- map_index: int = -1,
+ scope: TaskScopeDep,
Review Comment:
Moving `map_index` into `TaskScopeDep` changes its position in the generated
spec. FastAPI emits the route's own query params first, so `get_extra_links` is
now `try_number, map_index, ...` while the released Python client has
`get_extra_links(dag_id, dag_run_id, task_id, map_index, try_number)`. After a
client regen, a positional call like `get_extra_links("d", "r", "t", 3, 2)`
silently swaps slot and try, since both are optional ints. `get_xcom_entry`,
`get_log` and `clear_task_state_store` reorder too, though those fail on type
instead. Could `try_number` move into a dependency declared after `scope`, so
`map_index` stays ahead of it?
##########
airflow-core/src/airflow/api_fastapi/core_api/services/public/task_coordinates.py:
##########
@@ -111,3 +113,34 @@ def resolve_task_scope(
map_index=tasks[0].region_index,
region_id=tasks[0].region_id,
)
+
+
+def _coordinate_resolver(dag_bag: DagBagDep, session: SessionDep) ->
TaskCoordinateResolver:
+ return TaskCoordinateResolver(dag_bag, session)
+
+
+CoordinateResolverDep = Annotated[TaskCoordinateResolver,
Depends(_coordinate_resolver)]
+
+
+def _task_scope(
+ dag_id: str,
+ dag_run_id: str,
+ task_id: str,
+ resolver: CoordinateResolverDep,
+ map_index: int = -1,
+ region_id: Annotated[UUID | None, Query()] = None,
+ region_index: Annotated[int | None, Query(ge=-1)] = None,
+) -> TaskScope:
+ return resolve_task_scope(
+ dag_id=dag_id,
+ run_id=dag_run_id,
+ task_id=task_id,
+ session=resolver.session,
Review Comment:
`_task_scope` takes the request's `TaskCoordinateResolver` only to unpack
its session and dag bag, and `resolve_task_scope` then builds a second one.
Routes that also inject `CoordinateResolverDep` end up with two resolvers and
two sets of caches per request. Could `resolve_task_scope` take the resolver
instead?
##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_xcom.py:
##########
@@ -154,6 +172,130 @@ def _create_xcom(self, key, value, backend=None) -> None:
class TestGetXComEntry(TestXComEndpoint):
+ @pytest.mark.parametrize("mapped", [False, True])
+ def test_regional_xcom_projects_public_map_index(self, test_client,
dag_maker, session, mapped):
+ region = DynamicRegion(dag_id=TEST_DAG_ID, run_id=run_id,
node_id=TEST_TASK_ID if mapped else "loop")
+ session.add(region)
+ session.flush()
+ _create_regional_ti(dag_maker, session, region_id=region.id,
region_index=3)
+ XComModel.set(
+ dag_id=TEST_DAG_ID,
+ run_id=run_id,
+ task_id=TEST_TASK_ID,
+ region_id=region.id,
+ map_index=3,
+ key=TEST_XCOM_KEY,
+ value="live",
+ serialize=False,
+ session=session,
+ )
+ session.commit()
+ url =
f"/dags/{TEST_DAG_ID}/dagRuns/{run_id}/taskInstances/{TEST_TASK_ID}/xcomEntries"
+ expected_index = 3 if mapped else -1
+ response = test_client.get(
+ f"{url}/{TEST_XCOM_KEY}", params={"region_id": str(region.id),
"region_index": 3}
+ )
+ assert response.status_code == 200
+ assert response.json()["map_index"] == expected_index
+ assert response.json()["region_index"] == 3
+ response = test_client.get(url, params={"map_index_filter":
expected_index})
+ assert response.status_code == 200
+ assert response.json()["total_entries"] == 1
+ assert response.json()["xcom_entries"][0]["map_index"] ==
expected_index
+
+ @pytest.mark.parametrize(
+ ("params", "expected_indexes", "total"),
+ [
+ ({"map_index": -1}, [-1], 1),
+ ({"map_index_filter": 0}, [0], 1),
+ ({"order_by": "map_index", "limit": 1, "offset": 1}, [0], 3),
+ ],
+ )
+ def test_collection_projects_loop_and_mapped_group_before_pagination(
+ self, test_client, dag_maker, session, params, expected_indexes, total
+ ):
+ @task_group
+ def body():
+ EmptyOperator(task_id="member")
+
+ @task_group
+ def mapped_group(value):
+ EmptyOperator(task_id="member")
+
+ with dag_maker("regional-xcom", serialized=True):
+ loop = create_loop(body, max_iterations=7)
+ mapped_group.expand(value=[1, 2])
+ dr = dag_maker.create_dagrun()
+ regions = {}
+ for ti in dr.task_instances:
+ if ti.task_id == loop.gate_task_id:
+ continue
+ node_id = loop.group_id if ti.task_id == "body.member" else
ti.task_id
+ if node_id not in regions:
+ region = DynamicRegion(dag_id=dr.dag_id, run_id=dr.run_id,
node_id=node_id)
+ session.add(region)
+ session.flush()
+ regions[node_id] = region.id
+ ti.region_id = regions[node_id]
+ if ti.task_id == "body.member":
+ ti.map_index = 5
Review Comment:
`ti.map_index = 5` writes `region_index` through the synonym, and the test
then asserts `map_index == -1`. Writing `ti.region_index = 5` would say what's
meant, and moving the `region_index` expectation into a parametrize column
would drop the `if expected_indexes == [-1]` branch so every case checks it.
The regional test in `test_task_instances.py` builds its loop pass with
`map_index=2` the same way.
##########
airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_hitl.py:
##########
@@ -57,6 +61,61 @@
pytestmark = pytest.mark.db_test
+
[email protected]("operation", ["get", "history", "list", "respond"])
+def test_hitl_selects_loop_iteration(test_client, dag_maker, session,
operation):
+ @task_group
+ def body():
+ EmptyOperator(task_id="work")
+
+ with dag_maker(serialized=True) as dag:
+ loop = create_loop(body, max_iterations=4)
+ run = dag_maker.create_dagrun()
+ region = session.scalar(select(DynamicRegion).where(DynamicRegion.node_id
== loop.group_id))
+ later = TIModel(
+ task=dag.get_task("body.work"),
+ run_id=run.run_id,
+ dag_version_id=run.created_dag_version_id,
+ region_id=region.id,
Review Comment:
Only one region holds `body.work` here, so dropping `TI.region_id ==
scope.region_id` from the HITL readers would keep this green. The extra-links
tests are in the same spot, and `get_log` content, `/{map_index}` and delete
have no region test at all. A second region with the same `(task_id,
region_index)`, like
`test_exact_region_read_update_and_delete_preserve_sibling` in `test_xcom.py`,
would bind it. Splitting the four `operation` branches into separate tests
would also make each case easier to follow.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]