Files
getDomain/domain-api/tests/test_ops_agent_service.py
Your Name 7cbde2aa78 d
2026-04-22 14:13:21 +08:00

1426 lines
60 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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()