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

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. 

17 

18from __future__ import annotations 

19 

20from collections.abc import Iterable 

21from typing import TYPE_CHECKING 

22 

23from sqlalchemy import func, select 

24 

25from airflow.models.asset import AssetEvent, association_table 

26from airflow.models.dagrun import DagRun 

27 

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 

30 

31 

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. 

37 

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 }