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
« 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.
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.
25Per-task asset registration checks are intentionally not implemented here
26(deferred to AIP-93 — see TODO comment below).
27"""
29from __future__ import annotations
31import json
32from typing import Annotated
33from uuid import UUID
35from cadwyn import VersionedAPIRouter
36from fastapi import HTTPException, Query, status
37from sqlalchemy import select
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
52_TIWriterFields = tuple[str, str, str, int]
53NULL_UUID = UUID(int=0)
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
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)
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
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
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))
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
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)
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 )
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)
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)
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))
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 )
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)
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)