Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 41 additions & 15 deletions libs/code/deepagents_code/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -6663,6 +6663,7 @@ async def _run_session_start_sequence(self) -> None:

self._initial_session_started = True
self._startup_sequence_running = True
initial_submitted = False
try:
should_load_history = bool(self._lc_thread_id and self._agent) and (
self._resume_thread_intent is not None
Expand Down Expand Up @@ -6697,14 +6698,26 @@ async def _run_session_start_sequence(self) -> None:

if self._has_initial_submission():
await self._submit_initial_submission()
return
initial_submitted = True
finally:
self._startup_sequence_running = False

await self._drain_startup_backlog()
# Drain after the sequence completes. When an initial submission was
# dispatched it owns the session, so skip queued-input processing — but
# still drain deferred actions so an `/mcp login` queued during connect
# runs once the server is ready instead of being stranded.
await self._drain_startup_backlog(process_queue=not initial_submitted)

async def _drain_startup_backlog(self, *, process_queue: bool = True) -> None:
"""Drain deferred actions and queued input after server readiness.

async def _drain_startup_backlog(self) -> None:
"""Drain deferred actions and queued input after server readiness."""
Args:
process_queue: Whether to also process queued user input. Pass
`False` when an initial submission was just dispatched: it owns
the session, so queued input must wait, but deferred actions
(e.g. an `/mcp login` queued during connect) still need to run
rather than be stranded until the next agent turn.
"""
if self._agent_running or self._shell_running:
return

Expand All @@ -6722,7 +6735,7 @@ async def _drain_startup_backlog(self) -> None:
),
)

if self._pending_messages:
if process_queue and self._pending_messages:
await self._process_next_from_queue()

def _schedule_session_start_after_launch_init(
Expand Down Expand Up @@ -14736,11 +14749,12 @@ def on_worker_state_changed(self, event: Worker.StateChanged) -> None:
def _start_mcp_login(self, server_name: str) -> None:
"""Begin in-TUI OAuth login for `server_name`.

Guards against remote-server mode (no owned server to restart),
an absent local server, missing MCP config, an unknown server
name, and busy states. When the session is mid-run, the login
attempt is queued via `_defer_action` and runs once the user is
idle.
Rejects when MCP is disabled, in remote-server mode (no owned server
to restart), or while an agent switch is in progress. When the local
server is still connecting or the session is mid-run, the login is
queued via `_defer_action` and runs once the server is ready and the
user is idle. Config resolution and server-name validation happen
later, in `_run_mcp_login_worker`.

Args:
server_name: MCP server name from `mcpServers`.
Expand All @@ -14764,20 +14778,32 @@ def _start_mcp_login(self, server_name: str) -> None:
)
return

if self._connecting or self._server_proc is None:
if self._agent_switching:
self.notify(
"MCP login is unavailable until the local server is ready.",
"An agent switch is in progress; try again once it completes.",
severity="warning",
markup=False,
)
return

if self._agent_switching:
if self._connecting or self._server_proc is None:
# The server is still coming up (initial connect, a reconnect,
# or waiting on a model/credentials before startup). Queue the
# login instead of dropping the command so it runs once the
# server is ready; the queued action drains via
# `_maybe_drain_deferred` after connecting (and any startup
# submission) finishes.
self.notify(
"An agent switch is in progress; try again once it completes.",
severity="warning",
"MCP login will start once the local server is ready.",
timeout=5,
markup=False,
)
self._defer_action(
DeferredAction(
kind="mcp_login",
execute=lambda: self._run_mcp_login_worker(server_name),
),
)
return

if self._agent_running or self._shell_running:
Expand Down
142 changes: 137 additions & 5 deletions libs/code/tests/unit_tests/test_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -13519,6 +13519,102 @@ async def second_model() -> None: # noqa: RUF029
assert executed == ["thread", "second_model"]


class TestStartMcpLogin:
"""Tests for `_start_mcp_login` queueing before the server is connected."""

def _make_app(self) -> DeepAgentsApp:
"""Build an app that owns a server and has MCP enabled."""
app = DeepAgentsApp()
app._mcp_preload_kwargs = {"mcp_config_path": None} # ty: ignore
app._server_kwargs = {"model_name": "x"} # ty: ignore
app._agent_switching = False
app._agent_running = False
app._shell_running = False
app._defer_action = MagicMock() # ty: ignore
app.notify = MagicMock() # ty: ignore
app.run_worker = MagicMock() # ty: ignore
return app

async def test_queues_login_while_connecting(self) -> None:
"""A login sent while the server is connecting is deferred, not dropped."""
app = self._make_app()
app._connecting = True
app._server_proc = object() # ty: ignore

app._start_mcp_login("provider")

app._defer_action.assert_called_once() # ty: ignore
action = app._defer_action.call_args.args[0] # ty: ignore
assert action.kind == "mcp_login"
app.run_worker.assert_not_called() # ty: ignore
app.notify.assert_called_once() # ty: ignore
# The deferred-branch toast is informational (not a warning) and
# auto-dismisses, unlike the hard-reject guards above it.
notify_kwargs = app.notify.call_args.kwargs # ty: ignore
assert notify_kwargs.get("timeout") == 5
assert notify_kwargs.get("severity") != "warning"

async def test_rejects_login_during_agent_switch_even_while_connecting(
self,
) -> None:
"""Agent-swap reconnects warn instead of queueing a stranded login."""
app = self._make_app()
app._agent_switching = True
app._connecting = True
app._server_proc = object() # ty: ignore

app._start_mcp_login("provider")

app._defer_action.assert_not_called() # ty: ignore
app.run_worker.assert_not_called() # ty: ignore
app.notify.assert_called_once() # ty: ignore
assert "agent switch" in app.notify.call_args.args[0].lower() # ty: ignore

async def test_queues_login_when_server_proc_missing(self) -> None:
"""A login sent before the subprocess exists is deferred, not dropped."""
app = self._make_app()
app._connecting = False
app._server_proc = None # ty: ignore

app._start_mcp_login("provider")

app._defer_action.assert_called_once() # ty: ignore
assert app._defer_action.call_args.args[0].kind == "mcp_login" # ty: ignore
app.run_worker.assert_not_called() # ty: ignore

async def test_deferred_login_runs_target_worker(self) -> None:
"""The queued action invokes the login worker for the requested server."""
app = self._make_app()
app._connecting = True
app._server_proc = None # ty: ignore
app._run_mcp_login_worker = MagicMock() # ty: ignore

app._start_mcp_login("provider")

action = app._defer_action.call_args.args[0] # ty: ignore
action.execute()
app._run_mcp_login_worker.assert_called_once_with("provider") # ty: ignore

async def test_starts_worker_immediately_when_ready(self) -> None:
"""When the server is ready and idle, the login runs without deferral."""
app = self._make_app()
app._connecting = False
app._server_proc = object() # ty: ignore
# Mock the worker coroutine factory so building the `run_worker`
# argument doesn't create an un-awaited coroutine.
app._run_mcp_login_worker = MagicMock() # ty: ignore

app._start_mcp_login("provider")

app.run_worker.assert_called_once() # ty: ignore
app._defer_action.assert_not_called() # ty: ignore
# Each server logs in under its own worker group, and non-exclusively
# so a login never cancels an unrelated in-flight worker.
run_worker_kwargs = app.run_worker.call_args.kwargs # ty: ignore
assert run_worker_kwargs["group"] == "mcp-login-provider"
assert run_worker_kwargs["exclusive"] is False


class TestBuildModelSwitchErrorBody:
"""Tests for `_build_model_switch_error_body` link-aware formatting."""

Expand Down Expand Up @@ -18404,8 +18500,8 @@ async def test_toggle_disable_rejects_unknown_server_name(self) -> None:
assert kwargs.get("severity") == "warning"
assert kwargs.get("markup") is False

async def test_mcp_login_rejects_while_connecting(self) -> None:
"""`_connecting=True` prevents login until the server is ready."""
async def test_mcp_login_defers_while_connecting(self) -> None:
"""`_connecting=True` defers login until the server is ready."""
app = DeepAgentsApp(agent=MagicMock())
async with app.run_test() as pilot:
await pilot.pause()
Expand All @@ -18422,12 +18518,47 @@ async def test_mcp_login_rejects_while_connecting(self) -> None:
app._start_mcp_login("notion")
notify.assert_called_once()
assert "server is ready" in notify.call_args.args[0].lower()
assert not any(a.kind == "mcp_login" for a in app._deferred_actions)
assert any(a.kind == "mcp_login" for a in app._deferred_actions)
finally:
app._connecting = False

async def test_mcp_login_rejects_while_server_proc_is_none(self) -> None:
"""`_server_proc=None` prevents login until the server process exists."""
async def test_deferred_login_runs_after_server_ready(self) -> None:
"""A login queued during connect drains and runs once the server is ready.

Regression guard: an `/mcp login` deferred while connecting must
actually execute when the startup backlog drains, not sit in
`_deferred_actions` until some later agent turn.
"""
app = DeepAgentsApp(agent=MagicMock())
async with app.run_test() as pilot:
await pilot.pause()
app._mcp_preload_kwargs = {
"mcp_config_path": None,
"no_mcp": False,
"trust_project_mcp": None,
}
app._server_kwargs = {"some": "kwarg"}
app._agent_running = False
app._shell_running = False
app._startup_sequence_running = False
app._pending_messages.clear()

# Queue a login while the server is still connecting.
app._connecting = True
app._server_proc = None
with patch.object(app, "notify"):
app._start_mcp_login("notion")
assert any(a.kind == "mcp_login" for a in app._deferred_actions)

# Server becomes ready; draining the backlog runs the queued login.
app._connecting = False
with patch.object(app, "_run_mcp_login_worker", new=AsyncMock()) as worker:
await app._drain_startup_backlog()
worker.assert_awaited_once_with("notion")
assert not any(a.kind == "mcp_login" for a in app._deferred_actions)

async def test_mcp_login_defers_while_server_proc_is_none(self) -> None:
"""`_server_proc=None` defers login until the server process exists."""
app = DeepAgentsApp(agent=MagicMock())
async with app.run_test() as pilot:
await pilot.pause()
Expand All @@ -18447,6 +18578,7 @@ async def test_mcp_login_rejects_while_server_proc_is_none(self) -> None:
run_worker.assert_not_called()
notify.assert_called_once()
assert "server is ready" in notify.call_args.args[0].lower()
assert any(a.kind == "mcp_login" for a in app._deferred_actions)

async def test_mcp_login_rejects_while_agent_switching(self) -> None:
"""`_agent_switching=True` refuses login with a distinct message."""
Expand Down
Loading