From 6dc2803ccc3ff1904d95303efed323eeb86b368e Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Sun, 12 Jul 2026 18:52:32 +0800 Subject: [PATCH] feat(host-agent): execute cloud assignments --- .../host_agent/assignment.py | 71 +++++++++++ .../tests/test_assignment.py | 115 ++++++++++++++++++ .../cloud-control-plane-integration/tasks.md | 2 +- 3 files changed, 187 insertions(+), 1 deletion(-) create mode 100644 apps/device-host-agent/host_agent/assignment.py create mode 100644 apps/device-host-agent/tests/test_assignment.py diff --git a/apps/device-host-agent/host_agent/assignment.py b/apps/device-host-agent/host_agent/assignment.py new file mode 100644 index 0000000..6d5983e --- /dev/null +++ b/apps/device-host-agent/host_agent/assignment.py @@ -0,0 +1,71 @@ +from __future__ import annotations + +from dataclasses import dataclass, field +from typing import Any + +from cloud.internal_api.models import AssignmentModel +from core.models import Task +from host_agent.execution import ExecutionFactories + + +@dataclass(frozen=True) +class AssignmentExecutionResult: + status: str + failure_reason: str | None = None + metadata: dict[str, Any] = field(default_factory=dict) + + +class AssignmentExecutor: + def __init__(self, factories: ExecutionFactories) -> None: + self.factories = factories + + def execute(self, assignment: AssignmentModel) -> AssignmentExecutionResult: + if assignment.workflow_definition_id is not None: + return self._execute_workflow(assignment) + if assignment.goal is not None: + return self._execute_goal(assignment) + return AssignmentExecutionResult( + status="failed", + failure_reason="assignment has neither goal nor workflow definition", + ) + + def _execute_goal( + self, + assignment: AssignmentModel, + ) -> AssignmentExecutionResult: + task = Task(goal=assignment.goal or "", device_id=assignment.device_id) + completed = self.factories.task_runner_factory().run(task) + return AssignmentExecutionResult( + status="done" if completed.status == "completed" else "failed", + failure_reason=completed.failure_reason, + metadata={ + "runtime_task_id": completed.id, + "runtime_status": completed.status, + }, + ) + + def _execute_workflow( + self, + assignment: AssignmentModel, + ) -> AssignmentExecutionResult: + definition_id = assignment.workflow_definition_id or "" + definition = self.factories.workflow_store.get_definition(definition_id) + if definition is None: + return AssignmentExecutionResult( + status="failed", + failure_reason=f"unknown workflow definition {definition_id!r}", + ) + run = self.factories.workflow_runner_factory().run( + definition, + device_id=assignment.device_id, + ) + return AssignmentExecutionResult( + status="done" if run.status == "completed" else "failed", + failure_reason=( + None if run.status == "completed" else f"workflow ended as {run.status}" + ), + metadata={ + "workflow_run_id": run.id, + "workflow_status": run.status, + }, + ) diff --git a/apps/device-host-agent/tests/test_assignment.py b/apps/device-host-agent/tests/test_assignment.py new file mode 100644 index 0000000..9ce6750 --- /dev/null +++ b/apps/device-host-agent/tests/test_assignment.py @@ -0,0 +1,115 @@ +from __future__ import annotations + +from datetime import UTC, datetime +from types import SimpleNamespace + +from cloud.internal_api.models import AssignmentModel +from core.models import Task +from host_agent.assignment import AssignmentExecutor +from host_agent.execution import ExecutionFactories + + +def _assignment(**overrides) -> AssignmentModel: + values = { + "task_id": "cloud-task", + "attempt": 1, + "lease_id": "lease-a", + "lease_expires_at": datetime(2026, 7, 12, tzinfo=UTC), + "host_id": "host-a", + "device_id": "device-a", + "goal": "open settings", + } + values.update(overrides) + return AssignmentModel(**values) + + +def test_goal_assignment_executes_through_task_runner() -> None: + received: list[Task] = [] + + class FakeTaskRunner: + def run(self, task: Task) -> Task: + received.append(task) + task.status = "completed" + return task + + factories = ExecutionFactories( + task_runner_factory=lambda: FakeTaskRunner(), # type: ignore[arg-type,return-value] + workflow_runner_factory=lambda: object(), # type: ignore[arg-type,return-value] + workflow_store=object(), # type: ignore[arg-type] + ) + + result = AssignmentExecutor(factories).execute(_assignment()) + + assert result.status == "done" + assert received[0].goal == "open settings" + assert received[0].device_id == "device-a" + assert result.metadata["runtime_task_id"] == received[0].id + + +def test_goal_assignment_preserves_runtime_failure_reason() -> None: + class FakeTaskRunner: + def run(self, task: Task) -> Task: + task.status = "failed" + task.failure_reason = "planner unavailable" + return task + + factories = ExecutionFactories( + task_runner_factory=lambda: FakeTaskRunner(), # type: ignore[arg-type,return-value] + workflow_runner_factory=lambda: object(), # type: ignore[arg-type,return-value] + workflow_store=object(), # type: ignore[arg-type] + ) + + result = AssignmentExecutor(factories).execute(_assignment()) + + assert result.status == "failed" + assert result.failure_reason == "planner unavailable" + + +def test_workflow_assignment_loads_and_executes_definition() -> None: + definition = object() + calls: list[tuple[object, str]] = [] + + class FakeWorkflowStore: + def get_definition(self, definition_id: str): + return definition if definition_id == "workflow-a" else None + + class FakeWorkflowRunner: + def run(self, loaded_definition, device_id: str): + calls.append((loaded_definition, device_id)) + return SimpleNamespace(id="run-a", status="completed") + + factories = ExecutionFactories( + task_runner_factory=lambda: object(), # type: ignore[arg-type,return-value] + workflow_runner_factory=lambda: FakeWorkflowRunner(), # type: ignore[arg-type,return-value] + workflow_store=FakeWorkflowStore(), # type: ignore[arg-type] + ) + + result = AssignmentExecutor(factories).execute( + _assignment(goal=None, workflow_definition_id="workflow-a") + ) + + assert result.status == "done" + assert calls == [(definition, "device-a")] + assert result.metadata == { + "workflow_run_id": "run-a", + "workflow_status": "completed", + } + + +def test_unknown_workflow_fails_without_running() -> None: + class FakeWorkflowStore: + def get_definition(self, definition_id: str): + return None + + factories = ExecutionFactories( + task_runner_factory=lambda: object(), # type: ignore[arg-type,return-value] + workflow_runner_factory=lambda: object(), # type: ignore[arg-type,return-value] + workflow_store=FakeWorkflowStore(), # type: ignore[arg-type] + ) + + result = AssignmentExecutor(factories).execute( + _assignment(goal=None, workflow_definition_id="missing") + ) + + assert result.status == "failed" + assert result.failure_reason == "unknown workflow definition 'missing'" diff --git a/openspec/changes/cloud-control-plane-integration/tasks.md b/openspec/changes/cloud-control-plane-integration/tasks.md index a73538f..8243036 100644 --- a/openspec/changes/cloud-control-plane-integration/tasks.md +++ b/openspec/changes/cloud-control-plane-integration/tasks.md @@ -55,7 +55,7 @@ - [x] 7.1 Implement an authenticated Host Agent client for heartbeat, long-poll claim, renewal, and result operations with bounded retry/backoff. - [x] 7.2 Build complete device snapshots from the local `DeviceManager` and synchronize them at the configured interval. - [x] 7.3 Compose local `TaskRunner` and `WorkflowRunner` factories without importing cloud concerns into Runtime-owned packages. -- [ ] 7.4 Execute goal assignments through the configured Runtime Planner/Executor and workflow assignments through the existing workflow runner. +- [x] 7.4 Execute goal assignments through the configured Runtime Planner/Executor and workflow assignments through the existing workflow runner. - [ ] 7.5 Run lease renewal alongside active execution and stop further interruptible actions after confirmed lease loss. - [ ] 7.6 Normalize and report successful/failed terminal outcomes, including Runtime failure reasons, with idempotent retries after response loss. - [ ] 7.7 Implement graceful shutdown that stops polling, finishes or interrupts current work according to lease policy, and performs a final heartbeat when possible.