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
« 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"""
6from typing import TYPE_CHECKING, Generator, Optional, Sequence, Tuple
8from sqlalchemy import delete, select
9from sqlalchemy.ext.asyncio import AsyncSession
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
18# We have a limit of 32,767 parameters at a time for a single query...
19MAXIMUM_QUERY_PARAMETERS = 32_767
21# ...and logs have a certain number of fields...
22NUMBER_OF_LOG_FIELDS = len(schemas.core.Log.model_fields)
24# ...so we can only INSERT batches of a certain size at a time
25LOG_BATCH_SIZE = MAXIMUM_QUERY_PARAMETERS // NUMBER_OF_LOG_FIELDS
27if TYPE_CHECKING: 27 ↛ 28line 27 didn't jump to line 28 because the condition on line 27 was never true
28 import logging
30logger: "logging.Logger" = get_logger(__name__)
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
40def _sanitize_log_strings(log: LogCreate) -> LogCreate:
41 """Strip null bytes from log string fields.
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
56@db_injector
57async def create_logs(
58 db: PrefectDBInterface, session: AsyncSession, logs: Sequence[LogCreate]
59) -> None:
60 """
61 Creates new logs
63 Args:
64 session: a database session
65 logs: a list of log schemas
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)
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
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.
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
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)
115 if log_filter:
116 query = query.where(log_filter.as_sql_filter())
118 result = await session.execute(query)
119 return result.scalars().unique().all()
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