Skip to content

runtime

Application lifespan and connection supervision for custom transports.

ServerRuntime dataclass

Bases: Generic[LifespanT]

An active server returned by Server.serve() or MCPServer.serve().

Each connect() call serves a peer with independent protocol and request-ID state.

Source code in src/mcp/server/runtime.py
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
@dataclass
class ServerRuntime(Generic[LifespanT]):
    """An active server returned by `Server.serve()` or `MCPServer.serve()`.

    Each `connect()` call serves a peer with independent protocol and request-ID state.
    """

    _server: Server[LifespanT]
    _lifespan_state: LifespanT
    _task_group: anyio.abc.TaskGroup
    _limiter: anyio.CapacityLimiter
    _active: bool = field(default=True, init=False)

    @classmethod
    @asynccontextmanager
    async def open(
        cls, server: Server[LifespanT], *, max_connections: int = 100
    ) -> AsyncIterator[ServerRuntime[LifespanT]]:
        """Share one lifespan, closing connections before application cleanup.

        Cleanup has a five-second cancellation deadline per layer and must cooperate.
        """
        if max_connections < 1:
            raise ValueError("max_connections must be positive")
        body_error: BaseException | None = None
        with anyio.CancelScope() as lifespan_scope:
            async with server.lifespan(server) as state:
                try:
                    async with anyio.create_task_group() as tg:
                        runtime = cls(server, state, tg, anyio.CapacityLimiter(max_connections))
                        try:
                            yield runtime
                        finally:
                            runtime._active = False
                            tg.cancel_scope.cancel()
                except BaseException as exc:
                    body_error = exc
                    raise
                finally:
                    lifespan_scope.shield = True
                    lifespan_scope.deadline = anyio.current_time() + 5
        if lifespan_scope.cancelled_caught:
            logger.warning("Server lifespan cleanup exceeded five seconds")
            if body_error is not None:
                raise body_error
        await resync_tracer()

    async def connect(
        self,
        transport: Transport | DispatcherTransport,
        *,
        session_id: str | None = None,
        transport_builder: TransportContextBuilder | None = None,
    ) -> None:
        """Open and supervise one peer's transport until it disconnects.

        Waits for capacity, then owns the entered transport. Opening failures reach
        the caller; later failures are logged and isolated. After return, caller
        cancellation does not close the connection.
        Dispatcher transports signal readiness through `Dispatcher.run()`;
        message transports return before receiving the first MCP request.

        Args:
            transport: An unopened message transport or `DispatcherTransport`.
            session_id: Optional identity for a handshake-era connection.
            transport_builder: Builds handler metadata for each inbound message.

        Raises:
            RuntimeError: If this runtime has closed.
        """
        if not self._active:
            raise RuntimeError("Server runtime is closed")
        if isinstance(transport, DispatcherTransport) and (session_id is not None or transport_builder is not None):
            raise ValueError("Dispatcher transports supply their own context and do not use handshake-era sessions")

        async def serve(*, task_status: anyio.abc.TaskStatus[None]) -> None:
            ready = False
            run_error: BaseException | None = None

            class ReadyStatus:
                def started(self, value: None = None) -> None:
                    nonlocal ready
                    task_status.started()
                    ready = True

            status = ReadyStatus()
            try:
                async with self._limiter:
                    with anyio.CancelScope() as cleanup_scope:
                        async with AsyncExitStack() as stack:
                            if isinstance(transport, DispatcherTransport):
                                dispatcher = await stack.enter_async_context(transport.connection)
                                run = partial(
                                    serve_modern_dispatcher,
                                    self._server,
                                    dispatcher,
                                    lifespan_state=self._lifespan_state,
                                    task_status=status,
                                )
                            else:
                                read, write = await stack.enter_async_context(transport)
                                run = partial(
                                    serve_dual_era_loop,
                                    self._server,
                                    read,
                                    write,
                                    lifespan_state=self._lifespan_state,
                                    session_id=session_id,
                                    transport_builder=transport_builder,
                                )
                            try:
                                if not isinstance(transport, DispatcherTransport):
                                    status.started()
                                await run()
                            except BaseException as exc:
                                run_error = exc
                                raise
                            finally:
                                cleanup_scope.shield = True
                                cleanup_scope.deadline = anyio.current_time() + 5
                    if cleanup_scope.cancelled_caught:
                        logger.warning("Transport cleanup exceeded five seconds")
                        if run_error is not None:
                            raise run_error
                    return
            except Exception:
                if not ready:
                    raise
                logger.exception("Transport connection failed")

        await self._task_group.start(serve)

open async classmethod

open(
    server: Server[LifespanT], *, max_connections: int = 100
) -> AsyncIterator[ServerRuntime[LifespanT]]

