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

208 statements  

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

1""" 

2Functions for interacting with block schema ORM objects. 

3Intended for internal use by the Prefect REST API. 

4""" 

5 

6import json 

7from copy import copy 

8from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple, Union 

9from uuid import UUID 

10 

11import sqlalchemy as sa 

12from sqlalchemy import delete, select 

13from sqlalchemy.ext.asyncio import AsyncSession 

14 

15from prefect.server import schemas 

16from prefect.server.database import PrefectDBInterface, db_injector, orm_models 

17from prefect.server.models import storage_defaults 

18from prefect.server.models.block_types import read_block_type_by_slug 

19from prefect.server.schemas.actions import BlockSchemaCreate 

20from prefect.server.schemas.core import BlockSchema, BlockSchemaReference 

21 

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

23 from prefect.client.schemas.actions import ( 

24 BlockSchemaCreate as ClientBlockSchemaCreate, 

25 ) 

26 from prefect.client.schemas.objects import BlockSchema as ClientBlockSchema 

27 

28 

29class MissingBlockTypeException(Exception): 

30 """Raised when the block type corresponding to a block schema cannot be found""" 

31 

32 

33@db_injector 

34async def create_block_schema( 

35 db: PrefectDBInterface, 

36 session: AsyncSession, 

37 block_schema: Union[ 

38 schemas.actions.BlockSchemaCreate, 

39 schemas.core.BlockSchema, 

40 "ClientBlockSchemaCreate", 

41 "ClientBlockSchema", 

42 ], 

43 override: bool = False, 

44 definitions: Optional[dict[str, Any]] = None, 

45) -> Union[BlockSchema, orm_models.BlockSchema]: 

46 """ 

47 Create a new block schema. 

48 

49 Args: 

50 session: A database session 

51 block_schema: a block schema object 

52 definitions: Definitions of fields from block schema fields 

53 attribute. Used when recursively creating nested block schemas 

54 

55 Returns: 

56 block_schema: an ORM block schema model 

57 """ 

58 from prefect.blocks.core import Block, _get_non_block_reference_definitions 

59 

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

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

62 # the appropriate server model. 

63 if not isinstance(block_schema, schemas.actions.BlockSchemaCreate): 

64 block_schema = schemas.actions.BlockSchemaCreate.model_validate( 

65 block_schema.model_dump( 

66 mode="json", 

67 exclude={"id", "created", "updated", "checksum", "block_type"}, 

68 ) 

69 ) 

70 

71 insert_values = block_schema.model_dump_for_orm( 

72 exclude_unset=False, 

73 exclude={"block_type", "id", "created", "updated"}, 

74 ) 

75 

76 definitions = definitions or block_schema.fields.get("definitions") 

77 fields_for_checksum = insert_values["fields"] 

78 if definitions: 

79 # Ensure definitions are available if this is a nested schema 

80 # that is being registered 

81 fields_for_checksum["definitions"] = definitions 

82 checksum = Block._calculate_schema_checksum(fields_for_checksum) 

83 

84 # Check for existing block schema based on calculated checksum 

85 existing_block_schema = await read_block_schema_by_checksum( 

86 session=session, checksum=checksum, version=block_schema.version 

87 ) 

88 # Return existing block schema if it exists. Allows block schema creation to be called multiple 

89 # times for the same schema without errors. 

90 if existing_block_schema: 

91 return existing_block_schema 

92 

93 insert_values["checksum"] = checksum 

94 

95 if definitions: 

96 # Get non block definitions for saving to the DB. 

97 non_block_definitions = _get_non_block_reference_definitions( 

98 insert_values["fields"], definitions 

99 ) 

100 if non_block_definitions: 

101 insert_values["fields"]["definitions"] = ( 

102 _get_non_block_reference_definitions( 

103 insert_values["fields"], definitions 

104 ) 

105 ) 

106 else: 

107 # Prevent storing definitions for blocks. Those are reconstructed on read. 

108 insert_values["fields"].pop("definitions", None) 

109 

110 # Prevent saving block schema references in the block_schema table. They have 

111 # their own table. 

112 block_schema_references: Dict = insert_values["fields"].pop( 

113 "block_schema_references", {} 

114 ) 

115 

116 insert_stmt = db.queries.insert(db.BlockSchema).values(**insert_values) 

117 if override: 

