Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/common/db/assets.py: 89%
14 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
18from __future__ import annotations
20from collections.abc import Iterable
21from typing import TYPE_CHECKING
23from sqlalchemy import func, select
25from airflow.models.asset import AssetEvent, association_table
26from airflow.models.dagrun import DagRun
28if TYPE_CHECKING: 28 ↛ 29line 28 didn't jump to line 29 because the condition on line 28 was never true
29 from sqlalchemy.orm import Session
32def resolve_triggering_event_refs(
33 events: Iterable[AssetEvent], *, session: Session
34) -> set[tuple[int, str, str]]:
35 """
36 Return ``(event_id, dag_id, run_id)`` tuples where the asset event actually triggered the run.
38 A dag run consumes every asset event queued when it is created, but only the run's most recent
39 consumed event triggers it; earlier consumed events are merely included.
40 """
41 run_ids = {run.id for event in events for run in event.created_dagruns}
42 if not run_ids:
43 return set()
44 ranked_events = (
45 select(
46 association_table.c.event_id,
47 DagRun.dag_id,
48 DagRun.run_id,
49 func.row_number()
50 .over(
51 partition_by=association_table.c.dag_run_id,
52 order_by=(AssetEvent.timestamp.desc(), AssetEvent.id.desc()),
53 )
54 .label("rank"),
55 )
56 .join(AssetEvent, AssetEvent.id == association_table.c.event_id)
57 .join(DagRun, DagRun.id == association_table.c.dag_run_id)
58 .where(association_table.c.dag_run_id.in_(run_ids))
59 .subquery()
60 )
61 return {
62 (event_id, dag_id, run_id)
63 for event_id, dag_id, run_id in session.execute(
64 select(ranked_events.c.event_id, ranked_events.c.dag_id, ranked_events.c.run_id).where(
65 ranked_events.c.rank == 1
66 )
67 )
68 }