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

1from uuid import UUID 

2 

3import sqlalchemy as sa 

4from sqlalchemy.ext.asyncio import AsyncSession 

5 

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 

10 

11SERVER_DEFAULT_RESULT_STORAGE_CONFIGURATION_KEY = "server-default-result-storage" 

12 

13 

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.") 

24 

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 ) 

30 

31 

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 ) 

43 

44 

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() 

54 

55 return schemas.core.ServerDefaultResultStorage.model_validate(configured_value) 

56 

57 

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 ) 

63 

64 

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 

72 

73 return await clear_server_default_result_storage(session=session) 

74 

75 

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 

84 

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 

91 

92 

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 

101 

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