118 insert_stmt = insert_stmt.on_conflict_do_update( 

119 index_elements=db.orm.block_schema_unique_upsert_columns, 

120 set_=insert_values, 

121 ) 

122 await session.execute(insert_stmt) 

123 

124 query = ( 

125 sa.select(db.BlockSchema) 

126 .where( 

127 db.BlockSchema.checksum == insert_values["checksum"], 

128 ) 

129 .order_by(db.BlockSchema.created.desc()) 

130 .limit(1) 

131 .execution_options(populate_existing=True) 

132 ) 

133 

134 if block_schema.version is not None: 

135 query = query.where(db.BlockSchema.version == block_schema.version) 

136 

137 result = await session.execute(query) 

138 created_block_schema = copy(result.scalar_one()) 

139 

140 await _register_nested_block_schemas( 

141 db, 

142 session=session, 

143 parent_block_schema_id=created_block_schema.id, 

144 block_schema_references=block_schema_references, 

145 base_fields=insert_values["fields"], 

146 definitions=definitions, 

147 override=override, 

148 ) 

149 

150 created_block_schema.fields["block_schema_references"] = block_schema_references 

151 if definitions is not None: 

152 created_block_schema.fields["definitions"] = definitions 

153 

154 return created_block_schema 

155 

156 

157async def _register_nested_block_schemas( 

158 db: PrefectDBInterface, 

159 session: AsyncSession, 

160 parent_block_schema_id: UUID, 

161 block_schema_references: dict[str, Union[dict[str, str], List[dict[str, str]]]], 

162 base_fields: Dict, 

163 definitions: Optional[Dict], 

164 override: bool = False, 

165) -> None: 

166 """ 

167 Iterates through each of the block schema references declared on the block schema. 

168 Attempts to register each of the nested block schemas if they have not already been 

169 registered. An error is thrown if the corresponding block type for a block schema 

170 has not been registered. 

171 

172 Args: 

173 session: A database session. 

174 parent_block_schema_id: The ID of the parent block schema. 

175 block_schema_references: A dictionary containing the block schema references for 

176 the child block schemas of the parent block schema. 

177 base_fields: The field name and type declarations for the parent block schema. 

178 definitions: A dictionary of the field name and type declarations of each 

179 child block schema. 

180 override: Flag controlling if a block schema should updated in place. 

181 """ 

182 for reference_name, reference_values in block_schema_references.items(): 

183 # Operate on a list so that we can share the same code paths for union cases 

184 reference_values = ( 

185 reference_values 

186 if isinstance(reference_values, list) 

187 else [reference_values] 

188 ) 

189 for reference_values_entry in reference_values: 189 ↛ 182line 189 didn't jump to line 182 because the loop on line 189 didn't complete

190 # Check to make sure that associated block type exists 

191 reference_block_type = await read_block_type_by_slug( 

192 session=session, 

193 block_type_slug=reference_values_entry["block_type_slug"], 

194 ) 

195 if reference_block_type is None: 

196 raise MissingBlockTypeException( 

197 "Cannot create block schema because block type" 

198 f" {reference_values_entry['block_type_slug']!r} was not found.Did" 

199 " you forget to register the block type?" 

200 ) 

201 

202 reference_block_schema: Union[BlockSchema, orm_models.BlockSchema, None] 

203 

204 # Checks to see if the visited block schema has been previously created 

205 reference_block_schema = await read_block_schema_by_checksum( 

206 session=session, 

207 checksum=reference_values_entry["block_schema_checksum"], 

208 ) 

209 # Attempts to create block schema since it has not already been registered 

210 if reference_block_schema is None: 

211 if definitions is None: 

212 raise ValueError( 

213 "Unable to create nested block schema due to missing" 

214 " definitions in root block schema fields" 

215 ) 

216 sub_block_schema_fields = _get_fields_for_child_schema( 

217 db, definitions, base_fields, reference_name, reference_block_type 

218 ) 

219 

220 if sub_block_schema_fields is None: 

221 raise ValueError( 

222 "Unable to create nested block schema for block type" 

223 f" {reference_block_type.name!r} due to missing definition." 

224 ) 

225 

226 reference_block_schema = await create_block_schema( 

227 session=session, 

228 block_schema=BlockSchemaCreate( 

229 fields=sub_block_schema_fields, 

230 block_type_id=reference_block_type.id, 

231 ), 

232 override=override, 

233 definitions=definitions, 

234 ) 

235 # Create a block schema reference linking the nested block schema to its parent. 

236 await create_block_schema_reference( 

237 session=session, 

238 block_schema_reference=BlockSchemaReference( 

239 parent_block_schema_id=parent_block_schema_id, 

240 reference_block_schema_id=reference_block_schema.id, 

241 name=reference_name, 

242 ), 

243 ) 

244 

245 

246def _get_fields_for_child_schema( 

247 db: PrefectDBInterface, 

248 definitions: Dict, 

249 base_fields: Dict, 

250 reference_name: str, 

251 reference_block_type: orm_models.BlockType, 

252) -> dict[str, Any]: 

253 """ 

254 Returns the field definitions for a child schema. The fields definitions are pulled from the provided `definitions` 

255 dictionary based on the information extracted from `base_fields` using the `reference_name`. `reference_block_type` 

256 is used to disambiguate fields that have a union type. 

257 """ 

258 from prefect.blocks.core import _collect_nested_reference_strings 

259 

260 spec_reference = base_fields["properties"][reference_name] 

261 sub_block_schema_fields = None 

262 reference_strings = _collect_nested_reference_strings(spec_reference) 

263 if len(reference_strings) == 1: 

264 sub_block_schema_fields = definitions.get( 

265 reference_strings[0].replace("#/definitions/", "") 

266 ) 

267 else: 

268 for reference_string in reference_strings: 268 ↛ 283line 268 didn't jump to line 283 because the loop on line 268 didn't complete

269 definition_key = reference_string.replace("#/definitions/", "") 

270 potential_sub_block_schema_fields = definitions[definition_key] 

271 # Determines the definition to use when registering a child 

272 # block schema by verifying that the block type name stored in 

273 # the definition matches the name of the block type that we're 

274 # currently trying to register a block schema for. 

275 if ( 

276 definitions[definition_key]["block_type_slug"] 

277 == reference_block_type.slug 

278 ): 

279 # Once we've found the matching definition, we no longer 

280 # need to iterate 

281 sub_block_schema_fields = potential_sub_block_schema_fields 

282 break 

283 return sub_block_schema_fields # type: ignore 

284 

285 

286@db_injector 

287async def delete_block_schema( 

288 db: PrefectDBInterface, session: AsyncSession, block_schema_id: UUID 

289) -> bool: 

290 """ 

291 Delete a block schema by id. 

292 

293 Args: 

294 session: A database session 

295 block_schema_id: a block schema id 

296 

297 Returns: 

298 bool: whether or not the block schema was deleted 

299 """ 

300 

301 default_references_block_schema = ( 

302 await storage_defaults.server_default_result_storage_references_block_schema( 

303 session=session, 

304 block_schema_id=block_schema_id, 

305 ) 

306 ) 

307 

308 result = await session.execute( 

309 delete(db.BlockSchema).where(db.BlockSchema.id == block_schema_id) 

310 ) 

311 if result.rowcount <= 0: 311 ↛ 314line 311 didn't jump to line 314 because the condition on line 311 was always true

312 return False 

313 

314 if default_references_block_schema: 314 ↛ 317line 314 didn't jump to line 317 because the condition on line 314 was always true

315 await storage_defaults.clear_server_default_result_storage(session=session) 

316 

317 return True 

318 

319 

320@db_injector 

321async def read_block_schema( 

322 db: PrefectDBInterface, 

323 session: AsyncSession, 

324 block_schema_id: UUID, 

325) -> Union[BlockSchema, None]: 

326 """ 

327 Reads a block schema by id. Will reconstruct the block schema's fields attribute 

328 to include block schema references. 

329 

330 Args: 

331 session: A database session 

332 block_schema_id: a block_schema id 

333 

334 Returns: 

335 orm_models..BlockSchema: the block_schema 

336 """ 

337 

338 # Construction of a recursive query which returns the specified block schema 

339 # along with and nested block schemas coupled with the ID of their parent schema 

340 # the key that they reside under. 

341 block_schema_references_query = ( 

342 sa.select(db.BlockSchemaReference) 

343 .select_from(db.BlockSchemaReference) 

344 .filter_by(parent_block_schema_id=block_schema_id) 

345 .cte("block_schema_references", recursive=True) 

346 ) 

