Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/asset_state_store.py: 100%

44 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 

19import json 

20from typing import Annotated 

21 

22from fastapi import Depends, HTTPException, status 

23from sqlalchemy import select 

24 

25from airflow._shared.state import AssetScope, AssetStateStoreWriterKind 

26from airflow.api_fastapi.common.db.common import SessionDep, paginated_select 

27from airflow.api_fastapi.common.parameters import QueryLimit, QueryOffset 

28from airflow.api_fastapi.common.router import AirflowRouter 

29from airflow.api_fastapi.core_api.datamodels.asset_state_store import ( 

30 AssetStateStoreBody, 

31 AssetStateStoreCollectionResponse, 

32 AssetStateStoreLastUpdatedBy, 

33 AssetStateStoreResponse, 

34) 

35from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

36from airflow.api_fastapi.core_api.security import requires_access_asset 

37from airflow.models.asset import AssetModel 

38from airflow.models.asset_state_store import AssetStateStoreModel 

39from airflow.state.metastore import _get_db_backend 

40 

41asset_state_store_router = AirflowRouter( 

42 tags=["Asset State Store"], 

43 prefix="/assets/{asset_id}/state-store", 

44) 

45 

46 

47def _get_asset_or_404(asset_id: int, session: SessionDep) -> int: 

48 exists = session.scalar(select(AssetModel.id).where(AssetModel.id == asset_id)) 

49 if exists is None: 

50 raise HTTPException( 

51 status_code=status.HTTP_404_NOT_FOUND, 

52 detail=f"Asset with id {asset_id!r} not found", 

53 ) 

54 return asset_id 

55 

56 

57AssetIdDep = Annotated[int, Depends(_get_asset_or_404)] 

58 

59 

60@asset_state_store_router.get( 

61 "", 

62 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

63 dependencies=[Depends(requires_access_asset(method="GET"))], 

64) 

65def list_asset_state_store( 

66 asset_id: AssetIdDep, 

67 limit: QueryLimit, 

68 offset: QueryOffset, 

69 session: SessionDep, 

70) -> AssetStateStoreCollectionResponse: 

71 """List all state store entries for an asset.""" 

72 base = ( 

73 select( 

74 AssetStateStoreModel.key, 

75 AssetStateStoreModel.value, 

76 AssetStateStoreModel.updated_at, 

77 AssetStateStoreModel.last_updated_by_kind, 

78 AssetStateStoreModel.last_updated_by_dag_id, 

79 AssetStateStoreModel.last_updated_by_run_id, 

80 AssetStateStoreModel.last_updated_by_task_id, 

81 AssetStateStoreModel.last_updated_by_map_index, 

82 ) 

83 .where(AssetStateStoreModel.asset_id == asset_id) 

84 .order_by(AssetStateStoreModel.key.asc()) 

85 ) 

86 paginated, total_entries = paginated_select( 

87 statement=base, 

88 filters=None, 

89 order_by=None, 

90 offset=offset, 

91 limit=limit, 

92 session=session, 

93 ) 

94 rows = session.execute(paginated).all() 

95 entries = [ 

96 AssetStateStoreResponse( 

97 key=r.key, 

98 value=json.loads(r.value), 

99 updated_at=r.updated_at, 

100 last_updated_by=AssetStateStoreLastUpdatedBy( 

101 kind=r.last_updated_by_kind, 

102 dag_id=r.last_updated_by_dag_id, 

103 run_id=r.last_updated_by_run_id, 

104 task_id=r.last_updated_by_task_id, 

105 map_index=r.last_updated_by_map_index, 

106 ) 

107 if r.last_updated_by_kind is not None 

108 else None, 

109 ) 

110 for r in rows 

111 ] 

112 return AssetStateStoreCollectionResponse(asset_state_store=entries, total_entries=total_entries) 

113 

114 

115@asset_state_store_router.get( 

116 "/{key:path}", 

117 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

118 dependencies=[Depends(requires_access_asset(method="GET"))], 

119) 

120def get_asset_state_store( 

121 asset_id: AssetIdDep, 

122 key: str, 

123 session: SessionDep, 

124) -> AssetStateStoreResponse: 

125 """Get a single asset state store entry.""" 

126 row = session.execute( 

127 select( 

128 AssetStateStoreModel.key, 

129 AssetStateStoreModel.value, 

130 AssetStateStoreModel.updated_at, 

131 AssetStateStoreModel.last_updated_by_kind, 

132 AssetStateStoreModel.last_updated_by_dag_id, 

133 AssetStateStoreModel.last_updated_by_run_id, 

134 AssetStateStoreModel.last_updated_by_task_id, 

135 AssetStateStoreModel.last_updated_by_map_index, 

136 ).where( 

137 AssetStateStoreModel.asset_id == asset_id, 

138 AssetStateStoreModel.key == key, 

139 ) 

140 ).one_or_none() 

141 if row is None: 

142 raise HTTPException( 

143 status_code=status.HTTP_404_NOT_FOUND, 

144 detail=f"Asset state store key {key!r} not found", 

145 ) 

146 return AssetStateStoreResponse( 

147 key=row.key, 

148 value=json.loads(row.value), 

149 updated_at=row.updated_at, 

150 last_updated_by=AssetStateStoreLastUpdatedBy( 

151 kind=row.last_updated_by_kind, 

152 dag_id=row.last_updated_by_dag_id, 

153 run_id=row.last_updated_by_run_id, 

154 task_id=row.last_updated_by_task_id, 

155 map_index=row.last_updated_by_map_index, 

156 ) 

157 if row.last_updated_by_kind is not None 

158 else None, 

159 ) 

160 

161 

162@asset_state_store_router.put( 

163 "/{key:path}", 

164 status_code=status.HTTP_204_NO_CONTENT, 

165 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

166 dependencies=[Depends(requires_access_asset(method="PUT"))], 

167) 

168def set_asset_state_store( 

169 asset_id: AssetIdDep, 

170 key: str, 

171 body: AssetStateStoreBody, 

172 session: SessionDep, 

173) -> None: 

174 """Set an asset state store value. Creates or overwrites the key.""" 

175 _get_db_backend().set_asset_state_store( 

176 AssetScope(asset_id=asset_id), 

177 key, 

178 json.dumps(body.value), 

179 kind=AssetStateStoreWriterKind.API, 

180 session=session, 

181 ) 

182 

183 

184@asset_state_store_router.delete( 

185 "/{key:path}", 

186 status_code=status.HTTP_204_NO_CONTENT, 

187 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

188 dependencies=[Depends(requires_access_asset(method="DELETE"))], 

189) 

190def delete_asset_state_store( 

191 asset_id: AssetIdDep, 

192 key: str, 

193 session: SessionDep, 

194) -> None: 

195 """Delete a single asset state store key. No-op if the key does not exist.""" 

196 _get_db_backend().delete(AssetScope(asset_id=asset_id), key, session=session) 

197 

198 

199@asset_state_store_router.delete( 

200 "", 

201 status_code=status.HTTP_204_NO_CONTENT, 

202 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

203 dependencies=[Depends(requires_access_asset(method="DELETE"))], 

204) 

205def clear_asset_state_store( 

206 asset_id: AssetIdDep, 

207 session: SessionDep, 

208) -> None: 

209 """Delete all state store keys for an asset.""" 

210 _get_db_backend().clear(AssetScope(asset_id=asset_id), session=session)