Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/services/public/assets.py: 86%

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

17from __future__ import annotations 

18 

19from collections.abc import Sequence 

20from typing import TYPE_CHECKING 

21 

22from airflow.api_fastapi.common.db.assets import resolve_triggering_event_refs 

23from airflow.api_fastapi.core_api.datamodels.assets import AssetEventResponse 

24 

25if TYPE_CHECKING: 25 ↛ 26line 25 didn't jump to line 26 because the condition on line 25 was never true

26 from sqlalchemy.orm import Session 

27 

28 from airflow.models.asset import AssetEvent 

29 

30 

31def serialize_asset_events(events: Sequence[AssetEvent], *, session: Session) -> list[AssetEventResponse]: 

32 """ 

33 Serialize asset events, flagging on each created dag run whether this event triggered it. 

34 

35 Only a run's most recent consumed event triggers it; the rest were merely included. The flag is 

36 resolved per ``(event_id, dag_id, run_id)`` tuple and set on the run before validation so the response 

37 carries it. 

38 """ 

39 triggering_refs = resolve_triggering_event_refs(events, session=session) 

40 asset_events = [] 

41 for event in events: 

42 for run in event.created_dagruns: 

43 run.triggering = (event.id, run.dag_id, run.run_id) in triggering_refs 

44 asset_events.append(AssetEventResponse.model_validate(event)) 

45 return asset_events