347 block_schema_references_join = ( 

348 sa.select(db.BlockSchemaReference) 

349 .select_from(db.BlockSchemaReference) 

350 .join( 

351 block_schema_references_query, 

352 db.BlockSchemaReference.parent_block_schema_id 

353 == block_schema_references_query.c.reference_block_schema_id, 

354 ) 

355 ) 

356 recursive_block_schema_references_cte = block_schema_references_query.union_all( 

357 block_schema_references_join 

358 ) 

359 nested_block_schemas_query = ( 

360 sa.select( 

361 db.BlockSchema, 

362 recursive_block_schema_references_cte.c.name, 

363 recursive_block_schema_references_cte.c.parent_block_schema_id, 

364 ) 

365 .select_from(db.BlockSchema) 

366 .join( 

367 recursive_block_schema_references_cte, 

368 db.BlockSchema.id 

369 == recursive_block_schema_references_cte.c.reference_block_schema_id, 

370 isouter=True, 

371 ) 

372 .filter( 

373 sa.or_( 

374 db.BlockSchema.id == block_schema_id, 

375 recursive_block_schema_references_cte.c.parent_block_schema_id.is_not( 

376 None 

377 ), 

378 ) 

379 ) 

380 ) 

381 result = await session.execute(nested_block_schemas_query) 

382 

383 return _construct_full_block_schema(result.all()) # type: ignore[arg-type] 

384 

385 

386def _construct_full_block_schema( 

387 block_schemas_with_references: List[ 

388 Tuple[BlockSchema, Optional[str], Optional[UUID]] 

389 ], 

390 root_block_schema: Optional[BlockSchema] = None, 

391 checksum_index: Optional[Dict[str, BlockSchema]] = None, 

392) -> Optional[BlockSchema]: 

393 """ 

394 Reconstruct root_block_schema.fields with nested block schema references. 

395 

396 Args: 

397 block_schemas_with_references: Tuples of (BlockSchema, parent_name, parent_id). 

398 root_block_schema: Traversal root; inferred when omitted. 

399 checksum_index: Pre-built checksum -> BlockSchema mapping. When None, a 

400 local index is built from block_schemas_with_references before 

401 traversal begins (first-wins semantics, None checksums excluded). 

402 Pass an externally built index when reconstructing many schemas in a 

403 loop to share a single dict across all calls. 

404 

405 Returns: 

406 BlockSchema: A block schema with a fully reconstructed fields attribute 

407 """ 

408 if len(block_schemas_with_references) == 0: 

409 return None 

410 root_block_schema = ( 

411 copy(root_block_schema) 

412 if root_block_schema is not None 

413 else _find_root_block_schema(block_schemas_with_references) 

414 ) 

415 if root_block_schema is None: 415 ↛ 416line 415 didn't jump to line 416 because the condition on line 415 was never true

416 raise ValueError( 

417 "Unable to determine root block schema during schema reconstruction." 

418 ) 

419 root_block_schema.fields = _construct_block_schema_fields_with_block_references( 

420 root_block_schema, block_schemas_with_references 

421 ) 

422 if checksum_index is None: 

423 # Build a local index for single-schema callers that don't pre-build one. 

424 checksum_index = {} 

425 for bs, _, _ in block_schemas_with_references: 

426 if bs.checksum is not None and bs.checksum not in checksum_index: 426 ↛ 425line 426 didn't jump to line 425 because the condition on line 426 was always true

427 checksum_index[bs.checksum] = bs 

428 definitions = _construct_block_schema_spec_definitions( 

429 root_block_schema, 

430 block_schemas_with_references, 

431 checksum_index=checksum_index, 

432 ) 

433 # Definitions for non block object may already exist in the block schema OpenAPI 

434 # spec, so we need to combine block and non-block definitions. 

435 if definitions or root_block_schema.fields.get("definitions"): 

436 root_block_schema.fields["definitions"] = { 

437 **root_block_schema.fields.get("definitions", {}), 

438 **definitions, 

439 } 

440 return root_block_schema 

441 

442 

443def _find_root_block_schema( 

444 block_schemas_with_references: List[ 

445 Tuple[BlockSchema, Optional[str], Optional[UUID]] 

446 ], 

447) -> Union[BlockSchema, None]: 

448 """ 

449 Attempts to find the root block schema from a list of block schemas 

450 with references. Returns None if a root block schema is not found. 

451 Returns only the first potential root block schema if multiple are found. 

452 """ 

453 return next( 

454 ( 

455 copy(block_schema) 

456 for ( 

457 block_schema, 

458 _, 

459 parent_block_schema_id, 

460 ) in block_schemas_with_references 

461 if parent_block_schema_id is None 

462 ), 

463 None, 

464 ) 

465 

466 

467def _construct_block_schema_spec_definitions( 

468 root_block_schema: BlockSchema, 

469 block_schemas_with_references: List[ 

470 Tuple[BlockSchema, Optional[str], Optional[UUID]] 

471 ], 

472 checksum_index: Optional[Dict[str, BlockSchema]] = None, 

473) -> dict[str, Any]: 

474 """ 

475 Args: 

476 root_block_schema: The schema whose block_schema_references drive the traversal. 

477 block_schemas_with_references: The full flat result set returned by 

478 the recursive CTE query — shared across all recursive calls so 

479 that child lookups never require a second database round-trip. 

480 checksum_index: Pre-built checksum -> BlockSchema mapping. Forwarded 

481 unchanged to every recursive call to avoid rebuilding it at each 

482 level of nesting. 

483 

484 Returns: 

485 dict[str, Any]: A flat definitions dict mapping block schema titles to 

486 their reconstructed fields, suitable for merging into the root 

487 schema's `fields["definitions"]`. 

488 """ 

489 definitions: dict[str, Any] = {} 

490 for _, block_schema_references in root_block_schema.fields[ 

491 "block_schema_references" 

492 ].items(): 

493 block_schema_references = ( 

494 block_schema_references 

495 if isinstance(block_schema_references, list) 

496 else [block_schema_references] 

497 ) 

498 for block_schema_reference in block_schema_references: 

499 child_block_schema = _find_block_schema_via_checksum( 

500 block_schemas_with_references, 

501 block_schema_reference["block_schema_checksum"], 

502 checksum_index=checksum_index, 

503 ) 

504 

505 if child_block_schema is not None: 505 ↛ 498line 505 didn't jump to line 498 because the condition on line 505 was always true

506 child_block_schema = _construct_full_block_schema( 

507 block_schemas_with_references=block_schemas_with_references, 

508 root_block_schema=child_block_schema, 

509 checksum_index=checksum_index, 

510 ) 

511 assert child_block_schema 

512 definitions = _add_block_schemas_fields_to_definitions( 

513 definitions, child_block_schema 

514 ) 

515 return definitions 

516 

517 

518def _find_block_schema_via_checksum( 

519 block_schemas_with_references: List[ 

520 Tuple[BlockSchema, Optional[str], Optional[UUID]] 

521 ], 

522 checksum: str, 

523 checksum_index: Optional[Dict[str, BlockSchema]] = None, 

524) -> Optional[BlockSchema]: 

525 """ 

526 Return the block schema whose checksum matches, or None if not present. 

527 

528 When checksum_index is provided the lookup is O(1) and a miss returns None 

529 definitively — the linear scan is NOT consulted. Pass a complete index built 

530 from the same row set as block_schemas_with_references. When checksum_index 

531 is None, a linear scan of block_schemas_with_references is used instead. 

532 """ 

533 if checksum_index is not None: 533 ↛ 535line 533 didn't jump to line 535 because the condition on line 533 was always true

534 return checksum_index.get(checksum) 

535 return next( 

536 ( 

537 block_schema 

538 for block_schema, _, _ in block_schemas_with_references 

539 if block_schema.checksum == checksum 

540 ), 

541 None, 

542 ) 

543 

544 

545def _add_block_schemas_fields_to_definitions( 

546 definitions: Dict, child_block_schema: BlockSchema 

547) -> dict[str, Any]: 

548 """ 

549 Returns a new definitions dict with the fields of a block schema and it's child 

550 block schemas added to the existing definitions. 

551 """ 

552 block_schema_title = child_block_schema.fields.get("title") 

553 if block_schema_title is not None: 553 ↛ 563line 553 didn't jump to line 563 because the condition on line 553 was always true

554 # Definitions are declared as a flat dict, so we pop off definitions 

555 # from child schemas and add them to the parent definitions dict 

556 child_definitions = child_block_schema.fields.pop("definitions", {}) 

