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

47 statements  

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

1""" 

2Functions for interacting with log ORM objects. 

3Intended for internal use by the Prefect REST API. 

4""" 

5 

6from typing import TYPE_CHECKING, Generator, Optional, Sequence, Tuple 

7 

8from sqlalchemy import delete, select 

9from sqlalchemy.ext.asyncio import AsyncSession 

10 

11import prefect.server.schemas as schemas 

12from prefect.logging import get_logger 

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

14from prefect.server.logs import messaging 

15from prefect.server.schemas.actions import LogCreate 

16from prefect.utilities.collections import batched_iterable 

17 

18# We have a limit of 32,767 parameters at a time for a single query... 

19MAXIMUM_QUERY_PARAMETERS = 32_767 

20 

21# ...and logs have a certain number of fields... 

22NUMBER_OF_LOG_FIELDS = len(schemas.core.Log.model_fields) 

23 

24# ...so we can only INSERT batches of a certain size at a time 

25LOG_BATCH_SIZE = MAXIMUM_QUERY_PARAMETERS // NUMBER_OF_LOG_FIELDS 

26 

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

28 import logging 

29 

30logger: "logging.Logger" = get_logger(__name__) 

31 

32 

33def split_logs_into_batches( 

34 logs: Sequence[schemas.actions.LogCreate], 

35) -> Generator[Tuple[LogCreate, ...], None, None]: 

36 for batch in batched_iterable(logs, LOG_BATCH_SIZE): 

37 yield batch 

38 

39 

40def _sanitize_log_strings(log: LogCreate) -> LogCreate: 

41 """Strip null bytes from log string fields. 

42 

43 PostgreSQL rejects strings containing null bytes (0x00) with 

44 `CharacterNotInRepertoireError`. Rather than letting the INSERT 

45 fail, we replace them here so the rest of the log record is preserved. 

46 """ 

47 sanitized_message = log.message.replace("\x00", "") 

48 sanitized_name = log.name.replace("\x00", "") 

49 if sanitized_message != log.message or sanitized_name != log.name: 

50 return log.model_copy( 

51 update={"message": sanitized_message, "name": sanitized_name} 

52 ) 

53 return log 

54 

55 

56@db_injector 

57async def create_logs( 

58 db: PrefectDBInterface, session: AsyncSession, logs: Sequence[LogCreate] 

59) -> None: 

60 """ 

61 Creates new logs 

62 

63 Args: 

64 session: a database session 

65 logs: a list of log schemas 

66 

67 Returns: 

68 None 

69 """ 

70 try: 

71 logs = [_sanitize_log_strings(log) for log in logs] 

72 full_logs = [schemas.core.Log(**log.model_dump()) for log in logs] 

73 await session.execute( 

74 db.queries.insert(db.Log).values( 

75 [log.model_dump(exclude={"created", "updated"}) for log in full_logs] 

76 ) 

77 ) 

78 await messaging.publish_logs(full_logs) 

79 

80 except RuntimeError as exc: 

81 if "can't create new thread at interpreter shutdown" in str(exc): 

82 # Background logs sometimes fail to write when the interpreter is shutting down. 

83 # This is a known issue in Python 3.12.2 that can be ignored and is fixed in Python 3.12.3. 

84 # see e.g. https://github.com/python/cpython/issues/113964 

85 logger.debug("Received event during interpreter shutdown, ignoring") 

86 else: 

87 raise 

88 

89 

90@db_injector 

91async def read_logs( 

92 db: PrefectDBInterface, 

93 session: AsyncSession, 

94 log_filter: Optional[schemas.filters.LogFilter], 

95 offset: Optional[int] = None, 

96 limit: Optional[int] = None, 

97 sort: schemas.sorting.LogSort = schemas.sorting.LogSort.TIMESTAMP_ASC, 

98) -> Sequence[orm_models.Log]: 

99 """ 

100 Read logs. 

101 

102 Args: 

103 session: a database session 

104 db: the database interface 

105 log_filter: only select logs that match these filters 

106 offset: Query offset 

107 limit: Query limit 

108 sort: Query sort 

109 

110 Returns: 

111 List[orm_models.Log]: the matching logs 

112 """ 

113 query = select(db.Log).order_by(*sort.as_sql_sort()).offset(offset).limit(limit) 

114 

115 if log_filter: 

116 query = query.where(log_filter.as_sql_filter()) 

117 

118 result = await session.execute(query) 

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

120 

121 

122@db_injector 

123async def delete_logs( 

124 db: PrefectDBInterface, 

125 session: AsyncSession, 

126 log_filter: schemas.filters.LogFilter, 

127) -> int: 

128 query = delete(db.Log).where(log_filter.as_sql_filter()) 

129 result = await session.execute(query) 

130 return result.rowcount