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

447 lines
20 KiB
Python

import json
import os
import tempfile
import unittest
import urllib.error
from unittest.mock import patch
from app import node_agent
class NodeAgentDeliveryQueueTests(unittest.TestCase):
def test_base_payload_includes_delivery_queue_snapshot(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir):
node_agent._ensure_queue_dirs()
payload = node_agent._base_payload()
delivery_queue = payload["metadata"]["delivery_queue"]
self.assertEqual("healthy", delivery_queue["state"])
self.assertEqual(0, delivery_queue["pending_count"])
self.assertEqual(0, delivery_queue["dead_letter_count"])
def test_base_payload_prefers_non_loopback_identity(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir), patch.object(
node_agent,
"CONTROL_PLANE_BASE_URL",
"http://152.53.37.118:8100",
), patch.object(
node_agent.socket,
"gethostname",
return_value="localhost",
), patch.object(
node_agent.socket,
"getfqdn",
return_value="localhost.localdomain",
), patch.object(
node_agent.os,
"uname",
return_value=type("Uname", (), {"nodename": "localhost"})(),
), patch.object(
node_agent,
"NODE_CODE",
"mainland-controller-01",
), patch.object(
node_agent.socket,
"getaddrinfo",
return_value=[(None, None, None, None, ("127.0.0.1", 0))],
), patch.object(
node_agent.socket,
"gethostbyname",
return_value="127.0.0.1",
):
class FakeSocket:
def connect(self, target):
self.target = target
def getsockname(self):
return ("121.204.244.188", 12345)
def close(self):
return None
with patch.object(node_agent.socket, "socket", return_value=FakeSocket()):
payload = node_agent._base_payload()
self.assertNotIn(payload["hostname"], {"", "localhost", "localhost.localdomain"})
self.assertEqual("121.204.244.188", payload["ip"])
def test_detect_runtime_snapshot_degrades_to_worker_runtime_when_detect_status_fails(self) -> None:
with patch.dict(os.environ, {}, clear=False):
with patch(
"app.services.worker_control_service.detect_worker_runtime",
return_value={
"running": True,
"process_count": 1,
"latest_start_time": "2026-04-20 18:00:00",
"message": "active/running",
},
), patch(
"app.services.detect_service.get_detect_status",
side_effect=RuntimeError('connection to server at "127.0.0.1", port 5432 failed'),
):
snapshot = node_agent._detect_runtime_snapshot()
self.assertTrue(snapshot["worker_online"])
self.assertTrue(snapshot["service_running"])
self.assertFalse(snapshot["detecting"])
self.assertEqual("active/running", snapshot["phase_detail"])
self.assertEqual("2026-04-20 18:00:00", snapshot["updated_at"])
self.assertIn("127.0.0.1", snapshot["error"])
def test_detect_runtime_snapshot_infers_worker_online_from_active_threads(self) -> None:
with patch.dict(os.environ, {}, clear=False):
with patch(
"app.services.worker_control_service.detect_worker_runtime",
return_value={
"running": False,
"process_count": 0,
"latest_start_time": "",
"message": "",
},
), patch(
"app.services.detect_service.get_detect_status",
return_value={
"worker_online": False,
"detecting": False,
"active_thread_count": 19,
"max_thread_count": 1200,
"phase_label": "检测中",
"phase_detail": "Worker 正在处理 162 个检测任务",
"runtime_state": {
"service_running": False,
"updated_at": "2026-04-21 00:00:06",
},
"active_job": {"items_running": 0},
},
):
snapshot = node_agent._detect_runtime_snapshot()
self.assertTrue(snapshot["worker_online"])
self.assertTrue(snapshot["service_running"])
self.assertTrue(snapshot["detecting"])
self.assertTrue(snapshot["detect_participating"])
self.assertEqual(19, snapshot["active_threads"])
self.assertEqual(19, snapshot["current_load"])
def test_job_event_network_failure_is_queued_for_retry(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir), patch.object(
node_agent,
"_post",
side_effect=urllib.error.URLError("offline"),
):
result = node_agent._job_event(
7,
event_type="executor_received",
message="accepted",
client_event_id="evt-001",
summary_text="已接单",
focus_ref={"kind": "ops_job", "job_id": 7},
)
self.assertEqual("queued", result["state"])
pending_dir = os.path.join(temp_dir, "pending")
pending_files = sorted(os.listdir(pending_dir))
self.assertEqual(1, len(pending_files))
with open(os.path.join(pending_dir, pending_files[0]), "r", encoding="utf-8") as handle:
record = json.load(handle)
self.assertEqual("evt-001", record["payload"]["client_event_id"])
self.assertEqual("已接单", record["payload"]["summary_text"])
self.assertEqual("ops_job", record["payload"]["focus_ref"]["kind"])
self.assertEqual("job_event", record["kind"])
def test_job_complete_transport_failure_preserves_client_request_id(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir), patch.object(
node_agent,
"_post",
side_effect=urllib.error.URLError("offline"),
):
result = node_agent._job_complete(
9,
status="success",
stdout="done",
stderr="",
result={"ok": True},
client_request_id="complete-001",
duration_ms=321,
summary_text="执行成功",
focus_ref={"kind": "ops_job", "job_id": 9},
)
self.assertEqual("queued", result["state"])
pending_dir = os.path.join(temp_dir, "pending")
pending_files = sorted(os.listdir(pending_dir))
self.assertEqual(1, len(pending_files))
with open(os.path.join(pending_dir, pending_files[0]), "r", encoding="utf-8") as handle:
record = json.load(handle)
self.assertEqual("complete-001", record["payload"]["client_request_id"])
self.assertEqual(321, record["payload"]["duration_ms"])
self.assertEqual("执行成功", record["payload"]["summary_text"])
self.assertEqual("ops_job", record["payload"]["focus_ref"]["kind"])
self.assertEqual("job_complete", record["kind"])
def test_flush_delivery_queue_delivers_pending_records(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir):
record = node_agent._build_delivery_record(
"job_event",
"/api/v1/ops/agent/jobs/7/events",
{
"node_code": "mainland-worker-01",
"event_type": "executor_received",
"message": "accepted",
"client_event_id": "evt-002",
},
"evt-002",
)
node_agent._store_pending_delivery(record)
with patch.object(node_agent, "_post", return_value={"code": 0, "message": "ok"}):
summary = node_agent._flush_delivery_queue(limit=10)
self.assertEqual(1, summary["delivered"])
self.assertEqual([], os.listdir(os.path.join(temp_dir, "pending")))
def test_flush_delivery_queue_moves_semantic_failure_to_dead_letter(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir):
record = node_agent._build_delivery_record(
"job_complete",
"/api/v1/ops/agent/jobs/9/complete",
{
"node_code": "mainland-worker-01",
"status": "success",
"client_request_id": "complete-002",
},
"complete-002",
)
node_agent._store_pending_delivery(record)
with patch.object(
node_agent,
"_post",
return_value={
"code": 1,
"message": "任务不存在或不属于当前节点",
"detail_code": "ops_job_not_owned_by_agent",
},
):
summary = node_agent._flush_delivery_queue(limit=10)
self.assertEqual(1, summary["dead_letter"])
self.assertEqual([], os.listdir(os.path.join(temp_dir, "pending")))
self.assertEqual(1, len(os.listdir(os.path.join(temp_dir, "dead-letter"))))
def test_replay_dead_letter_records_moves_record_back_to_pending_and_flushes(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir):
record = node_agent._build_delivery_record(
"job_event",
"/api/v1/ops/agent/jobs/7/events",
{
"node_code": "mainland-worker-01",
"event_type": "executor_received",
"message": "accepted",
"client_event_id": "evt-003",
},
"evt-003",
)
node_agent._store_dead_letter_delivery(
record,
reason="temporary failure",
detail_code="transient_error",
)
with patch.object(node_agent, "_post", return_value={"code": 0, "message": "ok"}):
ok, message, data = node_agent._execute_action(
"delivery.queue.replay",
{
"selector": {"record_id": "evt-003"},
"limit": 10,
"flush_after_replay": True,
"reason": "恢复后重放",
},
)
self.assertTrue(ok)
self.assertEqual("已重放 1 条死信记录", message)
self.assertEqual(1, data["replayed_total"])
self.assertEqual(1, data["flush_summary"]["delivered"])
self.assertEqual([], os.listdir(os.path.join(temp_dir, "dead-letter")))
self.assertEqual([], os.listdir(os.path.join(temp_dir, "pending")))
def test_discard_dead_letter_records_archives_record(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
with patch.object(node_agent, "AGENT_QUEUE_DIR", temp_dir):
record = node_agent._build_delivery_record(
"job_complete",
"/api/v1/ops/agent/jobs/8/complete",
{
"node_code": "mainland-worker-01",
"status": "failed",
"client_request_id": "complete-003",
},
"complete-003",
)
node_agent._store_dead_letter_delivery(
record,
reason="permanent failure",
detail_code="ops_job_invalid_status",
)
ok, message, data = node_agent._execute_action(
"delivery.queue.discard",
{
"selector": {"record_id": "complete-003"},
"limit": 10,
"discarded_by": "web-ui",
"reason": "确认不需要再补发",
},
)
self.assertTrue(ok)
self.assertEqual("已丢弃 1 条死信记录", message)
self.assertEqual(1, data["discarded_total"])
self.assertEqual([], os.listdir(os.path.join(temp_dir, "dead-letter")))
discarded_files = os.listdir(os.path.join(temp_dir, "discarded"))
self.assertEqual(1, len(discarded_files))
with open(os.path.join(temp_dir, "discarded", discarded_files[0]), "r", encoding="utf-8") as handle:
archived = json.load(handle)
self.assertEqual("web-ui", archived["discarded_by"])
self.assertEqual("确认不需要再补发", archived["discard_reason"])
def test_process_job_passes_envelope_context_to_completion(self) -> None:
job = {
"id": 12,
"job_code": "ops-20260418-000012",
"job_type": "release_deploy",
"action": "logs.collect",
"payload": {"service_name": "domaincheck-worker", "lines": 120},
"step_key": "worker_logs",
"step_title": "收集 Worker 日志",
"focus_ref": {"kind": "ops_job", "job_id": 12, "job_code": "ops-20260418-000012"},
"release_context": {"release_version": "v2026.04.18", "rollout_id": 3},
"step_ref": {"step_id": 22, "step_key": "worker_logs", "step_title": "收集 Worker 日志"},
}
with patch.object(node_agent, "_job_start") as mock_job_start, patch.object(
node_agent,
"_job_event",
) as mock_job_event, patch.object(
node_agent,
"_execute_action",
return_value=(True, "logs collected", {"summary": "收集完成", "stdout": "ok", "stderr": ""}),
) as mock_execute_action, patch.object(
node_agent,
"_job_complete",
return_value={"state": "sent"},
) as mock_job_complete:
node_agent._process_job(job)
self.assertEqual("worker_logs", mock_job_start.call_args.kwargs["job"]["step_key"])
self.assertEqual("ops_job", mock_job_event.call_args.kwargs["focus_ref"]["kind"])
self.assertEqual("已接单 logs.collect", mock_job_event.call_args.kwargs["summary_text"])
self.assertEqual("logs.collect", mock_execute_action.call_args.args[0])
self.assertEqual("v2026.04.18", mock_job_complete.call_args.kwargs["release_context"]["release_version"])
self.assertEqual("ops_job", mock_job_complete.call_args.kwargs["focus_ref"]["kind"])
self.assertEqual("收集完成", mock_job_complete.call_args.kwargs["summary_text"])
self.assertGreaterEqual(mock_job_complete.call_args.kwargs["duration_ms"], 0)
def test_process_job_continues_when_job_start_delivery_fails(self) -> None:
job = {
"id": 18,
"job_code": "ops-20260418-000018",
"job_type": "ops_action",
"action": "logs.collect",
"payload": {"service_name": "domaincheck-worker", "lines": 50},
"step_key": "worker_logs",
"step_title": "收集 Worker 日志",
"focus_ref": {"kind": "ops_job", "job_id": 18, "job_code": "ops-20260418-000018"},
"release_context": {},
"step_ref": {"step_id": 28, "step_key": "worker_logs", "step_title": "收集 Worker 日志"},
}
with patch.object(
node_agent,
"_job_start",
side_effect=urllib.error.URLError("temporary offline"),
) as mock_job_start, patch.object(
node_agent,
"_job_event",
return_value={"state": "queued"},
) as mock_job_event, patch.object(
node_agent,
"_execute_action",
return_value=(True, "logs collected", {"summary": "收集完成", "stdout": "ok", "stderr": ""}),
) as mock_execute_action, patch.object(
node_agent,
"_job_complete",
return_value={"state": "sent"},
) as mock_job_complete:
node_agent._process_job(job)
mock_job_start.assert_called_once()
mock_execute_action.assert_called_once()
mock_job_complete.assert_called_once()
event_payload = mock_job_event.call_args.kwargs["payload"]
self.assertEqual("failed_local", event_payload["start_delivery_state"])
self.assertIn("temporary offline", event_payload["start_delivery_error"])
def test_register_heartbeat_and_pull_use_configured_timeouts(self) -> None:
with patch.object(node_agent, "_post", return_value={"code": 0, "message": "ok", "data": {"jobs": []}}) as mock_post, patch.object(
node_agent,
"AGENT_REGISTER_TIMEOUT_SECONDS",
91,
), patch.object(
node_agent,
"AGENT_HEARTBEAT_TIMEOUT_SECONDS",
92,
), patch.object(
node_agent,
"AGENT_PULL_TIMEOUT_SECONDS",
93,
):
node_agent._register()
node_agent._heartbeat()
jobs = node_agent._pull_jobs()
self.assertEqual([], jobs)
self.assertEqual(3, mock_post.call_count)
self.assertEqual(91, mock_post.call_args_list[0].kwargs["timeout"])
self.assertEqual(92, mock_post.call_args_list[1].kwargs["timeout"])
self.assertEqual(93, mock_post.call_args_list[2].kwargs["timeout"])
def test_job_delivery_uses_configured_timeouts(self) -> None:
with patch.object(
node_agent,
"_deliver_or_queue",
side_effect=lambda **kwargs: {"state": "queued", "timeout": kwargs["timeout"]},
) as mock_deliver, patch.object(
node_agent,
"AGENT_JOB_COMPLETE_TIMEOUT_SECONDS",
94,
), patch.object(
node_agent,
"AGENT_JOB_EVENT_TIMEOUT_SECONDS",
47,
):
complete_result = node_agent._job_complete(
11,
status="success",
stdout="ok",
stderr="",
result={"ok": True},
)
event_result = node_agent._job_event(
11,
event_type="executor_received",
message="accepted",
)
self.assertEqual("queued", complete_result["state"])
self.assertEqual("queued", event_result["state"])
self.assertEqual(2, mock_deliver.call_count)
self.assertEqual(94, mock_deliver.call_args_list[0].kwargs["timeout"])
self.assertEqual(47, mock_deliver.call_args_list[1].kwargs["timeout"])
if __name__ == "__main__":
unittest.main()