557 return { 

558 **definitions, 

559 **{block_schema_title: child_block_schema.fields}, 

560 **child_definitions, 

561 } 

562 else: 

563 return definitions 

564 

565 

566def _construct_block_schema_fields_with_block_references( 

567 parent_block_schema: BlockSchema, 

568 block_schemas_with_references: List[ 

569 Tuple[BlockSchema, Optional[str], Optional[UUID]] 

570 ], 

571) -> dict[str, Any]: 

572 """ 

573 Constructs the block_schema_references in a block schema's fields attributes. Returns 

574 a copy of the block schema with block_schema_references added. 

575 

576 Args: 

577 parent_block_schema: The block schema that needs block references populated. 

578 block_schema_with_references: A list of tuples with the structure: 

579 - A block schema object 

580 - The name the block schema lives under in the parent block schema 

581 - The ID of the block schema's parent block schema 

582 

583 Returns: 

584 Dict: Block schema fields with block schema references added. 

585 

586 """ 

587 block_schema_fields_copy = { 

588 **parent_block_schema.fields, 

589 "block_schema_references": {}, 

590 } 

591 for ( 

592 nested_block_schema, 

593 name, 

594 parent_block_schema_id, 

595 ) in block_schemas_with_references: 

596 if parent_block_schema_id == parent_block_schema.id: 

597 assert nested_block_schema.block_type, ( 

598 f"{nested_block_schema} has no block type" 

599 ) 

600 

601 new_block_schema_reference = { 

602 "block_schema_checksum": nested_block_schema.checksum, 

603 "block_type_slug": nested_block_schema.block_type.slug, 

604 } 

605 # A block reference for this key does not yet exist 

606 if name not in block_schema_fields_copy["block_schema_references"]: 

607 block_schema_fields_copy["block_schema_references"][name] = ( 

608 new_block_schema_reference 

609 ) 

610 else: 

611 # List of block references for this key already exist and the block 

612 # reference that we are attempting add isn't present 

613 if ( 

614 isinstance( 

615 block_schema_fields_copy["block_schema_references"][name], 

616 list, 

617 ) 

618 and new_block_schema_reference 

619 not in block_schema_fields_copy["block_schema_references"][name] 

620 ): 

621 block_schema_fields_copy["block_schema_references"][name].append( 

622 new_block_schema_reference 

623 ) 

624 # A single block reference for this key already exists and it does not 

625 # match the block reference that we are attempting to add 

626 elif ( 626 ↛ 591line 626 didn't jump to line 591 because the condition on line 626 was always true

627 block_schema_fields_copy["block_schema_references"][name] 

628 != new_block_schema_reference 

629 ): 

630 block_schema_fields_copy["block_schema_references"][name] = [ 

631 block_schema_fields_copy["block_schema_references"][name], 

632 new_block_schema_reference, 

633 ] 

634 return block_schema_fields_copy 

635 

636 

637@db_injector 

638async def read_block_schemas( 

639 db: PrefectDBInterface, 

640 session: AsyncSession, 

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

642 limit: Optional[int] = None, 

643 offset: Optional[int] = None, 

644) -> List[BlockSchema]: 

645 """ 

646 Reads block schemas, optionally filtered by type or name. 

647 

648 Args: 

649 session: A database session 

650 block_schema_filter: a block schema filter object 

651 limit (int): query limit 

652 offset (int): query offset 

653 

654 Returns: 

655 List[orm_models.BlockSchema]: the block_schemas 

656 """ 

657 # schemas are ordered by `created DESC` to get the most recently created 

658 # ones first (and to facilitate getting the newest one with `limit=1`). 

659 filtered_block_schemas_query = select(db.BlockSchema.id).order_by( 

660 db.BlockSchema.created.desc() 

661 ) 

662 

663 if block_schema_filter: 

664 filtered_block_schemas_query = filtered_block_schemas_query.where( 

665 block_schema_filter.as_sql_filter() 

666 ) 

667 

668 if offset is not None: 

669 filtered_block_schemas_query = filtered_block_schemas_query.offset(offset) 

670 if limit is not None: 

671 filtered_block_schemas_query = filtered_block_schemas_query.limit(limit) 

672 

673 filtered_block_schema_ids = ( 

674 (await session.execute(filtered_block_schemas_query)).scalars().unique().all() 

675 ) 

