Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/task_state_store.py: 55%
38 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
21from uuid import UUID
23from cadwyn import VersionedAPIRouter
24from fastapi import HTTPException, Path, Security, status
25from sqlalchemy.orm import Session
27from airflow._shared.state import TaskScope
28from airflow.api_fastapi.common.db.common import SessionDep
29from airflow.api_fastapi.execution_api.datamodels.task_state_store import (
30 TaskStateStorePutBody,
31 TaskStateStoreResponse,
32)
33from airflow.api_fastapi.execution_api.security import ExecutionAPIRoute, require_auth
34from airflow.models.taskinstance import TaskInstance as TI
35from airflow.state import get_state_backend
37router = VersionedAPIRouter(
38 route_class=ExecutionAPIRoute,
39 responses={
40 status.HTTP_401_UNAUTHORIZED: {"description": "Unauthorized"},
41 status.HTTP_403_FORBIDDEN: {"description": "Access denied"},
42 status.HTTP_404_NOT_FOUND: {"description": "Not found"},
43 },
44 dependencies=[Security(require_auth, scopes=["ti:self"])],
45)
48def _get_task_scope_for_ti(task_instance_id: UUID, session: Session) -> TaskScope:
49 ti = session.get(TI, task_instance_id)
50 if ti is None:
51 raise HTTPException(
52 status_code=status.HTTP_404_NOT_FOUND,
53 detail={
54 "reason": "not_found",
55 "message": f"Task instance {task_instance_id} not found",
56 },
57 )
58 return TaskScope(dag_id=ti.dag_id, run_id=ti.run_id, task_id=ti.task_id, map_index=ti.map_index)
61@router.get("/{task_instance_id}/{key:path}")
62def get_task_state_store(
63 task_instance_id: UUID,
64 key: Annotated[str, Path(min_length=1)],
65 session: SessionDep,
66) -> TaskStateStoreResponse:
67 """Get value for a task state store key."""
68 scope = _get_task_scope_for_ti(task_instance_id, session)
69 value = get_state_backend().get(scope, key, session=session)
70 if value is None:
71 raise HTTPException(
72 status_code=status.HTTP_404_NOT_FOUND,
73 detail={
74 "reason": "not_found",
75 "message": f"Task state store key {key!r} not found",
76 },
77 )
78 return TaskStateStoreResponse(value=json.loads(value))
81@router.put("/{task_instance_id}/{key:path}", status_code=status.HTTP_204_NO_CONTENT)
82def set_task_state_store(
83 task_instance_id: UUID,
84 key: Annotated[str, Path(min_length=1)],
85 body: TaskStateStorePutBody,
86 session: SessionDep,
87) -> None:
88 """Set a task state store key, creating or updating the row."""
89 scope = _get_task_scope_for_ti(task_instance_id, session)
90 get_state_backend().set(scope, key, json.dumps(body.value), expires_at=body.expires_at, session=session)
93@router.delete("/{task_instance_id}/{key:path}", status_code=status.HTTP_204_NO_CONTENT)
94def delete_task_state_store(
95 task_instance_id: UUID,
96 key: Annotated[str, Path(min_length=1)],
97 session: SessionDep,
98) -> None:
99 """Delete a single task state store key."""
100 scope = _get_task_scope_for_ti(task_instance_id, session)
101 get_state_backend().delete(scope, key, session=session)
104@router.delete("/{task_instance_id}", status_code=status.HTTP_204_NO_CONTENT)
105def clear_task_state_store(
106 task_instance_id: UUID,
107 session: SessionDep,
108) -> None:
109 """Delete all task state store keys for this task instance."""
110 scope = _get_task_scope_for_ti(task_instance_id, session)
111 get_state_backend().clear(scope, session=session)