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
« 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"""
6import json
7from copy import copy
8from typing import TYPE_CHECKING, Any, Dict, List, Optional, Tuple, Union
9from uuid import UUID
11import sqlalchemy as sa
12from sqlalchemy import delete, select
13from sqlalchemy.ext.asyncio import AsyncSession
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
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
29class MissingBlockTypeException(Exception):
30 """Raised when the block type corresponding to a block schema cannot be found"""
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.
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
55 Returns:
56 block_schema: an ORM block schema model
57 """
58 from prefect.blocks.core import Block, _get_non_block_reference_definitions
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 )
71 insert_values = block_schema.model_dump_for_orm(
72 exclude_unset=False,
73 exclude={"block_type", "id", "created", "updated"},
74 )
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)
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
93 insert_values["checksum"] = checksum
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)
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 )
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)
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 )
134 if block_schema.version is not None:
135 query = query.where(db.BlockSchema.version == block_schema.version)
137 result = await session.execute(query)
138 created_block_schema = copy(result.scalar_one())
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 )
150 created_block_schema.fields["block_schema_references"] = block_schema_references
151 if definitions is not None:
152 created_block_schema.fields["definitions"] = definitions
154 return created_block_schema
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.
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 )
202 reference_block_schema: Union[BlockSchema, orm_models.BlockSchema, None]
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 )
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 )
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 )
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
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
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.
293 Args:
294 session: A database session
295 block_schema_id: a block schema id
297 Returns:
298 bool: whether or not the block schema was deleted
299 """
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 )
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
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)
317 return True
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.
330 Args:
331 session: A database session
332 block_schema_id: a block_schema id
334 Returns:
335 orm_models..BlockSchema: the block_schema
336 """
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)
383 return _construct_full_block_schema(result.all()) # type: ignore[arg-type]
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.
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.
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
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 )
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.
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 )
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
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.
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 )
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
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.
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
583 Returns:
584 Dict: Block schema fields with block schema references added.
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 )
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
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.
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
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 )
663 if block_schema_filter:
664 filtered_block_schemas_query = filtered_block_schemas_query.where(
665 block_schema_filter.as_sql_filter()
666 )
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)
673 filtered_block_schema_ids = (
674 (await session.execute(filtered_block_schemas_query)).scalars().unique().all()
675 )
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 )
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 )
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)
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))
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.
766 Args:
767 session: A database session
768 checksum: a block_schema checksum
769 version: A block_schema version
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.
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 )
787 if version is not None:
788 root_block_schema_query = root_block_schema_query.filter_by(version=version)
790 root_block_schema_cte = root_block_schema_query.cte("root_block_schema")
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]
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.
844 Args:
845 session: A database session.
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})
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.
870 Args:
871 session: A database session.
872 block_schema_reference: A block schema reference object.
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 )
885 existing_reference = (await session.execute(query_stmt)).scalar()
886 if existing_reference:
887 return existing_reference
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)
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()