Share one lifespan, closing connections before application cleanup.

Cleanup has a five-second cancellation deadline per layer and must cooperate.

Source code in src/mcp/server/runtime.py
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
@classmethod
@asynccontextmanager
async def open(
    cls, server: Server[LifespanT], *, max_connections: int = 100
) -> AsyncIterator[ServerRuntime[LifespanT]]:
    """Share one lifespan, closing connections before application cleanup.

    Cleanup has a five-second cancellation deadline per layer and must cooperate.
    """
    if max_connections < 1:
        raise ValueError("max_connections must be positive")
    body_error: BaseException | None = None
    with anyio.CancelScope() as lifespan_scope:
        async with server.lifespan(server) as state:
            try:
                async with anyio.create_task_group() as tg:
                    runtime = cls(server, state, tg, anyio.CapacityLimiter(max_connections))
                    try:
                        yield runtime
                    finally:
                        runtime._active = False
                        tg.cancel_scope.cancel()
            except BaseException as exc:
                body_error = exc
                raise
            finally:
                lifespan_scope.shield = True
                lifespan_scope.deadline = anyio.current_time() + 5
    if lifespan_scope.cancelled_caught:
        logger.warning("Server lifespan cleanup exceeded five seconds")
        if body_error is not None:
            raise body_error
    await resync_tracer()

connect async

connect(
    transport: Transport | DispatcherTransport,
    *,
    session_id: str | None = None,
    transport_builder: TransportContextBuilder | None = None
) -> None

Open and supervise one peer's transport until it disconnects.

Waits for capacity, then owns the entered transport. Opening failures reach the caller; later failures are logged and isolated. After return, caller cancellation does not close the connection. Dispatcher transports signal readiness through Dispatcher.run(); message transports return before receiving the first MCP request.

Parameters:

Name Type Description Default
transport Transport | DispatcherTransport

An unopened message transport or DispatcherTransport.

required
session_id str | None

Optional identity for a handshake-era connection.

None
transport_builder TransportContextBuilder | None

Builds handler metadata for each inbound message.

None

Raises:

Type Description
RuntimeError

If this runtime has closed.

Source code in src/mcp/server/runtime.py
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
async def connect(
    self,
    transport: Transport | DispatcherTransport,
    *,
    session_id: str | None = None,
    transport_builder: TransportContextBuilder | None = None,
) -> None:
    """Open and supervise one peer's transport until it disconnects.

    Waits for capacity, then owns the entered transport. Opening failures reach
    the caller; later failures are logged and isolated. After return, caller
    cancellation does not close the connection.
    Dispatcher transports signal readiness through `Dispatcher.run()`;
    message transports return before receiving the first MCP request.

    Args:
        transport: An unopened message transport or `DispatcherTransport`.
        session_id: Optional identity for a handshake-era connection.
        transport_builder: Builds handler metadata for each inbound message.

    Raises:
        RuntimeError: If this runtime has closed.
    """
    if not self._active:
        raise RuntimeError("Server runtime is closed")
    if isinstance(transport, DispatcherTransport) and (session_id is not None or transport_builder is not None):
        raise ValueError("Dispatcher transports supply their own context and do not use handshake-era sessions")

    async def serve(*, task_status: anyio.abc.TaskStatus[None]) -> None:
        ready = False
        run_error: BaseException | None = None

        class ReadyStatus:
            def started(self, value: None = None) -> None:
                nonlocal ready
                task_status.started()
                ready = True

        status = ReadyStatus()
        try:
            async with self._limiter:
                with anyio.CancelScope() as cleanup_scope:
                    async with AsyncExitStack() as stack:
                        if isinstance(transport, DispatcherTransport):
                            dispatcher = await stack.enter_async_context(transport.connection)
                            run = partial(
                                serve_modern_dispatcher,
                                self._server,
                                dispatcher,
                                lifespan_state=self._lifespan_state,
                                task_status=status,
                            )
                        else:
                            read, write = await stack.enter_async_context(transport)
                            run = partial(
                                serve_dual_era_loop,
                                self._server,
                                read,
                                write,
                                lifespan_state=self._lifespan_state,
                                session_id=session_id,
                                transport_builder=transport_builder,
                            )
                        try:
                            if not isinstance(transport, DispatcherTransport):
                                status.started()
                            await run()
                        except BaseException as exc:
                            run_error = exc
                            raise
                        finally:
                            cleanup_scope.shield = True
                            cleanup_scope.deadline = anyio.current_time() + 5
                if cleanup_scope.cancelled_caught:
                    logger.warning("Transport cleanup exceeded five seconds")
                    if run_error is not None:
                        raise run_error
                return
        except Exception:
            if not ready:
                raise
            logger.exception("Transport connection failed")

    await self._task_group.start(serve)