Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/block_types.py: 78%

88 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-07 02:04 +0000

1""" 

2Functions for interacting with block type ORM objects. 

3Intended for internal use by the Prefect REST API. 

4""" 

5 

6import html 

7from typing import TYPE_CHECKING, Optional, Sequence, Union 

8from uuid import UUID 

9 

10import sqlalchemy as sa 

11from sqlalchemy.ext.asyncio import AsyncSession 

12 

13from prefect.server import schemas 

14from prefect.server.database import PrefectDBInterface, db_injector 

15from prefect.server.database.orm_models import BlockType 

16from prefect.server.events import clients 

17from prefect.server.events.schemas import lifecycle 

18from prefect.server.models import storage_defaults 

19from prefect.types._datetime import now 

20 

21if TYPE_CHECKING: 21 ↛ 22line 21 didn't jump to line 22 because the condition on line 21 was never true

22 from prefect.client.schemas import BlockType as ClientBlockType 

23 from prefect.client.schemas.actions import BlockTypeUpdate as ClientBlockTypeUpdate 

24 

25 

26async def emit_block_type_created_event(block_type: BlockType) -> None: 

27 """Emit an event when a block type is created.""" 

28 async with clients.PrefectServerEventsClient() as events_client: 

29 await events_client.emit( 

30 lifecycle.block_type_created_event(block_type, now("UTC")) 

31 ) 

32 

33 

34async def emit_block_type_updated_event(block_type: BlockType) -> None: 

35 """Emit an event when a block type is updated.""" 

36 async with clients.PrefectServerEventsClient() as events_client: 

37 await events_client.emit( 

38 lifecycle.block_type_updated_event(block_type, now("UTC")) 

39 ) 

40 

41 

42async def emit_block_type_deleted_event(block_type: BlockType) -> None: 

43 """Emit an event when a block type is deleted.""" 

44 async with clients.PrefectServerEventsClient() as events_client: 

45 await events_client.emit( 

46 lifecycle.block_type_deleted_event(block_type, now("UTC")) 

47 ) 

48 

49 

50@db_injector 

51async def create_block_type( 

52 db: PrefectDBInterface, 

53 session: AsyncSession, 

54 block_type: Union[schemas.core.BlockType, "ClientBlockType"], 

55 override: bool = False, 

56) -> Union[BlockType, None]: 

57 """ 

58 Create a new block type. 

59 

60 Args: 

61 session: A database session 

62 block_type: a block type object 

63 

64 Returns: 

65 block_type: an ORM block type model 

66 """ 

67 # We take a shortcut in many unit tests and in block registration to pass client 

68 # models directly to this function. We will support this by converting them to 

69 # the appropriate server model. 

70 if not isinstance(block_type, schemas.core.BlockType): 

71 block_type = schemas.core.BlockType.model_validate( 

72 block_type.model_dump(mode="json") 

73 ) 

74 

75 upsert_start = now("UTC") 

76 insert_values = block_type.model_dump_for_orm( 

77 exclude_unset=False, exclude={"created", "updated", "id"} 

78 ) 

79 if insert_values.get("description") is not None: 

80 insert_values["description"] = html.escape( 

81 insert_values["description"], quote=False 

82 ) 

83 if insert_values.get("code_example") is not None: 

84 insert_values["code_example"] = html.escape( 

85 insert_values["code_example"], quote=False 

86 ) 

87 insert_stmt = db.queries.insert(db.BlockType).values(**insert_values) 

88 if override: 

89 insert_stmt = insert_stmt.on_conflict_do_update( 

90 index_elements=db.orm.block_type_unique_upsert_columns, 

91 set_=insert_values, 

92 ) 

93 await session.execute(insert_stmt) 

94 

95 query = ( 

96 sa.select(db.BlockType) 

97 .where( 

98 sa.and_( 

99 db.BlockType.name == insert_values["name"], 

100 ) 

101 ) 

102 .execution_options(populate_existing=True) 

103 ) 

104 

105 result = await session.execute(query) 

106 model = result.scalar() 

107 

108 if model is not None and model.created >= upsert_start: 

109 await emit_block_type_created_event(model) 

110 

111 return model 

112 

113 

114@db_injector 

115async def read_block_type( 

116 db: PrefectDBInterface, 

117 session: AsyncSession, 

118 block_type_id: UUID, 

119) -> Union[BlockType, None]: 

120 """ 

121 Reads a block type by id. 

122 

123 Args: 

124 session: A database session 

125 block_type_id: a block_type id 

126 

127 Returns: 

128 BlockType: an ORM block type model 

129 """ 

130 return await session.get(db.BlockType, block_type_id) 

131 

132 

133@db_injector 

134async def read_block_type_by_slug( 

135 db: PrefectDBInterface, session: AsyncSession, block_type_slug: str 

136) -> Union[BlockType, None]: 