676 

677 block_schema_references_query = ( 

678 sa.select(db.BlockSchemaReference) 

679 .select_from(db.BlockSchemaReference) 

680 .filter( 

681 db.BlockSchemaReference.parent_block_schema_id.in_( 

682 filtered_block_schemas_query 

683 ) 

684 ) 

685 .cte("block_schema_references", recursive=True) 

686 ) 

687 block_schema_references_join = ( 

688 sa.select(db.BlockSchemaReference) 

689 .select_from(db.BlockSchemaReference) 

690 .join( 

691 block_schema_references_query, 

692 db.BlockSchemaReference.parent_block_schema_id 

693 == block_schema_references_query.c.reference_block_schema_id, 

694 ) 

695 ) 

696 recursive_block_schema_references_cte = block_schema_references_query.union_all( 

697 block_schema_references_join 

698 ) 

699 

700 nested_block_schemas_query = ( 

701 sa.select( 

702 db.BlockSchema, 

703 recursive_block_schema_references_cte.c.name, 

704 recursive_block_schema_references_cte.c.parent_block_schema_id, 

705 ) 

706 .select_from(db.BlockSchema) 

707 # in order to reconstruct nested block schemas efficiently, we need to visit them 

708 # in the order they were created (so that we guarantee that nested/referenced schemas) 

709 # have already been seen. Therefore this second query sorts by created ASC 

710 .order_by(db.BlockSchema.created.asc()) 

711 .join( 

712 recursive_block_schema_references_cte, 

713 db.BlockSchema.id 

714 == recursive_block_schema_references_cte.c.reference_block_schema_id, 

715 isouter=True, 

716 ) 

717 .filter( 

718 sa.or_( 

719 db.BlockSchema.id.in_(filtered_block_schemas_query), 

720 recursive_block_schema_references_cte.c.parent_block_schema_id.is_not( 

721 None 

722 ), 

723 ) 

724 ) 

725 ) 

726 

727 block_schemas_with_references = ( 

728 (await session.execute(nested_block_schemas_query)).unique().all() 

729 ) 

730 checksum_index: Dict[str, BlockSchema] = {} 

731 for bs, _, _ in block_schemas_with_references: 

732 if bs.checksum is not None and bs.checksum not in checksum_index: 

733 checksum_index[bs.checksum] = bs 

734 fully_constructed_block_schemas = [] 

735 visited_block_schema_ids = [] 

736 for root_block_schema, _, _ in block_schemas_with_references: 

737 if ( 

738 root_block_schema.id in filtered_block_schema_ids 

739 and root_block_schema.id not in visited_block_schema_ids 

740 ): 

741 constructed = _construct_full_block_schema( 

742 block_schemas_with_references=block_schemas_with_references, # type: ignore[arg-type] 

743 root_block_schema=root_block_schema, 

744 checksum_index=checksum_index, 

745 ) 

746 assert constructed 

747 fully_constructed_block_schemas.append(constructed) 

748 visited_block_schema_ids.append(root_block_schema.id) 

749 

750 # because we reconstructed schemas ordered by created ASC, we 

751 # reverse the final output to restore created DESC 

752 return list(reversed(fully_constructed_block_schemas)) 

753 

754 

755@db_injector 

756async def read_block_schema_by_checksum( 

757 db: PrefectDBInterface, 

758 session: AsyncSession, 

759 checksum: str, 

760 version: Optional[str] = None, 

761) -> Optional[BlockSchema]: 

762 """ 

763 Reads a block_schema by checksum. Will reconstruct the block schema's fields 

764 attribute to include block schema references. 

765 

766 Args: 

767 session: A database session 

768 checksum: a block_schema checksum 

769 version: A block_schema version 

770 

771 Returns: 

772 orm_models.BlockSchema: the block_schema 

773 """ 

774 # Construction of a recursive query which returns the specified block schema 

775 # along with and nested block schemas coupled with the ID of their parent schema 

776 # the key that they reside under. 

777 

778 # The same checksum with different versions can occur in the DB. Return only the 

779 # most recently created one. 

780 root_block_schema_query = ( 

781 sa.select(db.BlockSchema) 

782 .filter_by(checksum=checksum) 

783 .order_by(db.BlockSchema.created.desc()) 

784 .limit(1) 

785 ) 

786 

787 if version is not None: 

