Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/asset_state_store.py: 68%

79 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""" 

18Execution API routes for asset state store. 

19 

20Routes are split into ``/by-name`` and ``/by-uri`` sub-prefixes mirroring the 

21existing ``/assets/by-name`` and ``/assets/by-uri`` pattern. Callers pass 

22whichever identifier their inlet type carries: ``Asset``/``AssetNameRef`` use 

23the name routes, ``AssetUriRef`` uses the URI routes. 

24 

25Per-task asset registration checks are intentionally not implemented here 

26(deferred to AIP-93 — see TODO comment below). 

27""" 

28 

29from __future__ import annotations 

30 

31import json 

32from typing import Annotated 

33from uuid import UUID 

34 

35from cadwyn import VersionedAPIRouter 

36from fastapi import HTTPException, Query, status 

37from sqlalchemy import select 

38 

39from airflow._shared.state import AssetScope, AssetStateStoreWriterKind 

40from airflow.api_fastapi.common.db.common import SessionDep 

41from airflow.api_fastapi.execution_api.datamodels.asset_state_store import ( 

42 AssetStateStorePutBody, 

43 AssetStateStoreResponse, 

44) 

45from airflow.api_fastapi.execution_api.datamodels.token import TIToken 

46from airflow.api_fastapi.execution_api.security import CurrentTIToken, ExecutionAPIRoute 

47from airflow.models.asset import AssetModel 

48from airflow.models.taskinstance import TaskInstance 

49from airflow.state import get_state_backend 

50from airflow.state.metastore import MetastoreBackend 

51 

52_TIWriterFields = tuple[str, str, str, int] 

53NULL_UUID = UUID(int=0) 

54 

55 

56def _fetch_ti_writer_fields(token: TIToken, session: SessionDep) -> _TIWriterFields: 

57 """Return (dag_id, run_id, task_id, map_index) for the TI identified by the token.""" 

58 row = session.execute( 

59 select( 

60 TaskInstance.dag_id, 

61 TaskInstance.run_id, 

62 TaskInstance.task_id, 

63 TaskInstance.map_index, 

64 ).where(TaskInstance.id == token.id) 

65 ).one_or_none() 

66 if row is None: 66 ↛ 67line 66 didn't jump to line 67 because the condition on line 66 was never true

67 raise HTTPException( 

68 status_code=status.HTTP_404_NOT_FOUND, 

69 detail={"reason": "not_found", "message": f"Task instance {token.id!r} not found"}, 

70 ) 

71 return row.dag_id, row.run_id, row.task_id, row.map_index 

72 

73 

74# TODO(AIP-103): enforce that the requesting task is registered with the asset 

75# (via task_inlet_asset_reference or task_outlet_asset_reference) before 

76# allowing reads/writes. Currently any task with a valid execution token can 

77# access any asset's state store — the same gap exists in /assets and /asset-events. 

78# Proper fix is a unified asset-registration check across all asset routes, 

79# not just here. 

80router = VersionedAPIRouter( 

81 route_class=ExecutionAPIRoute, 

82 responses={ 

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

84 status.HTTP_404_NOT_FOUND: {"description": "Not found"}, 

85 }, 

86) 

87 

88 

89def _resolve_asset_id_by_name(name: str, session: SessionDep) -> int: 

90 asset_id = session.scalar(select(AssetModel.id).where(AssetModel.name == name, AssetModel.active.has())) 

91 if asset_id is None: 91 ↛ 92line 91 didn't jump to line 92 because the condition on line 91 was never true

92 raise HTTPException( 

93 status_code=status.HTTP_404_NOT_FOUND, 

94 detail={"reason": "not_found", "message": f"Asset with name={name!r} not found"}, 

95 ) 

96 return asset_id 

97 

98 

99def _resolve_asset_id_by_uri(uri: str, session: SessionDep) -> int: 

100 asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri == uri, AssetModel.active.has())) 

101 if asset_id is None: 

102 raise HTTPException( 

103 status_code=status.HTTP_404_NOT_FOUND, 

104 detail={"reason": "not_found", "message": f"Asset with uri={uri!r} not found"}, 

105 ) 

106 return asset_id 

107 

108 

109@router.get("/by-name/value") 

110def get_asset_state_store_by_name( 

111 name: Annotated[str, Query(min_length=1)], 

112 key: Annotated[str, Query(min_length=1)], 

113 session: SessionDep, 

114) -> AssetStateStoreResponse: 

115 """Get an asset state store value by asset name.""" 

116 asset_id = _resolve_asset_id_by_name(name, session) 

117 value = get_state_backend().get(AssetScope(asset_id=asset_id), key, session=session) 

118 if value is None: 

119 raise HTTPException( 

120 status_code=status.HTTP_404_NOT_FOUND, 

121 detail={"reason": "not_found", "message": f"Asset state store key {key!r} not found"}, 

122 ) 

123 return AssetStateStoreResponse(value=json.loads(value)) 

124 

125 

126def _put_asset_state_store( 

127 scope: AssetScope, 

128 key: str, 

129 body: AssetStateStorePutBody, 

130 token: TIToken, 

131 session: SessionDep, 

132) -> None: 

133 backend = get_state_backend() 

