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

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 typing import Annotated 

21 

22from fastapi import APIRouter, HTTPException, Query, status 

23from sqlalchemy import and_, select 

24 

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 

33 

34router = APIRouter( 

35 responses={ 

36 status.HTTP_404_NOT_FOUND: {"description": "Asset not found"}, 

37 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"}, 

38 }, 

39) 

40 

41 

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 ) 

74 

75 

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 ) 

100 

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) 

105 

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 ) 

113 

114 

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) 

129 

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 )