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

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 

21from uuid import UUID 

22 

23from cadwyn import VersionedAPIRouter 

24from fastapi import HTTPException, Path, Security, status 

25from sqlalchemy.orm import Session 

26 

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 

36 

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) 

46 

47 

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) 

59 

60 

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)) 

79 

80 

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) 

91 

92 

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) 

102 

103 

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)