134 if isinstance(backend, MetastoreBackend): 134 ↛ 160line 134 didn't jump to line 160 because the condition on line 134 was always true

135 if token.id == NULL_UUID: 135 ↛ 137line 135 didn't jump to line 137 because the condition on line 135 was never true

136 # Since the asset state store routes do not have `task_instance_id` in their path params, the default kicks in which is"00000000-0000-0000-0000-000000000000" 

137 backend.set_asset_state_store( 

138 scope, 

139 key, 

140 json.dumps(body.value), 

141 kind=AssetStateStoreWriterKind.WATCHER, 

142 session=session, 

143 ) 

144 else: 

145 ti_fields = _fetch_ti_writer_fields(token, session) 

146 dag_id, run_id, task_id, map_index = ti_fields 

147 

148 backend.set_asset_state_store( 

149 scope, 

150 key, 

151 json.dumps(body.value), 

152 kind=AssetStateStoreWriterKind.TASK, 

153 dag_id=dag_id, 

154 run_id=run_id, 

155 task_id=task_id, 

156 map_index=map_index, 

157 session=session, 

158 ) 

159 else: 

160 backend.set(scope, key, json.dumps(body.value), session=session) 

161 

162 

163@router.put("/by-name/value", status_code=status.HTTP_204_NO_CONTENT) 

164def set_asset_state_store_by_name( 

165 name: Annotated[str, Query(min_length=1)], 

166 key: Annotated[str, Query(min_length=1)], 

167 body: AssetStateStorePutBody, 

168 session: SessionDep, 

169 token: TIToken = CurrentTIToken, 

170) -> None: 

171 """Set an asset state store value by asset name.""" 

172 _put_asset_state_store( 

173 AssetScope(asset_id=_resolve_asset_id_by_name(name, session)), key, body, token, session 

174 ) 

175 

176 

177@router.delete("/by-name/value", status_code=status.HTTP_204_NO_CONTENT) 

178def delete_asset_state_store_by_name( 

179 name: Annotated[str, Query(min_length=1)], 

180 key: Annotated[str, Query(min_length=1)], 

181 session: SessionDep, 

182) -> None: 

183 """Delete a single asset state store key by asset name.""" 

184 asset_id = _resolve_asset_id_by_name(name, session) 

185 get_state_backend().delete(AssetScope(asset_id=asset_id), key, session=session) 

186 

187 

188@router.delete("/by-name/clear", status_code=status.HTTP_204_NO_CONTENT) 

189def clear_asset_state_store_by_name( 

190 name: Annotated[str, Query(min_length=1)], 

191 session: SessionDep, 

192) -> None: 

193 """Delete all state store keys for an asset by asset name.""" 

194 asset_id = _resolve_asset_id_by_name(name, session) 

195 get_state_backend().clear(AssetScope(asset_id=asset_id), session=session) 

196 

197 

198@router.get("/by-uri/value") 

199def get_asset_state_store_by_uri( 

200 uri: Annotated[str, Query(min_length=1)], 

201 key: Annotated[str, Query(min_length=1)], 

202 session: SessionDep, 

203) -> AssetStateStoreResponse: 

204 """Get an asset state store value by asset URI.""" 

205 asset_id = _resolve_asset_id_by_uri(uri, session) 

206 value = get_state_backend().get(AssetScope(asset_id=asset_id), key, session=session) 

207 if value is None: 

208 raise HTTPException( 

209 status_code=status.HTTP_404_NOT_FOUND, 

210 detail={"reason": "not_found", "message": f"Asset state store key {key!r} not found"}, 

211 ) 

212 return AssetStateStoreResponse(value=json.loads(value)) 

213 

214 

215@router.put("/by-uri/value", status_code=status.HTTP_204_NO_CONTENT) 

216def set_asset_state_store_by_uri( 

217 uri: Annotated[str, Query(min_length=1)], 

218 key: Annotated[str, Query(min_length=1)], 

219 body: AssetStateStorePutBody, 

220 session: SessionDep, 

221 token: TIToken = CurrentTIToken, 

222) -> None: 

223 """Set an asset state store value by asset URI.""" 

224 _put_asset_state_store( 

225 AssetScope(asset_id=_resolve_asset_id_by_uri(uri, session)), key, body, token, session 

226 ) 

227 

228 

229@router.delete("/by-uri/value", status_code=status.HTTP_204_NO_CONTENT) 

230def delete_asset_state_store_by_uri( 

231 uri: Annotated[str, Query(min_length=1)], 

232 key: Annotated[str, Query(min_length=1)], 

233 session: SessionDep, 

234) -> None: 

235 """Delete a single asset state store key by asset URI.""" 

236 asset_id = _resolve_asset_id_by_uri(uri, session) 

237 get_state_backend().delete(AssetScope(asset_id=asset_id), key, session=session) 

238 

239 

240@router.delete("/by-uri/clear", status_code=status.HTTP_204_NO_CONTENT) 

241def clear_asset_state_store_by_uri( 

242 uri: Annotated[str, Query(min_length=1)], 

243 session: SessionDep, 

244) -> None: 

245 """Delete all state store keys for an asset by asset URI.""" 

246 asset_id = _resolve_asset_id_by_uri(uri, session) 

247 get_state_backend().clear(AssetScope(asset_id=asset_id), session=session)