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
« 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 datetime import datetime
21from typing import Literal
23from pydantic import AwareDatetime, JsonValue, field_validator
25from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel
26from airflow.configuration import conf
29class TaskStateStoreResponse(BaseModel):
30 """A single task state store key/value pair with metadata."""
32 key: str
33 value: JsonValue
34 updated_at: datetime
35 expires_at: datetime | None
38class TaskStateStoreCollectionResponse(BaseModel):
39 """All task state store entries for a task instance."""
41 task_state_store: list[TaskStateStoreResponse]
42 total_entries: int
45class TaskStateStoreBody(StrictBaseModel):
46 """
47 Request body for setting a task state store value.
49 ``expires_at`` controls expiry:
51 - ``"default"``: apply the configured ``[state_store] default_retention_days``.
52 - ``null``: never expire.
53 - aware datetime: expire at that time.
54 """
56 value: JsonValue
57 expires_at: AwareDatetime | None | Literal["default"] = "default"
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
77class TaskStateStorePatchBody(StrictBaseModel):
78 """Request body for patching only the value of an existing task state store key."""
80 value: JsonValue
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