chore(cloud): verify integration change
This commit is contained in:
@@ -65,8 +65,7 @@ class ActiveAssignmentRunner:
|
|||||||
asyncio.to_thread(
|
asyncio.to_thread(
|
||||||
self.executor.execute,
|
self.executor.execute,
|
||||||
assignment,
|
assignment,
|
||||||
should_stop=lambda: guard.is_lost()
|
should_stop=lambda: guard.is_lost() or self._stop_requested.is_set(),
|
||||||
or self._stop_requested.is_set(),
|
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
renewal = asyncio.create_task(
|
renewal = asyncio.create_task(
|
||||||
|
|||||||
@@ -177,7 +177,9 @@ def test_one_host_executes_assignment_through_outbound_protocol(tmp_path) -> Non
|
|||||||
async with _control_plane(tmp_path / "one-host.sqlite3", "host-a") as app:
|
async with _control_plane(tmp_path / "one-host.sqlite3", "host-a") as app:
|
||||||
async with _host_client(app, "host-a", request_paths=paths) as client:
|
async with _host_client(app, "host-a", request_paths=paths) as client:
|
||||||
await _sync_fake_device(client, "host-a", "device-a")
|
await _sync_fake_device(client, "host-a", "device-a")
|
||||||
task_id = app.state.cloud_services.scheduler.submit(goal="open settings")
|
task_id = app.state.cloud_services.scheduler.submit(
|
||||||
|
goal="open settings"
|
||||||
|
)
|
||||||
app.state.cloud_services.scheduler.assign()
|
app.state.cloud_services.scheduler.assign()
|
||||||
assignment = await client.claim()
|
assignment = await client.claim()
|
||||||
assert assignment is not None
|
assert assignment is not None
|
||||||
@@ -287,7 +289,9 @@ def test_host_agent_restart_reuses_active_lease(tmp_path) -> None:
|
|||||||
assert (await restarted_client.renew(assignment)).status == "renewed"
|
assert (await restarted_client.renew(assignment)).status == "renewed"
|
||||||
await restarted_client.report_result(assignment, status="done")
|
await restarted_client.report_result(assignment, status="done")
|
||||||
|
|
||||||
assert app.state.cloud_services.repository.get_task(task_id).status == "done"
|
assert (
|
||||||
|
app.state.cloud_services.repository.get_task(task_id).status == "done"
|
||||||
|
)
|
||||||
|
|
||||||
asyncio.run(scenario())
|
asyncio.run(scenario())
|
||||||
|
|
||||||
|
|||||||
@@ -71,7 +71,7 @@
|
|||||||
|
|
||||||
## 9. Verification And Project Records
|
## 9. Verification And Project Records
|
||||||
|
|
||||||
- [ ] 9.1 Run formatting, static checks, all non-integration tests, and targeted PostgreSQL integration/concurrency tests.
|
- [x] 9.1 Run formatting, static checks, all non-integration tests, and targeted PostgreSQL integration/concurrency tests.
|
||||||
- [ ] 9.2 Run an end-to-end cloud submission through a Host Agent and fake device until the public SDK reports done and a failure case until it reports failed.
|
- [ ] 9.2 Run an end-to-end cloud submission through a Host Agent and fake device until the public SDK reports done and a failure case until it reports failed.
|
||||||
- [ ] 9.3 Verify existing local REST/MCP/console behavior and dependency-boundary tests remain unchanged.
|
- [ ] 9.3 Verify existing local REST/MCP/console behavior and dependency-boundary tests remain unchanged.
|
||||||
- [ ] 9.4 Run OpenSpec validation for `cloud-control-plane-integration` and map automated tests to every new or modified scenario.
|
- [ ] 9.4 Run OpenSpec validation for `cloud-control-plane-integration` and map automated tests to every new or modified scenario.
|
||||||
|
|||||||
@@ -19,8 +19,7 @@ target_metadata = Base.metadata
|
|||||||
|
|
||||||
def _database_url() -> str:
|
def _database_url() -> str:
|
||||||
return normalize_database_url(
|
return normalize_database_url(
|
||||||
os.environ.get("CLOUD_DATABASE_URL")
|
os.environ.get("CLOUD_DATABASE_URL") or config.get_main_option("sqlalchemy.url")
|
||||||
or config.get_main_option("sqlalchemy.url")
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@@ -38,7 +37,9 @@ def run_migrations_offline() -> None:
|
|||||||
def run_migrations_online() -> None:
|
def run_migrations_online() -> None:
|
||||||
supplied_connection = config.attributes.get("connection")
|
supplied_connection = config.attributes.get("connection")
|
||||||
if supplied_connection is not None:
|
if supplied_connection is not None:
|
||||||
context.configure(connection=supplied_connection, target_metadata=target_metadata)
|
context.configure(
|
||||||
|
connection=supplied_connection, target_metadata=target_metadata
|
||||||
|
)
|
||||||
with context.begin_transaction():
|
with context.begin_transaction():
|
||||||
context.run_migrations()
|
context.run_migrations()
|
||||||
return
|
return
|
||||||
|
|||||||
+11
-6
@@ -111,7 +111,9 @@ class TaskRunner:
|
|||||||
screenshot = self._planning_screenshot(task.device_id)
|
screenshot = self._planning_screenshot(task.device_id)
|
||||||
steps = self._plan(task.goal, scene, context, screenshot=screenshot)
|
steps = self._plan(task.goal, scene, context, screenshot=screenshot)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
reason = f"{type(exc).__name__}: {exc}" if str(exc) else type(exc).__name__
|
reason = (
|
||||||
|
f"{type(exc).__name__}: {exc}" if str(exc) else type(exc).__name__
|
||||||
|
)
|
||||||
self._update_task(
|
self._update_task(
|
||||||
task,
|
task,
|
||||||
status="failed",
|
status="failed",
|
||||||
@@ -135,7 +137,9 @@ class TaskRunner:
|
|||||||
context=context,
|
context=context,
|
||||||
)
|
)
|
||||||
context.add_step_result(result)
|
context.add_step_result(result)
|
||||||
self._record_step_result(world_handle, context, task, scene, step, result)
|
self._record_step_result(
|
||||||
|
world_handle, context, task, scene, step, result
|
||||||
|
)
|
||||||
if not result.success:
|
if not result.success:
|
||||||
self._update_task(
|
self._update_task(
|
||||||
task,
|
task,
|
||||||
@@ -251,11 +255,10 @@ class TaskRunner:
|
|||||||
def _planner_accepts(self, name: str) -> bool:
|
def _planner_accepts(self, name: str) -> bool:
|
||||||
try:
|
try:
|
||||||
parameters = signature(self.planner.plan).parameters
|
parameters = signature(self.planner.plan).parameters
|
||||||
except (TypeError, ValueError):
|
except TypeError, ValueError:
|
||||||
return True
|
return True
|
||||||
return name in parameters or any(
|
return name in parameters or any(
|
||||||
parameter.kind is Parameter.VAR_KEYWORD
|
parameter.kind is Parameter.VAR_KEYWORD for parameter in parameters.values()
|
||||||
for parameter in parameters.values()
|
|
||||||
)
|
)
|
||||||
|
|
||||||
def _planning_screenshot(self, device_id: str) -> bytes | None:
|
def _planning_screenshot(self, device_id: str) -> bytes | None:
|
||||||
@@ -309,7 +312,9 @@ class TaskRunner:
|
|||||||
"description": step.description,
|
"description": step.description,
|
||||||
"args": step.args,
|
"args": step.args,
|
||||||
},
|
},
|
||||||
result=result.to_dict() if hasattr(result, "to_dict") else {"result": result},
|
result=result.to_dict()
|
||||||
|
if hasattr(result, "to_dict")
|
||||||
|
else {"result": result},
|
||||||
screenshot=screenshot,
|
screenshot=screenshot,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -153,12 +153,15 @@ def test_client_applies_bearer_token_to_every_public_method(tmp_path) -> None:
|
|||||||
assert client.list_devices() == []
|
assert client.list_devices() == []
|
||||||
assert client.list_hosts() == []
|
assert client.list_hosts() == []
|
||||||
assert client.list_plugins() == []
|
assert client.list_plugins() == []
|
||||||
assert client.register_plugin(
|
assert (
|
||||||
name="authenticated-plugin",
|
client.register_plugin(
|
||||||
version="1.0.0",
|
name="authenticated-plugin",
|
||||||
entry_point_kind="tool",
|
version="1.0.0",
|
||||||
target="cloud.store:CloudStore",
|
entry_point_kind="tool",
|
||||||
)["name"] == "authenticated-plugin"
|
target="cloud.store:CloudStore",
|
||||||
|
)["name"]
|
||||||
|
== "authenticated-plugin"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_client_raises_typed_authorization_error_without_exposing_token(
|
def test_client_raises_typed_authorization_error_without_exposing_token(
|
||||||
|
|||||||
@@ -74,7 +74,11 @@ def test_legacy_data_survives_upgrade_and_downgrade(tmp_path) -> None:
|
|||||||
"('task-a', 'goal', null, :constraints, 'queued', null, null, "
|
"('task-a', 'goal', null, :constraints, 'queued', null, null, "
|
||||||
"'2026-01-01T00:00:00+00:00')"
|
"'2026-01-01T00:00:00+00:00')"
|
||||||
),
|
),
|
||||||
{"constraints": json.dumps({"driver_type": None, "capability_tags": []})},
|
{
|
||||||
|
"constraints": json.dumps(
|
||||||
|
{"driver_type": None, "capability_tags": []}
|
||||||
|
)
|
||||||
|
},
|
||||||
)
|
)
|
||||||
finally:
|
finally:
|
||||||
engine.dispose()
|
engine.dispose()
|
||||||
@@ -83,7 +87,9 @@ def test_legacy_data_survives_upgrade_and_downgrade(tmp_path) -> None:
|
|||||||
engine = create_engine(database_url)
|
engine = create_engine(database_url)
|
||||||
try:
|
try:
|
||||||
with engine.connect() as connection:
|
with engine.connect() as connection:
|
||||||
assert connection.scalar(text("select count(*) from host_registrations")) == 1
|
assert (
|
||||||
|
connection.scalar(text("select count(*) from host_registrations")) == 1
|
||||||
|
)
|
||||||
assert connection.scalar(text("select count(*) from scheduled_tasks")) == 1
|
assert connection.scalar(text("select count(*) from scheduled_tasks")) == 1
|
||||||
task_columns = {
|
task_columns = {
|
||||||
column["name"] for column in inspect(engine).get_columns("scheduled_tasks")
|
column["name"] for column in inspect(engine).get_columns("scheduled_tasks")
|
||||||
@@ -96,7 +102,9 @@ def test_legacy_data_survives_upgrade_and_downgrade(tmp_path) -> None:
|
|||||||
engine = create_engine(database_url)
|
engine = create_engine(database_url)
|
||||||
try:
|
try:
|
||||||
with engine.connect() as connection:
|
with engine.connect() as connection:
|
||||||
assert connection.scalar(text("select count(*) from host_registrations")) == 1
|
assert (
|
||||||
|
connection.scalar(text("select count(*) from host_registrations")) == 1
|
||||||
|
)
|
||||||
assert connection.scalar(text("select count(*) from scheduled_tasks")) == 1
|
assert connection.scalar(text("select count(*) from scheduled_tasks")) == 1
|
||||||
task_columns = {
|
task_columns = {
|
||||||
column["name"] for column in inspect(engine).get_columns("scheduled_tasks")
|
column["name"] for column in inspect(engine).get_columns("scheduled_tasks")
|
||||||
|
|||||||
@@ -25,9 +25,10 @@ def test_compose_defines_database_control_plane_and_outbound_host_agent() -> Non
|
|||||||
assert services["host-agent"]["volumes"] == [
|
assert services["host-agent"]["volumes"] == [
|
||||||
"${HOST_AGENT_TASKS_PATH:-./tasks}:/app/tasks"
|
"${HOST_AGENT_TASKS_PATH:-./tasks}:/app/tasks"
|
||||||
]
|
]
|
||||||
assert services["host-agent"]["environment"][
|
assert (
|
||||||
"HOST_AGENT_CONTROL_PLANE_URL"
|
services["host-agent"]["environment"]["HOST_AGENT_CONTROL_PLANE_URL"]
|
||||||
] == "http://cloud-api:8001"
|
== "http://cloud-api:8001"
|
||||||
|
)
|
||||||
assert services["host-agent"]["environment"]["AI_PLANNER_ENABLED"] == (
|
assert services["host-agent"]["environment"]["AI_PLANNER_ENABLED"] == (
|
||||||
"${AI_PLANNER_ENABLED:-false}"
|
"${AI_PLANNER_ENABLED:-false}"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -152,7 +152,9 @@ def test_workflow_runner_failing_planned_goal_marks_run_failed(tmp_path) -> None
|
|||||||
assert run.step_results[0].detail["failure_reason"] == "device offline"
|
assert run.step_results[0].detail["failure_reason"] == "device offline"
|
||||||
|
|
||||||
|
|
||||||
def test_workflow_runner_skill_invocation_executes_resolved_tool_calls(tmp_path) -> None:
|
def test_workflow_runner_skill_invocation_executes_resolved_tool_calls(
|
||||||
|
tmp_path,
|
||||||
|
) -> None:
|
||||||
skill_store, skill = _skill_store()
|
skill_store, skill = _skill_store()
|
||||||
calls: list[dict[str, object]] = []
|
calls: list[dict[str, object]] = []
|
||||||
definition = WorkflowDefinition(
|
definition = WorkflowDefinition(
|
||||||
@@ -227,7 +229,9 @@ def test_workflow_runner_branch_paths_and_sequential_advancement(tmp_path) -> No
|
|||||||
steps=[
|
steps=[
|
||||||
BranchStep(
|
BranchStep(
|
||||||
"branch",
|
"branch",
|
||||||
ConditionSpec("world_variable_equals", {"name": "ready", "value": True}),
|
ConditionSpec(
|
||||||
|
"world_variable_equals", {"name": "ready", "value": True}
|
||||||
|
),
|
||||||
on_true="true-step",
|
on_true="true-step",
|
||||||
on_false="false-step",
|
on_false="false-step",
|
||||||
),
|
),
|
||||||
@@ -320,7 +324,9 @@ def test_workflow_runner_resume_after_branch_crash_advances_to_recorded_target(
|
|||||||
steps=[
|
steps=[
|
||||||
BranchStep(
|
BranchStep(
|
||||||
"branch",
|
"branch",
|
||||||
ConditionSpec("world_variable_equals", {"name": "ready", "value": True}),
|
ConditionSpec(
|
||||||
|
"world_variable_equals", {"name": "ready", "value": True}
|
||||||
|
),
|
||||||
on_true="true-step",
|
on_true="true-step",
|
||||||
on_false="false-step",
|
on_false="false-step",
|
||||||
),
|
),
|
||||||
@@ -469,7 +475,9 @@ def test_workflow_runner_branch_wait_and_skill_combination(tmp_path) -> None:
|
|||||||
steps=[
|
steps=[
|
||||||
BranchStep(
|
BranchStep(
|
||||||
"branch",
|
"branch",
|
||||||
ConditionSpec("world_variable_equals", {"name": "ready", "value": True}),
|
ConditionSpec(
|
||||||
|
"world_variable_equals", {"name": "ready", "value": True}
|
||||||
|
),
|
||||||
on_true="wait",
|
on_true="wait",
|
||||||
on_false="skip",
|
on_false="skip",
|
||||||
),
|
),
|
||||||
|
|||||||
+1
-3
@@ -302,9 +302,7 @@ class WorkflowRunner:
|
|||||||
detail={"reason": f"unknown condition kind: {exc}"},
|
detail={"reason": f"unknown condition kind: {exc}"},
|
||||||
)
|
)
|
||||||
|
|
||||||
elapsed = (
|
elapsed = (datetime.now(started_at.tzinfo) - started_at).total_seconds()
|
||||||
datetime.now(started_at.tzinfo) - started_at
|
|
||||||
).total_seconds()
|
|
||||||
if elapsed >= timeout_seconds:
|
if elapsed >= timeout_seconds:
|
||||||
return WorkflowStepResult(
|
return WorkflowStepResult(
|
||||||
step_id=step.step_id,
|
step_id=step.step_id,
|
||||||
|
|||||||
Reference in New Issue
Block a user