Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/datamodels/asset_state_store.py: 90%
36 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
22from pydantic import JsonValue, field_validator
24from airflow._shared.state import AssetStateStoreWriterKind
25from airflow.api_fastapi.core_api.base import BaseModel, StrictBaseModel
26from airflow.configuration import conf
29class AssetStateStoreLastUpdatedBy(BaseModel):
30 """Writer info for the last write to an asset state store entry."""
32 kind: AssetStateStoreWriterKind
33 dag_id: str | None = None
34 run_id: str | None = None
35 task_id: str | None = None
36 map_index: int | None = None
39class AssetStateStoreResponse(BaseModel):
40 """A single asset state store key/value pair with metadata."""
42 key: str
43 value: JsonValue
44 updated_at: datetime
45 last_updated_by: AssetStateStoreLastUpdatedBy | None = None
48class AssetStateStoreCollectionResponse(BaseModel):
49 """All asset state store entries for an asset."""
51 asset_state_store: list[AssetStateStoreResponse]
52 total_entries: int
55class AssetStateStoreBody(StrictBaseModel):
56 """Request body for setting an asset state store value."""
58 value: JsonValue
60 @field_validator("value")
61 @classmethod
62 def value_is_json_representable(cls, v: JsonValue) -> JsonValue:
63 if v is None:
64 raise ValueError("value cannot be null")
65 try:
66 serialized = json.dumps(v, allow_nan=False)
67 except ValueError:
68 raise ValueError("value contains non-finite numbers; NaN and Inf are not JSON representable")
69 limit = conf.getint("state_store", "max_value_storage_bytes")
70 if limit > 0 and len(serialized) > limit: 70 ↛ 71line 70 didn't jump to line 71 because the condition on line 70 was never true
71 raise ValueError(
72 f"value exceeds max_value_storage_bytes ({limit}); "
73 "raise [state_store] max_value_storage_bytes or set it to 0 to disable the limit"
74 )
75 return v