kaxil commented on code in PR #74222:
URL: https://github.com/apache/airflow/pull/74222#discussion_r4196073265


##########
airflow-core/src/airflow/api_fastapi/common/db/dag_runs.py:
##########
@@ -96,29 +95,17 @@ def attach_dag_versions_to_runs(dag_runs: Sequence[DagRun], 
*, session: Session)
 
     run_key_values = [(dr.dag_id, dr.run_id) for dr in runs_needing_versions]
 
-    ti_sub = (
+    rows = session.execute(

Review Comment:
   Same cause as the Gantt query: this used to union `TaskInstanceHistory`, and 
the hook now limits it to the working set. When an earlier try of a run ran on 
an older Dag version, that version drops out of `dag_versions` in the dag run 
list and grid responses. The non-prefetched `DagRun.dag_versions` path still 
picks it up through `historical_task_instances`, so the two paths now disagree. 
`include_all_attempts=True` here would fix it.



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/ui/gantt.py:
##########
@@ -61,7 +60,7 @@ def get_gantt_data(
     session: SessionDep,
 ) -> GanttResponse:
     """Get all task instance tries for Gantt chart."""
-    # Exclude mapped tasks (use grid summaries) and UP_FOR_RETRY (already in 
history)
+    # Pending retries retain timing for backoff; only the archived attempt 
belongs on the chart.

Review Comment:
   This select now goes through the `do_orm_execute` hook, so it only returns 
working-set rows. After a retry the chart shows only the latest try. During 
backoff the only live row is the UP_FOR_RETRY successor, which the clause below 
drops, so the task disappears from the chart entirely. 3.3.2 unioned 
`TaskInstanceHistory` here. Adding 
`.execution_options(include_all_attempts=True)` should be enough: `archive()` 
already moves unfinished states to FAILED, and the existing UP_FOR_RETRY clause 
still drops the pending successor. `test_gantt.py` has no retried-task case, 
which would have caught this. Fine as a follow-up.



##########
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_state_store.py:
##########
@@ -47,7 +47,7 @@
 
 def _get_task_scope_for_ti(task_instance_id: UUID, session: Session) -> 
TaskScope:
     ti = session.get(TI, task_instance_id)
-    if ti is None:
+    if ti is None or ti.working_set is not True:

Review Comment:
   A GET from an archived attempt gets 404 here, and the SDK maps any 404 to 
`TASK_STORE_NOT_FOUND`, so `task_state_store.get()` hands back the default. 
`ResumableJobMixin` reads that as no saved job id. If the attempt was archived 
while its worker was still running (heartbeat timeout through `handle_failure`, 
or adopt-or-reset), it submits the job again. At 2026-10-30 the `set` after 
that gets 410, so the new job id is never stored. Before this change the read 
went to the shared (dag_id, run_id, task_id, map_index) scope and returned the 
stored id. 3.3 also returned 404 here because the old UUID was gone. `get_xcom` 
already returns 410 for an archived attempt when 
`IdentifyArchivedTaskStateUpdates.is_applied`. Doing that here too would make 
the SDK raise. The `Airflow-API-Version` header in 
`test_archived_attempt_returns_404` has no effect because the `client` fixture 
stubs out `require_auth`. The real-auth setup from 
`test_archived_attempt_read_response_by_version` would cover both
  versions.



##########
airflow-core/src/airflow/migrations/versions/0142_3_4_0_unify_task_attempt_ownership.py:
##########
@@ -0,0 +1,484 @@
+#
+# 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.
+
+"""
+Unify task attempt ownership without rewriting legacy XCom data.

Review Comment:
   This drops `task_instance_history` and `hitl_detail_history` and renames 
`xcom` to `xcom_v1` and `rendered_task_instance_fields` to `rtif_v1`, but the 
PR has no newsfragment. Anyone with SQL, BI dashboards or `db clean --tables` 
scripts against those names only finds out at upgrade time. A `significant` 
newsfragment listing the old and new names would cover it.



##########
airflow-core/src/airflow/models/renderedtifields.py:
##########
@@ -335,15 +353,83 @@ def _do_delete_old_records(
         run_ids_to_keep: list[str] | ScalarSelect[str],
         session: Session,
     ) -> None:
-        # This query might deadlock occasionally and it should be retried if 
fails (see decorator)
-        stmt = (
+        from airflow.models.taskinstance import TaskInstance
+
+        legacy = LegacyRenderedTaskInstanceFields.__table__
+        session.execute(
+            delete(legacy).where(
+                legacy.c.dag_id == dag_id,
+                legacy.c.task_id == task_id,
+                legacy.c.run_id.not_in(run_ids_to_keep),
+            )
+        )
+        session.execute(
             delete(cls)

Review Comment:
   This prune runs on every rendered-fields write 
(`num_dag_runs_to_retain_rendered_fields` defaults to 30). `rtif_v2` has no 
dag_id/task_id/run_id columns, so the task scope only comes from the EXISTS 
into `task_instance`, and that matches every attempt of every run of this task, 
not only the retained window. The old delete was a range on the RTIF primary 
key. Since 0142 is unreleased, bounding the subquery to the runs that just left 
the window (or putting the coordinates on `rtif_v2`) would keep the per-write 
cost flat as history grows. `synchronize_session` also went from `False` to 
`"fetch"`, which adds a pre-select on MySQL with nothing in the session to sync.



##########
airflow-core/src/airflow/models/xcom.py:
##########
@@ -86,35 +99,123 @@ class XComModel(TaskInstanceDependencies):
         # separately, and enforce uniqueness with DagRun.id instead.
         Index("idx_xcom_key", key),
         Index("idx_xcom_task_instance", dag_id, task_id, run_id, map_index),
-        CheckConstraint(mapped_length >= 0, name="mapped_length_not_negative"),
+        CheckConstraint(mapped_length >= 0, 
name=conv("ck_xcom_mapped_length_not_negative")),
         PrimaryKeyConstraint("dag_run_id", "task_id", "map_index", "key", 
name="xcom_pkey"),
         ForeignKeyConstraint(
             [dag_id, task_id, run_id, map_index],
             [
-                "task_instance.dag_id",
-                "task_instance.task_id",
-                "task_instance.run_id",
-                "task_instance.map_index",
+                "legacy_task_data_owner.dag_id",
+                "legacy_task_data_owner.task_id",
+                "legacy_task_data_owner.run_id",
+                "legacy_task_data_owner.map_index",
             ],
             name="xcom_task_instance_fkey",
             ondelete="CASCADE",
         ),
     )
 
-    dag_run = relationship(
-        "DagRun",
-        primaryjoin="XComModel.dag_run_id == foreign(DagRun.id)",
-        uselist=False,
-        lazy="joined",
-        passive_deletes="all",
-    )
-    logical_date = association_proxy("dag_run", "logical_date")
 
-    task = relationship(
-        "TaskInstance",
-        viewonly=True,
-        lazy="raise",
+class XComModelV2(Base):
+    """XCom values stored per task attempt, keyed by the attempt UUID."""
+
+    __tablename__ = "xcom_v2"
+
+    id: Mapped[UUID] = mapped_column(Uuid(), primary_key=True, 
default=uuid6.uuid7)
+    task_instance_id: Mapped[UUID] = mapped_column(Uuid(), nullable=False)
+    key: Mapped[str] = mapped_column(String(512, **COLLATION_ARGS), 
nullable=False)
+    value: Mapped[Any] = mapped_column(JSON().with_variant(postgresql.JSONB, 
"postgresql"), nullable=True)
+    timestamp: Mapped[datetime] = mapped_column(UtcDateTime, 
default=timezone.utcnow, nullable=False)
+    dag_result: Mapped[bool | None] = mapped_column(Boolean, nullable=True, 
default=False)
+    mapped_length: Mapped[int | None] = mapped_column(Integer, nullable=True)
+
+    __table_args__ = (
+        PrimaryKeyConstraint("id", name="xcom_v2_pkey"),
+        UniqueConstraint("task_instance_id", "key", name="xcom_v2_ti_key_uq"),

Review Comment:
   `xcom_v1` keeps `idx_xcom_key`, but `xcom_v2` has no index on `key` alone. 
The global XCom list (`/dags/~/dagRuns/~/taskInstances/~/xcomEntries` with a 
key or key-prefix filter, which the Browse > XComs page sends) used that index, 
and every row written after the upgrade lands here. Adding the same index to 
`xcom_v2` in 0142 would keep that search on an index.



##########
airflow-core/src/airflow/models/taskinstance.py:
##########
@@ -639,7 +701,9 @@ class TaskInstance(Base, LoggingMixin, BaseWorkload):
     duration: Mapped[float | None] = mapped_column(Float, nullable=True)
     state: Mapped[str | None] = mapped_column(String(20), nullable=True)
     try_number: Mapped[int] = mapped_column(Integer, default=0)
-    max_tries: Mapped[int] = mapped_column(Integer, server_default="-1")
+    max_tries: Mapped[int] = mapped_column(Integer, server_default="-1", 
nullable=False)
+    working_set: Mapped[bool | None] = mapped_column(Boolean, default=True, 
server_default=true())
+    archived_reason: Mapped[str | None] = mapped_column(String(50))

Review Comment:
   `prepare_db_for_next_try` is the only caller of `archive()` and always 
passes `reason="retry"`, but it also runs for a user clear, restart completion, 
the orphan reset in the scheduler and the restore of a removed task. So every 
archived row says "retry" (or "legacy" from the migration), and nothing reads 
the column yet. Could `prepare_db_for_next_try` take the reason (as a 
`Literal`) so the value means something before AIP-111 or the UI starts relying 
on it?



-- 
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]

Reply via email to