788 root_block_schema_query = root_block_schema_query.filter_by(version=version) 

789 

790 root_block_schema_cte = root_block_schema_query.cte("root_block_schema") 

791 

792 block_schema_references_query = ( 

793 sa.select(db.BlockSchemaReference) 

794 .select_from(db.BlockSchemaReference) 

795 .filter_by(parent_block_schema_id=root_block_schema_cte.c.id) 

796 .cte("block_schema_references", recursive=True) 

797 ) 

798 block_schema_references_join = ( 

799 sa.select(db.BlockSchemaReference) 

800 .select_from(db.BlockSchemaReference) 

801 .join( 

802 block_schema_references_query, 

803 db.BlockSchemaReference.parent_block_schema_id 

804 == block_schema_references_query.c.reference_block_schema_id, 

805 ) 

806 ) 

807 recursive_block_schema_references_cte = block_schema_references_query.union_all( 

808 block_schema_references_join 

809 ) 

810 nested_block_schemas_query = ( 

811 sa.select( 

812 db.BlockSchema, 

813 recursive_block_schema_references_cte.c.name, 

814 recursive_block_schema_references_cte.c.parent_block_schema_id, 

815 ) 

816 .select_from(db.BlockSchema) 

817 .join( 

818 recursive_block_schema_references_cte, 

819 db.BlockSchema.id 

820 == recursive_block_schema_references_cte.c.reference_block_schema_id, 

821 isouter=True, 

822 ) 

823 .filter( 

824 sa.or_( 

825 db.BlockSchema.id == root_block_schema_cte.c.id, 

826 recursive_block_schema_references_cte.c.parent_block_schema_id.is_not( 

827 None 

828 ), 

829 ) 

830 ) 

831 ) 

832 result = await session.execute(nested_block_schemas_query) 

833 return _construct_full_block_schema(result.all()) # type: ignore[arg-type] 

834 

835 

836@db_injector 

837async def read_available_block_capabilities( 

838 db: PrefectDBInterface, 

839 session: AsyncSession, 

840) -> List[str]: 

841 """ 

842 Retrieves a list of all available block capabilities. 

843 

844 Args: 

845 session: A database session. 

846 

847 Returns: 

848 List[str]: List of all available block capabilities. 

849 """ 

850 query = sa.select( 

851 db.queries.json_arr_agg( 

852 db.queries.cast_to_json(db.BlockSchema.capabilities.distinct()) 

853 ) 

854 ) 

855 capability_combinations = (await session.execute(query)).scalars().first() or list() 

856 if db.queries.uses_json_strings and isinstance(capability_combinations, str): 

857 capability_combinations = json.loads(capability_combinations) 

858 return list({c for capabilities in capability_combinations for c in capabilities}) 

859 

860 

861@db_injector 

862async def create_block_schema_reference( 

863 db: PrefectDBInterface, 

864 session: AsyncSession, 

865 block_schema_reference: schemas.core.BlockSchemaReference, 

866) -> Union[orm_models.BlockSchemaReference, None]: 

867 """ 

868 Retrieves a list of all available block capabilities. 

869 

870 Args: 

871 session: A database session. 

872 block_schema_reference: A block schema reference object. 

873 

874 Returns: 

875 orm_models.BlockSchemaReference: The created BlockSchemaReference 

876 """ 

877 query_stmt = sa.select(db.BlockSchemaReference).where( 

878 db.BlockSchemaReference.name == block_schema_reference.name, 

879 db.BlockSchemaReference.parent_block_schema_id 

880 == block_schema_reference.parent_block_schema_id, 

881 db.BlockSchemaReference.reference_block_schema_id 

882 == block_schema_reference.reference_block_schema_id, 

883 ) 

884 

885 existing_reference = (await session.execute(query_stmt)).scalar() 

886 if existing_reference: 

887 return existing_reference 

888 

889 insert_stmt = db.queries.insert(db.BlockSchemaReference).values( 

890 **block_schema_reference.model_dump_for_orm( 

891 exclude_unset=True, exclude={"created", "updated"} 

892 ) 

893 ) 

894 await session.execute(insert_stmt) 

895 

896 result = await session.execute( 

897 sa.select(db.BlockSchemaReference).where( 

898 db.BlockSchemaReference.id == block_schema_reference.id 

899 ) 

900 ) 

901 return result.scalar()