From 08cef7ca3cb8c940dc54374590f4bdaa3fb62a9f Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Tue, 14 Jul 2026 16:35:48 +0800 Subject: [PATCH] fix(host-agent): thread device manager into task runner observer/screenshot create_task_runner() built TaskRunner's observer/screenshot_provider by calling describe_screen(device_id)/take_screenshot(device_id) without manager=, so both silently fell back to the process-global DEFAULT_MANAGER singleton instead of the Host Agent's real, device-populated DeviceManager. DEFAULT_MANAGER never has any device registered, so every task's first step raised DeviceNotFoundError even though the console (which does pass manager=) showed the same device as connected. Deterministic on every task, independent of process count. Add regression tests confirming both lambdas now resolve devices via the configured manager; verified each fails with the original DeviceNotFoundError symptom when the fix is reverted. Co-Authored-By: Claude Sonnet 5 --- .../device-host-agent/host_agent/execution.py | 6 ++ .../device-host-agent/tests/test_execution.py | 98 +++++++++++++++++++ 2 files changed, 104 insertions(+) diff --git a/apps/device-host-agent/host_agent/execution.py b/apps/device-host-agent/host_agent/execution.py index 732d44d..d43f21b 100644 --- a/apps/device-host-agent/host_agent/execution.py +++ b/apps/device-host-agent/host_agent/execution.py @@ -14,6 +14,8 @@ from runtime.planner_config import PlannerConfig, load_config as load_planner_co from runtime.task import TaskRunner from storage.task_metadata import TaskMetadataStore from storage.timeline import Timeline +from tools.describe_screen import describe_screen +from tools.screenshot import take_screenshot from workflow.runner import WorkflowRunner from workflow.store import WorkflowStore @@ -39,6 +41,10 @@ def create_execution_factories( def create_task_runner() -> TaskRunner: return TaskRunner( executor=Executor(tools=default_tool_registry(manager=manager)), + observer=lambda device_id: describe_screen(device_id, manager=manager), + screenshot_provider=lambda device_id: take_screenshot( + device_id, manager=manager + ), metadata_store=metadata_store, timeline=timeline, planner=_host_agent_planner(resolved_host_agent_config), diff --git a/apps/device-host-agent/tests/test_execution.py b/apps/device-host-agent/tests/test_execution.py index cfd4194..d815eca 100644 --- a/apps/device-host-agent/tests/test_execution.py +++ b/apps/device-host-agent/tests/test_execution.py @@ -5,6 +5,8 @@ from pathlib import Path import pytest from device.manager import DeviceManager +from driver.base import Driver +from host_agent import execution from host_agent.cloud_planner_client import CloudProxyToolCallingClient from host_agent.config import HostAgentConfig from host_agent.execution import create_execution_factories @@ -15,6 +17,54 @@ from workflow.runner import WorkflowRunner from workflow.store import WorkflowStore +class FakeDriver(Driver): + def __init__(self) -> None: + self.calls: list[tuple[str, tuple[object, ...]]] = [] + + def connect(self) -> None: + self.calls.append(("connect", ())) + + def disconnect(self) -> None: + return None + + def screenshot(self) -> bytes: + return b"fake-screenshot-bytes" + + def tap(self, x: float, y: float) -> None: + return None + + def swipe( + self, + start_x: float, + start_y: float, + end_x: float, + end_y: float, + duration_ms: int = 500, + ) -> None: + return None + + def input(self, text: str) -> None: + return None + + def launch(self, app_id: str) -> None: + return None + + def terminate(self, app_id: str) -> None: + return None + + def tree(self): + return None + + def home(self) -> None: + return None + + def lock(self) -> None: + return None + + def unlock(self) -> None: + return None + + def test_execution_factories_compose_existing_runtime_and_workflow(tmp_path) -> None: manager = DeviceManager() workflow_store = WorkflowStore(tmp_path / "workflows.sqlite3") @@ -32,6 +82,54 @@ def test_execution_factories_compose_existing_runtime_and_workflow(tmp_path) -> assert isinstance(workflow_runner.task_runner_factory(), TaskRunner) +def test_created_task_runner_screenshot_provider_uses_configured_manager( + tmp_path, +) -> None: + """Regression test: `create_task_runner()` must thread the Host Agent's + own `manager` into `screenshot_provider`. Omitting `manager=` makes the + tool fall back to the process-global `DEFAULT_MANAGER` singleton, which + never has this device registered, so it raises `DeviceNotFoundError` even + though the device is connected on the manager actually in use. + """ + manager = DeviceManager() + manager.register_device("phone-1", lambda: FakeDriver()) + manager.connect("phone-1") + factories = create_execution_factories( + manager, workflow_store=WorkflowStore(tmp_path / "workflows.sqlite3") + ) + + task_runner = factories.task_runner_factory() + + assert task_runner.screenshot_provider("phone-1") == b"fake-screenshot-bytes" + + +def test_created_task_runner_observer_uses_configured_manager( + tmp_path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Regression test: `create_task_runner()` must thread the Host Agent's + own `manager` into `observer` the same way it does for + `screenshot_provider` -- see the test above for the failure mode this + guards against. + """ + manager = DeviceManager() + seen: dict[str, object] = {} + + def fake_describe_screen(device_id, *, manager=None): + seen["device_id"] = device_id + seen["manager"] = manager + return "scene-stub" + + monkeypatch.setattr(execution, "describe_screen", fake_describe_screen) + factories = create_execution_factories( + manager, workflow_store=WorkflowStore(tmp_path / "workflows.sqlite3") + ) + + task_runner = factories.task_runner_factory() + + assert task_runner.observer("phone-1") == "scene-stub" + assert seen == {"device_id": "phone-1", "manager": manager} + + def test_created_task_runner_defaults_to_ai_planner( tmp_path, monkeypatch: pytest.MonkeyPatch ) -> None: