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


##########
airflow-core/src/airflow/api_fastapi/core_api/services/public/event_logs.py:
##########
@@ -16,14 +16,27 @@
 # under the License.
 from __future__ import annotations
 
-from sqlalchemy import inspect
+from sqlalchemy import case, inspect, select
 from sqlalchemy.orm.attributes import set_committed_value
 
 from airflow.api_fastapi.core_api.datamodels.event_logs import EventLogResponse
 from airflow.models import Log
+from airflow.models.task_coordinates import public_map_index_expression
+from airflow.models.taskinstance import TaskInstance
 
 
-def event_log_to_response(event_log: Log) -> EventLogResponse:
+def event_log_public_map_index():
+    """Resolve an attributed execution's public index without guessing once 
its row has been purged."""
+    attributed = (
+        select(public_map_index_expression(TaskInstance))

Review Comment:
   `event_log_public_map_index()` builds its subquery as an ORM select on 
`TaskInstance`, so `_restrict_to_current_attempts` adds `working_set IS TRUE` 
to it (the outer statement isn't a PK lookup, and `with_loader_criteria` 
reaches into scalar subqueries). Once an attempt is archived, which 
`prepare_db_for_next_try` does on every retry, the subquery finds nothing and 
the CASE yields NULL. So a plain task that failed try 1 and retried now shows 
`map_index: null` for try 1's audit rows, and `GET /eventLogs?...&map_index=-1` 
drops them from both the page and `total_entries`. That's the opposite of the 
"keeps doing so after that attempt is retired" behaviour in the description, 
and I'd expect the `archived=True` cases of 
`test_projects_public_mapping_index_for_exact_execution` to fail on it. 
Building the subquery against `TaskInstance.__table__` (which the hook leaves 
alone) should fix it. Setting `include_all_attempts` on the outer select alone 
wouldn't, because the count wraps the sta
 tement in a new select that doesn't carry the option.



##########
providers/openlineage/src/airflow/providers/openlineage/utils/utils.py:
##########
@@ -1652,6 +1703,39 @@ def _extract_ol_info_from_asset_event(asset_event: 
AssetEvent) -> dict[str, str]
     return None
 
 
+def _get_asset_event_source_regions(events: list[AssetEvent]) -> dict[UUID, 
UUID]:
+    if not AIRFLOW_V_3_4_PLUS:
+        return {}
+    source_ids = {
+        source_id
+        for event in events
+        if isinstance(source_id := getattr(event, "source_task_instance_id", 
None), UUID)
+    }
+    if not source_ids:
+        return {}
+
+    from sqlalchemy import select
+    from sqlalchemy.orm import object_session
+
+    from airflow.models.asset import AssetEvent
+    from airflow.utils.session import create_session
+
+    session = None
+    for event in events:
+        if isinstance(event, AssetEvent):
+            session = object_session(event)
+            if session is not None:
+                break
+    with nullcontext(session) if session is not None else create_session() as 
session:
+        return dict(
+            session.execute(
+                select(TaskInstance.id, 
TaskInstance.region_id).where(TaskInstance.id.in_(source_ids))

Review Comment:
   I don't think this lookup can change the result. It's an ORM select on 
`TaskInstance` with an `in_()` bind, so the attempt hook filters it to 
`working_set IS TRUE` and a retried or cleared source never shows up here. For 
live sources the caller already joinedloads `source_task_instance`, and 
`get_regional_task_instance_run_id(ti)` returns the same `str(ti.id)`. So in 
production it's one extra session and query per asset-triggered run event with 
no effect on the output. 
`test_cleared_regional_asset_source_keeps_retiring_uuid` passes because the 
relationship fallback resolves the retired TI from the session identity map, 
not through this branch. Could we drop `_get_asset_event_source_regions` and 
the `source_regions` shortcut and rely on the TI branch?



##########
airflow-core/src/airflow/models/log.py:
##########
@@ -73,15 +73,16 @@ class Log(Base):
     task_instance: Mapped[TaskInstance | None] = relationship(
         "TaskInstance",
         viewonly=True,
-        foreign_keys=[dag_id, task_id, run_id, map_index],
-        primaryjoin="and_(Log.dag_id == TaskInstance.dag_id, Log.task_id == 
TaskInstance.task_id, Log.run_id == TaskInstance.run_id, Log.map_index == 
TaskInstance.region_index, TaskInstance.working_set.is_(True))",
+        foreign_keys=[task_instance_id],
+        primaryjoin="Log.task_instance_id == TaskInstance.id",

Review Comment:
   With the join now only on `task_instance_id`, every audit row written before 
the upgrade (the column came in 0142 with no backfill) loses 
`task_display_name` in `GET /eventLogs`, where 3.3 returned it through the 
coordinate join. Is that intended? The display name is task-level, so a 
coordinate fallback for `task_instance_id IS NULL` rows wouldn't run into the 
wrong-attempt problem the description is avoiding. If it is intended, a 
newsfragment line would help API users who read that field.



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py:
##########
@@ -130,8 +136,14 @@ def get_event_logs(
     # Exact match filters (for backward compatibility)
     dag_id: Annotated[FilterParam[str | None], 
Depends(filter_param_factory(Log.dag_id, str | None))],
     task_id: Annotated[FilterParam[str | None], 
Depends(filter_param_factory(Log.task_id, str | None))],
+    task_instance_id: Annotated[

Review Comment:
   This adds a `task_instance_id` filter, but `EventLogResponse` doesn't return 
the field, and the same goes for `source_task_instance_id` on 
`AssetEventResponse`. Two loop passes of an unmapped task on try 1 come back 
with identical `task_id`, `map_index=-1` and `try_number=1`, so a client 
listing events can't tell which pass wrote which row without one request per 
attempt. `AssetStateStoreLastUpdatedBy` already exposes `task_instance_id` in 
this PR. Could the other two responses carry it as an optional field too?



##########
providers/openlineage/src/airflow/providers/openlineage/utils/utils.py:
##########
@@ -1466,6 +1495,14 @@ def is_dag_run_asset_triggered(
     return dag_run.run_type == DagRunType.DATASET_TRIGGERED  # type: 
ignore[attr-defined]  # This attr is available on AF2, but mypy can't see it
 
 
+def get_regional_task_instance_run_id(task_instance: TaskInstance | 
RuntimeTaskInstance) -> str | None:
+    if AIRFLOW_V_3_4_PLUS:
+        region_id = getattr(task_instance, "region_id", None)
+        if isinstance(region_id, UUID) and region_id.int != 0:

Review Comment:
   Since mapped expansions got their own region, this returns the TI UUID for 
every plain `.expand()` task too, not only loop work. Their OpenLineage run ids 
move from the deterministic `build_task_instance_run_id` value to `str(ti.id)` 
whenever the provider runs on a core that has regions, even though dag, task, 
try, logical date and map index already identified them uniquely. Is that 
intended? If so, it's a visible change for lineage consumers that reproduce run 
ids and probably deserves a provider changelog note. If not, the check would 
need to know whether the region sits inside a loop.



##########
airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py:
##########
@@ -204,7 +216,9 @@ def get_event_logs(
         # that bypass Log.__init__ (which always sets dttm = 
timezone.utcnow()).
         # Making EventLogResponse.when nullable would be a breaking API 
contract change for
         # clients that currently rely on `when` always being present.
-        
select(Log).where(Log.dttm.is_not(None)).options(*_eager_load_display_names())
+        select(Log, event_log_public_map_index())

Review Comment:
   Putting `event_log_public_map_index()` in the paginated statement means the 
count query (`select(count()).select_from(stmt.subquery())`) carries the 
correlated subquery too. Postgres and SQLite should flatten it away, but MySQL 
won't merge a derived table that has a subquery in its select list, so it 
materializes it and runs the TI lookup for every attributed log row on each 
audit-log page load. Adding the column after `paginated_select` 
(`event_logs_select.add_columns(...)`) would keep the count as cheap as before. 
I haven't run EXPLAIN on MySQL, so worth a quick check.



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