Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/models/storage_defaults.py: 39%
44 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
1from uuid import UUID
3import sqlalchemy as sa
4from sqlalchemy.ext.asyncio import AsyncSession
6from prefect.server import schemas
7from prefect.server.database import orm_models
8from prefect.server.exceptions import ObjectNotFoundError
9from prefect.server.models import block_documents, configuration
11SERVER_DEFAULT_RESULT_STORAGE_CONFIGURATION_KEY = "server-default-result-storage"
14async def validate_server_default_result_storage_block(
15 session: AsyncSession,
16 block_document_id: UUID,
17) -> None:
18 block_document = await block_documents.read_block_document_by_id(
19 session=session,
20 block_document_id=block_document_id,
21 )
22 if block_document is None:
23 raise ObjectNotFoundError(f"Block document {block_document_id!s} not found.")
25 block_schema = block_document.block_schema
26 if block_schema is None or "write-path" not in block_schema.capabilities:
27 raise ValueError(
28 f"Block document {block_document_id!s} cannot be used for result storage."
29 )
32async def write_server_default_result_storage(
33 session: AsyncSession,
34 storage_default: schemas.core.ServerDefaultResultStorage,
35) -> orm_models.Configuration:
36 return await configuration.write_configuration(
37 session=session,
38 configuration=schemas.core.Configuration(
39 key=SERVER_DEFAULT_RESULT_STORAGE_CONFIGURATION_KEY,
40 value=storage_default.model_dump(mode="json"),
41 ),
42 )
45async def read_server_default_result_storage(
46 session: AsyncSession,
47) -> schemas.core.ServerDefaultResultStorage:
48 query = sa.select(orm_models.Configuration.value).where(
49 orm_models.Configuration.key == SERVER_DEFAULT_RESULT_STORAGE_CONFIGURATION_KEY
50 )
51 configured_value = await session.scalar(query)
52 if configured_value is None:
53 return schemas.core.ServerDefaultResultStorage()
55 return schemas.core.ServerDefaultResultStorage.model_validate(configured_value)
58async def clear_server_default_result_storage(session: AsyncSession) -> bool:
59 return await configuration.delete_configuration(
60 session=session,
61 key=SERVER_DEFAULT_RESULT_STORAGE_CONFIGURATION_KEY,
62 )
65async def clear_server_default_result_storage_for_block(
66 session: AsyncSession,
67 block_document_id: UUID,
68) -> bool:
69 storage_default = await read_server_default_result_storage(session=session)
70 if storage_default.default_result_storage_block_id != block_document_id:
71 return False
73 return await clear_server_default_result_storage(session=session)
76async def server_default_result_storage_references_block_schema(
77 session: AsyncSession,
78 block_schema_id: UUID,
79) -> bool:
80 storage_default = await read_server_default_result_storage(session=session)
81 block_document_id = storage_default.default_result_storage_block_id
82 if block_document_id is None:
83 return False
85 query = (
86 sa.select(orm_models.BlockDocument.id)
87 .where(orm_models.BlockDocument.id == block_document_id)
88 .where(orm_models.BlockDocument.block_schema_id == block_schema_id)
89 )
90 return await session.scalar(query) is not None
93async def server_default_result_storage_references_block_type(
94 session: AsyncSession,
95 block_type_id: UUID,
96) -> bool:
97 storage_default = await read_server_default_result_storage(session=session)
98 block_document_id = storage_default.default_result_storage_block_id
99 if block_document_id is None:
100 return False
102 query = (
103 sa.select(orm_models.BlockDocument.id)
104 .where(orm_models.BlockDocument.id == block_document_id)
105 .where(orm_models.BlockDocument.block_type_id == block_type_id)
106 )
107 return await session.scalar(query) is not None