137 """ 

138 Reads a block type by slug. 

139 

140 Args: 

141 session: A database session 

142 block_type_slug: a block type slug 

143 

144 Returns: 

145 BlockType: an ORM block type model 

146 

147 """ 

148 result = await session.execute( 

149 sa.select(db.BlockType).where(db.BlockType.slug == block_type_slug) 

150 ) 

151 return result.scalar() 

152 

153 

154@db_injector 

155async def read_block_types( 

156 db: PrefectDBInterface, 

157 session: AsyncSession, 

158 block_type_filter: Optional[schemas.filters.BlockTypeFilter] = None, 

159 block_schema_filter: Optional[schemas.filters.BlockSchemaFilter] = None, 

160 limit: Optional[int] = None, 

161 offset: Optional[int] = None, 

162) -> Sequence[BlockType]: 

163 """ 

164 Reads block types with an optional limit and offset 

165 

166 Args: 

167 

168 Returns: 

169 List[BlockType]: List of 

170 """ 

171 query = sa.select(db.BlockType).order_by(db.BlockType.name) 

172 

173 if block_type_filter is not None: 

174 query = query.where(block_type_filter.as_sql_filter()) 

175 

176 if block_schema_filter is not None: 

177 exists_clause = sa.select(db.BlockSchema).where( 

178 db.BlockSchema.block_type_id == db.BlockType.id, 

179 block_schema_filter.as_sql_filter(), 

180 ) 

181 query = query.where(exists_clause.exists()) 

182 

183 if offset is not None: 183 ↛ 186line 183 didn't jump to line 186 because the condition on line 183 was always true

184 query = query.offset(offset) 

185 

186 if limit is not None: 186 ↛ 189line 186 didn't jump to line 189 because the condition on line 186 was always true

187 query = query.limit(limit) 

188 

189 result = await session.execute(query) 

190 return result.scalars().unique().all() 

191 

192 

193@db_injector 

194async def update_block_type( 

195 db: PrefectDBInterface, 

196 session: AsyncSession, 

197 block_type_id: Union[str, UUID], 

198 block_type: Union[ 

199 schemas.actions.BlockTypeUpdate, 

200 schemas.core.BlockType, 

201 "ClientBlockTypeUpdate", 

202 "ClientBlockType", 

203 ], 

204) -> bool: 

205 """ 

206 Update a block type by id. 

207 

208 Args: 

209 session: A database session 

210 block_type_id: Data to update block type with 

211 block_type: A block type id 

212 

213 Returns: 

214 bool: True if the block type was updated 

215 """ 

216 

217 # We take a shortcut in many unit tests and in block registration to pass client 

218 # models directly to this function. We will support this by converting them to 

219 # the appropriate server model. 

220 if not isinstance(block_type, schemas.actions.BlockTypeUpdate): 

221 block_type = schemas.actions.BlockTypeUpdate.model_validate( 

222 block_type.model_dump( 

223 mode="json", 

224 exclude={"id", "created", "updated", "name", "slug", "is_protected"}, 

225 ) 

226 ) 

227 

228 existing = await session.get(db.BlockType, block_type_id) 

229 if existing is None: 229 ↛ 230line 229 didn't jump to line 230 because the condition on line 229 was never true

230 return False 

231 

232 update_statement = ( 

233 sa.update(db.BlockType) 

234 .where(db.BlockType.id == block_type_id) 

235 .values(**block_type.model_dump_for_orm(exclude_unset=True, exclude={"id"})) 

236 ) 

237 await session.execute(update_statement) 

238 

239 await session.refresh(existing) 

240 await emit_block_type_updated_event(existing) 

241 return True 

242 

243 

244@db_injector 

245async def delete_block_type( 

246 db: PrefectDBInterface, session: AsyncSession, block_type_id: UUID 

247) -> bool: 

248 """ 

249 Delete a block type by id. 

250 

251 Args: 

252 session: A database session 

253 block_type_id: A block type id 

254 

255 Returns: 

256 bool: True if the block type was updated 

257 """ 

258 

259 existing = await session.get(db.BlockType, block_type_id) 

260 if existing is None: 260 ↛ 261line 260 didn't jump to line 261 because the condition on line 260 was never true

261 return False 

262 

263 default_references_block_type = ( 

264 await storage_defaults.server_default_result_storage_references_block_type( 

265 session=session, 

266 block_type_id=block_type_id, 

267 ) 

268 ) 

269 

270 await emit_block_type_deleted_event(existing) 

271 

272 await session.execute( 

273 sa.delete(db.BlockType).where(db.BlockType.id == block_type_id) 

274 ) 

275 

276 if default_references_block_type: 

277 await storage_defaults.clear_server_default_result_storage(session=session) 

278 

279 return True