Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/api/logs.py: 48%
48 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"""
2Routes for interacting with log objects.
3"""
5from typing import Optional, Sequence
7from fastapi import Body, Depends, WebSocket, status
8from pydantic import TypeAdapter
9from starlette.status import WS_1002_PROTOCOL_ERROR
11import prefect.server.api.dependencies as dependencies
12import prefect.server.models as models
13from prefect.server.database import PrefectDBInterface, provide_database_interface
14from prefect.server.logs import stream
15from prefect.server.schemas.actions import LogCreate
16from prefect.server.schemas.core import Log
17from prefect.server.schemas.filters import LogFilter
18from prefect.server.schemas.sorting import LogSort
19from prefect.server.utilities import subscriptions
20from prefect.server.utilities.server import PrefectRouter
22router: PrefectRouter = PrefectRouter(prefix="/logs", tags=["Logs"])
25@router.post("/", status_code=status.HTTP_201_CREATED)
26async def create_logs(
27 logs: Sequence[LogCreate],
28 db: PrefectDBInterface = Depends(provide_database_interface),
29) -> None:
30 """
31 Create new logs from the provided schema.
33 For more information, see https://docs.prefect.io/v3/how-to-guides/workflows/add-logging.
34 """
35 for batch in models.logs.split_logs_into_batches(logs):
36 async with db.session_context(begin_transaction=True) as session:
37 await models.logs.create_logs(session=session, logs=batch)
40logs_adapter: TypeAdapter[Sequence[Log]] = TypeAdapter(Sequence[Log])
43@router.post("/filter")
44async def read_logs(
45 limit: int = dependencies.LimitBody(),
46 offset: int = Body(0, ge=0),
47 logs: Optional[LogFilter] = None,
48 sort: LogSort = Body(LogSort.TIMESTAMP_ASC),
49 db: PrefectDBInterface = Depends(provide_database_interface),
50) -> Sequence[Log]:
51 """
52 Query for logs.
53 """
54 async with db.session_context() as session:
55 return logs_adapter.validate_python(
56 await models.logs.read_logs(
57 session=session, log_filter=logs, offset=offset, limit=limit, sort=sort
58 )
59 )
62@router.websocket("/out")
63async def stream_logs_out(websocket: WebSocket) -> None:
64 """Serve a WebSocket to stream live logs"""
65 websocket = await subscriptions.accept_prefect_socket(websocket)
66 if not websocket:
67 return
69 try:
70 # After authentication, the next message is expected to be a filter message, any
71 # other type of message will close the connection.
72 message = await websocket.receive_json()
74 if message["type"] != "filter":
75 return await websocket.close(
76 WS_1002_PROTOCOL_ERROR, reason="Expected 'filter' message"
77 )
79 try:
80 filter = LogFilter.model_validate(message["filter"])
81 except Exception as e:
82 return await websocket.close(
83 WS_1002_PROTOCOL_ERROR, reason=f"Invalid filter: {e}"
84 )
86 # No backfill support for logs - only live streaming
87 # Subscribe to the ongoing log stream
88 async with stream.logs(filter) as log_stream:
89 async for log in log_stream:
90 if not log:
91 if await subscriptions.still_connected(websocket):
92 continue
93 break
95 await websocket.send_json(
96 {"type": "log", "log": log.model_dump(mode="json")}
97 )
99 except subscriptions.NORMAL_DISCONNECT_EXCEPTIONS: # pragma: no cover
100 pass # it's fine if a client disconnects either normally or abnormally
102 return None