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

44 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. 

17""" 

18Common Data Models for Airflow REST API. 

19 

20:meta private: 

21""" 

22 

23from __future__ import annotations 

24 

25import enum 

26from typing import Annotated, Any, Generic, Literal, TypeVar, Union 

27 

28from pydantic import Discriminator, Field, Tag 

29 

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

31 

32# Common Bulk Data Models 

33T = TypeVar("T") 

34K = TypeVar("K") 

35 

36 

37class BulkAction(str, enum.Enum): 

38 """Bulk Action to be performed on the used model.""" 

39 

40 CREATE = "create" 

41 DELETE = "delete" 

42 UPDATE = "update" 

43 

44 

45class BulkActionOnExistence(enum.Enum): 

46 """Bulk Action to be taken if the entity already exists or not.""" 

47 

48 FAIL = "fail" 

49 SKIP = "skip" 

50 OVERWRITE = "overwrite" 

51 

52 

53class BulkActionNotOnExistence(enum.Enum): 

54 """Bulk Action to be taken if the entity does not exist.""" 

55 

56 FAIL = "fail" 

57 SKIP = "skip" 

58 

59 

60class BulkBaseAction(StrictBaseModel, Generic[T]): 

61 """Base class for bulk actions.""" 

62 

63 action: BulkAction = Field(..., description="The action to be performed on the entities.") 

64 

65 

66class BulkCreateAction(BulkBaseAction[T]): 

67 """Bulk Create entity serializer for request bodies.""" 

68 

69 action: Literal[BulkAction.CREATE] = Field(description="The action to be performed on the entities.") 

70 entities: list[T] = Field(..., description="A list of entities to be created.") 

71 action_on_existence: BulkActionOnExistence = BulkActionOnExistence.FAIL 

72 

73 

74class BulkUpdateAction(BulkBaseAction[T]): 

75 """Bulk Update entity serializer for request bodies.""" 

76 

77 action: Literal[BulkAction.UPDATE] = Field(description="The action to be performed on the entities.") 

78 entities: list[T] = Field(..., description="A list of entities to be updated.") 

79 update_mask: list[str] | None = Field( 

80 default=None, 

81 description=( 

82 "A list of field names to update for each entity." 

83 "Only these fields will be applied from the request body to the database model." 

84 "Any extra fields provided will be ignored." 

85 ), 

86 ) 

87 action_on_non_existence: BulkActionNotOnExistence = BulkActionNotOnExistence.FAIL 

88 

89 

90class BulkDeleteAction(BulkBaseAction[T]): 

91 """Bulk Delete entity serializer for request bodies.""" 

92 

93 action: Literal[BulkAction.DELETE] = Field(description="The action to be performed on the entities.") 

94 entities: list[Union[str, T]] = Field( 

95 ..., 

96 description="A list of entity id/key or entity objects to be deleted.", 

97 ) 

98 action_on_non_existence: BulkActionNotOnExistence = BulkActionNotOnExistence.FAIL 

99 

100 

101def _action_discriminator(action: Any) -> str: 

102 return BulkAction(action["action"]).value 

103 

104 

105class BulkBody(StrictBaseModel, Generic[T]): 

106 """Serializer for bulk entity operations.""" 

107 

108 actions: list[ 

109 Annotated[ 

110 Union[ 

111 Annotated[BulkCreateAction[T], Tag(BulkAction.CREATE.value)], 

112 Annotated[BulkUpdateAction[T], Tag(BulkAction.UPDATE.value)], 

113 Annotated[BulkDeleteAction[T], Tag(BulkAction.DELETE.value)], 

114 ], 

115 Discriminator(_action_discriminator), 

116 ] 

117 ] 

118 

119 

120class BulkActionResponse(BaseModel): 

121 """ 

122 Serializer for individual bulk action responses. 

123 

124 Represents the outcome of a single bulk operation (create, update, or delete). 

125 The response includes a list of successful keys and any errors encountered during the operation. 

126 This structure helps users understand which key actions succeeded and which failed. 

127 """ 

128 

129 success: list[str] = Field( 

130 default=[], description="A list of unique id/key representing successful operations." 

131 ) 

132 errors: list[dict[str, Any]] = Field( 

133 default=[], 

134 description="A list of errors encountered during the operation, each containing details about the issue.", 

135 ) 

136 

137 

138class BulkResponse(BaseModel): 

139 """ 

140 Serializer for responses to bulk entity operations. 

141 

142 This represents the results of create, update, and delete actions performed on entity in bulk. 

143 Each action (if requested) is represented as a field containing details about successful keys and any encountered errors. 

144 Fields are populated in the response only if the respective action was part of the request, else are set None. 

145 """ 

146 

147 create: BulkActionResponse | None = Field( 

148 default=None, 

149 description="Details of the bulk create operation, including successful keys and errors.", 

150 ) 

151 update: BulkActionResponse | None = Field( 

152 default=None, 

153 description="Details of the bulk update operation, including successful keys and errors.", 

154 ) 

155 delete: BulkActionResponse | None = Field( 

156 default=None, 

157 description="Details of the bulk delete operation, including successful keys and errors.", 

158 )