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
« 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
19import json
20from typing import Annotated
22from fastapi import Depends, HTTPException, status
23from sqlalchemy import select
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
41asset_state_store_router = AirflowRouter(
42 tags=["Asset State Store"],
43 prefix="/assets/{asset_id}/state-store",
44)
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
57AssetIdDep = Annotated[int, Depends(_get_asset_or_404)]
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)
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 )
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 )
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)
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)