1426 lines
60 KiB
Python
1426 lines
60 KiB
Python
import unittest
|
||
from datetime import datetime
|
||
from unittest.mock import patch
|
||
|
||
from app.api.routes.ops_agent import _build_agent_response
|
||
from app.services.ops_agent_service import (
|
||
_build_agent_runtime_config_bundle,
|
||
_runtime_config_bundle_hash_payload,
|
||
agent_append_job_event,
|
||
agent_complete_job,
|
||
agent_mark_job_started,
|
||
agent_pull_jobs,
|
||
agent_register,
|
||
build_node_agent_bootstrap_plan,
|
||
build_managed_node_handover_bootstrap_plan,
|
||
execute_managed_node_onboarding_acceptance,
|
||
execute_managed_node_onboarding_bootstrap,
|
||
execute_managed_node_onboarding_recovery,
|
||
get_managed_node_delivery_queue,
|
||
get_managed_node_handover,
|
||
get_managed_node_onboarding,
|
||
list_ops_job_events,
|
||
list_managed_node_delivery_queue_records,
|
||
list_managed_nodes_with_agent_state,
|
||
preview_managed_node_onboarding_acceptance,
|
||
preview_managed_node_onboarding_bootstrap,
|
||
preview_managed_node_onboarding_recovery,
|
||
request_managed_node_delivery_queue_action,
|
||
)
|
||
|
||
|
||
class _EmptyCursor:
|
||
def __enter__(self):
|
||
return self
|
||
|
||
def __exit__(self, exc_type, exc, tb):
|
||
return False
|
||
|
||
def execute(self, query, params=None) -> None:
|
||
self.query = query
|
||
self.params = params
|
||
|
||
def fetchall(self):
|
||
return []
|
||
|
||
|
||
class _EmptyConnection:
|
||
def __enter__(self):
|
||
return self
|
||
|
||
def __exit__(self, exc_type, exc, tb):
|
||
return False
|
||
|
||
def cursor(self):
|
||
return _EmptyCursor()
|
||
|
||
|
||
class _SequenceCursor:
|
||
def __init__(self, *, fetchone_values=None, fetchall_values=None):
|
||
self.fetchone_values = list(fetchone_values or [])
|
||
self.fetchall_values = list(fetchall_values or [])
|
||
self.executed = []
|
||
|
||
def __enter__(self):
|
||
return self
|
||
|
||
def __exit__(self, exc_type, exc, tb):
|
||
return False
|
||
|
||
def execute(self, query, params=None) -> None:
|
||
self.executed.append((query, params))
|
||
|
||
def fetchone(self):
|
||
if self.fetchone_values:
|
||
return self.fetchone_values.pop(0)
|
||
return None
|
||
|
||
def fetchall(self):
|
||
if self.fetchall_values:
|
||
return self.fetchall_values.pop(0)
|
||
return []
|
||
|
||
|
||
class _SequenceConnection:
|
||
def __init__(self, cursor):
|
||
self.cursor_obj = cursor
|
||
self.autocommit = True
|
||
self.committed = False
|
||
self.rolled_back = False
|
||
|
||
def __enter__(self):
|
||
return self
|
||
|
||
def __exit__(self, exc_type, exc, tb):
|
||
return False
|
||
|
||
def cursor(self):
|
||
return self.cursor_obj
|
||
|
||
def commit(self) -> None:
|
||
self.committed = True
|
||
|
||
def rollback(self) -> None:
|
||
self.rolled_back = True
|
||
|
||
|
||
class OpsAgentServiceTests(unittest.TestCase):
|
||
def test_runtime_config_bundle_hash_ignores_generated_at(self) -> None:
|
||
bundle_a = {
|
||
"node_code": "mainland-worker-01",
|
||
"thread_count": 100,
|
||
"node_thread_counts": {"mainland-worker-01": 100},
|
||
"runtime_settings": {"worker_log_sync_enabled": True},
|
||
"generated_at": "2026-04-19 17:30:00",
|
||
"config_hash": "stale-hash",
|
||
}
|
||
bundle_b = {
|
||
"node_code": "mainland-worker-01",
|
||
"thread_count": 100,
|
||
"node_thread_counts": {"mainland-worker-01": 100},
|
||
"runtime_settings": {"worker_log_sync_enabled": True},
|
||
"generated_at": "2026-04-19 17:35:00",
|
||
"config_hash": "other-stale-hash",
|
||
}
|
||
|
||
self.assertEqual(
|
||
_runtime_config_bundle_hash_payload(bundle_a),
|
||
_runtime_config_bundle_hash_payload(bundle_b),
|
||
)
|
||
|
||
@patch("app.services.ops_agent_service.get_sensitive_words_payload")
|
||
@patch("app.services.ops_agent_service.get_runtime_settings")
|
||
@patch("app.services.ops_agent_service.get_settings_payload")
|
||
def test_build_agent_runtime_config_bundle_keeps_hash_stable_for_same_config(
|
||
self,
|
||
mock_get_settings_payload,
|
||
mock_get_runtime_settings,
|
||
mock_get_sensitive_words_payload,
|
||
) -> None:
|
||
mock_get_settings_payload.return_value = {
|
||
"detect_options": {"detect_wayback": True},
|
||
"proxy_config": {"enabled": True},
|
||
"thread_count": 100,
|
||
"node_thread_counts": {"mainland-worker-01": 100},
|
||
}
|
||
mock_get_runtime_settings.return_value = {
|
||
"worker_log_sync_enabled": True,
|
||
"worker_log_sync_mode": "full",
|
||
}
|
||
mock_get_sensitive_words_payload.return_value = {
|
||
"text": "foo\nbar",
|
||
"total": 2,
|
||
"items": ["foo", "bar"],
|
||
}
|
||
|
||
with patch(
|
||
"app.services.ops_agent_service._format_time",
|
||
side_effect=["2026-04-19 17:30:00", "2026-04-19 17:35:00"],
|
||
):
|
||
first_bundle = _build_agent_runtime_config_bundle("mainland-worker-01")
|
||
second_bundle = _build_agent_runtime_config_bundle("mainland-worker-01")
|
||
|
||
self.assertNotEqual(first_bundle["generated_at"], second_bundle["generated_at"])
|
||
self.assertEqual(first_bundle["config_hash"], second_bundle["config_hash"])
|
||
|
||
def test_build_agent_response_extracts_detail_code(self) -> None:
|
||
response = _build_agent_response(False, "bad request", {"detail_code": "agent_token_invalid", "foo": "bar"})
|
||
|
||
self.assertEqual(1, response.code)
|
||
self.assertEqual("agent_token_invalid", response.detail_code)
|
||
self.assertEqual("bar", response.data["foo"])
|
||
|
||
@patch("app.services.ops_agent_service.issue_node_agent_token")
|
||
def test_build_node_agent_bootstrap_plan_includes_script_artifacts(self, mock_issue_node_agent_token) -> None:
|
||
mock_issue_node_agent_token.return_value = (
|
||
True,
|
||
"ok",
|
||
{
|
||
"token": "agent-token-example",
|
||
"token_preview": "agent-...mple",
|
||
"record_id": 12,
|
||
"node_code": "mainland-worker-02",
|
||
"expires_at": "2026-04-20 10:00:00",
|
||
"created_at": "2026-04-17 10:00:00",
|
||
},
|
||
)
|
||
|
||
ok, message, data = build_node_agent_bootstrap_plan(
|
||
node_code="mainland-worker-02",
|
||
node_region="mainland",
|
||
node_role="worker",
|
||
issued_by="test",
|
||
control_plane_base_url="http://152.53.37.118:8100/api/v1",
|
||
root_dir="/opt/domaincheck",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("节点 Agent 接入方案已生成", message)
|
||
plan = data["bootstrap_plan"]
|
||
self.assertEqual("http://152.53.37.118:8100", data["control_plane_base_url"])
|
||
self.assertEqual("bootstrap-node-agent-mainland-worker-02.sh", plan["bootstrap_script_name"])
|
||
self.assertEqual("/tmp/bootstrap-node-agent-mainland-worker-02.sh", plan["bootstrap_script_path"])
|
||
self.assertIn("OPS_AGENT_TOKEN=agent-token-example", plan["env_content"])
|
||
self.assertIn("bash \"${INSTALL_SCRIPT}\" \"${ROOT_DIR}\"", plan["bootstrap_script_content"])
|
||
self.assertIn("systemctl restart \"${SERVICE_NAME}\"", plan["bootstrap_script_content"])
|
||
self.assertIn("cat >/tmp/bootstrap-node-agent-mainland-worker-02.sh <<'EOF'", plan["bootstrap_run_script_block"])
|
||
self.assertIn("bash /tmp/bootstrap-node-agent-mainland-worker-02.sh", plan["bootstrap_run_script_block"])
|
||
self.assertIn("curl -s http://152.53.37.118:8100/api/v1/ops/overview", plan["health_check_block"])
|
||
|
||
@patch("app.services.ops_agent_service.list_managed_nodes_with_agent_state")
|
||
def test_get_managed_node_handover_summarizes_pending_bootstrap_node(self, mock_list_managed_nodes_with_agent_state) -> None:
|
||
mock_list_managed_nodes_with_agent_state.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"is_managed": True,
|
||
"is_enabled": True,
|
||
"has_ssh_access": True,
|
||
"ssh_host": "121.204.244.248",
|
||
"ssh_port": 22,
|
||
"ssh_user": "root",
|
||
"auth_mode": "key",
|
||
"agent_state": "pending_bootstrap",
|
||
"agent_state_label": "待接入",
|
||
"agent_state_reason": "已签发有效 Agent Token,但节点尚未完成 register/heartbeat。",
|
||
"remote_access_state": "ssh_ready",
|
||
"remote_access_label": "SSH 可执行",
|
||
"remote_access_reason": "SSH 已备好,可先生成接入方案。",
|
||
"delivery_queue_state": "healthy",
|
||
"delivery_queue_label": "正常",
|
||
"delivery_queue_reason": "当前没有待重试回执。",
|
||
"delivery_queue_pending_count": 0,
|
||
"delivery_queue_dead_letter_count": 0,
|
||
"cluster_status": "online",
|
||
"cluster_current_load": 0,
|
||
"cluster_last_heartbeat_at": "2026-04-18 12:00:00",
|
||
"latest_token": {
|
||
"token_status": "active",
|
||
"expires_at": "2026-04-21 12:00:00",
|
||
},
|
||
"latest_job": {
|
||
"job_code": "ops-job-demo",
|
||
"status": "success",
|
||
},
|
||
"metadata": {},
|
||
}
|
||
]
|
||
}
|
||
|
||
payload = get_managed_node_handover(
|
||
"mainland-worker-01",
|
||
control_plane_base_url="http://152.53.37.118:8100/api/v1",
|
||
)
|
||
|
||
self.assertEqual("mainland-worker-01", payload["node_code"])
|
||
self.assertEqual("pending_bootstrap", payload["stage"]["code"])
|
||
self.assertEqual("execute_bootstrap_plan", payload["stage"]["next_step"]["code"])
|
||
self.assertEqual("http://152.53.37.118:8100", payload["control_plane_base_url"])
|
||
self.assertEqual("/api/v1/ops/nodes/mainland-worker-01/handover/bootstrap-plan", payload["endpoints"]["bootstrap_plan"])
|
||
self.assertTrue(payload["readiness"]["has_ssh_access"])
|
||
self.assertIn("已签发有效 Token", "".join(payload["blocking_items"]))
|
||
|
||
@patch("app.services.ops_agent_service.get_managed_node_handover")
|
||
def test_get_managed_node_onboarding_emits_bootstrap_recovery_decision(self, mock_get_managed_node_handover) -> None:
|
||
mock_get_managed_node_handover.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"summary": "已签发有效 Token,但节点尚未完成 register/heartbeat。",
|
||
"node": {
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
},
|
||
"stage": {
|
||
"code": "pending_bootstrap",
|
||
"label": "待接入",
|
||
"summary": "已签发有效 Token,但节点尚未完成 register/heartbeat。",
|
||
},
|
||
"readiness": {
|
||
"has_ssh_access": True,
|
||
"is_agent_online": False,
|
||
"is_remote_access_ready": True,
|
||
},
|
||
"control_plane_base_url": "http://127.0.0.1:8100",
|
||
"root_dir": "/opt/domaincheck",
|
||
}
|
||
|
||
payload = get_managed_node_onboarding("mainland-worker-01")
|
||
|
||
self.assertEqual("pending_bootstrap", payload["onboarding_stage"]["code"])
|
||
self.assertEqual("bootstrap_run", payload["recovery_decision"]["action"])
|
||
self.assertEqual("bootstrap", payload["recovery_decision"]["window"])
|
||
self.assertIn("node-bootstrap-run", payload["recovery_decision"]["command_hint"])
|
||
|
||
@patch("app.services.ops_agent_service.get_managed_node_handover")
|
||
def test_get_managed_node_onboarding_emits_acceptance_recovery_decision(self, mock_get_managed_node_handover) -> None:
|
||
mock_get_managed_node_handover.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"summary": "Node Agent 已在线,可以继续跑接管验收。",
|
||
"node": {
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
},
|
||
"stage": {
|
||
"code": "ready",
|
||
"label": "已接入",
|
||
"summary": "Node Agent 已在线,可以继续跑接管验收。",
|
||
},
|
||
"readiness": {
|
||
"has_ssh_access": True,
|
||
"is_agent_online": True,
|
||
"is_remote_access_ready": True,
|
||
},
|
||
"control_plane_base_url": "http://127.0.0.1:8100",
|
||
"root_dir": "/opt/domaincheck",
|
||
}
|
||
|
||
payload = get_managed_node_onboarding("mainland-worker-01")
|
||
|
||
self.assertEqual("acceptance_ready", payload["onboarding_stage"]["code"])
|
||
self.assertEqual("run_acceptance", payload["recovery_decision"]["action"])
|
||
self.assertEqual("acceptance", payload["recovery_decision"]["window"])
|
||
self.assertIn("node-acceptance-run", payload["recovery_decision"]["command_hint"])
|
||
|
||
@patch("app.services.ops_agent_service.build_node_agent_bootstrap_plan")
|
||
@patch("app.services.ops_agent_service.get_managed_node_handover")
|
||
def test_build_managed_node_handover_bootstrap_plan_reuses_node_defaults(
|
||
self,
|
||
mock_get_managed_node_handover,
|
||
mock_build_node_agent_bootstrap_plan,
|
||
) -> None:
|
||
mock_get_managed_node_handover.side_effect = [
|
||
{
|
||
"node": {
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
},
|
||
"control_plane_base_url": "http://152.53.37.118:8100",
|
||
"root_dir": "/srv/domaincheck",
|
||
"bootstrap_defaults": {
|
||
"node_region": "mainland",
|
||
"node_role": "worker",
|
||
"control_plane_base_url": "http://152.53.37.118:8100",
|
||
"root_dir": "/srv/domaincheck",
|
||
"expires_in_hours": 72,
|
||
"metadata": {
|
||
"issued_from": "ops-managed-node-handover",
|
||
},
|
||
},
|
||
},
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"stage": {
|
||
"code": "pending_bootstrap",
|
||
},
|
||
},
|
||
]
|
||
mock_build_node_agent_bootstrap_plan.return_value = (
|
||
True,
|
||
"节点 Agent 接入方案已生成",
|
||
{
|
||
"token_preview": "agent-...mple",
|
||
"bootstrap_plan": {
|
||
"bootstrap_script_name": "bootstrap-node-agent-mainland-worker-01.sh",
|
||
},
|
||
},
|
||
)
|
||
|
||
ok, message, data = build_managed_node_handover_bootstrap_plan(
|
||
node_code="mainland-worker-01",
|
||
issued_by="web-ui",
|
||
expires_in_hours=48,
|
||
metadata={"issued_from": "ops-center-manual"},
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("节点 Agent 接入方案已生成", message)
|
||
mock_build_node_agent_bootstrap_plan.assert_called_once()
|
||
called_kwargs = mock_build_node_agent_bootstrap_plan.call_args.kwargs
|
||
self.assertEqual("mainland-worker-01", called_kwargs["node_code"])
|
||
self.assertEqual("mainland", called_kwargs["node_region"])
|
||
self.assertEqual("worker", called_kwargs["node_role"])
|
||
self.assertEqual("http://152.53.37.118:8100", called_kwargs["control_plane_base_url"])
|
||
self.assertEqual("/srv/domaincheck", called_kwargs["root_dir"])
|
||
self.assertEqual("ops-center-manual", called_kwargs["metadata"]["issued_from"])
|
||
self.assertEqual("managed-node-handover", data["requested_via"])
|
||
self.assertEqual("pending_bootstrap", data["handover"]["stage"]["code"])
|
||
|
||
@patch("app.services.ops_playbook_service.preview_ops_playbook")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_preview_managed_node_onboarding_bootstrap_uses_control_plane_playbook(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_preview_ops_playbook,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-controller-01",
|
||
"bootstrap_plan_request": {
|
||
"node_code": "mainland-controller-01",
|
||
"node_region": "mainland",
|
||
"node_role": "control",
|
||
"control_plane_base_url": "http://127.0.0.1:8100",
|
||
"root_dir": "/opt/domaincheck",
|
||
},
|
||
"acceptance": {"status_code": "blocked"},
|
||
}
|
||
mock_preview_ops_playbook.return_value = (True, "ok", {"playbook": {"key": "onboarding.bootstrap"}})
|
||
|
||
ok, message, data = preview_managed_node_onboarding_bootstrap(
|
||
node_code="mainland-controller-01",
|
||
requested_by="tester",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("ok", message)
|
||
self.assertEqual("mainland-controller-01", data["node_code"])
|
||
self.assertEqual("control-plane", data["execution_mode"])
|
||
mock_preview_ops_playbook.assert_called_once_with(
|
||
{
|
||
"playbook_key": "onboarding.bootstrap",
|
||
"node_codes": ["mainland-controller-01"],
|
||
"execution_mode": "control-plane",
|
||
"auto_approve": True,
|
||
"requested_by": "tester",
|
||
"payload": {
|
||
"node_code": "mainland-controller-01",
|
||
"node_region": "mainland",
|
||
"node_role": "control",
|
||
"control_plane_base_url": "http://127.0.0.1:8100",
|
||
"root_dir": "/opt/domaincheck",
|
||
},
|
||
}
|
||
)
|
||
|
||
@patch("app.services.ops_playbook_service.execute_ops_playbook")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_execute_managed_node_onboarding_bootstrap_uses_control_plane_playbook(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_execute_ops_playbook,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-controller-01",
|
||
"bootstrap_plan_request": {
|
||
"node_code": "mainland-controller-01",
|
||
"node_region": "mainland",
|
||
"node_role": "control",
|
||
"control_plane_base_url": "http://127.0.0.1:8100",
|
||
"root_dir": "/opt/domaincheck",
|
||
},
|
||
"acceptance": {"status_code": "blocked"},
|
||
}
|
||
mock_execute_ops_playbook.return_value = (
|
||
True,
|
||
"created",
|
||
{
|
||
"run_code": "pbr-demo",
|
||
"playbook_run": {
|
||
"run_code": "pbr-demo",
|
||
"focus_summary": "接入工单已生成,节点待执行 bootstrap 脚本并回连 Agent。",
|
||
"focus_ref": {"kind": "playbook_run", "run_code": "pbr-demo"},
|
||
},
|
||
},
|
||
)
|
||
|
||
ok, message, data = execute_managed_node_onboarding_bootstrap(
|
||
node_code="mainland-controller-01",
|
||
requested_by="tester",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("created", message)
|
||
self.assertEqual("control-plane", data["execution_mode"])
|
||
self.assertEqual("acceptance_ready", data["follow_up"]["next_stage_code"])
|
||
self.assertIn("node-recover", data["follow_up"]["recommended_commands"][0])
|
||
self.assertIn("待执行 bootstrap 脚本", data["follow_up"]["summary"])
|
||
mock_execute_ops_playbook.assert_called_once()
|
||
called_payload = mock_execute_ops_playbook.call_args.args[0]
|
||
self.assertEqual("onboarding.bootstrap", called_payload["playbook_key"])
|
||
self.assertEqual(["mainland-controller-01"], called_payload["node_codes"])
|
||
self.assertEqual("control-plane", called_payload["execution_mode"])
|
||
self.assertTrue(called_payload["auto_approve"])
|
||
self.assertEqual("tester", called_payload["requested_by"])
|
||
|
||
@patch("app.services.ops_playbook_service.preview_ops_playbook")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_preview_managed_node_onboarding_acceptance_uses_acceptance_execution_mode(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_preview_ops_playbook,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"acceptance": {
|
||
"status_code": "ssh_only",
|
||
"execution_mode": "ssh",
|
||
},
|
||
}
|
||
mock_preview_ops_playbook.return_value = (True, "ok", {"playbook": {"key": "onboarding.acceptance"}})
|
||
|
||
ok, message, data = preview_managed_node_onboarding_acceptance(
|
||
node_code="mainland-worker-01",
|
||
requested_by="tester",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("ok", message)
|
||
self.assertEqual("ssh", data["execution_mode"])
|
||
mock_preview_ops_playbook.assert_called_once_with(
|
||
{
|
||
"playbook_key": "onboarding.acceptance",
|
||
"node_codes": ["mainland-worker-01"],
|
||
"execution_mode": "ssh",
|
||
"auto_approve": True,
|
||
"requested_by": "tester",
|
||
}
|
||
)
|
||
|
||
@patch("app.services.ops_playbook_service.execute_ops_playbook")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_execute_managed_node_onboarding_acceptance_blocks_when_not_ready(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_execute_ops_playbook,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"acceptance": {
|
||
"status_code": "blocked",
|
||
"summary": "当前还没进入标准接管就绪状态",
|
||
},
|
||
}
|
||
|
||
ok, message, data = execute_managed_node_onboarding_acceptance(
|
||
node_code="mainland-worker-01",
|
||
requested_by="tester",
|
||
)
|
||
|
||
self.assertFalse(ok)
|
||
self.assertIn("当前还没进入标准接管就绪状态", message)
|
||
self.assertIn("onboarding", data)
|
||
mock_execute_ops_playbook.assert_not_called()
|
||
|
||
@patch("app.services.ops_agent_service.preview_managed_node_onboarding_bootstrap")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_preview_managed_node_onboarding_recovery_routes_to_bootstrap(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_preview_bootstrap,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"recovery_decision": {
|
||
"action": "bootstrap_run",
|
||
"summary": "建议先签发 onboarding.bootstrap。",
|
||
},
|
||
}
|
||
mock_preview_bootstrap.return_value = (True, "ok", {"playbook_preview": {"playbook": {"key": "onboarding.bootstrap"}}})
|
||
|
||
ok, message, data = preview_managed_node_onboarding_recovery(node_code="mainland-worker-01", requested_by="tester")
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("bootstrap_run", data["recovery_action"])
|
||
self.assertEqual("ok", message)
|
||
mock_preview_bootstrap.assert_called_once()
|
||
|
||
@patch("app.services.ops_agent_service.execute_managed_node_onboarding_acceptance")
|
||
@patch("app.services.ops_agent_service.get_managed_node_onboarding")
|
||
def test_execute_managed_node_onboarding_recovery_routes_to_acceptance(
|
||
self,
|
||
mock_get_managed_node_onboarding,
|
||
mock_execute_acceptance,
|
||
) -> None:
|
||
mock_get_managed_node_onboarding.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"recovery_decision": {
|
||
"action": "run_acceptance",
|
||
"summary": "建议直接执行 onboarding.acceptance。",
|
||
},
|
||
}
|
||
mock_execute_acceptance.return_value = (
|
||
True,
|
||
"created",
|
||
{
|
||
"playbook_run": {"run_code": "pbr-acceptance"},
|
||
"follow_up": {
|
||
"next_stage_code": "ready",
|
||
"summary": "接管验收已发起;接下来应继续确认 health / node-agent / worker 三段验收全部通过。",
|
||
"recommended_commands": [
|
||
"bash domain-api/deploy/multi-region/drive_ops_center.sh node-onboarding http://127.0.0.1:8100 mainland-worker-01",
|
||
"bash domain-api/deploy/multi-region/drive_ops_center.sh doctor http://127.0.0.1:8100",
|
||
],
|
||
},
|
||
},
|
||
)
|
||
|
||
ok, message, data = execute_managed_node_onboarding_recovery(node_code="mainland-worker-01", requested_by="tester")
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("run_acceptance", data["recovery_action"])
|
||
self.assertEqual("created", message)
|
||
self.assertEqual("ready", data["selected_run"]["follow_up"]["next_stage_code"])
|
||
mock_execute_acceptance.assert_called_once()
|
||
|
||
@patch("app.services.runtime_status_service.get_runtime_status")
|
||
@patch("app.services.cluster_runtime_service.get_cluster_snapshot")
|
||
@patch("app.services.ops_job_service.list_managed_nodes")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
@patch("app.services.ops_agent_service.ensure_ops_agent_schema")
|
||
def test_list_managed_nodes_with_agent_state_merges_detect_participation(
|
||
self,
|
||
mock_ensure_ops_agent_schema,
|
||
mock_get_db,
|
||
mock_list_managed_nodes,
|
||
mock_get_cluster_snapshot,
|
||
mock_get_runtime_status,
|
||
) -> None:
|
||
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
mock_ensure_ops_agent_schema.return_value = None
|
||
mock_get_db.return_value = _EmptyConnection()
|
||
mock_list_managed_nodes.return_value = [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"title": "Mainland Worker 01",
|
||
"is_enabled": True,
|
||
"metadata": {},
|
||
"last_seen_at": now,
|
||
},
|
||
{
|
||
"node_code": "mainland-controller-01",
|
||
"region": "mainland",
|
||
"role": "control",
|
||
"title": "Mainland Controller 01",
|
||
"is_enabled": True,
|
||
"metadata": {},
|
||
"last_seen_at": now,
|
||
},
|
||
]
|
||
mock_get_cluster_snapshot.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"status": "busy",
|
||
"current_load": 4,
|
||
"last_heartbeat_at": now,
|
||
"is_effective_worker": True,
|
||
"detect_participating": True,
|
||
},
|
||
{
|
||
"node_code": "mainland-controller-01",
|
||
"region": "mainland",
|
||
"role": "control",
|
||
"status": "online",
|
||
"current_load": 0,
|
||
"last_heartbeat_at": now,
|
||
"is_effective_worker": True,
|
||
"detect_participating": False,
|
||
},
|
||
]
|
||
}
|
||
mock_get_runtime_status.return_value = {
|
||
"detect": {
|
||
"participating_nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"participation_state": "running",
|
||
"participation_label": "执行中",
|
||
"participation_reason": "当前正有 4 个任务线程在跑",
|
||
"is_current_participant": True,
|
||
"is_dispatch_active": True,
|
||
"items_total": 50,
|
||
"items_claimed": 12,
|
||
"items_running": 4,
|
||
"items_completed": 8,
|
||
"items_failed": 1,
|
||
"processed_recent": 9,
|
||
"processed_per_minute": 0.6,
|
||
}
|
||
],
|
||
"non_participating_nodes": [
|
||
{
|
||
"node_code": "mainland-controller-01",
|
||
"participation_state": "standby",
|
||
"participation_label": "在线待命",
|
||
"participation_reason": "节点在线但当前没有领任务",
|
||
"is_current_participant": False,
|
||
"is_dispatch_active": False,
|
||
"items_total": 0,
|
||
"items_claimed": 0,
|
||
"items_running": 0,
|
||
"items_completed": 0,
|
||
"items_failed": 0,
|
||
"processed_recent": 0,
|
||
"processed_per_minute": 0,
|
||
}
|
||
],
|
||
}
|
||
}
|
||
|
||
payload = list_managed_nodes_with_agent_state()
|
||
|
||
self.assertEqual(2, len(payload["nodes"]))
|
||
node_map = {item["node_code"]: item for item in payload["nodes"]}
|
||
self.assertTrue(node_map["mainland-worker-01"]["is_current_participant"])
|
||
self.assertTrue(node_map["mainland-worker-01"]["cluster_detect_participating"])
|
||
self.assertEqual("running", node_map["mainland-worker-01"]["participation_state"])
|
||
self.assertEqual("执行中", node_map["mainland-worker-01"]["participation_label"])
|
||
self.assertTrue(node_map["mainland-worker-01"]["is_dispatch_active"])
|
||
self.assertEqual(50, node_map["mainland-worker-01"]["items_total"])
|
||
self.assertEqual(4, node_map["mainland-worker-01"]["items_running"])
|
||
self.assertEqual("standby", node_map["mainland-controller-01"]["participation_state"])
|
||
self.assertEqual("在线待命", node_map["mainland-controller-01"]["participation_label"])
|
||
self.assertFalse(node_map["mainland-controller-01"]["is_current_participant"])
|
||
self.assertEqual(1, payload["summary"]["participating"])
|
||
self.assertEqual(1, payload["summary"]["dispatch_active"])
|
||
self.assertEqual(1, payload["summary"]["standby"])
|
||
self.assertEqual(1, payload["summary"]["non_participating"])
|
||
|
||
@patch("app.services.runtime_status_service.get_runtime_status")
|
||
@patch("app.services.cluster_runtime_service.get_cluster_snapshot")
|
||
@patch("app.services.ops_job_service.list_managed_nodes")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
@patch("app.services.ops_agent_service.ensure_ops_agent_schema")
|
||
def test_list_managed_nodes_with_agent_state_marks_ssh_ready_without_agent(
|
||
self,
|
||
mock_ensure_ops_agent_schema,
|
||
mock_get_db,
|
||
mock_list_managed_nodes,
|
||
mock_get_cluster_snapshot,
|
||
mock_get_runtime_status,
|
||
) -> None:
|
||
mock_ensure_ops_agent_schema.return_value = None
|
||
mock_get_db.return_value = _EmptyConnection()
|
||
mock_list_managed_nodes.return_value = [
|
||
{
|
||
"node_code": "mainland-worker-02",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"title": "Mainland Worker 02",
|
||
"ssh_host": "121.204.244.249",
|
||
"ssh_port": 22,
|
||
"ssh_user": "root",
|
||
"auth_mode": "key",
|
||
"deploy_channel": "stable",
|
||
"is_enabled": True,
|
||
"metadata": {},
|
||
"last_seen_at": "",
|
||
}
|
||
]
|
||
mock_get_cluster_snapshot.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-02",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"status": "online",
|
||
"current_load": 0,
|
||
"hostname": "worker-02",
|
||
"ip": "121.204.244.249",
|
||
"last_heartbeat_at": "2026-04-18 10:00:00",
|
||
"is_effective_worker": True,
|
||
"detect_participating": False,
|
||
}
|
||
]
|
||
}
|
||
mock_get_runtime_status.return_value = {
|
||
"detect": {
|
||
"participating_nodes": [],
|
||
"non_participating_nodes": [],
|
||
}
|
||
}
|
||
|
||
payload = list_managed_nodes_with_agent_state()
|
||
|
||
self.assertEqual(1, len(payload["nodes"]))
|
||
row = payload["nodes"][0]
|
||
self.assertFalse(row["is_agent_online"])
|
||
self.assertTrue(row["has_ssh_access"])
|
||
self.assertEqual("SSH 已备好", row["ssh_access_label"])
|
||
self.assertEqual("ssh_ready", row["remote_access_state"])
|
||
self.assertEqual("SSH 可执行", row["remote_access_label"])
|
||
self.assertTrue(row["is_remote_access_ready"])
|
||
self.assertEqual(1, payload["summary"]["ssh_ready"])
|
||
self.assertEqual(0, payload["summary"]["agent_ready"])
|
||
self.assertEqual(1, payload["summary"]["remote_access_ready"])
|
||
|
||
@patch("app.services.runtime_status_service.get_runtime_status")
|
||
@patch("app.services.cluster_runtime_service.get_cluster_snapshot")
|
||
@patch("app.services.ops_job_service.list_managed_nodes")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
@patch("app.services.ops_agent_service.ensure_ops_agent_schema")
|
||
def test_list_managed_nodes_with_agent_state_prefers_non_loopback_cluster_identity_and_infers_worker_online(
|
||
self,
|
||
mock_ensure_ops_agent_schema,
|
||
mock_get_db,
|
||
mock_list_managed_nodes,
|
||
mock_get_cluster_snapshot,
|
||
mock_get_runtime_status,
|
||
) -> None:
|
||
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
mock_ensure_ops_agent_schema.return_value = None
|
||
mock_get_db.return_value = _EmptyConnection()
|
||
mock_list_managed_nodes.return_value = [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"title": "Mainland Worker 01",
|
||
"is_enabled": True,
|
||
"metadata": {
|
||
"hostname": "S244-248",
|
||
"ip": "121.204.244.248",
|
||
},
|
||
"last_seen_at": now,
|
||
}
|
||
]
|
||
mock_get_cluster_snapshot.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"status": "busy",
|
||
"current_load": 19,
|
||
"hostname": "localhost",
|
||
"ip": "127.0.0.1",
|
||
"last_heartbeat_at": now,
|
||
"is_effective_worker": True,
|
||
"detect_participating": False,
|
||
"metadata": {
|
||
"worker_online": False,
|
||
"active_threads": 19,
|
||
"max_threads": 1200,
|
||
"phase_label": "检测中",
|
||
"phase_detail": "Worker 正在处理 162 个检测任务",
|
||
},
|
||
}
|
||
]
|
||
}
|
||
mock_get_runtime_status.return_value = {
|
||
"detect": {
|
||
"participating_nodes": [],
|
||
"non_participating_nodes": [],
|
||
}
|
||
}
|
||
|
||
payload = list_managed_nodes_with_agent_state()
|
||
|
||
self.assertEqual(1, len(payload["nodes"]))
|
||
row = payload["nodes"][0]
|
||
self.assertEqual("S244-248", row["cluster_hostname"])
|
||
self.assertEqual("121.204.244.248", row["cluster_ip"])
|
||
self.assertTrue(row["detect_runtime"]["worker_online"])
|
||
self.assertTrue(row["detect_runtime"]["detect_participating"])
|
||
self.assertEqual(19, row["detect_runtime"]["active_threads"])
|
||
|
||
def test_agent_register_requires_node_code_with_detail_code(self) -> None:
|
||
ok, message, data = agent_register({}, token="ignored")
|
||
|
||
self.assertFalse(ok)
|
||
self.assertEqual("node_code 不能为空", message)
|
||
self.assertEqual("agent_node_code_required", data["detail_code"])
|
||
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_ops_job_event")
|
||
@patch("app.services.ops_agent_service.append_ops_job_event")
|
||
def test_agent_append_job_event_preserves_client_event_id(
|
||
self,
|
||
mock_append_ops_job_event,
|
||
mock_get_ops_job_event,
|
||
mock_authenticate_agent_token,
|
||
) -> None:
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_append_ops_job_event.return_value = 88
|
||
mock_get_ops_job_event.return_value = {
|
||
"id": 88,
|
||
"event_key": "job-event:88",
|
||
"focus_ref": {"kind": "ops_job_event", "job_id": 12, "focus_event_id": 88},
|
||
}
|
||
|
||
ok, message, data = agent_append_job_event(
|
||
12,
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"event_type": "executor_received",
|
||
"message": "accepted",
|
||
"client_event_id": "evt-001",
|
||
"summary_text": "已接单",
|
||
"focus_ref": {"kind": "ops_job", "job_id": 12},
|
||
},
|
||
token="token",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("事件已写入", message)
|
||
self.assertEqual("evt-001", data["client_event_id"])
|
||
self.assertEqual(88, data["event"]["id"])
|
||
self.assertEqual("ops_job_event", data["event"]["focus_ref"]["kind"])
|
||
self.assertEqual("evt-001", mock_append_ops_job_event.call_args.kwargs["client_event_id"])
|
||
self.assertEqual("已接单", mock_append_ops_job_event.call_args.kwargs["payload"]["summary_text"])
|
||
self.assertEqual("ops_job", mock_append_ops_job_event.call_args.kwargs["payload"]["focus_ref"]["kind"])
|
||
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
@patch("app.services.ops_agent_service.ensure_ops_agent_schema")
|
||
def test_list_ops_job_events_adds_contract_aliases(
|
||
self,
|
||
mock_ensure_ops_agent_schema,
|
||
mock_get_db,
|
||
) -> None:
|
||
mock_ensure_ops_agent_schema.return_value = None
|
||
cursor = _SequenceCursor(
|
||
fetchall_values=[
|
||
[
|
||
(
|
||
301,
|
||
101,
|
||
201,
|
||
"mainland-worker-01",
|
||
"evt-001",
|
||
"job_dispatch_requested",
|
||
"warning",
|
||
"任务等待 node-agent 拉取: 101",
|
||
'{"dispatch": true, "summary_text": "等待 Agent 拉取", "occurred_at": "2026-04-18 06:45:01", "focus_ref": {"kind": "ops_job", "job_id": 101}}',
|
||
datetime(2026, 4, 18, 6, 45, 3),
|
||
)
|
||
]
|
||
]
|
||
)
|
||
mock_get_db.return_value = _SequenceConnection(cursor)
|
||
|
||
events = list_ops_job_events(101, limit=20)
|
||
|
||
self.assertEqual(1, len(events))
|
||
row = events[0]
|
||
self.assertEqual("job-event:301", row["event_key"])
|
||
self.assertEqual("warning", row["level"])
|
||
self.assertEqual("警告", row["level_label"])
|
||
self.assertEqual("任务等待 node-agent 拉取: 101", row["summary"])
|
||
self.assertEqual("等待 Agent 拉取", row["summary_text"])
|
||
self.assertTrue(row["has_payload"])
|
||
self.assertEqual("2026-04-18 06:45:01", row["occurred_at"])
|
||
self.assertEqual("ops_job", row["focus_ref"]["kind"])
|
||
self.assertEqual("job_detail", row["ui_intent"]["kind"])
|
||
self.assertEqual(101, row["ui_intent"]["job_id"])
|
||
self.assertEqual(301, row["ui_intent"]["focus_event_id"])
|
||
|
||
@patch("app.services.ops_agent_service.get_ops_job")
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
def test_agent_complete_job_is_idempotent_for_same_client_request_id(
|
||
self,
|
||
mock_get_db,
|
||
mock_authenticate_agent_token,
|
||
mock_get_ops_job,
|
||
) -> None:
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_get_ops_job.return_value = {"id": 12, "status": "success"}
|
||
cursor = _SequenceCursor(fetchone_values=[(12, "complete-001")])
|
||
connection = _SequenceConnection(cursor)
|
||
mock_get_db.return_value = connection
|
||
|
||
ok, message, data = agent_complete_job(
|
||
12,
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"status": "success",
|
||
"result": {"ok": True},
|
||
"client_request_id": "complete-001",
|
||
},
|
||
token="token",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("任务结果已幂等回写", message)
|
||
self.assertTrue(data["deduplicated"])
|
||
self.assertEqual("complete-001", data["client_request_id"])
|
||
self.assertEqual(12, data["completion_summary"]["job_id"])
|
||
self.assertTrue(data["completion_summary"]["deduplicated"])
|
||
self.assertTrue(connection.rolled_back)
|
||
|
||
@patch("app.services.ops_agent_service.append_ops_job_event")
|
||
@patch("app.services.ops_agent_service.get_ops_job")
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
def test_agent_pull_jobs_returns_agent_job_envelope(
|
||
self,
|
||
mock_get_db,
|
||
mock_authenticate_agent_token,
|
||
mock_get_ops_job,
|
||
mock_append_ops_job_event,
|
||
) -> None:
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_append_ops_job_event.return_value = 501
|
||
mock_get_ops_job.return_value = {
|
||
"id": 101,
|
||
"job_code": "ops-20260418-000101",
|
||
"action": "logs.collect",
|
||
"target_node_code": "mainland-worker-01",
|
||
"status": "dispatching",
|
||
"execution_mode": "remote-agent",
|
||
"execution_mode_label": "Node Agent",
|
||
"requested_by": "web-ui",
|
||
"payload": {"service_name": "domaincheck-worker", "lines": 120},
|
||
"metadata": {"source": "ops-center"},
|
||
"policy": {"timeout_seconds": 120, "auto_approve": True},
|
||
"summary": "收集 Worker 日志",
|
||
"rollout_id": 0,
|
||
"focus_ref": {"kind": "ops_job", "job_id": 101, "job_code": "ops-20260418-000101"},
|
||
"steps": [
|
||
{
|
||
"id": 201,
|
||
"step_key": "worker_logs",
|
||
"title": "收集 Worker 日志",
|
||
"node_code": "mainland-worker-01",
|
||
"status": "dispatching",
|
||
}
|
||
],
|
||
}
|
||
cursor = _SequenceCursor(fetchall_values=[[(101,)]])
|
||
connection = _SequenceConnection(cursor)
|
||
mock_get_db.return_value = connection
|
||
|
||
ok, message, data = agent_pull_jobs(
|
||
{"node_code": "mainland-worker-01"},
|
||
token="token",
|
||
limit=2,
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("ok", message)
|
||
self.assertEqual("ops-agent/v1", data["protocol_version"])
|
||
self.assertEqual("agent_job", data["envelope_type"])
|
||
self.assertEqual(1, data["count"])
|
||
envelope = data["jobs"][0]
|
||
self.assertEqual(101, envelope["job_id"])
|
||
self.assertEqual("inspection", envelope["job_type"])
|
||
self.assertEqual("worker_logs", envelope["step_key"])
|
||
self.assertEqual("收集 Worker 日志", envelope["step_title"])
|
||
self.assertEqual(120, envelope["policy"]["timeout_seconds"])
|
||
self.assertEqual("ops_job", envelope["focus_ref"]["kind"])
|
||
self.assertEqual("logs.collect", envelope["job_ref"]["action"])
|
||
self.assertEqual(201, envelope["step_ref"]["step_id"])
|
||
self.assertEqual("", envelope["release_context"]["release_version"])
|
||
self.assertTrue(connection.committed)
|
||
|
||
@patch("app.services.runtime_status_service.get_runtime_status")
|
||
@patch("app.services.cluster_runtime_service.get_cluster_snapshot")
|
||
@patch("app.services.ops_job_service.list_managed_nodes")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
@patch("app.services.ops_agent_service.ensure_ops_agent_schema")
|
||
def test_list_managed_nodes_with_agent_state_includes_delivery_queue_snapshot(
|
||
self,
|
||
mock_ensure_ops_agent_schema,
|
||
mock_get_db,
|
||
mock_list_managed_nodes,
|
||
mock_get_cluster_snapshot,
|
||
mock_get_runtime_status,
|
||
) -> None:
|
||
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
mock_ensure_ops_agent_schema.return_value = None
|
||
mock_get_db.return_value = _EmptyConnection()
|
||
mock_list_managed_nodes.return_value = [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"title": "Mainland Worker 01",
|
||
"is_enabled": True,
|
||
"metadata": {
|
||
"delivery_queue": {
|
||
"state": "dead_letter",
|
||
"label": "死信 2",
|
||
"reason": "当前存在 2 条死信记录",
|
||
"pending_count": 3,
|
||
"dead_letter_count": 2,
|
||
"oldest_pending_at": "2026-04-18 05:00:00",
|
||
"oldest_dead_letter_at": "2026-04-18 05:01:00",
|
||
"last_flush_at": "2026-04-18 05:02:00",
|
||
}
|
||
},
|
||
"last_seen_at": now,
|
||
}
|
||
]
|
||
mock_get_cluster_snapshot.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"status": "online",
|
||
"current_load": 0,
|
||
"last_heartbeat_at": now,
|
||
"is_effective_worker": True,
|
||
"detect_participating": False,
|
||
}
|
||
]
|
||
}
|
||
mock_get_runtime_status.return_value = {
|
||
"detect": {
|
||
"participating_nodes": [],
|
||
"non_participating_nodes": [],
|
||
}
|
||
}
|
||
|
||
payload = list_managed_nodes_with_agent_state()
|
||
|
||
self.assertEqual(1, len(payload["nodes"]))
|
||
row = payload["nodes"][0]
|
||
self.assertEqual("dead_letter", row["delivery_queue_state"])
|
||
self.assertEqual("死信 2", row["delivery_queue_label"])
|
||
self.assertEqual(3, row["delivery_queue_pending_count"])
|
||
self.assertEqual(2, row["delivery_queue_dead_letter_count"])
|
||
self.assertEqual("2026-04-18 05:01:00", row["delivery_queue_oldest_dead_letter_at"])
|
||
self.assertEqual(1, payload["summary"]["queue_dead_letter_nodes"])
|
||
self.assertEqual(3, payload["summary"]["queue_pending_records"])
|
||
self.assertEqual(2, payload["summary"]["queue_dead_letter_records"])
|
||
|
||
@patch("app.services.ops_agent_service.list_managed_nodes_with_agent_state")
|
||
def test_get_managed_node_delivery_queue_returns_summary_and_head_records(
|
||
self,
|
||
mock_list_managed_nodes_with_agent_state,
|
||
) -> None:
|
||
mock_list_managed_nodes_with_agent_state.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"title": "Mainland Worker 01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"is_managed": True,
|
||
"is_enabled": True,
|
||
"is_agent_online": True,
|
||
"agent_state": "online",
|
||
"agent_state_label": "已接管",
|
||
"remote_access_state": "hybrid_ready",
|
||
"remote_access_label": "Agent + SSH",
|
||
"delivery_queue_state": "dead_letter",
|
||
"delivery_queue_label": "死信 2",
|
||
"delivery_queue_reason": "当前存在 2 条死信记录,建议优先查看节点日志或诊断编排。",
|
||
"delivery_queue_pending_count": 3,
|
||
"delivery_queue_dead_letter_count": 2,
|
||
"delivery_queue_last_flush_at": "2026-04-18 06:02:00",
|
||
"delivery_queue_oldest_pending_at": "2026-04-18 06:00:00",
|
||
"delivery_queue_oldest_pending_request_id": "complete-ops-123-1",
|
||
"delivery_queue_oldest_pending_kind": "job_complete",
|
||
"delivery_queue_oldest_dead_letter_at": "2026-04-18 06:01:00",
|
||
"delivery_queue_oldest_dead_letter_request_id": "event-ops-123-2",
|
||
"delivery_queue_oldest_dead_letter_kind": "job_event",
|
||
}
|
||
]
|
||
}
|
||
|
||
payload = get_managed_node_delivery_queue("mainland-worker-01")
|
||
|
||
self.assertEqual("mainland-worker-01", payload["node_code"])
|
||
self.assertEqual("dead_letter", payload["summary"]["state"])
|
||
self.assertEqual(3, payload["summary"]["pending_count"])
|
||
self.assertEqual(2, payload["summary"]["dead_letter_count"])
|
||
self.assertEqual("head_only", payload["record_visibility"])
|
||
self.assertTrue(payload["capabilities"]["can_flush"])
|
||
self.assertTrue(payload["capabilities"]["can_replay"])
|
||
self.assertTrue(payload["capabilities"]["can_discard"])
|
||
self.assertEqual("head_only", payload["records_capability"]["record_visibility"])
|
||
self.assertFalse(payload["records_capability"]["full_records_available"])
|
||
self.assertEqual("job_complete", payload["head_records"]["pending"]["request_kind"])
|
||
self.assertEqual("job_event", payload["head_records"]["dead_letter"]["request_kind"])
|
||
|
||
@patch("app.services.ops_agent_service.list_managed_nodes_with_agent_state")
|
||
def test_list_managed_node_delivery_queue_records_filters_head_records_by_state(
|
||
self,
|
||
mock_list_managed_nodes_with_agent_state,
|
||
) -> None:
|
||
mock_list_managed_nodes_with_agent_state.return_value = {
|
||
"nodes": [
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"title": "Mainland Worker 01",
|
||
"region": "mainland",
|
||
"role": "worker",
|
||
"is_managed": True,
|
||
"is_enabled": True,
|
||
"is_agent_online": True,
|
||
"agent_state": "online",
|
||
"agent_state_label": "已接管",
|
||
"remote_access_state": "hybrid_ready",
|
||
"remote_access_label": "Agent + SSH",
|
||
"delivery_queue_state": "dead_letter",
|
||
"delivery_queue_label": "死信 1",
|
||
"delivery_queue_reason": "当前存在 1 条死信记录。",
|
||
"delivery_queue_pending_count": 1,
|
||
"delivery_queue_dead_letter_count": 1,
|
||
"delivery_queue_last_flush_at": "2026-04-18 06:02:00",
|
||
"delivery_queue_oldest_pending_at": "2026-04-18 06:00:00",
|
||
"delivery_queue_oldest_pending_request_id": "complete-ops-123-1",
|
||
"delivery_queue_oldest_pending_kind": "job_complete",
|
||
"delivery_queue_oldest_dead_letter_at": "2026-04-18 06:01:00",
|
||
"delivery_queue_oldest_dead_letter_request_id": "event-ops-123-2",
|
||
"delivery_queue_oldest_dead_letter_kind": "job_event",
|
||
}
|
||
]
|
||
}
|
||
|
||
payload = list_managed_node_delivery_queue_records("mainland-worker-01", state="dead_letter", limit=20)
|
||
|
||
self.assertEqual("mainland-worker-01", payload["node_code"])
|
||
self.assertEqual(1, payload["summary"]["total"])
|
||
self.assertEqual(1, payload["summary"]["returned_records"])
|
||
self.assertEqual("head_only", payload["summary"]["record_visibility"])
|
||
self.assertEqual("head_only", payload["record_visibility"])
|
||
self.assertTrue(payload["capabilities"]["can_replay"])
|
||
self.assertEqual("dead_letter", payload["records"][0]["state"])
|
||
self.assertEqual("event-ops-123-2", payload["records"][0]["client_request_id"])
|
||
|
||
@patch("app.services.ops_agent_service.create_ops_job")
|
||
@patch("app.services.ops_agent_service.get_ops_action_template")
|
||
@patch("app.services.ops_agent_service.build_ops_template_payload")
|
||
@patch("app.services.ops_agent_service.get_managed_node_delivery_queue")
|
||
def test_request_managed_node_delivery_queue_action_creates_remote_agent_job_for_flush(
|
||
self,
|
||
mock_get_managed_node_delivery_queue,
|
||
mock_build_ops_template_payload,
|
||
mock_get_ops_action_template,
|
||
mock_create_ops_job,
|
||
) -> None:
|
||
mock_get_managed_node_delivery_queue.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"summary": {"pending_count": 3, "dead_letter_count": 1},
|
||
"records_capability": {"record_visibility": "head_only"},
|
||
"head_records": {"pending": {}, "dead_letter": {}},
|
||
}
|
||
mock_get_ops_action_template.return_value = {
|
||
"key": "delivery.queue.flush",
|
||
"title": "立即冲刷回执队列",
|
||
"default_auto_approve": True,
|
||
}
|
||
mock_build_ops_template_payload.return_value = (True, "ok", {"limit": 20})
|
||
mock_create_ops_job.return_value = (
|
||
True,
|
||
"ok",
|
||
{"job": {"id": 91, "action": "delivery.queue.flush", "status": "queued"}},
|
||
)
|
||
|
||
ok, message, data = request_managed_node_delivery_queue_action(
|
||
"mainland-worker-01",
|
||
"delivery.queue.flush",
|
||
payload={"requested_by": "web-ui", "limit": 20},
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("立即冲刷回执队列任务已创建", message)
|
||
called_payload = mock_create_ops_job.call_args.args[0]
|
||
self.assertEqual("delivery.queue.flush", called_payload["action"])
|
||
self.assertEqual("mainland-worker-01", called_payload["target_node_code"])
|
||
self.assertEqual("remote-agent", called_payload["execution_mode"])
|
||
self.assertEqual({"limit": 20}, called_payload["payload"])
|
||
self.assertEqual("delivery.queue.flush", data["action_key"])
|
||
self.assertEqual(91, data["job"]["id"])
|
||
|
||
@patch("app.services.ops_agent_service.get_ops_action_template")
|
||
@patch("app.services.ops_agent_service.build_ops_template_payload")
|
||
@patch("app.services.ops_agent_service.get_managed_node_delivery_queue")
|
||
def test_request_managed_node_delivery_queue_action_rejects_invisible_record(
|
||
self,
|
||
mock_get_managed_node_delivery_queue,
|
||
mock_build_ops_template_payload,
|
||
mock_get_ops_action_template,
|
||
) -> None:
|
||
mock_get_managed_node_delivery_queue.return_value = {
|
||
"node_code": "mainland-worker-01",
|
||
"summary": {"pending_count": 0, "dead_letter_count": 1},
|
||
"records_capability": {"record_visibility": "head_only"},
|
||
"head_records": {
|
||
"pending": {},
|
||
"dead_letter": {
|
||
"record_id": "evt-visible",
|
||
"state": "dead_letter",
|
||
"client_request_id": "evt-visible",
|
||
},
|
||
},
|
||
}
|
||
mock_get_ops_action_template.return_value = {
|
||
"key": "delivery.queue.discard",
|
||
"title": "丢弃死信回执",
|
||
"default_auto_approve": False,
|
||
}
|
||
mock_build_ops_template_payload.return_value = (
|
||
True,
|
||
"ok",
|
||
{"record_id": "evt-hidden", "limit": 1, "reason": "ignore"},
|
||
)
|
||
|
||
ok, message, data = request_managed_node_delivery_queue_action(
|
||
"mainland-worker-01",
|
||
"delivery.queue.discard",
|
||
payload={"requested_by": "web-ui", "reason": "ignore"},
|
||
record_id="evt-hidden",
|
||
)
|
||
|
||
self.assertFalse(ok)
|
||
self.assertEqual("当前控制面仅支持对可见头部记录执行单条动作", message)
|
||
self.assertEqual("evt-hidden", data["record_id"])
|
||
|
||
@patch("app.services.ops_agent_service.get_ops_job")
|
||
@patch("app.services.ops_agent_service.append_ops_job_event")
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
def test_agent_pull_jobs_commits_before_appending_event(
|
||
self,
|
||
mock_get_db,
|
||
mock_authenticate_agent_token,
|
||
mock_append_ops_job_event,
|
||
mock_get_ops_job,
|
||
) -> None:
|
||
cursor = _SequenceCursor(fetchall_values=[[(36,)]] )
|
||
connection = _SequenceConnection(cursor)
|
||
mock_get_db.return_value = connection
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_get_ops_job.return_value = {"id": 36, "job_code": "ops-demo", "action": "health.snapshot"}
|
||
|
||
def _assert_after_commit(**kwargs):
|
||
self.assertTrue(connection.committed)
|
||
self.assertEqual(36, kwargs["job_id"])
|
||
self.assertEqual("agent_dispatched", kwargs["event_type"])
|
||
|
||
mock_append_ops_job_event.side_effect = _assert_after_commit
|
||
|
||
ok, message, data = agent_pull_jobs({"node_code": "mainland-worker-01"}, token="agent-token", limit=1)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("ok", message)
|
||
self.assertEqual(1, data["count"])
|
||
self.assertTrue(connection.committed)
|
||
mock_append_ops_job_event.assert_called_once()
|
||
|
||
@patch("app.services.ops_agent_service.get_ops_job")
|
||
@patch("app.services.ops_agent_service.append_ops_job_event")
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
def test_agent_mark_job_started_commits_before_appending_event(
|
||
self,
|
||
mock_get_db,
|
||
mock_authenticate_agent_token,
|
||
mock_append_ops_job_event,
|
||
mock_get_ops_job,
|
||
) -> None:
|
||
cursor = _SequenceCursor(fetchone_values=[(36,)], fetchall_values=[[(101,), (102,)]])
|
||
connection = _SequenceConnection(cursor)
|
||
mock_get_db.return_value = connection
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_get_ops_job.return_value = {"id": 36, "job_code": "ops-demo", "action": "health.snapshot"}
|
||
|
||
def _assert_after_commit(**kwargs):
|
||
self.assertTrue(connection.committed)
|
||
self.assertEqual(36, kwargs["job_id"])
|
||
self.assertEqual("agent_started", kwargs["event_type"])
|
||
self.assertEqual({"step_ids": [101, 102]}, kwargs["payload"])
|
||
|
||
mock_append_ops_job_event.side_effect = _assert_after_commit
|
||
|
||
ok, message, data = agent_mark_job_started(36, {"node_code": "mainland-worker-01"}, token="agent-token")
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("任务已标记为运行中", message)
|
||
self.assertEqual(36, data["job"]["id"])
|
||
self.assertTrue(connection.committed)
|
||
mock_append_ops_job_event.assert_called_once()
|
||
|
||
@patch("app.services.ops_release_service.refresh_release_rollout_for_job")
|
||
@patch("app.services.ops_agent_service.get_ops_job")
|
||
@patch("app.services.ops_agent_service.append_ops_job_event")
|
||
@patch("app.services.ops_agent_service._authenticate_agent_token")
|
||
@patch("app.services.ops_agent_service.get_db")
|
||
def test_agent_complete_job_commits_before_appending_event(
|
||
self,
|
||
mock_get_db,
|
||
mock_authenticate_agent_token,
|
||
mock_append_ops_job_event,
|
||
mock_get_ops_job,
|
||
mock_refresh_release_rollout_for_job,
|
||
) -> None:
|
||
cursor = _SequenceCursor(fetchone_values=[(36, "")], fetchall_values=[[(101,)]] )
|
||
connection = _SequenceConnection(cursor)
|
||
mock_get_db.return_value = connection
|
||
mock_authenticate_agent_token.return_value = (True, "ok", {"node_code": "mainland-worker-01"})
|
||
mock_get_ops_job.return_value = {"id": 36, "job_code": "ops-demo", "status": "success"}
|
||
|
||
def _assert_after_commit(**kwargs):
|
||
self.assertTrue(connection.committed)
|
||
self.assertEqual(36, kwargs["job_id"])
|
||
self.assertEqual("agent_completed", kwargs["event_type"])
|
||
self.assertEqual("info", kwargs["level"])
|
||
self.assertEqual([101], kwargs["payload"]["step_ids"])
|
||
|
||
mock_append_ops_job_event.side_effect = _assert_after_commit
|
||
|
||
ok, message, data = agent_complete_job(
|
||
36,
|
||
{
|
||
"node_code": "mainland-worker-01",
|
||
"status": "success",
|
||
"stdout": "ok",
|
||
"stderr": "",
|
||
"result": {"summary": "done"},
|
||
"duration_ms": 120,
|
||
},
|
||
token="agent-token",
|
||
)
|
||
|
||
self.assertTrue(ok)
|
||
self.assertEqual("任务结果已回写", message)
|
||
self.assertEqual(36, data["job"]["id"])
|
||
self.assertTrue(connection.committed)
|
||
mock_append_ops_job_event.assert_called_once()
|
||
mock_refresh_release_rollout_for_job.assert_called_once_with(36)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
unittest.main()
|