Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/asset_events.py: 56%
39 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 typing import Annotated
22from fastapi import APIRouter, HTTPException, Query, status
23from sqlalchemy import and_, select
25from airflow.api_fastapi.common.db.common import SessionDep
26from airflow.api_fastapi.common.types import UtcDateTime
27from airflow.api_fastapi.execution_api.datamodels.asset import AssetResponse
28from airflow.api_fastapi.execution_api.datamodels.asset_event import (
29 AssetEventResponse,
30 AssetEventsResponse,
31)
32from airflow.models.asset import AssetAliasModel, AssetEvent, AssetModel
34router = APIRouter(
35 responses={
36 status.HTTP_404_NOT_FOUND: {"description": "Asset not found"},
37 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"},
38 },
39)
42def _get_asset_events_through_sql_clauses(
43 *, join_clause, where_clause, session: SessionDep, ascending: bool = True, limit: int | None = None
44) -> AssetEventsResponse:
45 order_by_clause = AssetEvent.timestamp.asc() if ascending else AssetEvent.timestamp.desc()
46 asset_events_query = select(AssetEvent).join(join_clause).where(where_clause).order_by(order_by_clause)
47 if limit: 47 ↛ 48line 47 didn't jump to line 48 because the condition on line 47 was never true
48 asset_events_query = asset_events_query.limit(limit)
49 asset_events = session.scalars(asset_events_query)
50 return AssetEventsResponse.model_validate(
51 {
52 "asset_events": [
53 AssetEventResponse(
54 id=event.id,
55 timestamp=event.timestamp,
56 extra=event.extra,
57 asset=AssetResponse(
58 name=event.asset.name,
59 uri=event.asset.uri,
60 group=event.asset.group,
61 extra=event.asset.extra,
62 ),
63 created_dagruns=event.created_dagruns,
64 source_task_id=event.source_task_id,
65 source_dag_id=event.source_dag_id,
66 source_run_id=event.source_run_id,
67 source_map_index=event.source_map_index,
68 partition_key=event.partition_key,
69 )
70 for event in asset_events
71 ]
72 }
73 )
76@router.get("/by-asset")
77def get_asset_event_by_asset_name_uri(
78 name: Annotated[str | None, Query(description="The name of the Asset")],
79 uri: Annotated[str | None, Query(description="The URI of the Asset")],
80 session: SessionDep,
81 after: Annotated[UtcDateTime | None, Query(description="The start of the time range")] = None,
82 before: Annotated[UtcDateTime | None, Query(description="The end of the time range")] = None,
83 ascending: Annotated[bool, Query(description="Whether to sort results in ascending order")] = True,
84 limit: Annotated[int | None, Query(description="The maximum number of results to return")] = None,
85) -> AssetEventsResponse:
86 if name and uri: 86 ↛ 87line 86 didn't jump to line 87 because the condition on line 86 was never true
87 where_clause = and_(AssetModel.name == name, AssetModel.uri == uri)
88 elif uri: 88 ↛ 90line 88 didn't jump to line 90 because the condition on line 88 was always true
89 where_clause = and_(AssetModel.uri == uri, AssetModel.active.has())
90 elif name:
91 where_clause = and_(AssetModel.name == name, AssetModel.active.has())
92 else:
93 raise HTTPException(
94 status_code=status.HTTP_400_BAD_REQUEST,
95 detail={
96 "reason": "Missing parameter",
97 "message": "name and uri cannot both be None",
98 },
99 )
101 if after: 101 ↛ 102line 101 didn't jump to line 102 because the condition on line 101 was never true
102 where_clause = and_(where_clause, AssetEvent.timestamp >= after)
103 if before: 103 ↛ 104line 103 didn't jump to line 104 because the condition on line 103 was never true
104 where_clause = and_(where_clause, AssetEvent.timestamp <= before)
106 return _get_asset_events_through_sql_clauses(
107 join_clause=AssetEvent.asset,
108 where_clause=where_clause,
109 session=session,
110 ascending=ascending,
111 limit=limit,
112 )
115@router.get("/by-asset-alias")
116def get_asset_event_by_asset_alias(
117 name: Annotated[str, Query(description="The name of the Asset Alias")],
118 session: SessionDep,
119 after: Annotated[UtcDateTime | None, Query(description="The start of the time range")] = None,
120 before: Annotated[UtcDateTime | None, Query(description="The end of the time range")] = None,
121 ascending: Annotated[bool, Query(description="Whether to sort results in ascending order")] = True,
122 limit: Annotated[int | None, Query(description="The maximum number of results to return")] = None,
123) -> AssetEventsResponse:
124 where_clause = AssetAliasModel.name == name
125 if after:
126 where_clause = and_(where_clause, AssetEvent.timestamp >= after)
127 if before:
128 where_clause = and_(where_clause, AssetEvent.timestamp <= before)
130 return _get_asset_events_through_sql_clauses(
131 join_clause=AssetEvent.source_aliases,
132 where_clause=where_clause,
133 session=session,
134 ascending=ascending,
135 limit=limit,
136 )