feat(host-agent): execute cloud assignments
This commit is contained in:
@@ -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,
|
||||||
|
},
|
||||||
|
)
|
||||||
@@ -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'"
|
||||||
@@ -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.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.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.
|
- [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.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.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.
|
- [ ] 7.7 Implement graceful shutdown that stops polling, finishes or interrupts current work according to lease policy, and performs a final heartbeat when possible.
|
||||||
|
|||||||
Reference in New Issue
Block a user