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
« 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
19from collections.abc import Sequence
20from typing import TYPE_CHECKING
22from airflow.api_fastapi.common.db.assets import resolve_triggering_event_refs
23from airflow.api_fastapi.core_api.datamodels.assets import AssetEventResponse
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
28 from airflow.models.asset import AssetEvent
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.
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