From 28dccb908c3e58d26a8d372e3779c28a61ed4c8d Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Tue, 21 Jul 2026 13:46:27 +0800 Subject: [PATCH] docs(superpowers): add host-agent MCP server implementation plan 15-task TDD plan implementing the spec committed in 3b62195. Covers the four new host-agent modules (mcp_token, mcp_lock, web/mcp_auth, web/mcp), cloud heartbeat + scheduler coordination, console mount wiring, CLI subcommand, and documentation. Co-Authored-By: Claude Opus 4.6 --- .../plans/2026-07-21-host-agent-mcp-server.md | 2663 +++++++++++++++++ 1 file changed, 2663 insertions(+) create mode 100644 docs/superpowers/plans/2026-07-21-host-agent-mcp-server.md diff --git a/docs/superpowers/plans/2026-07-21-host-agent-mcp-server.md b/docs/superpowers/plans/2026-07-21-host-agent-mcp-server.md new file mode 100644 index 0000000..900f412 --- /dev/null +++ b/docs/superpowers/plans/2026-07-21-host-agent-mcp-server.md @@ -0,0 +1,2663 @@ +# Host-Agent MCP Server Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Mount a Streamable HTTP MCP server inside the host-agent process so external MCP clients (Hermes Agent) can drive devices directly, coexisting with the Cloud Control Plane worker path. + +**Architecture:** Reuse the existing host-agent console FastAPI + uvicorn on port 8765. Add `app.mount("/mcp", auth_wrapped(FastMCP.streamable_http_app()))`. The FastMCP server wraps `api.mcp.tool_handlers(manager=...)` with a busy-check + lazy-lock decorator. Per-device session-level locks with 60s TTL. Cloud scheduler learns about MCP-busy devices via a new heartbeat field. + +**Tech Stack:** Python 3.13, FastAPI, uvicorn, FastMCP (`mcp>=1.28,<2`), Pydantic v2, pytest, secrets/tempfile/os.replace for atomic token files, threading.Lock for in-process locks. + +## Global Constraints + +- **Python**: `>=3.13,<3.14` (workspace pin, see `apps/device-host-agent/pyproject.toml`) +- **Package boundary**: `runtime/` and `api/` packages MUST NOT import `cloud` or `host_agent`. Enforced by `apps/device-host-agent/tests/test_execution.py::test_runtime_owned_packages_do_not_import_host_or_cloud_concerns`. +- **`tool_handlers` signature**: After Task 2, `tool_handlers(*, manager: DeviceManager)` is required (no None default). The single `api/mcp.py::tool_handlers` wrapper in this repo MUST get an explicit manager — never fall back to `DEFAULT_MANAGER`. +- **Network binding**: MCP server inherits console's loopback-only bind. `HOST_AGENT_CONSOLE_ALLOW_NON_LOOPBACK=true` is the only escape hatch (existing config). +- **Token storage**: `tempfile.NamedTemporaryFile` + `os.replace` atomic rename; file mode `0o600` on POSIX; JSON schema `{"version": 1, "token": "...", "created_at": ""}`. +- **Lock semantics**: per-device, session-level, lazy acquire on first device-touching tool call, renew on every call, release on session end, TTL sweep on read. MVP callers all use try-acquire (`acquire()` returning False on conflict); `wait_until_usable` is implemented + unit-tested but **not** invoked by any production caller. +- **Test commands**: + - Single package: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_.py -v` + - Cloud package: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_.py -v` + - Full non-integration: `uv run --all-packages pytest -m "not integration"` + - Ruff: `uv run --with ruff ruff check ` then `uv run --with ruff ruff format --check ` + - Compile: `uv run --all-packages python -m compileall ` + +## File Structure + +### New files +- `apps/device-host-agent/host_agent/mcp_token.py` — `McpTokenStore`, `McpToken`, `McpTokenStoreError` +- `apps/device-host-agent/host_agent/mcp_lock.py` — `McpBusyTracker`, `McpDeviceLease` +- `apps/device-host-agent/host_agent/web/mcp_auth.py` — `BearerAuthMiddleware` +- `apps/device-host-agent/host_agent/web/mcp.py` — `build_mcp_server(*, manager, mcp_busy_tracker, status_tracker) -> FastMCP` +- `apps/device-host-agent/tests/test_mcp_token.py` +- `apps/device-host-agent/tests/test_mcp_lock.py` +- `apps/device-host-agent/tests/test_mcp_auth.py` +- `apps/device-host-agent/tests/test_web_mcp.py` +- `docs/MCP_INTEGRATION.md` + +### Modified files +- `apps/device-host-agent/pyproject.toml` — add `mcp>=1.28,<2` direct dep +- `api/mcp.py` — `tool_handlers(*, manager: DeviceManager)` required-arg tighten +- `apps/device-host-agent/host_agent/app.py` — construct new components in `create_application()`, pass to `create_console_app()` +- `apps/device-host-agent/host_agent/web/app.py` — `create_console_app(...)` accepts new params, mounts `/mcp`, adds `mcp_busy_devices` to `/api/status` +- `apps/device-host-agent/host_agent/web/templates/dashboard.html` — one MCP status row +- `apps/device-host-agent/host_agent/cli.py` — new `mcp-token` subcommand +- `apps/device-host-agent/host_agent/heartbeat.py` — `HeartbeatSynchronizer` accepts optional `mcp_busy_tracker`, `build_device_snapshot` adds `mcp_busy_device_ids` +- `apps/device-host-agent/host_agent/client.py` — `heartbeat()` accepts and sends `mcp_busy_device_ids` +- `apps/device-host-agent/host_agent/assignment.py` — `AssignmentExecutor` accepts optional `mcp_busy_tracker`; `execute()` does fail-fast check +- `packages/cloud-platform/cloud/internal_api/models.py` — `HeartbeatRequest.mcp_busy_device_ids: list[str] = []` +- `packages/cloud-platform/cloud/pool.py` — `PooledDevice.mcp_busy: bool = False`; `sync_host_devices()` accepts `mcp_busy_device_ids`; `_to_pooled()` sets flag +- `packages/cloud-platform/cloud/scheduler.py` — `_matches()` rejects `device.mcp_busy` +- `packages/cloud-platform/cloud/internal_api/api.py` — heartbeat handler passes `mcp_busy_device_ids` to `pool.sync_host_devices()` +- `docs/MACOS_IPHONE_SETUP.md` — add "MCP Integration" section + +--- + +## Task 1: Add `mcp` as direct dependency of device-host-agent + +**Files:** +- Modify: `apps/device-host-agent/pyproject.toml:6-14` + +**Interfaces:** +- Produces: device-host-agent now declares `mcp>=1.28,<2` as a direct dep (it was previously transitively available via device-agent-runtime). + +- [ ] **Step 1: Read current pyproject.toml** + +Run: `cat apps/device-host-agent/pyproject.toml` +Expected output: see lines 6-14 are the `dependencies = [...]` list with 7 entries (`device-agent-runtime`, `device-cloud-platform`, `fastapi`, `filelock`, `httpx`, `jinja2`, `uvicorn[standard]`). + +- [ ] **Step 2: Add `mcp>=1.28,<2` to dependencies** + +Edit `apps/device-host-agent/pyproject.toml`. The final `dependencies` block should be: + +```toml +dependencies = [ + "device-agent-runtime==0.1.0", + "device-cloud-platform==0.1.0", + "fastapi>=0.115.0", + "filelock>=3.0", + "httpx>=0.27.0", + "jinja2>=3.1", + "mcp>=1.28,<2", + "uvicorn[standard]>=0.30.0", +] +``` + +(Note: keep entries alphabetically ordered per existing convention.) + +- [ ] **Step 3: Re-lock the workspace** + +Run: `uv lock` +Expected: uv reports "Resolved X packages" and updates `uv.lock` with `mcp` listed under `device-host-agent`'s dependencies. `mcp` version remains `1.28.1` (already in lock as transitive — now dual-listed as direct). + +- [ ] **Step 4: Verify sync works** + +Run: `uv sync --locked --all-packages` +Expected: no error; "Resolved X packages" + "Installed X packages" or "Audited X packages" with no diffs. + +- [ ] **Step 5: Verify import works in device-host-agent context** + +Run: `uv run --package device-host-agent python -c "from mcp.server.fastmcp import FastMCP; print('ok')"` +Expected output: `ok` + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/pyproject.toml uv.lock +git commit -m "build(host-agent): add mcp as direct dependency" +``` + +--- + +## Task 2: Tighten `tool_handlers` signature (D12) + +Eliminate the silent fallback to `DEFAULT_MANAGER` that caused the `DeviceNotFoundError` incident (see `apps/device-host-agent/host_agent/execution.py` regression tests for the same trap in task runners). + +**Files:** +- Modify: `api/mcp.py:18-95` +- Modify: `api/mcp.py:98-110` (`create_mcp_server` now must receive non-None manager or raise) +- Test: `tests/test_mcp.py` (extend) + +**Interfaces:** +- Produces: `tool_handlers(*, manager: DeviceManager) -> dict[str, Callable[..., Any]]` (manager is keyword-required, no default). +- Produces: `create_mcp_server(*, manager: DeviceManager, ...)` raises `ValueError` if `manager is None`. + +- [ ] **Step 1: Write failing test that `tool_handlers()` without manager raises** + +Append to `tests/test_mcp.py`: + +```python +def test_tool_handlers_requires_manager() -> None: + """D12: tool_handlers must not silently fall back to DEFAULT_MANAGER.""" + from api.mcp import tool_handlers + import pytest + + with pytest.raises(TypeError): + tool_handlers() # type: ignore[call-arg] +``` + +- [ ] **Step 2: Run test to verify it fails (signature still allows None)** + +Run: `uv run --all-packages pytest tests/test_mcp.py::test_tool_handlers_requires_manager -v` +Expected: FAIL — currently `tool_handlers(manager=None)` succeeds because `None` is the default. + +- [ ] **Step 3: Tighten signature in `api/mcp.py`** + +Replace `api/mcp.py:18-22` (the existing `def tool_handlers(...)` declaration) with: + +```python +def tool_handlers( + *, + manager: DeviceManager, +) -> dict[str, Callable[..., Any]]: +``` + +Remove the line `device_manager = manager or DEFAULT_MANAGER` (was line 22 in original) and rename all uses of `device_manager` inside the function body to `manager`. The `DEFAULT_MANAGER` import on line 6 (`from device.manager import DeviceManager, DEFAULT_MANAGER`) becomes unused — remove `DEFAULT_MANAGER` from the import, keeping only `DeviceManager`. + +In `create_mcp_server` (around line 98-110), if `manager is None: raise ValueError("create_mcp_server requires a non-None manager")` before constructing handlers. Update the signature `manager: DeviceManager | None = None` to `manager: DeviceManager` (required keyword). + +- [ ] **Step 4: Run test to verify it passes** + +Run: `uv run --all-packages pytest tests/test_mcp.py::test_tool_handlers_requires_manager -v` +Expected: PASS. + +- [ ] **Step 5: Verify all existing callers still work** + +The grep in spec §10 Q1 found these callers, all of which already pass `manager=manager`: +- `api/mcp.py:110` — internal call inside `create_mcp_server` +- `tests/test_mcp.py:13, 56` +- `tests/test_skill_catalog_e2e.py:287` + +Run: `uv run --all-packages pytest tests/test_mcp.py tests/test_skill_catalog_e2e.py -v` +Expected: all PASS. + +- [ ] **Step 6: Verify `test_runtime_owned_packages_do_not_import_host_or_cloud_concerns` still passes** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_execution.py::test_runtime_owned_packages_do_not_import_host_or_cloud_concerns -v` +Expected: PASS (this task doesn't change package boundaries). + +- [ ] **Step 7: Commit** + +```bash +git add api/mcp.py tests/test_mcp.py +git commit -m "refactor(api): make tool_handlers require a DeviceManager + +Eliminates the silent fallback to DEFAULT_MANAGER that produced the +DeviceNotFoundError incident. All existing callers already pass +manager explicitly." +``` + +--- + +## Task 3: `McpTokenStore` — token file persistence + +**Files:** +- Create: `apps/device-host-agent/host_agent/mcp_token.py` +- Test: `apps/device-host-agent/tests/test_mcp_token.py` + +**Interfaces:** +- Produces: `McpTokenStore(path: Path, *, now: Callable[[], datetime] | None = None)` with methods: + - `load_or_create() -> McpToken` — atomic; reads if exists, else generates + writes + - `verify(presented: str) -> bool` — `hmac.compare_digest` against loaded token +- Produces: `McpToken(version: int, token: str, created_at: datetime)` (frozen dataclass) +- Produces: `McpTokenStoreError(RuntimeError)` + +- [ ] **Step 1: Write failing tests for `McpTokenStore`** + +Create `apps/device-host-agent/tests/test_mcp_token.py`: + +```python +from __future__ import annotations + +import json +import os +import stat +import sys +from datetime import datetime, timezone +from pathlib import Path + +import pytest + +from host_agent.mcp_token import McpToken, McpTokenStore, McpTokenStoreError + + +def test_load_or_create_generates_when_missing(tmp_path: Path) -> None: + store = McpTokenStore(tmp_path / "host_mcp_token.json") + token = store.load_or_create() + assert token.version == 1 + assert len(token.token) >= 40 # secrets.token_urlsafe(32) -> ~43 chars + assert isinstance(token.created_at, datetime) + # File now exists. + assert (tmp_path / "host_mcp_token.json").exists() + + +def test_load_or_create_is_idempotent(tmp_path: Path) -> None: + store = McpTokenStore(tmp_path / "host_mcp_token.json") + first = store.load_or_create() + second = McpTokenStore(tmp_path / "host_mcp_token.json").load_or_create() + assert first.token == second.token + + +def test_load_or_create_writes_json_schema(tmp_path: Path) -> None: + path = tmp_path / "host_mcp_token.json" + McpTokenStore(path).load_or_create() + data = json.loads(path.read_text()) + assert set(data) == {"version", "token", "created_at"} + assert data["version"] == 1 + assert isinstance(data["token"], str) + # created_at is ISO 8601. + datetime.fromisoformat(data["created_at"]) + + +@pytest.mark.skipif(sys.platform == "win32", reason="POSIX perms only") +def test_load_or_create_sets_posix_permissions(tmp_path: Path) -> None: + path = tmp_path / "host_mcp_token.json" + McpTokenStore(path).load_or_create() + mode = stat.S_IMODE(os.fstat(os.open(path, os.O_RDONLY)).st_mode) + assert mode == 0o600 + + +def test_verify_accepts_correct_token(tmp_path: Path) -> None: + store = McpTokenStore(tmp_path / "host_mcp_token.json") + token = store.load_or_create() + assert store.verify(token.token) is True + + +def test_verify_rejects_wrong_token(tmp_path: Path) -> None: + store = McpTokenStore(tmp_path / "host_mcp_token.json") + store.load_or_create() + assert store.verify("wrong") is False + + +def test_load_or_create_raises_on_corrupt_json(tmp_path: Path) -> None: + path = tmp_path / "host_mcp_token.json" + path.write_text("{not valid json") + with pytest.raises(McpTokenStoreError): + McpTokenStore(path).load_or_create() + + +def test_load_or_create_raises_on_unwritable_dir(tmp_path: Path) -> None: + unwritable = tmp_path / "ro" + unwritable.mkdir() + os.chmod(unwritable, 0o500) # r-x for owner + try: + with pytest.raises(McpTokenStoreError): + McpTokenStore(unwritable / "host_mcp_token.json").load_or_create() + finally: + os.chmod(unwritable, 0o700) # restore so cleanup works + + +def test_load_or_create_concurrent_calls_do_not_corrupt( + tmp_path: Path, +) -> None: + """Two store instances racing to create: both end up reading the same token.""" + import threading + + path = tmp_path / "host_mcp_token.json" + results: list[McpToken] = [] + barrier = threading.Barrier(2) + + def worker() -> None: + barrier.wait() + store = McpTokenStore(path) + results.append(store.load_or_create()) + + threads = [threading.Thread(target=worker) for _ in range(2)] + for t in threads: + t.start() + for t in threads: + t.join() + assert len(results) == 2 + assert results[0].token == results[1].token +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_token.py -v` +Expected: FAIL — `host_agent.mcp_token` doesn't exist. + +- [ ] **Step 3: Implement `McpTokenStore`** + +Create `apps/device-host-agent/host_agent/mcp_token.py`: + +```python +"""Bearer-token persistence for the host-agent MCP server. + +The token is generated on first start and persisted to a JSON file with +0o600 permissions (POSIX) alongside the host identity. Rotation = delete +the file and restart host-agent. +""" + +from __future__ import annotations + +import json +import os +import secrets +import tempfile +from dataclasses import dataclass +from datetime import UTC, datetime +from pathlib import Path +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from collections.abc import Callable + + +_TOKEN_BYTES = 32 + + +class McpTokenStoreError(RuntimeError): + """Raised when the MCP token file cannot be read or written.""" + + +@dataclass(frozen=True) +class McpToken: + version: int + token: str + created_at: datetime + + +class McpTokenStore: + def __init__( + self, + path: Path, + *, + now: Callable[[], datetime] | None = None, + ) -> None: + self._path = Path(path) + self._now = now or (lambda: datetime.now(UTC)) + + def load_or_create(self) -> McpToken: + if self._path.exists(): + return self._read_existing() + return self._generate_and_write() + + def verify(self, presented: str) -> bool: + try: + token = self.load_or_create() + except McpTokenStoreError: + return False + import hmac + + return hmac.compare_digest(token.token, presented) + + def _read_existing(self) -> McpToken: + try: + data = json.loads(self._path.read_text()) + except (OSError, json.JSONDecodeError) as exc: + raise McpTokenStoreError( + f"cannot read MCP token file {self._path}: {exc}" + ) from exc + if not isinstance(data, dict): + raise McpTokenStoreError("MCP token file is not a JSON object") + try: + return McpToken( + version=int(data["version"]), + token=str(data["token"]), + created_at=datetime.fromisoformat(str(data["created_at"])), + ) + except (KeyError, TypeError, ValueError) as exc: + raise McpTokenStoreError( + f"MCP token file schema invalid: {exc}" + ) from exc + + def _generate_and_write(self) -> McpToken: + token = McpToken( + version=1, + token=secrets.token_urlsafe(_TOKEN_BYTES), + created_at=self._now(), + ) + payload = { + "version": token.version, + "token": token.token, + "created_at": token.created_at.isoformat(), + } + try: + self._atomic_write(json.dumps(payload, indent=2)) + except OSError as exc: + raise McpTokenStoreError( + f"cannot write MCP token file {self._path}: {exc}" + ) from exc + return token + + def _atomic_write(self, content: str) -> None: + self._path.parent.mkdir(parents=True, exist_ok=True) + # Atomic on POSIX; on Windows os.replace is also atomic per docs. + fd, tmp_name = tempfile.mkstemp( + prefix=".host_mcp_token.", + suffix=".tmp", + dir=str(self._path.parent), + ) + try: + with os.fdopen(fd, "w", encoding="utf-8") as fh: + fh.write(content) + os.chmod(tmp_name, 0o600) + os.replace(tmp_name, self._path) + except BaseException: + try: + os.unlink(tmp_name) + except OSError: + pass + raise +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_token.py -v` +Expected: all 8 tests PASS (Windows skips the POSIX perm test). + +- [ ] **Step 5: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/mcp_token.py apps/device-host-agent/tests/test_mcp_token.py` +Expected: clean. + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/host_agent/mcp_token.py apps/device-host-agent/tests/test_mcp_token.py +git commit -m "feat(host-agent): add McpTokenStore for MCP bearer token" +``` + +--- + +## Task 4: `McpBusyTracker` — per-device session-level lock with TTL + +**Files:** +- Create: `apps/device-host-agent/host_agent/mcp_lock.py` +- Test: `apps/device-host-agent/tests/test_mcp_lock.py` + +**Interfaces:** +- Produces: `McpBusyTracker(*, ttl_seconds: float = 60.0, now: Callable[[], datetime] | None = None)` with methods: + - `acquire(device_id: str, session_id: str) -> bool` + - `renew(device_id: str, session_id: str) -> bool` + - `release(session_id: str) -> list[str]` — returns freed device_ids + - `release_device(device_id: str, session_id: str) -> bool` + - `busy_device_ids() -> list[str]` — lazy-sweep expired leases + - `snapshot() -> list[McpDeviceLease]` + - `wait_until_usable(device_id, session_id, *, timeout, poll_interval=1.0, cloud_busy_check=None) -> bool` — atomic acquire on success +- Produces: `McpDeviceLease(device_id: str, session_id: str, acquired_at: datetime, last_seen_at: datetime)` (frozen) + +- [ ] **Step 1: Write failing tests** + +Create `apps/device-host-agent/tests/test_mcp_lock.py`: + +```python +from __future__ import annotations + +import threading +from datetime import UTC, datetime, timedelta +from typing import Any + +import pytest + +from host_agent.mcp_lock import McpBusyTracker, McpDeviceLease + + +def _tracker_with_now() -> tuple[McpBusyTracker, list[datetime]]: + times: list[datetime] = [] + + def now() -> datetime: + return times[-1] if times else datetime(2026, 1, 1, tzinfo=UTC) + + tracker = McpBusyTracker(ttl_seconds=60.0, now=now) + return tracker, times + + +def test_acquire_succeeds_on_empty() -> None: + tracker, _ = _tracker_with_now() + assert tracker.acquire("phone-1", "sess-a") is True + assert "phone-1" in tracker.busy_device_ids() + + +def test_acquire_fails_when_held_by_other_session() -> None: + tracker, _ = _tracker_with_now() + assert tracker.acquire("phone-1", "sess-a") is True + assert tracker.acquire("phone-1", "sess-b") is False + + +def test_acquire_is_idempotent_for_same_session() -> None: + tracker, _ = _tracker_with_now() + assert tracker.acquire("phone-1", "sess-a") is True + # Same session re-acquiring is allowed (acts as renew). + assert tracker.acquire("phone-1", "sess-a") is True + + +def test_renew_refreshes_last_seen() -> None: + tracker, times = _tracker_with_now() + times.append(datetime(2026, 1, 1, 12, 0, tzinfo=UTC)) + tracker.acquire("phone-1", "sess-a") + initial = tracker.snapshot()[0] + times.append(datetime(2026, 1, 1, 12, 0, 30, tzinfo=UTC)) + assert tracker.renew("phone-1", "sess-a") is True + refreshed = tracker.snapshot()[0] + assert refreshed.last_seen_at > initial.last_seen_at + + +def test_renew_fails_when_held_by_other() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + assert tracker.renew("phone-1", "sess-b") is False + + +def test_release_returns_freed_device_ids() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + tracker.acquire("phone-2", "sess-a") + freed = tracker.release("sess-a") + assert sorted(freed) == ["phone-1", "phone-2"] + assert tracker.busy_device_ids() == [] + + +def test_release_only_frees_caller_session() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + tracker.acquire("phone-1", "sess-b") # fails + freed = tracker.release("sess-b") + assert freed == [] + assert "phone-1" in tracker.busy_device_ids() + + +def test_ttl_sweeps_expired_leases() -> None: + tracker, times = _tracker_with_now() + times.append(datetime(2026, 1, 1, 12, 0, tzinfo=UTC)) + tracker.acquire("phone-1", "sess-a") + # Advance past TTL without renew. + times.append(datetime(2026, 1, 1, 12, 1, 1, tzinfo=UTC)) # 61s later + assert tracker.busy_device_ids() == [] + + +def test_renew_after_ttl_tolerates_same_session() -> None: + """Scene 10: lease expired but session_id matches -> re-acquire.""" + tracker, times = _tracker_with_now() + times.append(datetime(2026, 1, 1, 12, 0, tzinfo=UTC)) + tracker.acquire("phone-1", "sess-a") + times.append(datetime(2026, 1, 1, 12, 1, 1, tzinfo=UTC)) # expired + # renew from the same session should succeed (re-acquire). + assert tracker.renew("phone-1", "sess-a") is True + assert "phone-1" in tracker.busy_device_ids() + + +def test_snapshot_matches_busy_device_ids() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + tracker.acquire("phone-2", "sess-a") + snap = tracker.snapshot() + assert {lease.device_id for lease in snap} == set(tracker.busy_device_ids()) + + +def test_wait_until_usable_succeeds_when_free() -> None: + tracker, _ = _tracker_with_now() + ok = tracker.wait_until_usable( + "phone-1", "sess-a", timeout=1.0, poll_interval=0.01 + ) + assert ok is True + assert "phone-1" in tracker.busy_device_ids() + + +def test_wait_until_usable_returns_false_on_timeout() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + ok = tracker.wait_until_usable( + "phone-1", "sess-b", timeout=0.1, poll_interval=0.02 + ) + assert ok is False + + +def test_wait_until_usable_blocks_then_succeeds_when_released() -> None: + tracker, _ = _tracker_with_now() + tracker.acquire("phone-1", "sess-a") + + def releaser() -> None: + import time + + time.sleep(0.05) + tracker.release("sess-a") + + t = threading.Thread(target=releaser) + t.start() + try: + ok = tracker.wait_until_usable( + "phone-1", "sess-b", timeout=2.0, poll_interval=0.02 + ) + assert ok is True + finally: + t.join() + + +def test_wait_until_usable_blocks_then_fails_when_cloud_remains_busy() -> None: + tracker, _ = _tracker_with_now() + ok = tracker.wait_until_usable( + "phone-1", + "sess-a", + timeout=0.1, + poll_interval=0.02, + cloud_busy_check=lambda: True, + ) + assert ok is False + assert tracker.busy_device_ids() == [] +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_lock.py -v` +Expected: FAIL — `host_agent.mcp_lock` doesn't exist. + +- [ ] **Step 3: Implement `McpBusyTracker`** + +Create `apps/device-host-agent/host_agent/mcp_lock.py`: + +```python +"""Per-device MCP session-level busy tracker. + +The cloud-side assignment path and the MCP-driven path both drive devices +through the same in-process ``DeviceManager``. This tracker records which +devices are currently held by an MCP session so that: + +- MCP tool calls against a device held by another session (or by a cloud + assignment — checked separately by the caller via ``AgentStatusTracker``) + can fail fast with a busy error. +- The heartbeat payload can advertise ``mcp_busy_device_ids`` so the cloud + scheduler won't dispatch conflicting assignments to the same device. + +Leases expire ``ttl_seconds`` after the last ``renew()`` call (set on every +tool call from the holding session). Expired leases are lazy-swept on read. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass +from datetime import UTC, datetime +from threading import Lock +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from collections.abc import Callable + + +@dataclass(frozen=True) +class McpDeviceLease: + device_id: str + session_id: str + acquired_at: datetime + last_seen_at: datetime + + +class McpBusyTracker: + def __init__( + self, + *, + ttl_seconds: float = 60.0, + now: Callable[[], datetime] | None = None, + ) -> None: + self._ttl = float(ttl_seconds) + self._now = now or (lambda: datetime.now(UTC)) + self._lock = Lock() + # device_id -> McpDeviceLease + self._leases: dict[str, McpDeviceLease] = {} + + def acquire(self, device_id: str, session_id: str) -> bool: + with self._lock: + self._sweep_locked() + existing = self._leases.get(device_id) + if existing is not None and existing.session_id != session_id: + return False + now = self._now() + lease = McpDeviceLease( + device_id=device_id, + session_id=session_id, + acquired_at=( + existing.acquired_at if existing is not None else now + ), + last_seen_at=now, + ) + self._leases[device_id] = lease + return True + + def renew(self, device_id: str, session_id: str) -> bool: + with self._lock: + self._sweep_locked() + existing = self._leases.get(device_id) + # Tolerate boundary: lease may have been swept, but if the caller + # is the legitimate previous holder, re-acquire on their behalf. + if existing is None: + now = self._now() + self._leases[device_id] = McpDeviceLease( + device_id=device_id, + session_id=session_id, + acquired_at=now, + last_seen_at=now, + ) + return True + if existing.session_id != session_id: + return False + self._leases[device_id] = McpDeviceLease( + device_id=device_id, + session_id=session_id, + acquired_at=existing.acquired_at, + last_seen_at=self._now(), + ) + return True + + def release(self, session_id: str) -> list[str]: + with self._lock: + freed = [ + device_id + for device_id, lease in self._leases.items() + if lease.session_id == session_id + ] + for device_id in freed: + del self._leases[device_id] + return freed + + def release_device(self, device_id: str, session_id: str) -> bool: + with self._lock: + existing = self._leases.get(device_id) + if existing is None or existing.session_id != session_id: + return False + del self._leases[device_id] + return True + + def busy_device_ids(self) -> list[str]: + with self._lock: + self._sweep_locked() + return sorted(self._leases) + + def snapshot(self) -> list[McpDeviceLease]: + with self._lock: + self._sweep_locked() + return sorted(self._leases.values(), key=lambda l: l.device_id) + + def wait_until_usable( + self, + device_id: str, + session_id: str, + *, + timeout: float, + poll_interval: float = 1.0, + cloud_busy_check: Callable[[], bool] | None = None, + ) -> bool: + """Block until ``device_id`` is acquirable by ``session_id`` or timeout. + + Reserved capability. MVP callers use try-acquire (``acquire`` -> False + means busy). This method exists for future wiring where the cloud + assignment path or an explicit MCP tool may opt to wait. + """ + deadline = time.monotonic() + timeout + while True: + cloud_busy = cloud_busy_check() if cloud_busy_check else False + if not cloud_busy: + if self.acquire(device_id, session_id): + return True + if time.monotonic() >= deadline: + return False + remaining = deadline - time.monotonic() + time.sleep(max(0.0, min(poll_interval, remaining))) + + def _sweep_locked(self) -> None: + """Caller holds ``self._lock``. Drops leases past their TTL.""" + cutoff = self._now() + expired = [ + device_id + for device_id, lease in self._leases.items() + if (cutoff - lease.last_seen_at).total_seconds() > self._ttl + ] + for device_id in expired: + del self._leases[device_id] +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_lock.py -v` +Expected: all 14 tests PASS. + +- [ ] **Step 5: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/mcp_lock.py apps/device-host-agent/tests/test_mcp_lock.py` +Expected: clean. + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/host_agent/mcp_lock.py apps/device-host-agent/tests/test_mcp_lock.py +git commit -m "feat(host-agent): add McpBusyTracker for per-device session locks" +``` + +--- + +## Task 5: `BearerAuthMiddleware` + +**Files:** +- Create: `apps/device-host-agent/host_agent/web/mcp_auth.py` +- Test: `apps/device-host-agent/tests/test_mcp_auth.py` + +**Interfaces:** +- Produces: `BearerAuthMiddleware(app: ASGIApp, token_store: McpTokenStore)` — Starlette `BaseHTTPMiddleware` subclass. +- Behavior: 401 + `WWW-Authenticate: Bearer` + JSON `{"error": "invalid token"}` on missing/wrong token. + +- [ ] **Step 1: Write failing tests** + +Create `apps/device-host-agent/tests/test_mcp_auth.py`: + +```python +from __future__ import annotations + +from pathlib import Path + +from starlette.applications import Starlette +from starlette.responses import JSONResponse, Response +from starlette.testclient import TestClient + +from host_agent.mcp_token import McpTokenStore +from host_agent.web.mcp_auth import BearerAuthMiddleware + + +def _make_client(tmp_path: Path) -> tuple[TestClient, str]: + store = McpTokenStore(tmp_path / "host_mcp_token.json") + token = store.load_or_create().token + + async def hello(request): # type: ignore[no-untyped-def] + return JSONResponse({"ok": True}) + + inner = Starlette(routes=[]) + inner.router.add_route("/", hello, methods=["GET"]) + wrapped = Starlette() + wrapped.router.add_middleware(BearerAuthMiddleware, token_store=store) + wrapped.mount("/", inner) + return TestClient(wrapped), token + + +def test_no_header_returns_401(tmp_path: Path) -> None: + client, _ = _make_client(tmp_path) + resp = client.get("/") + assert resp.status_code == 401 + assert resp.headers["WWW-Authenticate"] == "Bearer" + assert resp.json() == {"error": "invalid token"} + + +def test_wrong_token_returns_401(tmp_path: Path) -> None: + client, _ = _make_client(tmp_path) + resp = client.get("/", headers={"Authorization": "Bearer wrong"}) + assert resp.status_code == 401 + + +def test_correct_token_passes_through(tmp_path: Path) -> None: + client, token = _make_client(tmp_path) + resp = client.get("/", headers={"Authorization": f"Bearer {token}"}) + assert resp.status_code == 200 + assert resp.json() == {"ok": True} + + +def test_non_bearer_scheme_returns_401(tmp_path: Path) -> None: + client, token = _make_client(tmp_path) + resp = client.get("/", headers={"Authorization": f"Basic {token}"}) + assert resp.status_code == 401 + + +def test_header_case_insensitive(tmp_path: Path) -> None: + client, token = _make_client(tmp_path) + resp = client.get("/", headers={"authorization": f"Bearer {token}"}) + assert resp.status_code == 200 +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_auth.py -v` +Expected: FAIL — `host_agent.web.mcp_auth` doesn't exist. + +- [ ] **Step 3: Implement `BearerAuthMiddleware`** + +Create `apps/device-host-agent/host_agent/web/mcp_auth.py`: + +```python +"""Bearer-token auth middleware for the MCP sub-app. + +Mounted on the FastMCP ``streamable_http_app()`` (NOT the console FastAPI), +so cookie-session auth on console routes is unaffected. +""" + +from __future__ import annotations + +from starlette.middleware.base import BaseHTTPMiddleware +from starlette.requests import Request +from starlette.responses import JSONResponse, Response + +from host_agent.mcp_token import McpTokenStore + + +class BearerAuthMiddleware(BaseHTTPMiddleware): + def __init__(self, app, token_store: McpTokenStore) -> None: + super().__init__(app) + self._store = token_store + + async def dispatch(self, request: Request, call_next) -> Response: # type: ignore[no-untyped-def] + header = request.headers.get("Authorization") + if not header or not header.lower().startswith("bearer "): + return _unauthorized() + presented = header.split(" ", 1)[1].strip() + if not self._store.verify(presented): + return _unauthorized() + return await call_next(request) + + +def _unauthorized() -> JSONResponse: + return JSONResponse( + status_code=401, + content={"error": "invalid token"}, + headers={"WWW-Authenticate": "Bearer"}, + ) +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_mcp_auth.py -v` +Expected: all 5 tests PASS. + +- [ ] **Step 5: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/web/mcp_auth.py apps/device-host-agent/tests/test_mcp_auth.py` +Expected: clean. + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/host_agent/web/mcp_auth.py apps/device-host-agent/tests/test_mcp_auth.py +git commit -m "feat(host-agent): add BearerAuthMiddleware for MCP server" +``` + +--- + +## Task 6: `build_mcp_server` — wrapped FastMCP server + +The heart of the integration. Wraps each of the 11 tools from `api.mcp.tool_handlers(manager=...)` with: +1. session_id extraction from FastMCP context +2. busy check (cloud via `AgentStatusTracker`, MCP via `McpBusyTracker`) +3. lazy acquire on first call against a device +4. renew on every call +5. status mapping for `list_devices` / `device_status` (display status) + +**Files:** +- Create: `apps/device-host-agent/host_agent/web/mcp.py` +- Test: `apps/device-host-agent/tests/test_web_mcp.py` + +**Interfaces:** +- Produces: `build_mcp_server(*, manager: DeviceManager, mcp_busy_tracker: McpBusyTracker, status_tracker: AgentStatusTracker) -> FastMCP` +- Produces: `McpDeviceBusyError(Exception)` with `.device_id`, `.busy_owner` attributes for callers that want structured info. +- Produces: helpers `_cloud_busy_device_id(status_tracker: AgentStatusTracker) -> str | None` and `_display_status(device, busy_device_id) -> str` (extracted from `host_agent/web/app.py:64` for reuse; do NOT remove the original from `app.py` yet — Task 11 will consolidate). + +**Session ID extraction note:** The mcp SDK 1.28.1 exposes session_id via `mcp.server.fastmcp.Context`. The exact attribute path should be verified by inspecting the installed SDK at runtime; a fallback contextvars-based middleware can be added if `Context.session_id` isn't reliable. + +- [ ] **Step 1: Verify session_id extraction from mcp SDK** + +Run: `uv run --package device-host-agent python -c " +import inspect +from mcp.server.fastmcp import Context +print([a for a in dir(Context) if 'session' in a.lower() or 'id' in a.lower()]) +"` +Expected output: shows attributes including something like `session_id`, `request_id`, or similar. Record the attribute name for use in Step 3. + +If `session_id` (or equivalent) is NOT accessible via `Context`, fall back to a request-scoped `contextvars.ContextVar[str]` populated by a Starlette middleware on the mounted MCP app that reads `request.scope["mcp_session_id"]` (the mcp SDK sets this in ASGI scope). + +- [ ] **Step 2: Write failing tests for `build_mcp_server`** + +Create `apps/device-host-agent/tests/test_web_mcp.py`: + +```python +from __future__ import annotations + +from typing import Any + +import pytest + +from device.manager import DeviceManager +from driver.base import Driver +from host_agent.mcp_lock import McpBusyTracker +from host_agent.status import AgentStatusTracker +from host_agent.web.mcp import ( + McpDeviceBusyError, + build_mcp_server, + _call_tool_sync, +) + + +class _FakeDriver(Driver): + """Minimal driver: connect/screenshot/tap are the only methods exercised.""" + + def __init__(self) -> None: + self.taps: list[tuple[float, float]] = [] + + def connect(self) -> None: + return None + + def disconnect(self) -> None: + return None + + def screenshot(self) -> bytes: + return b"fake" + + def tap(self, x: float, y: float) -> None: + self.taps.append((x, y)) + + # All other Driver methods left as no-ops/defaults — add stubs as needed + # for type-checking. (long_press, swipe, swipe_path, double_tap, input, + # launch, tree, home, get_active_app_metadata, press_keycode, etc.) + # Use `def __getattr__(self, name): return lambda *a, **k: None` if + # needed during early development, then replace with explicit stubs. + + +def _make_manager_with_device(device_id: str = "phone-1") -> DeviceManager: + manager = DeviceManager() + manager.register( + device_id=device_id, + driver_type="wda", + factory=lambda: _FakeDriver(), + name=device_id, + ) + manager.connect(device_id) + return manager + + +def test_build_mcp_server_returns_fastmcp_instance() -> None: + from mcp.server.fastmcp import FastMCP + + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + assert isinstance(server, FastMCP) + + +def test_call_tool_succeeds_when_device_is_free() -> None: + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + result = _call_tool_sync( + server, "take_screenshot", {"device_id": "phone-1"}, session_id="sess-a" + ) + assert result["ok"] is True + assert "phone-1" in tracker.busy_device_ids() + + +def test_call_tool_fails_when_cloud_uses_device() -> None: + """AgentStatusTracker.current_assignment.device_id matches -> busy.""" + from cloud.internal_api.models import AssignmentModel + from datetime import datetime, UTC + + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + status.mark_assignment_started( + AssignmentModel( + task_id="t1", + attempt=1, + lease_id="l1", + lease_expires_at=datetime.now(UTC), + host_id="h1", + device_id="phone-1", + goal="cloud task", + ) + ) + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + with pytest.raises(McpDeviceBusyError) as exc: + _call_tool_sync( + server, + "take_screenshot", + {"device_id": "phone-1"}, + session_id="sess-a", + ) + assert exc.value.device_id == "phone-1" + assert exc.value.busy_owner == "cloud_assignment" + + +def test_call_tool_fails_when_another_mcp_session_holds_device() -> None: + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + # Pre-acquire as a different session. + tracker.acquire("phone-1", "sess-other") + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + with pytest.raises(McpDeviceBusyError) as exc: + _call_tool_sync( + server, + "take_screenshot", + {"device_id": "phone-1"}, + session_id="sess-a", + ) + assert exc.value.busy_owner.startswith("mcp_session:") + + +def test_call_tool_renews_when_same_session_already_holds() -> None: + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + _call_tool_sync( + server, "take_screenshot", {"device_id": "phone-1"}, session_id="sess-a" + ) + # Second call from the same session should succeed. + result = _call_tool_sync( + server, "take_screenshot", {"device_id": "phone-1"}, session_id="sess-a" + ) + assert result["ok"] is True + + +def test_list_devices_uses_display_status() -> None: + """Connected-but-idle devices report as 'connected', not 'busy'.""" + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + result = _call_tool_sync(server, "list_devices", {}, session_id="sess-a") + assert isinstance(result, list) + assert result[0]["status"] == "connected" + + +def test_unknown_device_returns_value_error() -> None: + manager = _make_manager_with_device() + tracker = McpBusyTracker() + status = AgentStatusTracker() + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + with pytest.raises(Exception) as exc: # api.errors wraps as semantic error + _call_tool_sync( + server, + "take_screenshot", + {"device_id": "does-not-exist"}, + session_id="sess-a", + ) + # The error should mention the device somehow. + assert "does-not-exist" in str(exc.value) + + +def test_manager_required_for_tool_handlers() -> None: + """Reinforces D12 — build_mcp_server itself requires non-None manager.""" + # build_mcp_server signature already requires manager as keyword-only. + # This is a documentation test. + import inspect + + sig = inspect.signature(build_mcp_server) + assert sig.parameters["manager"].kind == inspect.Parameter.KEYWORD_ONLY +``` + +- [ ] **Step 3: Implement `host_agent/web/mcp.py`** + +Create `apps/device-host-agent/host_agent/web/mcp.py`: + +```python +"""FastMCP server builder for the host-agent MCP endpoint. + +Wraps ``api.mcp.tool_handlers(manager=...)`` with: +- Cloud-busy and MCP-busy checks (per-device, fail-fast on conflict) +- Lazy session-level device lock acquire / renew +- Display-status mapping for list_devices / device_status so connected-but- + idle devices don't appear "busy" (which they do at the DeviceManager + layer because an Appium/WDA session is open). + +The builder returns a ``FastMCP`` instance. The caller (``create_console_app``) +is responsible for wrapping it in ``BearerAuthMiddleware`` and mounting at +``/mcp``. +""" + +from __future__ import annotations + +from collections.abc import Callable +from typing import Any + +from device.manager import DeviceManager +from host_agent.mcp_lock import McpBusyTracker +from host_agent.status import AgentStatusTracker +from mcp.server.fastmcp import FastMCP + + +# Tool names that don't target a specific device — skip busy check. +_NON_DEVICE_TOOLS = frozenset({"list_devices", "device_status"}) +# Tools that report status and should use the display-status mapping. +_STATUS_TOOLS = frozenset({"list_devices", "device_status"}) + + +class McpDeviceBusyError(Exception): + """Raised by the wrapper when the target device is held by cloud or + another MCP session.""" + + def __init__(self, device_id: str, busy_owner: str) -> None: + super().__init__(f"device {device_id} is busy (held by {busy_owner})") + self.device_id = device_id + self.busy_owner = busy_owner + + +def build_mcp_server( + *, + manager: DeviceManager, + mcp_busy_tracker: McpBusyTracker, + status_tracker: AgentStatusTracker, +) -> FastMCP: + """Construct the FastMCP server wrapping ``tool_handlers``.""" + # Imported lazily to keep package import graph flat. + from api.mcp import tool_handlers + + handlers = tool_handlers(manager=manager) + server = FastMCP("apex-host-agent") + + for tool_name, raw_handler in handlers.items(): + wrapped = _wrap_tool( + tool_name, + raw_handler, + mcp_busy_tracker=mcp_busy_tracker, + status_tracker=status_tracker, + ) + + def _make_tool(name: str, fn: Callable[..., Any]) -> None: + @server.tool(name=name) + def _tool(*args: Any, **kwargs: Any) -> Any: # noqa: ANN202 + return fn(*args, **kwargs) + + _make_tool(tool_name, wrapped) + + return server + + +def _wrap_tool( + tool_name: str, + handler: Callable[..., Any], + *, + mcp_busy_tracker: McpBusyTracker, + status_tracker: AgentStatusTracker, +) -> Callable[..., Any]: + def wrapped(*args: Any, **kwargs: Any) -> Any: + session_id = _current_session_id() + device_id = kwargs.get("device_id") + + if tool_name in _STATUS_TOOLS: + return _with_display_status(handler, status_tracker, *args, **kwargs) + + if device_id is not None and tool_name not in _NON_DEVICE_TOOLS: + _check_and_acquire(device_id, session_id, mcp_busy_tracker, status_tracker) + + return handler(*args, **kwargs) + + return wrapped + + +def _check_and_acquire( + device_id: str, + session_id: str, + mcp_busy_tracker: McpBusyTracker, + status_tracker: AgentStatusTracker, +) -> None: + cloud_busy = _cloud_busy_device_id(status_tracker) + if cloud_busy == device_id: + raise McpDeviceBusyError(device_id, "cloud_assignment") + if device_id in mcp_busy_tracker.busy_device_ids(): + existing = next( + ( + lease + for lease in mcp_busy_tracker.snapshot() + if lease.device_id == device_id + ), + None, + ) + if existing is not None and existing.session_id != session_id: + prefix = existing.session_id[:8] + raise McpDeviceBusyError(device_id, f"mcp_session:{prefix}") + if not mcp_busy_tracker.acquire(device_id, session_id): + # Race: someone else got it between check and acquire. + raise McpDeviceBusyError(device_id, "another_session") + mcp_busy_tracker.renew(device_id, session_id) + + +def _cloud_busy_device_id(status_tracker: AgentStatusTracker) -> str | None: + snap = status_tracker.snapshot() + current = snap.get("current_assignment") + if not isinstance(current, dict): + return None + device_id = current.get("device_id") + return device_id if isinstance(device_id, str) else None + + +def _with_display_status( + handler: Callable[..., Any], + status_tracker: AgentStatusTracker, + *args: Any, + **kwargs: Any, +) -> Any: + busy_device_id = _cloud_busy_device_id(status_tracker) + result = handler(*args, **kwargs) + if isinstance(result, list): + for item in result: + if isinstance(item, dict) and "status" in item: + item["status"] = _display_status(item["status"], item.get("id"), busy_device_id) + return result + if isinstance(result, dict) and "status" in result: + result["status"] = _display_status( + result["status"], result.get("device_id"), busy_device_id + ) + return result + + +def _display_status(raw: str, device_id: Any, busy_device_id: str | None) -> str: + """Mirror host_agent.web.app._device_display_status semantics: + a device that's locally 'busy' because it's connected-but-idle reports + 'connected', unless it's the device currently running a cloud assignment. + """ + if raw == "busy" and device_id != busy_device_id: + return "connected" + return raw + + +def _current_session_id() -> str: + """Extract session_id from the current FastMCP tool-call context. + + The mcp SDK exposes session_id via Context; for sync handlers invoked + outside a request lifecycle (e.g. tests), fall back to a contextvars + value set by the test harness. + """ + # Try the FastMCP context first. + try: + from mcp.server.fastmcp import get_context + + ctx = get_context() + # SDK 1.28 exposes session_id via the underlying server transport. + session_id = getattr(ctx, "session_id", None) + if isinstance(session_id, str) and session_id: + return session_id + request_id = getattr(ctx, "request_id", None) + if isinstance(request_id, str) and request_id: + return request_id + except Exception: + pass + # Test fallback. + return _TEST_SESSION_ID.get("") + + +# contextvars fallback for tests; production code uses FastMCP context. +import contextvars + +_TEST_SESSION_ID: contextvars.ContextVar[str] = contextvars.ContextVar( + "_TEST_SESSION_ID", default="" +) + + +def _call_tool_sync( + server: FastMCP, + tool_name: str, + arguments: dict[str, Any], + *, + session_id: str, +) -> Any: + """Test helper: invoke a registered tool synchronously with a forced + session_id. Bypasses the HTTP/MCP transport layer to keep tests fast.""" + token = _TEST_SESSION_ID.set(session_id) + try: + # Walk FastMCP's tool registry to find the underlying callable. + tool = server._tool_manager.get_tool(tool_name) # type: ignore[attr-defined] + if tool is None: + raise KeyError(f"tool {tool_name!r} not registered") + # FastMCP Tool wraps a coroutine; our wrappers are sync, so unwrap. + return tool.fn(**arguments) + finally: + _TEST_SESSION_ID.reset(token) +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_web_mcp.py -v` +Expected: 8 tests PASS. + +If `_tool_manager.get_tool(tool_name).fn` access pattern doesn't match mcp 1.28.1 internals, replace with introspection: + +```python +# Fallback for different FastMCP internal layouts: +tools = getattr(server, "_tool_manager", None) +if tools is not None: + registry = getattr(tools, "_tools", None) or getattr(tools, "tools", None) + if isinstance(registry, dict): + tool = registry.get(tool_name) +``` + +- [ ] **Step 5: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/web/mcp.py apps/device-host-agent/tests/test_web_mcp.py` +Expected: clean (or only minor style notes — fix what's flagged). + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/host_agent/web/mcp.py apps/device-host-agent/tests/test_web_mcp.py +git commit -m "feat(host-agent): wrap tool_handlers with busy check + status mapping" +``` + +--- + +## Task 7: Cloud `HeartbeatRequest` schema — add `mcp_busy_device_ids` + +**Files:** +- Modify: `packages/cloud-platform/cloud/internal_api/models.py:37-42` +- Test: `packages/cloud-platform/tests/test_internal_api_models.py` (create or extend) + +**Interfaces:** +- Produces: `HeartbeatRequest` now has `mcp_busy_device_ids: list[str] = Field(default_factory=list)`. +- Backward compat: requests without the field still validate (default `[]`). + +- [ ] **Step 1: Locate or create the test file** + +Run: `ls packages/cloud-platform/tests/` +If `test_internal_api_models.py` exists, extend it; otherwise create it. + +- [ ] **Step 2: Write failing tests** + +Create or extend `packages/cloud-platform/tests/test_internal_api_models.py`: + +```python +from __future__ import annotations + +from cloud.internal_api.models import HeartbeatRequest + + +def test_heartbeat_request_defaults_mcp_busy_device_ids_to_empty() -> None: + req = HeartbeatRequest(host_id="h1") + assert req.mcp_busy_device_ids == [] + + +def test_heartbeat_request_accepts_mcp_busy_device_ids() -> None: + req = HeartbeatRequest(host_id="h1", mcp_busy_device_ids=["phone-1"]) + assert req.mcp_busy_device_ids == ["phone-1"] + + +def test_heartbeat_request_omitting_field_is_backward_compatible() -> None: + """Old host-agents that don't send the field must still validate.""" + raw = {"host_id": "h1", "devices": []} + req = HeartbeatRequest.model_validate(raw) + assert req.mcp_busy_device_ids == [] +``` + +- [ ] **Step 3: Run tests to verify they fail** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_internal_api_models.py -v` +Expected: FAIL — `mcp_busy_device_ids` not yet on the model. + +- [ ] **Step 4: Add the field** + +In `packages/cloud-platform/cloud/internal_api/models.py`, modify the `HeartbeatRequest` class (around lines 37-42): + +```python +class HeartbeatRequest(BaseModel): + host_id: str = Field(min_length=1) + address: str | None = None + devices: list[DeviceSnapshotModel] = Field(default_factory=list) + policy_revision: int = Field(default=0, ge=0) + planner_transport: Literal["direct", "cloud"] = "direct" + mcp_busy_device_ids: list[str] = Field(default_factory=list) +``` + +- [ ] **Step 5: Run tests to verify they pass** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_internal_api_models.py -v` +Expected: all 3 PASS. + +- [ ] **Step 6: Lint check** + +Run: `uv run --with ruff ruff check packages/cloud-platform/cloud/internal_api/models.py packages/cloud-platform/tests/test_internal_api_models.py` +Expected: clean. + +- [ ] **Step 7: Commit** + +```bash +git add packages/cloud-platform/cloud/internal_api/models.py packages/cloud-platform/tests/test_internal_api_models.py +git commit -m "feat(cloud): accept mcp_busy_device_ids in heartbeat payload" +``` + +--- + +## Task 8: Cloud scheduler skips MCP-busy devices + +**Files:** +- Modify: `packages/cloud-platform/cloud/pool.py:40-49` (`PooledDevice` gets `mcp_busy: bool`) +- Modify: `packages/cloud-platform/cloud/pool.py:59-89` (`sync_host_devices` accepts `mcp_busy_device_ids`) +- Modify: `packages/cloud-platform/cloud/pool.py:117-134` (`_to_pooled` sets flag) +- Modify: `packages/cloud-platform/cloud/scheduler.py:206-220` (`_matches` checks `device.mcp_busy`) +- Modify: `packages/cloud-platform/cloud/internal_api/api.py:182-188` (heartbeat handler passes field) +- Test: `packages/cloud-platform/tests/test_pool.py` (extend), `packages/cloud-platform/tests/test_scheduler.py` (extend), `packages/cloud-platform/tests/test_internal_api.py` (extend if exists) + +**Interfaces:** +- Produces: `PooledDevice.mcp_busy: bool = False` +- Produces: `DevicePool.sync_host_devices(..., mcp_busy_device_ids: list[str] | None = None)` +- Produces: scheduler's `_matches()` returns False if `device.mcp_busy is True`. + +- [ ] **Step 1: Locate existing tests** + +Run: `ls packages/cloud-platform/tests/ | grep -E "pool|scheduler|internal_api"` + +- [ ] **Step 2: Write failing tests for pool** + +Append to `packages/cloud-platform/tests/test_pool.py`: + +```python +def test_sync_host_devices_marks_mcp_busy_devices() -> None: + """When a host reports device-1 as MCP-busy, the pool PooledDevice for + device-1 has mcp_busy=True.""" + # Use the existing test harness in this file for constructing a pool. + # (Adjust to match the file's existing fixture style — see other tests + # in the same file for the exact setup pattern.) + pool = _build_pool() # _build_pool is a helper already in test_pool.py + pool.sync_host_devices( + "host-1", + [_device("device-1", status="idle")], + mcp_busy_device_ids=["device-1"], + ) + devices = pool.list_devices() + busy = [d for d in devices if d.device_id == "device-1"] + assert len(busy) == 1 + assert busy[0].mcp_busy is True + + +def test_sync_host_devices_default_mcp_busy_is_false() -> None: + pool = _build_pool() + pool.sync_host_devices( + "host-1", [_device("device-1", status="idle")] + ) + devices = pool.list_devices() + assert devices[0].mcp_busy is False + + +def test_sync_host_devices_clears_mcp_busy_on_next_sync() -> None: + """MCP releases device → next heartbeat without device in + mcp_busy_device_ids → pool reflects mcp_busy=False.""" + pool = _build_pool() + pool.sync_host_devices( + "host-1", + [_device("device-1", status="idle")], + mcp_busy_device_ids=["device-1"], + ) + pool.sync_host_devices( + "host-1", [_device("device-1", status="idle")] + ) + devices = pool.list_devices() + assert devices[0].mcp_busy is False +``` + +If `_build_pool` / `_device` helpers don't exist, look at other tests in `test_pool.py` for the actual fixture names and reuse them. + +- [ ] **Step 3: Implement `PooledDevice.mcp_busy` and pool plumbing** + +In `packages/cloud-platform/cloud/pool.py`: + +1. Add field to `PooledDevice` (after line 49): + +```python +@dataclass(frozen=True) +class PooledDevice: + device_id: str + host_id: str + driver_type: str + status: PooledDeviceStatus + capability_tags: list[str] = field(default_factory=list) + synced_at: datetime | None = None + mcp_busy: bool = False +``` + +2. Update `sync_host_devices` signature (around line 59-67) — add `mcp_busy_device_ids` parameter: + +```python +def sync_host_devices( + self, + host_id: str, + snapshot: list[Device], + *, + address: str | None = None, + planner_transport: Literal["direct", "cloud"] = "direct", + allow_device_takeover: bool = False, + mcp_busy_device_ids: list[str] | None = None, +) -> None: + now = utc_now() + self.store.upsert_host( + host_id, + address=address, + last_seen_at=now, + planner_transport=planner_transport, + ) + busy_set = set(mcp_busy_device_ids or []) + devices = [ + self._to_pooled(device, host_id, now, mcp_busy=device.id in busy_set) + for device in snapshot + ] + # ... rest unchanged +``` + +3. Update `_to_pooled` (around line 117-134) — accept and pass through `mcp_busy`: + +```python +def _to_pooled( + self, + device: Device, + host_id: str, + synced_at: datetime, + *, + mcp_busy: bool = False, +) -> PooledDevice: + raw_status = ( + device.status if device.status in _HOST_REPORTED_STATUSES else "idle" + ) + tags = list(device.capability_tags or []) + return PooledDevice( + device_id=device.id, + host_id=host_id, + driver_type=device.driver_type, + status=raw_status, # type: ignore[arg-type] + capability_tags=tags, + synced_at=synced_at, + mcp_busy=mcp_busy, + ) +``` + +4. In `packages/cloud-platform/cloud/internal_api/api.py` heartbeat handler (around line 182-188), pass the new field: + +```python +pool.sync_host_devices( + host_id, + devices, + address=payload.address, + allow_device_takeover=allow_device_takeover, + planner_transport=payload.planner_transport, + mcp_busy_device_ids=payload.mcp_busy_device_ids, +) +``` + +- [ ] **Step 4: Run pool tests** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_pool.py -v` +Expected: all PASS (new + existing). + +- [ ] **Step 5: Write failing scheduler test** + +Append to `packages/cloud-platform/tests/test_scheduler.py`: + +```python +def test_mcp_busy_device_is_skipped_by_scheduler() -> None: + """A device with mcp_busy=True is not selected for assignment.""" + # Use existing test fixtures; create a pool with one mcp-busy device and + # one idle device, submit a task, verify the idle one is selected. + pool = _build_pool_with_devices( + idle_device="dev-idle", + mcp_busy_device="dev-busy", + ) + scheduler = _build_scheduler(pool) + task_id = scheduler.submit(goal="test", constraints=TaskConstraints()) + assignments = scheduler.assign() + if assignments: + assert all(a.device_id != "dev-busy" for a in assignments) +``` + +Use the existing fixture pattern in `test_scheduler.py` (look at other tests to copy the setup style). If the existing fixtures don't fit, construct the scenario manually by calling `pool.sync_host_devices(..., mcp_busy_device_ids=["dev-busy"])`. + +- [ ] **Step 6: Implement `_matches` check** + +In `packages/cloud-platform/cloud/scheduler.py`, modify `_matches` (lines 206-220). At the start of the function: + +```python +def _matches(device: "PooledDevice", constraints: TaskConstraints) -> bool: + if device.mcp_busy: + return False + # ... rest unchanged +``` + +Also update `scheduler.assign()` (around line 172-178) to filter on `mcp_busy`. Actually, the existing `device.status == "idle"` filter doesn't catch mcp-busy devices (they're still `status="idle"` from the host). The cleanest fix is in `_matches()`, so do it there. + +- [ ] **Step 7: Run scheduler tests** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_scheduler.py -v` +Expected: all PASS. + +- [ ] **Step 8: Run cloud-internal-api integration tests (if any)** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/ -v -k heartbeat` +Expected: all PASS. + +- [ ] **Step 9: Lint check** + +Run: `uv run --with ruff ruff check packages/cloud-platform/cloud/pool.py packages/cloud-platform/cloud/scheduler.py packages/cloud-platform/cloud/internal_api/api.py packages/cloud-platform/tests/test_pool.py packages/cloud-platform/tests/test_scheduler.py` +Expected: clean. + +- [ ] **Step 10: Commit** + +```bash +git add packages/cloud-platform/cloud/pool.py packages/cloud-platform/cloud/scheduler.py packages/cloud-platform/cloud/internal_api/api.py packages/cloud-platform/tests/test_pool.py packages/cloud-platform/tests/test_scheduler.py +git commit -m "feat(cloud): skip MCP-busy devices in scheduler" +``` + +--- + +## Task 9: Host-agent heartbeat sends `mcp_busy_device_ids` + +**Files:** +- Modify: `apps/device-host-agent/host_agent/heartbeat.py:18-27` (`build_device_snapshot` doesn't return the busy ids; add a sibling helper or change the API) +- Modify: `apps/device-host-agent/host_agent/heartbeat.py:30-53` (`HeartbeatSynchronizer.__init__` accepts `mcp_busy_tracker`) +- Modify: `apps/device-host-agent/host_agent/heartbeat.py:58-85` (`sync_once` reads tracker and passes to client) +- Modify: `apps/device-host-agent/host_agent/client.py:162+` (`heartbeat()` accepts and sends `mcp_busy_device_ids`) +- Test: `apps/device-host-agent/tests/test_heartbeat.py` (extend), `apps/device-host-agent/tests/test_client.py` (extend) + +**Interfaces:** +- Consumes: `McpBusyTracker.busy_device_ids()` from Task 4. +- Produces: `HeartbeatSynchronizer.__init__(..., mcp_busy_tracker: McpBusyTracker | None = None)`. +- Produces: `HostAgentClient.heartbeat(snapshot, *, address=..., policy_revision=..., mcp_busy_device_ids: list[str] | None = None)`. + +- [ ] **Step 1: Read existing `HostAgentClient.heartbeat` signature** + +Run: `uv run --package device-host-agent python -c "import inspect; from host_agent.client import HostAgentClient; print(inspect.signature(HostAgentClient.heartbeat))"` +Record the current parameter list — you'll add `mcp_busy_device_ids` to the end of the kwargs. + +- [ ] **Step 2: Write failing test for client** + +Append to `apps/device-host-agent/tests/test_client.py`: + +```python +def test_heartbeat_includes_mcp_busy_device_ids_in_payload() -> None: + """When mcp_busy_device_ids is passed, the client sends it in the request.""" + # Use the existing mock-transport pattern from this file. + # See test_heartbeat_* or test_claim_* in the same file for the pattern. + client = _build_client_with_mock_transport() + captured = _capture_request_payload(client) + + client.heartbeat( + [], + mcp_busy_device_ids=["phone-1"], + ) + assert captured()["mcp_busy_device_ids"] == ["phone-1"] + + +def test_heartbeat_omits_mcp_busy_device_ids_when_empty() -> None: + """Backward compat: empty list still goes over the wire (or is omitted, + depending on what the cloud model accepts — see Task 7 default_factory + accepts both). Verify whatever the implementation does.""" + client = _build_client_with_mock_transport() + captured = _capture_request_payload(client) + client.heartbeat([], mcp_busy_device_ids=[]) + # Either the key is present with [] or absent; both are valid per cloud model. + assert captured().get("mcp_busy_device_ids", []) == [] +``` + +If `_build_client_with_mock_transport` / `_capture_request_payload` helpers don't exist, look at existing tests in `test_client.py` to copy the pattern. + +- [ ] **Step 3: Modify `HostAgentClient.heartbeat`** + +Read `apps/device-host-agent/host_agent/client.py` around line 162 to find the `heartbeat()` method. Add `mcp_busy_device_ids: list[str] | None = None` parameter and include it in the POST body: + +```python +async def heartbeat( + self, + snapshot: list[DeviceSnapshotModel], + *, + address: str | None = None, + policy_revision: int = 0, + mcp_busy_device_ids: list[str] | None = None, +) -> HeartbeatResponse: + payload = { + "host_id": self._config.host_id, + "devices": [device.model_dump() for device in snapshot], + "policy_revision": policy_revision, + "planner_transport": self._config.ai_planner_transport, + } + if address is not None: + payload["address"] = address + if mcp_busy_device_ids: + payload["mcp_busy_device_ids"] = list(mcp_busy_device_ids) + # ... rest of method (POST + parse) unchanged +``` + +(The exact existing payload structure may differ slightly — match it.) + +- [ ] **Step 4: Run client tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_client.py -v` +Expected: all PASS. + +- [ ] **Step 5: Write failing test for `HeartbeatSynchronizer`** + +Append to `apps/device-host-agent/tests/test_heartbeat.py`: + +```python +async def test_sync_once_passes_mcp_busy_device_ids_to_client() -> None: + """When mcp_busy_tracker has a lease, sync_once relays the device_ids.""" + # Use existing test fixtures (mock manager, mock client) from this file. + manager = _build_manager_with_no_devices() + fake_client = _FakeClient() + tracker = McpBusyTracker() + tracker.acquire("phone-1", "sess-a") + sync = HeartbeatSynchronizer( + manager, + fake_client, + _config(), + mcp_busy_tracker=tracker, + ) + await sync.sync_once() + assert fake_client.last_heartbeat_kwargs.get("mcp_busy_device_ids") == ["phone-1"] + + +async def test_sync_once_passes_empty_when_tracker_is_none() -> None: + """Default: no tracker → no mcp_busy_device_ids kwarg (or empty).""" + manager = _build_manager_with_no_devices() + fake_client = _FakeClient() + sync = HeartbeatSynchronizer(manager, fake_client, _config()) + await sync.sync_once() + # If the client method is called without the kwarg, that's fine — verify + # the call didn't pass a non-empty list. + assert not fake_client.last_heartbeat_kwargs.get("mcp_busy_device_ids") +``` + +Look at existing tests in `test_heartbeat.py` for the actual fixture/helper names. + +- [ ] **Step 6: Modify `HeartbeatSynchronizer`** + +In `apps/device-host-agent/host_agent/heartbeat.py`: + +1. Add `mcp_busy_tracker: McpBusyTracker | None = None` to `__init__` (line 30-43) and store as `self.mcp_busy_tracker`. + +2. Import at top: `from host_agent.mcp_lock import McpBusyTracker` under `TYPE_CHECKING`. + +3. In `sync_once` (around line 58-85), compute `mcp_busy_device_ids` and pass to client: + +```python +async def sync_once(self) -> HeartbeatResponse: + snapshot = build_device_snapshot(self.manager) + mcp_busy_ids = ( + self.mcp_busy_tracker.busy_device_ids() + if self.mcp_busy_tracker is not None + else [] + ) + response = await self.client.heartbeat( + snapshot, + address=self.address, + policy_revision=self.policy_revision, + mcp_busy_device_ids=mcp_busy_ids, + ) + # ... rest unchanged +``` + +- [ ] **Step 7: Run heartbeat tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_heartbeat.py -v` +Expected: all PASS. + +- [ ] **Step 8: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/heartbeat.py apps/device-host-agent/host_agent/client.py apps/device-host-agent/tests/test_heartbeat.py apps/device-host-agent/tests/test_client.py` +Expected: clean. + +- [ ] **Step 9: Commit** + +```bash +git add apps/device-host-agent/host_agent/heartbeat.py apps/device-host-agent/host_agent/client.py apps/device-host-agent/tests/test_heartbeat.py apps/device-host-agent/tests/test_client.py +git commit -m "feat(host-agent): include mcp_busy_device_ids in heartbeat payload" +``` + +--- + +## Task 10: `AssignmentExecutor` fail-fast on MCP-held device + +**Files:** +- Modify: `apps/device-host-agent/host_agent/assignment.py:22-30` (`AssignmentExecutor.__init__` accepts optional `mcp_busy_tracker`) +- Modify: `apps/device-host-agent/host_agent/assignment.py:31-57` (`execute` does the check at entry) +- Test: `apps/device-host-agent/tests/test_assignment.py` (extend) + +**Interfaces:** +- Produces: `AssignmentExecutor(factories, *, mcp_busy_tracker: McpBusyTracker | None = None)`. +- Behavior: if `mcp_busy_tracker is not None and assignment.device_id in tracker.busy_device_ids()` at execute entry → return `AssignmentExecutionResult(status="failed", failure_reason="device held by active MCP session")`. + +- [ ] **Step 1: Read current `AssignmentExecutor`** + +Already read in earlier context — `assignment.py:22-57`. Confirm signature is still `def __init__(self, factories: ExecutionFactories)` before editing. + +- [ ] **Step 2: Write failing test** + +Append to `apps/device-host-agent/tests/test_assignment.py`: + +```python +def test_execute_fails_fast_when_mcp_session_holds_device() -> None: + """Cloud assignment arriving for a device currently held by an MCP + session must fail immediately rather than fight for the device.""" + from cloud.internal_api.models import AssignmentModel + from datetime import datetime, UTC + + factories = _build_factories() # use existing fixture from this file + tracker = McpBusyTracker() + tracker.acquire("phone-1", "sess-mcp") + executor = AssignmentExecutor(factories, mcp_busy_tracker=tracker) + assignment = AssignmentModel( + task_id="t1", + attempt=1, + lease_id="l1", + lease_expires_at=datetime.now(UTC), + host_id="h1", + device_id="phone-1", + goal="some goal", + ) + result = executor.execute(assignment) + assert result.status == "failed" + assert "MCP" in (result.failure_reason or "") + + +def test_execute_skips_check_when_tracker_is_none() -> None: + """Default backward-compat: no tracker → no fail-fast.""" + factories = _build_factories() + executor = AssignmentExecutor(factories) # no tracker + # Without a real workflow store / task runner this test gets harder; + # use a goal + a mock runner factory, or verify the entry-point path + # doesn't raise on the mcp_busy check. The minimum we need to prove is + # that None tracker doesn't crash — see other tests in this file for + # the existing pattern for asserting execute() runs through. +``` + +Look at existing `test_assignment.py` to copy `_build_factories` and any mock-task-runner patterns. + +- [ ] **Step 3: Modify `AssignmentExecutor`** + +In `apps/device-host-agent/host_agent/assignment.py`: + +```python +class AssignmentExecutor: + def __init__( + self, + factories: ExecutionFactories, + *, + mcp_busy_tracker: McpBusyTracker | None = None, + ) -> None: + self.factories = factories + self._progress = TaskProgressHolder() + self._mcp_busy_tracker = mcp_busy_tracker + + # ... latest_progress() unchanged ... + + def execute( + self, + assignment: AssignmentModel, + *, + should_stop: Callable[[], bool] | None = None, + stop_reason: Callable[[], str | None] = None, + ) -> AssignmentExecutionResult: + self._progress.clear() + with bind_planner_execution_context(assignment): + if should_stop is not None and should_stop(): + # existing logic + ... + if self._mcp_busy_tracker is not None and ( + assignment.device_id in self._mcp_busy_tracker.busy_device_ids() + ): + return AssignmentExecutionResult( + status="failed", + failure_reason=( + f"device {assignment.device_id} is held by an active " + "MCP session" + ), + ) + # existing workflow / goal dispatch unchanged + ... +``` + +Add `from host_agent.mcp_lock import McpBusyTracker` under `TYPE_CHECKING`. + +- [ ] **Step 4: Run tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_assignment.py -v` +Expected: all PASS (new + existing). + +- [ ] **Step 5: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/assignment.py apps/device-host-agent/tests/test_assignment.py` +Expected: clean. + +- [ ] **Step 6: Commit** + +```bash +git add apps/device-host-agent/host_agent/assignment.py apps/device-host-agent/tests/test_assignment.py +git commit -m "feat(host-agent): fail-fast cloud assignment when MCP holds device" +``` + +--- + +## Task 11: Console mount wiring + `/api/status` field + dashboard row + +**Files:** +- Modify: `apps/device-host-agent/host_agent/web/app.py:209-266` (`create_console_app` signature + body) +- Modify: `apps/device-host-agent/host_agent/web/app.py:353-379` (`/api/status` payload gets `mcp_busy_devices` field) +- Modify: `apps/device-host-agent/host_agent/web/templates/dashboard.html` (one MCP status row) +- Test: `apps/device-host-agent/tests/test_web_app.py` (extend) + +**Interfaces:** +- Produces: `create_console_app(..., mcp_server: FastMCP | None = None, mcp_token_store: McpTokenStore | None = None, mcp_busy_tracker: McpBusyTracker | None = None)`. When all three are provided, the app mounts `/mcp`. +- Produces: `GET /api/status` returns `mcp_busy_devices: list[str]` and `mcp_endpoint: str | None` (the latter is None when MCP not mounted). +- Produces: dashboard HTML includes `MCP{...}` status row. + +- [ ] **Step 1: Read current `/api/status` handler** + +Read `apps/device-host-agent/host_agent/web/app.py` around line 353 (`@app.get("/api/status")`). Record the response dict structure. + +- [ ] **Step 2: Write failing tests** + +Append to `apps/device-host-agent/tests/test_web_app.py`: + +```python +def test_console_app_mounts_mcp_when_all_components_provided(tmp_path) -> None: + from host_agent.mcp_lock import McpBusyTracker + from host_agent.mcp_token import McpTokenStore + from host_agent.web.mcp import build_mcp_server + from device.manager import DeviceManager + from host_agent.status import AgentStatusTracker + + manager = DeviceManager() + status = AgentStatusTracker() + tracker = McpBusyTracker() + token_store = McpTokenStore(tmp_path / "host_mcp_token.json") + server = build_mcp_server( + manager=manager, mcp_busy_tracker=tracker, status_tracker=status + ) + app = _build_console_app( # use existing helper from this file + manager=manager, + status_tracker=status, + mcp_server=server, + mcp_token_store=token_store, + mcp_busy_tracker=tracker, + ) + client = TestClient(app) + # /mcp exists (not 404). Without auth it returns 401. + resp = client.post("/mcp/", headers={"Content-Type": "application/json"}) + assert resp.status_code != 404 + + +def test_console_app_does_not_mount_mcp_when_components_missing() -> None: + app = _build_console_app() # no mcp_* kwargs + client = TestClient(app) + resp = client.post("/mcp/", headers={"Content-Type": "application/json"}) + assert resp.status_code == 404 + + +def test_api_status_includes_mcp_busy_devices(tmp_path) -> None: + # Mount with a tracker that has one lease; verify /api/status shows it. + ... + client = TestClient(app_with_mcp) + # Login first if the test helper requires session; see existing tests. + resp = client.get("/api/status") + assert resp.status_code == 200 + body = resp.json() + assert "mcp_busy_devices" in body + assert "phone-1" in body["mcp_busy_devices"] + + +def test_dashboard_renders_mcp_status_row(tmp_path) -> None: + client = TestClient(app_with_mcp) + # login + resp = client.get("/") + assert resp.status_code == 200 + assert b"MCP" in resp.content # the status row label +``` + +Use the existing helpers in `test_web_app.py` (look at other tests for the exact fixture style). + +- [ ] **Step 3: Modify `create_console_app`** + +In `apps/device-host-agent/host_agent/web/app.py:209-266`: + +1. Add three new kwargs to signature: `mcp_server`, `mcp_token_store`, `mcp_busy_tracker` (all default None). +2. After `app = FastAPI(...)` and existing route setup, conditionally mount MCP: + +```python +if mcp_server is not None and mcp_token_store is not None: + from host_agent.web.mcp_auth import BearerAuthMiddleware + from starlette.middleware import Middleware + + mcp_asgi = mcp_server.streamable_http_app() + # Wrap with bearer auth via a sub-Starlette app: + from starlette.applications import Starlette + + authed = Starlette( + routes=[], + middleware=[Middleware(BearerAuthMiddleware, token_store=mcp_token_store)], + ) + authed.router.mount("/", mcp_asgi) # bearer check wraps all mcp traffic + app.mount("/mcp", authed) +``` + +3. In the `/api/status` handler (around line 353), add fields: + +```python +return { + # ... existing fields ... + "mcp_endpoint": "/mcp" if mcp_server is not None else None, + "mcp_busy_devices": ( + mcp_busy_tracker.busy_device_ids() if mcp_busy_tracker is not None else [] + ), +} +``` + +- [ ] **Step 4: Modify dashboard template** + +Edit `apps/device-host-agent/host_agent/web/templates/dashboard.html`. Find the existing status rows (heartbeat, host policy, etc.) and add after them: + +```html + + MCP + + {% if mcp_endpoint %} + endpoint {{ mcp_endpoint }}; + {% if mcp_busy_devices %}busy: {{ mcp_busy_devices|join(", ") }}{% else %}idle{% endif %} + {% else %} + not configured + {% endif %} + + +``` + +Make sure the template context passes `mcp_endpoint` and `mcp_busy_devices` from the dashboard view handler. Look at the existing dashboard route in `app.py` (around line 320 `@app.get("/", response_class=HTMLResponse)`) to add these to the render context. + +- [ ] **Step 5: Run tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_web_app.py -v` +Expected: all PASS. + +- [ ] **Step 6: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/web/app.py apps/device-host-agent/tests/test_web_app.py` +Expected: clean. + +- [ ] **Step 7: Commit** + +```bash +git add apps/device-host-agent/host_agent/web/app.py apps/device-host-agent/host_agent/web/templates/dashboard.html apps/device-host-agent/tests/test_web_app.py +git commit -m "feat(host-agent): mount /mcp + surface MCP status in console" +``` + +--- + +## Task 12: Wire all components in `create_application` + +**Files:** +- Modify: `apps/device-host-agent/host_agent/app.py:155-288` (`create_application` body) + +**Interfaces:** +- Consumes: Tasks 3, 4, 5, 6, 9, 10, 11 (all the new components + plumbing). +- Produces: `create_application()` now constructs `McpTokenStore`, `McpBusyTracker`, `build_mcp_server`, passes them to `create_console_app`, `HeartbeatSynchronizer`, `AssignmentExecutor`. + +- [ ] **Step 1: Read current `create_application`** + +Read `apps/device-host-agent/host_agent/app.py:155-288` again. Identify the insertion points: +- After `resolved_manager` is constructed (around line 182-186) — needed for `build_mcp_server` +- After `TaskMetadataStore` / `Timeline` / `AssignmentExecutor` construction (around line 200-211) — `AssignmentExecutor` needs `mcp_busy_tracker` +- Before `create_console_app` call (around line 212-228) — passes new kwargs +- Before `HeartbeatSynchronizer` construction (around line 238-252) — passes `mcp_busy_tracker` + +- [ ] **Step 2: Add an integration test** + +Append to `apps/device-host-agent/tests/test_app.py`: + +```python +def test_create_application_wires_mcp_components(tmp_path, monkeypatch) -> None: + """create_application produces an app whose console is mounted at /mcp + and whose heartbeat reads from the in-process McpBusyTracker.""" + # Use the existing pattern in test_app.py for monkeypatching config to + # point at tmp_path. Look at other tests in the file for the exact setup. + app = create_application(config=_test_config(tmp_path)) + # /mcp is reachable (401 without auth, not 404). + from starlette.testclient import TestClient + + with TestClient(app.console_app) as client: + resp = client.post("/mcp/") + assert resp.status_code == 401 # auth required, not 404 + # Token file exists. + assert (tmp_path / "host_mcp_token.json").exists() +``` + +Note: if `HostAgentApplication` doesn't expose `console_app` as an attribute, you'll need to add it in step 3. + +- [ ] **Step 3: Modify `create_application`** + +In `apps/device-host-agent/host_agent/app.py`: + +1. After `resolved_manager = ...` is set (around line 182), add construction of new components: + +```python +mcp_token_store = McpTokenStore( + resolved_config.identity_path.parent / "host_mcp_token.json" +) +mcp_token = mcp_token_store.load_or_create() # eager; logs on first creation +if mcp_token_store._path.exists() and not mcp_token_store._path.stat().st_size: + import logging + + logging.getLogger(__name__).info( + "MCP token generated at %s", mcp_token_store._path + ) +# Actually log on first creation only — adjust by storing previous existence. +``` + +A cleaner approach for the "log on first creation" requirement: + +```python +token_path = resolved_config.identity_path.parent / "host_mcp_token.json" +token_existed_before = token_path.exists() +mcp_token_store = McpTokenStore(token_path) +mcp_token = mcp_token_store.load_or_create() +if not token_existed_before: + import logging + + logging.getLogger(__name__).info( + "MCP token generated at %s", token_path + ) +``` + +2. After `executor = AssignmentExecutor(...)` construction (around line 204-211), pass `mcp_busy_tracker`: + +```python +mcp_busy_tracker = McpBusyTracker(ttl_seconds=60.0) +executor = AssignmentExecutor( + create_execution_factories(...), + mcp_busy_tracker=mcp_busy_tracker, +) +``` + +3. Build the MCP server before `create_console_app`: + +```python +mcp_server = build_mcp_server( + manager=resolved_manager, + mcp_busy_tracker=mcp_busy_tracker, + status_tracker=status_tracker, +) +``` + +4. Pass new kwargs to `create_console_app` (around line 212-228): + +```python +console_app = create_console_app( + config=resolved_config, + manager=resolved_manager, + config_store=config_store, + local_account_store=..., + identity_store=..., + history_store=history_store, + status_tracker=status_tracker, + session_manager=..., + enrollment_client=..., + host_client=client, + metadata_store=metadata_store, + timeline=timeline, + executor=executor, + mcp_server=mcp_server, + mcp_token_store=mcp_token_store, + mcp_busy_tracker=mcp_busy_tracker, +) +``` + +5. Pass `mcp_busy_tracker` to `HeartbeatSynchronizer` (around line 238-252): + +```python +heartbeat = HeartbeatSynchronizer( + resolved_manager, + client, + resolved_config, + status_tracker=status_tracker, + mcp_busy_tracker=mcp_busy_tracker, + on_sync=..., + policy_cache=..., + on_policy_sync=..., +) +``` + +6. Add `console_app` to the returned `HostAgentApplication` dataclass if not already there (it currently only stores `console_server`). The test in Step 2 needs access to the FastAPI app for TestClient. Look at how the existing embedded uvicorn server is built (around line 229-236) and either: + - Add `console_app: FastAPI | None = None` to `HostAgentApplication` dataclass, OR + - Have the test build its own app via `create_console_app` directly. + +Option B (test builds its own app) is less invasive. Adjust the test in Step 2 accordingly. + +- [ ] **Step 4: Run app tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_app.py -v` +Expected: all PASS (new + existing 9 `create_application()` calls). + +- [ ] **Step 5: Run all host-agent tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/ -v` +Expected: all PASS. + +- [ ] **Step 6: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/app.py apps/device-host-agent/tests/test_app.py` +Expected: clean. + +- [ ] **Step 7: Commit** + +```bash +git add apps/device-host-agent/host_agent/app.py apps/device-host-agent/tests/test_app.py +git commit -m "feat(host-agent): wire MCP server into create_application" +``` + +--- + +## Task 13: CLI `mcp-token` subcommand + +**Files:** +- Modify: `apps/device-host-agent/host_agent/cli.py:19-39` +- Test: `apps/device-host-agent/tests/test_cli.py` (extend) + +**Interfaces:** +- Produces: `device-host-agent mcp-token` prints the current MCP token to stdout (generating if missing), then exits 0. + +- [ ] **Step 1: Write failing test** + +Append to `apps/device-host-agent/tests/test_cli.py`: + +```python +def test_mcp_token_subcommand_prints_token(tmp_path, capsys, monkeypatch) -> None: + monkeypatch.setenv("HOST_AGENT_IDENTITY_PATH", str(tmp_path / "host_identity.json")) + monkeypatch.setenv("HOST_AGENT_LOCAL_ACCOUNT_PATH", str(tmp_path / "host_local_account.json")) + # Possibly need to set other env vars per existing test pattern. + from host_agent.cli import main + + main(["mcp-token"]) + out = capsys.readouterr().out.strip() + assert len(out) >= 40 # token is ~43 chars + # Subsequent invocation prints the same token (idempotent). + main(["mcp-token"]) + out2 = capsys.readouterr().out.strip() + assert out == out2 +``` + +Look at existing `test_cli.py` for env setup patterns. + +- [ ] **Step 2: Implement subcommand** + +In `apps/device-host-agent/host_agent/cli.py`, extend the argparse subparsers and dispatch: + +```python +def main(argv: Sequence[str] | None = None) -> None: + parser = argparse.ArgumentParser(description="Run the Device Host Agent") + subparsers = parser.add_subparsers(dest="command") + subparsers.add_parser("setup", help="Create the local operator account") + subparsers.add_parser( + "mcp-token", + help="Print the MCP server bearer token (generating if missing)", + ) + args = parser.parse_args(argv) + + if args.command == "mcp-token": + _print_mcp_token() + return + + try: + if args.command == "setup": + _run_setup() + return + config = _resolve_config_with_local_account() + except LocalAccountSetupError as exc: + print(f"error: {exc}", file=sys.stderr) + raise SystemExit(1) from exc + + try: + create_application(config=config).run() + except InstanceAlreadyRunningError as exc: + print(f"error: {exc}", file=sys.stderr) + raise SystemExit(1) from exc + + +def _print_mcp_token() -> None: + config = load_host_agent_config() + store = McpTokenStore(config.identity_path.parent / "host_mcp_token.json") + print(store.load_or_create().token) +``` + +Add imports: `from host_agent.mcp_token import McpTokenStore`. + +- [ ] **Step 3: Run tests** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_cli.py -v` +Expected: all PASS. + +- [ ] **Step 4: Lint check** + +Run: `uv run --with ruff ruff check apps/device-host-agent/host_agent/cli.py apps/device-host-agent/tests/test_cli.py` +Expected: clean. + +- [ ] **Step 5: Commit** + +```bash +git add apps/device-host-agent/host_agent/cli.py apps/device-host-agent/tests/test_cli.py +git commit -m "feat(host-agent): add mcp-token CLI subcommand" +``` + +--- + +## Task 14: Documentation + +**Files:** +- Create: `docs/MCP_INTEGRATION.md` +- Modify: `docs/MACOS_IPHONE_SETUP.md` (add a new section near the bottom) + +No TDD for docs. Content is concrete and verifiable. + +- [ ] **Step 1: Create `docs/MCP_INTEGRATION.md`** + +```markdown +# Host-Agent MCP Server Integration + +The host-agent process exposes a Streamable HTTP MCP server on the same +port as the local console (default `127.0.0.1:8765`), at path `/mcp`. This +lets any MCP-compatible client — Hermes Agent, Claude Desktop, custom +scripts using the `mcp` Python SDK — drive devices directly through the +same `DeviceManager` the cloud worker uses. + +## Prerequisites + +- Host-agent built from this repo (see `docs/MACOS_IPHONE_SETUP.md`). +- An MCP client that supports the Streamable HTTP transport (mcp SDK + 1.20+ on the client side). + +## Get the bearer token + +The first time host-agent starts after this feature ships, it generates +a random bearer token and writes it to: + + /host_mcp_token.json + +(Default: `tasks/host_mcp_token.json` next to `host_identity.json`.) + +To print it for copy/paste: + + device-host-agent mcp-token + +To rotate: delete the file and restart host-agent. Old tokens stop +working immediately. + +## Hermes Agent configuration + +Add to `~/.hermes/config.yaml`: + +```yaml +mcp_servers: + apex_device: + url: "http://127.0.0.1:8765/mcp" + headers: + Authorization: "Bearer " +``` + +Start (or restart) Hermes. Verify by asking Hermes to list devices: + +> Use the apex_device MCP to list connected devices. + +## Tools exposed + +All 11 device tools from `api/mcp.py`: + +- `take_screenshot(device_id?)` +- `tap(x, y, device_id?)` +- `swipe(start_x, start_y, end_x, end_y, duration_ms?, device_id?)` +- `input_text(text, device_id?)` +- `launch_app(app_id, device_id?)` +- `find_text(query, device_id?)` +- `find_icon(name, device_id?)` +- `get_ui_tree(device_id?, include_app_info?)` +- `describe_screen(device_id?)` +- `list_devices()` +- `device_status(device_id)` + +## Concurrency model + +- The cloud worker and MCP clients share the same `DeviceManager`. +- Per-device, session-level locking: the first caller (cloud or MCP) to + touch a device holds it; the other side sees a busy error. +- MCP sessions hold their lock until the session ends OR 60 seconds of + inactivity. Cloud assignments hold theirs until the assignment + terminates. +- The cloud scheduler is told about MCP-held devices via the heartbeat + `mcp_busy_device_ids` field, so it normally won't even try to dispatch + to them. A 30-second window exists between an MCP acquire and the next + heartbeat; during that window cloud may dispatch, and the host-agent + will fail-fast the assignment with `failure_reason="device held by an + active MCP session"`. + +## Network binding + +The MCP endpoint is bound to the same address as the local console. By +default this is `127.0.0.1` (loopback only). To expose on a different +interface, set `HOST_AGENT_CONSOLE_BIND_HOST` AND +`HOST_AGENT_CONSOLE_ALLOW_NON_LOOPBACK=true` — both are required. This +is the same escape hatch the local console uses; there is no MCP-only +override. + +## Error responses + +| Condition | HTTP / JSON-RPC | Body | +|---|---|---| +| Missing/wrong bearer token | HTTP 401 | `{"error": "invalid token"}` + `WWW-Authenticate: Bearer` | +| Device busy (cloud) | JSON-RPC `-32000` | `"device X is busy (held by cloud assignment)"`, `data.busy_owner = "cloud_assignment"` | +| Device busy (other MCP) | JSON-RPC `-32000` | `"..."`, `data.busy_owner = "mcp_session:"` | +| Unknown device | JSON-RPC `-32602` | `"unknown device: X"` | +| Tool error | JSON-RPC `-32000` | Original exception message | + +## Troubleshooting + +- **`list_devices` returns `[]`**: no devices registered. Use the local + console at `http://127.0.0.1:8765/` to add one (Login → Devices). +- **`device X is busy` even when cloud console says device is idle**: + check whether another MCP session is holding it. The local console + dashboard shows active MCP sessions and held device_ids. +- **Token verification fails after restart**: confirm you copied the + token from the current `host_mcp_token.json`, not an older one. + Rotation = delete file + restart. + +## Out of scope (current version) + +- `wait_until_usable` MCP tool: implemented internally but not exposed. + MVP callers must handle busy errors themselves. +- MCP call history in the local console: only current state is surfaced, + not a call log. +- Token rotation CLI: use delete-and-restart for now. +- Non-loopback binding without explicit opt-in. +``` + +- [ ] **Step 2: Add a section to `docs/MACOS_IPHONE_SETUP.md`** + +Read the existing file structure first. Add a new section near the bottom (before any "Troubleshooting" appendix): + +```markdown +## MCP server (Hermes Agent integration) + +Host-agent now exposes an MCP server on the same port as the local +console (`127.0.0.1:8765/mcp`). To drive your iPhone from Hermes Agent +or any MCP-compatible client: + +1. Start host-agent normally. +2. Get the bearer token: `device-host-agent mcp-token`. +3. Configure Hermes per `docs/MCP_INTEGRATION.md`. + +The MCP path reuses the same WDA session that the cloud worker uses. +Per-device locking prevents both sides from driving the same device at +once; see `docs/MCP_INTEGRATION.md` for the full concurrency model. +``` + +- [ ] **Step 3: Commit** + +```bash +git add docs/MCP_INTEGRATION.md docs/MACOS_IPHONE_SETUP.md +git commit -m "docs: add MCP integration guide" +``` + +--- + +## Task 15: Final validation + +**Files:** none (verification only) + +- [ ] **Step 1: Run full non-integration test suite** + +Run: `uv run --all-packages pytest -m "not integration"` +Expected: ALL PASS. Compare against the pre-change baseline (per project memory: ~687 passed / 54 deselected as of recent runs). Any new failures must be explained. + +- [ ] **Step 2: Run host-agent subpackage tests in isolation** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/` +Expected: ALL PASS. + +- [ ] **Step 3: Run cloud-platform subpackage tests** + +Run: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/` +Expected: ALL PASS. + +- [ ] **Step 4: Ruff check + format check on all touched paths** + +Run: `uv run --with ruff ruff check apps/device-host-agent/ api/mcp.py packages/cloud-platform/ tests/ docs/MCP_INTEGRATION.md docs/MACOS_IPHONE_SETUP.md` +Expected: clean. + +Run: `uv run --with ruff ruff format --check apps/device-host-agent/ api/mcp.py packages/cloud-platform/ tests/` +Expected: clean. + +- [ ] **Step 5: compileall on touched Python files** + +Run: `uv run --all-packages python -m compileall apps/device-host-agent/host_agent/ api/mcp.py packages/cloud-platform/cloud/` +Expected: no errors. + +- [ ] **Step 6: OpenSpec strict validation (regression only — we did NOT add a new proposal)** + +Run: `openspec validate --strict` +Expected: PASS (no openspec changes in this feature; validates existing specs still pass). + +- [ ] **Step 7: Architecture boundary regression** + +Run: `uv run --package device-host-agent pytest apps/device-host-agent/tests/test_execution.py::test_runtime_owned_packages_do_not_import_host_or_cloud_concerns -v` +Expected: PASS — `runtime/` and `api/` packages still do not import `host_agent` or `cloud`. + +- [ ] **Step 8: Final manual sanity check** + +Run: `uv run --package device-host-agent python -c " +from host_agent.app import create_application +# Just exercise that import graph composes at module level. +print('imports ok') +"` +Expected: `imports ok`. + +- [ ] **Step 9: Final commit (if anything was reformatted)** + +If ruff format fixed anything: + +```bash +git add . +git commit -m "style: ruff format after MCP server integration" +``` + +Otherwise, no commit needed — the feature is fully committed across Tasks 1-14. + +--- + +## Open Items Flagged for User Decision + +Per the brainstorming spec, this implementation does NOT go through the openspec proposal process. If you (the user) prefer to retroactively create an openspec change for traceability, do so as a follow-up — the spec at `docs/superpowers/specs/2026-07-21-host-agent-mcp-server-design.md` is structured to translate cleanly into an openspec proposal if desired.