Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/datamodels/task_state_store.py: 85%

46 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 datetime import datetime 

21from typing import Literal 

22 

23from pydantic import AwareDatetime, JsonValue, field_validator 

24 

25from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel 

26from airflow.configuration import conf 

27 

28 

29class TaskStateStoreResponse(BaseModel): 

30 """A single task state store key/value pair with metadata.""" 

31 

32 key: str 

33 value: JsonValue 

34 updated_at: datetime 

35 expires_at: datetime | None 

36 

37 

38class TaskStateStoreCollectionResponse(BaseModel): 

39 """All task state store entries for a task instance.""" 

40 

41 task_state_store: list[TaskStateStoreResponse] 

42 total_entries: int 

43 

44 

45class TaskStateStoreBody(StrictBaseModel): 

46 """ 

47 Request body for setting a task state store value. 

48 

49 ``expires_at`` controls expiry: 

50 

51 - ``"default"``: apply the configured ``[state_store] default_retention_days``. 

52 - ``null``: never expire. 

53 - aware datetime: expire at that time. 

54 """ 

55 

56 value: JsonValue 

57 expires_at: AwareDatetime | None | Literal["default"] = "default" 

58 

59 @field_validator("value") 

60 @classmethod 

61 def value_is_json_representable(cls, v: JsonValue) -> JsonValue: 

62 if v is None: 

63 raise ValueError("value cannot be null") 

64 try: 

65 serialized = json.dumps(v, allow_nan=False) 

66 except ValueError: 

67 raise ValueError("value contains non-finite numbers; NaN and Inf are not JSON representable") 

68 limit = conf.getint("state_store", "max_value_storage_bytes") 

69 if limit > 0 and len(serialized) > limit: 69 ↛ 70line 69 didn't jump to line 70 because the condition on line 69 was never true

70 raise ValueError( 

71 f"value exceeds max_value_storage_bytes ({limit}); " 

72 "raise [state_store] max_value_storage_bytes or set it to 0 to disable the limit" 

73 ) 

74 return v 

75 

76 

77class TaskStateStorePatchBody(StrictBaseModel): 

78 """Request body for patching only the value of an existing task state store key.""" 

79 

80 value: JsonValue 

81 

82 @field_validator("value") 

83 @classmethod 

84 def value_is_json_representable(cls, v: JsonValue) -> JsonValue: 

85 if v is None: 

86 raise ValueError("value cannot be null") 

87 try: 

88 serialized = json.dumps(v, allow_nan=False) 

89 except ValueError: 

90 raise ValueError("value contains non-finite numbers; NaN and Inf are not JSON representable") 

91 limit = conf.getint("state_store", "max_value_storage_bytes") 

92 if limit > 0 and len(serialized) > limit: 92 ↛ 93line 92 didn't jump to line 93 because the condition on line 92 was never true

93 raise ValueError( 

94 f"value exceeds max_value_storage_bytes ({limit}); " 

95 "raise [state_store] max_value_storage_bytes or set it to 0 to disable the limit" 

96 ) 

97 return v