mirror of
https://github.com/PrefectHQ/fastmcp.git
synced 2026-08-23 22:14:18 +02:00
Remove debug logging, keep minimal ContextVar fix
Removes all the verbose debug logging added during diagnosis while preserving the essential fix: Context.__aenter__ sets _current_docket and _current_worker from server instance attributes. This ensures ContextVars work in ASGI environments where lifespan and request handlers run in sibling async contexts. Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
parent
7b24bd58c8
commit
c127dd979f
5 changed files with 33 additions and 343 deletions
|
|
@ -175,25 +175,14 @@ class Context:
|
|||
|
||||
async def __aenter__(self) -> Context:
|
||||
"""Enter the context manager and set this context as the current context."""
|
||||
import socket
|
||||
|
||||
from fastmcp.utilities.logging import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
instance_id = f"{socket.gethostname()}#{id(self.fastmcp)}"
|
||||
|
||||
logger.info(f"[{instance_id}] Context.__aenter__ ENTERING")
|
||||
|
||||
parent_context = _current_context.get(None)
|
||||
if parent_context is not None:
|
||||
# Inherit state from parent context
|
||||
self._state = copy.deepcopy(parent_context._state)
|
||||
logger.info(f"[{instance_id}] Context: inherited state from parent context")
|
||||
|
||||
# Always set this context and save the token
|
||||
token = _current_context.set(self)
|
||||
self._tokens.append(token)
|
||||
logger.info(f"[{instance_id}] Context: _current_context SET (token={token})")
|
||||
|
||||
# Set current server for dependency injection (use weakref to avoid reference cycles)
|
||||
from fastmcp.server.dependencies import (
|
||||
|
|
@ -203,55 +192,21 @@ class Context:
|
|||
)
|
||||
|
||||
self._server_token = _current_server.set(weakref.ref(self.fastmcp))
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: _current_server SET (token={self._server_token})"
|
||||
)
|
||||
|
||||
# Set docket/worker from server instance for this request's context.
|
||||
# This ensures ContextVars work even in environments (like Lambda) where
|
||||
# lifespan ContextVars don't propagate to request handlers.
|
||||
server = self.fastmcp
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: server._docket={server._docket}, "
|
||||
f"server._worker={server._worker}, "
|
||||
f"_current_docket.get() BEFORE={_current_docket.get()}"
|
||||
)
|
||||
if server._docket is not None:
|
||||
self._docket_token = _current_docket.set(server._docket)
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: _current_docket SET from server._docket "
|
||||
f"(token={self._docket_token}), _current_docket.get() AFTER={_current_docket.get()}"
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
f"[{instance_id}] Context: server._docket is None, skipping _current_docket.set()"
|
||||
)
|
||||
|
||||
if server._worker is not None:
|
||||
self._worker_token = _current_worker.set(server._worker)
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: _current_worker SET from server._worker "
|
||||
f"(token={self._worker_token})"
|
||||
)
|
||||
else:
|
||||
logger.warning(
|
||||
f"[{instance_id}] Context: server._worker is None, skipping _current_worker.set()"
|
||||
)
|
||||
|
||||
logger.info(f"[{instance_id}] Context.__aenter__ COMPLETE")
|
||||
return self
|
||||
|
||||
async def __aexit__(self, exc_type, exc_val, exc_tb) -> None:
|
||||
"""Exit the context manager and reset the most recent token."""
|
||||
import socket
|
||||
|
||||
from fastmcp.utilities.logging import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
instance_id = f"{socket.gethostname()}#{id(self.fastmcp)}"
|
||||
|
||||
logger.info(f"[{instance_id}] Context.__aexit__ ENTERING (exc_type={exc_type})")
|
||||
|
||||
# Flush any remaining notifications before exiting
|
||||
await self._flush_notifications()
|
||||
|
||||
|
|
@ -263,34 +218,20 @@ class Context:
|
|||
)
|
||||
|
||||
if hasattr(self, "_worker_token"):
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: resetting _current_worker (token={self._worker_token})"
|
||||
)
|
||||
_current_worker.reset(self._worker_token)
|
||||
delattr(self, "_worker_token")
|
||||
if hasattr(self, "_docket_token"):
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: resetting _current_docket (token={self._docket_token})"
|
||||
)
|
||||
_current_docket.reset(self._docket_token)
|
||||
delattr(self, "_docket_token")
|
||||
if hasattr(self, "_server_token"):
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: resetting _current_server (token={self._server_token})"
|
||||
)
|
||||
_current_server.reset(self._server_token)
|
||||
delattr(self, "_server_token")
|
||||
|
||||
# Reset context token
|
||||
if self._tokens:
|
||||
token = self._tokens.pop()
|
||||
logger.info(
|
||||
f"[{instance_id}] Context: resetting _current_context (token={token})"
|
||||
)
|
||||
_current_context.reset(token)
|
||||
|
||||
logger.info(f"[{instance_id}] Context.__aexit__ COMPLETE")
|
||||
|
||||
@property
|
||||
def request_context(self) -> RequestContext[ServerSession, Any, Request] | None:
|
||||
"""Access to the underlying request context.
|
||||
|
|
|
|||
|
|
@ -386,8 +386,6 @@ 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)
|
||||
|
|
@ -397,26 +395,16 @@ 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,
|
||||
|
|
@ -469,9 +457,6 @@ 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] = {
|
||||
|
|
@ -483,55 +468,27 @@ 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]
|
||||
# Store on server instance for cross-context access
|
||||
self._worker = worker
|
||||
# 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)
|
||||
self._worker = None
|
||||
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,
|
||||
|
|
@ -607,33 +564,18 @@ 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:
|
||||
# Lifespan already ran - ContextVars will be set by Context.__aenter__
|
||||
# at request time, so we just yield here.
|
||||
logger.info(
|
||||
f"[{instance_id}] _lifespan_manager: already set, yielding "
|
||||
f"(ContextVars managed by Context at request time)"
|
||||
)
|
||||
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:
|
||||
|
|
@ -642,21 +584,13 @@ 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,13 +44,6 @@ 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())
|
||||
|
|
@ -65,20 +58,11 @@ async def handle_tool_as_task(
|
|||
|
||||
# Get Docket from ContextVar (set by Context.__aenter__ at request time)
|
||||
docket = _current_docket.get()
|
||||
_logger.info(
|
||||
f"[{instance_id}] handle_tool_as_task: tool={tool_name}, "
|
||||
f"_current_docket.get()={docket}"
|
||||
)
|
||||
if docket is None:
|
||||
_logger.error(
|
||||
f"[{instance_id}] handle_tool_as_task FAILED: no Docket available! "
|
||||
f"server._started.is_set()={server._started.is_set()}"
|
||||
)
|
||||
raise McpError(
|
||||
ErrorData(
|
||||
code=INTERNAL_ERROR,
|
||||
message=f"Background tasks require a running FastMCP server context "
|
||||
f"(instance={instance_id})",
|
||||
message="Background tasks require a running FastMCP server context",
|
||||
)
|
||||
)
|
||||
|
||||
|
|
@ -96,23 +80,9 @@ async def handle_tool_as_task(
|
|||
ttl_seconds = int(
|
||||
docket.execution_ttl.total_seconds() + TASK_MAPPING_TTL_BUFFER_SECONDS
|
||||
)
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to write to Redis: task_meta_key={task_meta_key}, "
|
||||
f"created_at_key={created_at_key}, ttl={ttl_seconds}"
|
||||
)
|
||||
try:
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
_logger.info(f"[{instance_id}] Redis write successful")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] Redis write FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
|
||||
# Send notifications/tasks/created per SEP-1686 (mandatory)
|
||||
# Send BEFORE queuing to avoid race where task completes before notification
|
||||
|
|
@ -134,38 +104,18 @@ async def handle_tool_as_task(
|
|||
|
||||
# Queue function to Docket by name (result storage via execution_ttl)
|
||||
# Use tool.key which matches what was registered - prefixed for mounted tools
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to call docket.add: tool.key={tool.key}, "
|
||||
f"task_key={task_key}, arguments={arguments}"
|
||||
)
|
||||
try:
|
||||
await docket.add(
|
||||
tool.key,
|
||||
key=task_key,
|
||||
)(**arguments)
|
||||
_logger.info(f"[{instance_id}] docket.add completed successfully")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] docket.add FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
await docket.add(
|
||||
tool.key,
|
||||
key=task_key,
|
||||
)(**arguments)
|
||||
|
||||
# Spawn subscription task to send status notifications (SEP-1686 optional feature)
|
||||
from fastmcp.server.tasks.subscriptions import subscribe_to_task_updates
|
||||
|
||||
# Start subscription in session's task group (persists for connection lifetime)
|
||||
_logger.info(
|
||||
f"[{instance_id}] Checking for subscription task group: "
|
||||
f"hasattr={hasattr(ctx.session, '_subscription_task_group')}"
|
||||
)
|
||||
if hasattr(ctx.session, "_subscription_task_group"):
|
||||
tg = ctx.session._subscription_task_group # type: ignore[attr-defined]
|
||||
_logger.info(f"[{instance_id}] Task group: {tg}")
|
||||
if tg:
|
||||
_logger.info(f"[{instance_id}] Starting subscription task")
|
||||
tg.start_soon( # type: ignore[union-attr]
|
||||
subscribe_to_task_updates,
|
||||
server_task_id,
|
||||
|
|
@ -173,9 +123,7 @@ async def handle_tool_as_task(
|
|||
ctx.session,
|
||||
docket,
|
||||
)
|
||||
_logger.info(f"[{instance_id}] Subscription task started")
|
||||
|
||||
_logger.info(f"[{instance_id}] About to return task stub")
|
||||
# Return task stub
|
||||
# Tasks MUST begin in "working" status per SEP-1686 final spec (line 381)
|
||||
return mcp.types.CallToolResult(
|
||||
|
|
@ -208,13 +156,6 @@ async def handle_prompt_as_task(
|
|||
Returns:
|
||||
GetPromptResult: 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())
|
||||
|
|
@ -229,20 +170,11 @@ async def handle_prompt_as_task(
|
|||
|
||||
# Get Docket from ContextVar (set by Context.__aenter__ at request time)
|
||||
docket = _current_docket.get()
|
||||
_logger.info(
|
||||
f"[{instance_id}] handle_prompt_as_task: prompt={prompt_name}, "
|
||||
f"_current_docket.get()={docket}"
|
||||
)
|
||||
if docket is None:
|
||||
_logger.error(
|
||||
f"[{instance_id}] handle_prompt_as_task FAILED: no Docket available! "
|
||||
f"server._started.is_set()={server._started.is_set()}"
|
||||
)
|
||||
raise McpError(
|
||||
ErrorData(
|
||||
code=INTERNAL_ERROR,
|
||||
message=f"Background tasks require a running FastMCP server context "
|
||||
f"(instance={instance_id})",
|
||||
message="Background tasks require a running FastMCP server context",
|
||||
)
|
||||
)
|
||||
|
||||
|
|
@ -260,23 +192,9 @@ async def handle_prompt_as_task(
|
|||
ttl_seconds = int(
|
||||
docket.execution_ttl.total_seconds() + TASK_MAPPING_TTL_BUFFER_SECONDS
|
||||
)
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to write to Redis: task_meta_key={task_meta_key}, "
|
||||
f"created_at_key={created_at_key}, ttl={ttl_seconds}"
|
||||
)
|
||||
try:
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
_logger.info(f"[{instance_id}] Redis write successful")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] Redis write FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
|
||||
# Send notifications/tasks/created per SEP-1686 (mandatory)
|
||||
# Send BEFORE queuing to avoid race where task completes before notification
|
||||
|
|
@ -295,38 +213,18 @@ async def handle_prompt_as_task(
|
|||
|
||||
# Queue function to Docket by name (result storage via execution_ttl)
|
||||
# Use prompt.key which matches what was registered - prefixed for mounted prompts
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to call docket.add: prompt.key={prompt.key}, "
|
||||
f"task_key={task_key}, arguments={arguments}"
|
||||
)
|
||||
try:
|
||||
await docket.add(
|
||||
prompt.key,
|
||||
key=task_key,
|
||||
)(**(arguments or {}))
|
||||
_logger.info(f"[{instance_id}] docket.add completed successfully")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] docket.add FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
await docket.add(
|
||||
prompt.key,
|
||||
key=task_key,
|
||||
)(**(arguments or {}))
|
||||
|
||||
# Spawn subscription task to send status notifications (SEP-1686 optional feature)
|
||||
from fastmcp.server.tasks.subscriptions import subscribe_to_task_updates
|
||||
|
||||
# Start subscription in session's task group (persists for connection lifetime)
|
||||
_logger.info(
|
||||
f"[{instance_id}] Checking for subscription task group: "
|
||||
f"hasattr={hasattr(ctx.session, '_subscription_task_group')}"
|
||||
)
|
||||
if hasattr(ctx.session, "_subscription_task_group"):
|
||||
tg = ctx.session._subscription_task_group # type: ignore[attr-defined]
|
||||
_logger.info(f"[{instance_id}] Task group: {tg}")
|
||||
if tg:
|
||||
_logger.info(f"[{instance_id}] Starting subscription task")
|
||||
tg.start_soon( # type: ignore[union-attr]
|
||||
subscribe_to_task_updates,
|
||||
server_task_id,
|
||||
|
|
@ -334,9 +232,7 @@ async def handle_prompt_as_task(
|
|||
ctx.session,
|
||||
docket,
|
||||
)
|
||||
_logger.info(f"[{instance_id}] Subscription task started")
|
||||
|
||||
_logger.info(f"[{instance_id}] About to return task stub")
|
||||
# Return task stub
|
||||
# Tasks MUST begin in "working" status per SEP-1686 final spec (line 381)
|
||||
return mcp.types.GetPromptResult(
|
||||
|
|
@ -370,13 +266,6 @@ async def handle_resource_as_task(
|
|||
Returns:
|
||||
ServerResult with ReadResourceResult stub
|
||||
"""
|
||||
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())
|
||||
|
|
@ -391,20 +280,11 @@ async def handle_resource_as_task(
|
|||
|
||||
# Get Docket from ContextVar (set by Context.__aenter__ at request time)
|
||||
docket = _current_docket.get()
|
||||
_logger.info(
|
||||
f"[{instance_id}] handle_resource_as_task: uri={uri}, "
|
||||
f"_current_docket.get()={docket}"
|
||||
)
|
||||
if docket is None:
|
||||
_logger.error(
|
||||
f"[{instance_id}] handle_resource_as_task FAILED: no Docket available! "
|
||||
f"server._started.is_set()={server._started.is_set()}"
|
||||
)
|
||||
raise McpError(
|
||||
ErrorData(
|
||||
code=INTERNAL_ERROR,
|
||||
message=f"Background tasks require a running FastMCP server context "
|
||||
f"(instance={instance_id})",
|
||||
message="Background tasks require a running FastMCP server context",
|
||||
)
|
||||
)
|
||||
|
||||
|
|
@ -419,23 +299,9 @@ async def handle_resource_as_task(
|
|||
ttl_seconds = int(
|
||||
docket.execution_ttl.total_seconds() + TASK_MAPPING_TTL_BUFFER_SECONDS
|
||||
)
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to write to Redis: task_meta_key={task_meta_key}, "
|
||||
f"created_at_key={created_at_key}, ttl={ttl_seconds}"
|
||||
)
|
||||
try:
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
_logger.info(f"[{instance_id}] Redis write successful")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] Redis write FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
async with docket.redis() as redis:
|
||||
await redis.set(task_meta_key, task_key, ex=ttl_seconds)
|
||||
await redis.set(created_at_key, created_at, ex=ttl_seconds)
|
||||
|
||||
# Send notifications/tasks/created per SEP-1686 (mandatory)
|
||||
# Send BEFORE queuing to avoid race where task completes before notification
|
||||
|
|
@ -459,57 +325,23 @@ async def handle_resource_as_task(
|
|||
|
||||
if isinstance(resource, FunctionResourceTemplate):
|
||||
params = match_uri_template(uri, resource.uri_template) or {}
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to call docket.add: resource.name={resource.name}, "
|
||||
f"task_key={task_key}, params={params} (template)"
|
||||
)
|
||||
try:
|
||||
await docket.add(
|
||||
resource.name,
|
||||
key=task_key,
|
||||
)(**params)
|
||||
_logger.info(f"[{instance_id}] docket.add completed successfully")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] docket.add FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
await docket.add(
|
||||
resource.name,
|
||||
key=task_key,
|
||||
)(**params)
|
||||
else:
|
||||
_logger.info(
|
||||
f"[{instance_id}] About to call docket.add: resource.name={resource.name}, "
|
||||
f"task_key={task_key} (static resource)"
|
||||
)
|
||||
try:
|
||||
await docket.add(
|
||||
resource.name,
|
||||
key=task_key,
|
||||
)()
|
||||
_logger.info(f"[{instance_id}] docket.add completed successfully")
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
_logger.error(
|
||||
f"[{instance_id}] docket.add FAILED: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
raise
|
||||
await docket.add(
|
||||
resource.name,
|
||||
key=task_key,
|
||||
)()
|
||||
|
||||
# Spawn subscription task to send status notifications (SEP-1686 optional feature)
|
||||
from fastmcp.server.tasks.subscriptions import subscribe_to_task_updates
|
||||
|
||||
# Start subscription in session's task group (persists for connection lifetime)
|
||||
_logger.info(
|
||||
f"[{instance_id}] Checking for subscription task group: "
|
||||
f"hasattr={hasattr(ctx.session, '_subscription_task_group')}"
|
||||
)
|
||||
if hasattr(ctx.session, "_subscription_task_group"):
|
||||
tg = ctx.session._subscription_task_group # type: ignore[attr-defined]
|
||||
_logger.info(f"[{instance_id}] Task group: {tg}")
|
||||
if tg:
|
||||
_logger.info(f"[{instance_id}] Starting subscription task")
|
||||
tg.start_soon( # type: ignore[union-attr]
|
||||
subscribe_to_task_updates,
|
||||
server_task_id,
|
||||
|
|
@ -517,9 +349,7 @@ async def handle_resource_as_task(
|
|||
ctx.session,
|
||||
docket,
|
||||
)
|
||||
_logger.info(f"[{instance_id}] Subscription task started")
|
||||
|
||||
_logger.info(f"[{instance_id}] About to return task stub")
|
||||
# Return task stub
|
||||
# Tasks MUST begin in "working" status per SEP-1686 final spec (line 381)
|
||||
return mcp.types.ServerResult(
|
||||
|
|
|
|||
|
|
@ -42,23 +42,13 @@ async def subscribe_to_task_updates(
|
|||
session: MCP ServerSession for sending notifications
|
||||
docket: Docket instance for subscribing to execution events
|
||||
"""
|
||||
logger.info(
|
||||
f"subscribe_to_task_updates STARTING: task_id={task_id}, task_key={task_key}"
|
||||
)
|
||||
try:
|
||||
logger.info(
|
||||
f"subscribe_to_task_updates: About to call docket.get_execution({task_key})"
|
||||
)
|
||||
execution = await docket.get_execution(task_key)
|
||||
logger.info(
|
||||
f"subscribe_to_task_updates: docket.get_execution returned {execution}"
|
||||
)
|
||||
if execution is None:
|
||||
logger.warning(f"No execution found for task {task_id}")
|
||||
return
|
||||
|
||||
# Subscribe to state and progress events from Docket
|
||||
logger.info("subscribe_to_task_updates: About to subscribe to execution events")
|
||||
async for event in execution.subscribe():
|
||||
if event["type"] == "state":
|
||||
# Send notifications/tasks/status when state changes
|
||||
|
|
@ -80,12 +70,7 @@ async def subscribe_to_task_updates(
|
|||
)
|
||||
|
||||
except Exception as e:
|
||||
import traceback
|
||||
|
||||
logger.error(
|
||||
f"subscribe_to_task_updates FAILED for {task_id}: {type(e).__name__}: {e}\n"
|
||||
f"Traceback:\n{traceback.format_exc()}"
|
||||
)
|
||||
logger.error(f"subscribe_to_task_updates failed for {task_id}: {e}")
|
||||
|
||||
|
||||
async def _send_status_notification(
|
||||
|
|
|
|||
8
uv.lock
generated
8
uv.lock
generated
|
|
@ -752,7 +752,7 @@ requires-dist = [
|
|||
{ name = "platformdirs", specifier = ">=4.0.0" },
|
||||
{ name = "py-key-value-aio", extras = ["disk", "keyring", "memory"], specifier = ">=0.3.0,<0.4.0" },
|
||||
{ name = "pydantic", extras = ["email"], specifier = ">=2.11.7" },
|
||||
{ name = "pydocket", specifier = ">=0.16.4" },
|
||||
{ name = "pydocket", specifier = ">=0.16.6" },
|
||||
{ name = "pyperclip", specifier = ">=1.9.0" },
|
||||
{ name = "python-dotenv", specifier = ">=1.1.0" },
|
||||
{ name = "rich", specifier = ">=13.9.4" },
|
||||
|
|
@ -1775,7 +1775,7 @@ wheels = [
|
|||
|
||||
[[package]]
|
||||
name = "pydocket"
|
||||
version = "0.16.4"
|
||||
version = "0.16.6"
|
||||
source = { registry = "https://pypi.org/simple" }
|
||||
dependencies = [
|
||||
{ name = "cloudpickle" },
|
||||
|
|
@ -1792,9 +1792,9 @@ dependencies = [
|
|||
{ name = "typer" },
|
||||
{ name = "typing-extensions" },
|
||||
]
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/4d/c6/eb7f3af72fa5c04b52a3f9390ff0c948987441987f9526dd992d2a6b3524/pydocket-0.16.4.tar.gz", hash = "sha256:d034d1ac75877560d86329fb3643e7b862fcbcdac407d876a62f5d9e386e8753", size = 297949, upload-time = "2026-01-08T21:58:31.637Z" }
|
||||
sdist = { url = "https://files.pythonhosted.org/packages/72/00/26befe5f58df7cd1aeda4a8d10bc7d1908ffd86b80fd995e57a2a7b3f7bd/pydocket-0.16.6.tar.gz", hash = "sha256:b96c96ad7692827214ed4ff25fcf941ec38371314db5dcc1ae792b3e9d3a0294", size = 299054, upload-time = "2026-01-09T22:09:15.405Z" }
|
||||
wheels = [
|
||||
{ url = "https://files.pythonhosted.org/packages/74/c5/e6ffed3902ead6cb906758749c42211f7f24ea6d8fdd772f531f5a81c9fa/pydocket-0.16.4-py3-none-any.whl", hash = "sha256:cdcdf74b987c2cd5d03c7353d15f8dd2ac9bd43f2a91f7441748d5a8ebd617c9", size = 67374, upload-time = "2026-01-08T21:58:30.01Z" },
|
||||
{ url = "https://files.pythonhosted.org/packages/0a/3f/7483e5a6dc6326b6e0c640619b5c5bd1d6e3c20e54d58f5fb86267cef00e/pydocket-0.16.6-py3-none-any.whl", hash = "sha256:683d21e2e846aa5106274e7d59210331b242d7fb0dce5b08d3b82065663ed183", size = 67697, upload-time = "2026-01-09T22:09:13.436Z" },
|
||||
]
|
||||
|
||||
[[package]]
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue