Files
agentic-mobile-control/docs/superpowers/plans/2026-07-21-host-agent-mcp-server.md
T
q792602257andClaude Opus 4.6 47eac0f2a7 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 <noreply@anthropic.com>
2026-07-21 13:47:15 +08:00

2664 lines
95 KiB
Markdown

# 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": "<isoformat>"}`.
- **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_<file>.py -v`
- Cloud package: `uv run --package device-cloud-platform pytest packages/cloud-platform/tests/test_<file>.py -v`
- Full non-integration: `uv run --all-packages pytest -m "not integration"`
- Ruff: `uv run --with ruff ruff check <paths>` then `uv run --with ruff ruff format --check <paths>`
- Compile: `uv run --all-packages python -m compileall <paths>`
## 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 `<tr><td>MCP</td><td>{...}</td></tr>` 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
<tr>
<td>MCP</td>
<td>
{% if mcp_endpoint %}
endpoint <code>{{ mcp_endpoint }}</code>;
{% if mcp_busy_devices %}busy: {{ mcp_busy_devices|join(", ") }}{% else %}idle{% endif %}
{% else %}
not configured
{% endif %}
</td>
</tr>
```
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:
<identity_path.parent>/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 <paste-token-here>"
```
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:<prefix>"` |
| 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.