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

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 

21 

22from pydantic import JsonValue, field_validator 

23 

24from airflow._shared.state import AssetStateStoreWriterKind 

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

26from airflow.configuration import conf 

27 

28 

29class AssetStateStoreLastUpdatedBy(BaseModel): 

30 """Writer info for the last write to an asset state store entry.""" 

31 

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 

37 

38 

39class AssetStateStoreResponse(BaseModel): 

40 """A single asset state store key/value pair with metadata.""" 

41 

42 key: str 

43 value: JsonValue 

44 updated_at: datetime 

45 last_updated_by: AssetStateStoreLastUpdatedBy | None = None 

46 

47 

48class AssetStateStoreCollectionResponse(BaseModel): 

49 """All asset state store entries for an asset.""" 

50 

51 asset_state_store: list[AssetStateStoreResponse] 

52 total_entries: int 

53 

54 

55class AssetStateStoreBody(StrictBaseModel): 

56 """Request body for setting an asset state store value.""" 

57 

58 value: JsonValue 

59 

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