mirror of
https://github.com/PrefectHQ/fastmcp.git
synced 2026-08-23 14:04:18 +02:00
Add debug logging for task lifecycle
This commit is contained in:
parent
3df88260eb
commit
8e927fbfed
2 changed files with 79 additions and 1 deletions
|
|
@ -385,6 +385,8 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
@asynccontextmanager
|
||||
async def _docket_lifespan(self) -> AsyncIterator[None]:
|
||||
"""Manage Docket instance and Worker for background task execution."""
|
||||
import socket
|
||||
|
||||
from fastmcp import settings
|
||||
|
||||
# Set FastMCP server in ContextVar so CurrentFastMCP can access it (use weakref to avoid reference cycles)
|
||||
|
|
@ -394,16 +396,26 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
_current_worker,
|
||||
)
|
||||
|
||||
# Get instance identifier for debugging Lambda lifecycle issues
|
||||
instance_id = f"{socket.gethostname()}#{id(self)}"
|
||||
logger.info(f"[{instance_id}] _docket_lifespan ENTERING")
|
||||
|
||||
server_token = _current_server.set(weakref.ref(self))
|
||||
logger.info(f"[{instance_id}] _current_server SET (token={server_token})")
|
||||
|
||||
try:
|
||||
# For directly mounted servers, the parent's Docket/Worker handles all
|
||||
# task execution. Skip creating our own to avoid race conditions with
|
||||
# multiple workers competing for tasks from the same queue.
|
||||
if self._is_mounted:
|
||||
logger.info(f"[{instance_id}] Server is mounted, skipping Docket setup")
|
||||
yield
|
||||
return
|
||||
|
||||
logger.info(
|
||||
f"[{instance_id}] Creating Docket with name={settings.docket.name}, "
|
||||
f"url={settings.docket.url}"
|
||||
)
|
||||
# Create Docket instance using configured name and URL
|
||||
async with Docket(
|
||||
name=settings.docket.name,
|
||||
|
|
@ -456,6 +468,9 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
|
||||
# Set Docket in ContextVar so CurrentDocket can access it
|
||||
docket_token = _current_docket.set(docket)
|
||||
logger.info(
|
||||
f"[{instance_id}] _current_docket SET (token={docket_token})"
|
||||
)
|
||||
try:
|
||||
# Build worker kwargs from settings
|
||||
worker_kwargs: dict[str, Any] = {
|
||||
|
|
@ -467,27 +482,52 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
worker_kwargs["name"] = settings.docket.worker_name
|
||||
|
||||
# Create and start Worker
|
||||
logger.info(
|
||||
f"[{instance_id}] Creating Worker with kwargs={worker_kwargs}"
|
||||
)
|
||||
async with Worker(docket, **worker_kwargs) as worker: # type: ignore[arg-type]
|
||||
# Set Worker in ContextVar so CurrentWorker can access it
|
||||
worker_token = _current_worker.set(worker)
|
||||
logger.info(
|
||||
f"[{instance_id}] _current_worker SET (token={worker_token}), "
|
||||
f"worker.name={worker.name}"
|
||||
)
|
||||
try:
|
||||
worker_task = asyncio.create_task(worker.run_forever())
|
||||
logger.info(
|
||||
f"[{instance_id}] Worker task started, "
|
||||
f"_docket_lifespan YIELDING (lifespan active)"
|
||||
)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
logger.info(
|
||||
f"[{instance_id}] _docket_lifespan EXITING "
|
||||
"(after yield, cancelling worker)"
|
||||
)
|
||||
worker_task.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await worker_task
|
||||
logger.info(f"[{instance_id}] Worker task cancelled")
|
||||
finally:
|
||||
logger.info(
|
||||
f"[{instance_id}] _current_worker RESET (token={worker_token})"
|
||||
)
|
||||
_current_worker.reset(worker_token)
|
||||
finally:
|
||||
# Reset ContextVar
|
||||
logger.info(
|
||||
f"[{instance_id}] _current_docket RESET (token={docket_token})"
|
||||
)
|
||||
_current_docket.reset(docket_token)
|
||||
# Clear instance attribute
|
||||
self._docket = None
|
||||
logger.info(f"[{instance_id}] self._docket = None")
|
||||
finally:
|
||||
# Reset server ContextVar
|
||||
logger.info(f"[{instance_id}] _current_server RESET (token={server_token})")
|
||||
_current_server.reset(server_token)
|
||||
logger.info(f"[{instance_id}] _docket_lifespan EXITED")
|
||||
|
||||
async def _register_mounted_server_functions(
|
||||
self,
|
||||
|
|
@ -563,16 +603,30 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
|
||||
@asynccontextmanager
|
||||
async def _lifespan_manager(self) -> AsyncIterator[None]:
|
||||
import socket
|
||||
|
||||
instance_id = f"{socket.gethostname()}#{id(self)}"
|
||||
logger.info(f"[{instance_id}] _lifespan_manager ENTERING")
|
||||
|
||||
if self._lifespan_result_set:
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager: already set, yielding without setup"
|
||||
)
|
||||
yield
|
||||
return
|
||||
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager: entering user lifespan and docket lifespan"
|
||||
)
|
||||
async with (
|
||||
self._lifespan(self) as user_lifespan_result,
|
||||
self._docket_lifespan(),
|
||||
):
|
||||
self._lifespan_result = user_lifespan_result
|
||||
self._lifespan_result_set = True
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager: _lifespan_result_set = True"
|
||||
)
|
||||
|
||||
async with AsyncExitStack[bool | None]() as stack:
|
||||
for server in self._mounted_servers:
|
||||
|
|
@ -581,13 +635,21 @@ class FastMCP(Generic[LifespanResultT]):
|
|||
)
|
||||
|
||||
self._started.set()
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager YIELDING (server started, "
|
||||
f"_started.is_set={self._started.is_set()})"
|
||||
)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager EXITING (after yield)"
|
||||
)
|
||||
self._started.clear()
|
||||
|
||||
self._lifespan_result_set = False
|
||||
self._lifespan_result = None
|
||||
logger.info(f"[{instance_id}] _lifespan_manager EXITED")
|
||||
|
||||
async def run_async(
|
||||
self,
|
||||
|
|
|
|||
|
|
@ -44,6 +44,13 @@ async def handle_tool_as_task(
|
|||
Returns:
|
||||
CallToolResult: Task stub with task metadata in _meta
|
||||
"""
|
||||
import socket
|
||||
|
||||
from fastmcp.utilities.logging import get_logger
|
||||
|
||||
_logger = get_logger(__name__)
|
||||
instance_id = f"{socket.gethostname()}#{id(server)}"
|
||||
|
||||
# Generate server-side task ID per SEP-1686 final spec (line 375-377)
|
||||
# Server MUST generate task IDs, clients no longer provide them
|
||||
server_task_id = str(uuid.uuid4())
|
||||
|
|
@ -57,11 +64,20 @@ async def handle_tool_as_task(
|
|||
session_id = ctx.session_id
|
||||
|
||||
docket = _current_docket.get()
|
||||
_logger.info(
|
||||
f"[{instance_id}] handle_tool_as_task: tool={tool_name}, "
|
||||
f"_current_docket.get()={docket}, server._docket={server._docket}"
|
||||
)
|
||||
if docket is None:
|
||||
_logger.error(
|
||||
f"[{instance_id}] handle_tool_as_task FAILED: _current_docket is None! "
|
||||
f"server._started.is_set()={server._started.is_set()}"
|
||||
)
|
||||
raise McpError(
|
||||
ErrorData(
|
||||
code=INTERNAL_ERROR,
|
||||
message="Background tasks require a running FastMCP server context",
|
||||
message=f"Background tasks require a running FastMCP server context "
|
||||
f"(instance={instance_id}, server._docket={server._docket})",
|
||||
)
|
||||
)
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue