feat: add ops center and node onboarding flow

This commit is contained in:
Your Name
2026-04-18 23:52:51 +08:00
parent 246838ae4c
commit b9c29481b5
142 changed files with 89727 additions and 186 deletions

View File

@@ -0,0 +1,991 @@
from __future__ import annotations
import json
import os
import socket
import subprocess
import time
import urllib.error
import urllib.request
from datetime import datetime
from uuid import uuid4
from app.services.ops_action_executor_core import (
build_service_name_map,
execute_structured_action,
supports_structured_action,
)
from app.services.ops_release_executor_core import (
execute_release_action,
normalize_release_health_check_services as _release_normalize_health_check_services,
normalize_text_list as _release_normalize_text_list,
)
CONTROL_PLANE_BASE_URL = str(os.getenv("OPS_CONTROL_PLANE_BASE_URL", "")).strip().rstrip("/")
AGENT_TOKEN = str(os.getenv("OPS_AGENT_TOKEN", "")).strip()
NODE_CODE = str(os.getenv("NODE_CODE", "")).strip()
NODE_REGION = str(os.getenv("NODE_REGION", "mainland")).strip() or "mainland"
NODE_ROLE = str(os.getenv("NODE_ROLE", "worker")).strip() or "worker"
AGENT_POLL_INTERVAL_SECONDS = max(2, int(os.getenv("OPS_AGENT_POLL_INTERVAL_SECONDS", "5") or 5))
WORKER_SERVICE_NAME = str(os.getenv("WORKER_SERVICE_NAME", os.getenv("WORKER_SERVICE", "domaincheck-worker"))).strip() or "domaincheck-worker"
API_SERVICE_NAME = str(os.getenv("API_SERVICE_NAME", "domaincheck-api")).strip() or "domaincheck-api"
SYNC_AGENT_SERVICE_NAME = str(os.getenv("SYNC_AGENT_SERVICE_NAME", "domaincheck-sync-agent")).strip() or "domaincheck-sync-agent"
NODE_AGENT_SERVICE_NAME = str(os.getenv("NODE_AGENT_SERVICE_NAME", "domaincheck-node-agent")).strip() or "domaincheck-node-agent"
AGENT_VERSION = "0.1.0"
_APP_DIR = os.path.dirname(os.path.abspath(__file__))
_PROJECT_DIR = os.path.dirname(_APP_DIR)
_DEFAULT_QUEUE_DIR = os.path.join(_PROJECT_DIR, "runtime", "node-agent-queue", NODE_CODE or "unbound")
AGENT_QUEUE_DIR = str(os.getenv("OPS_AGENT_QUEUE_DIR", _DEFAULT_QUEUE_DIR)).strip() or _DEFAULT_QUEUE_DIR
AGENT_QUEUE_FLUSH_LIMIT = max(1, int(os.getenv("OPS_AGENT_QUEUE_FLUSH_LIMIT", "20") or 20))
_PERMANENT_DELIVERY_DETAIL_CODES = {
"agent_node_code_required",
"ops_job_not_owned_by_agent",
"ops_job_invalid_status",
}
_LAST_QUEUE_FLUSH_SUMMARY = {
"scanned": 0,
"delivered": 0,
"deferred": 0,
"dead_letter": 0,
"last_flush_at": "",
}
def _normalize_text_list(raw_value: object) -> list[str]:
"""Keep node-agent payload normalization aligned with release executor helpers."""
return _release_normalize_text_list(raw_value)
def _normalize_release_health_check_services(
payload: dict | None,
restart_services: list[str] | None,
) -> list[str]:
"""Backward-compatible wrapper used by node-agent tests and payload shaping."""
return _release_normalize_health_check_services(
dict(payload or {}),
list(restart_services or []),
default_api_service_name=API_SERVICE_NAME,
)
def _log(message: str) -> None:
print(f"{datetime.now().isoformat(sep=' ', timespec='seconds')} [node-agent] {message}", flush=True)
def _json_env(name: str, fallback: object) -> object:
raw = str(os.getenv(name, "")).strip()
if not raw:
return fallback
try:
return json.loads(raw)
except Exception:
return fallback
AGENT_CAPABILITIES = _json_env(
"OPS_AGENT_CAPABILITIES",
[
"service.start",
"service.stop",
"service.restart",
"service.status",
"runtime.start_worker",
"runtime.stop_worker",
"runtime.restart_api",
"runtime.start_sync_agent",
"runtime.stop_sync_agent",
"health.snapshot",
"logs.collect",
"diagnostics.collect",
"delivery.queue.flush",
"delivery.queue.replay",
"delivery.queue.discard",
"deploy.release",
],
)
AGENT_LABELS = _json_env("OPS_AGENT_LABELS", {})
class AgentResponseError(RuntimeError):
def __init__(self, message: str, *, detail_code: str = "", payload: dict | None = None):
super().__init__(message)
self.detail_code = str(detail_code or "").strip()
self.payload = dict(payload or {})
def _now_text() -> str:
return datetime.now().isoformat(sep=" ", timespec="seconds")
def _queue_pending_dir() -> str:
return os.path.join(AGENT_QUEUE_DIR, "pending")
def _queue_dead_letter_dir() -> str:
return os.path.join(AGENT_QUEUE_DIR, "dead-letter")
def _queue_discarded_dir() -> str:
return os.path.join(AGENT_QUEUE_DIR, "discarded")
def _ensure_queue_dirs() -> None:
for path in (AGENT_QUEUE_DIR, _queue_pending_dir(), _queue_dead_letter_dir(), _queue_discarded_dir()):
os.makedirs(path, exist_ok=True)
def _write_json_file(path: str, payload: dict) -> None:
temp_path = f"{path}.tmp"
with open(temp_path, "w", encoding="utf-8") as handle:
json.dump(payload, handle, ensure_ascii=False, indent=2)
os.replace(temp_path, path)
def _read_json_file(path: str) -> dict:
with open(path, "r", encoding="utf-8") as handle:
return json.load(handle)
def _new_request_id(prefix: str, job_id: int) -> str:
normalized_prefix = str(prefix or "request").strip() or "request"
return f"{normalized_prefix}-{int(job_id or 0)}-{int(time.time() * 1000)}-{uuid4().hex[:10]}"
def _queue_file_names(directory: str) -> list[str]:
_ensure_queue_dirs()
try:
return sorted(file_name for file_name in os.listdir(directory) if file_name.endswith(".json"))
except FileNotFoundError:
return []
def _queue_head_record(directory: str) -> dict:
file_names = _queue_file_names(directory)
if not file_names:
return {}
file_path = os.path.join(directory, file_names[0])
try:
return _read_json_file(file_path)
except Exception:
return {}
def _queue_record_id(record: dict, *, file_name: str = "") -> str:
request_id = str((record or {}).get("request_id") or "").strip()
if request_id:
return request_id
normalized_file_name = str(file_name or "").strip()
if normalized_file_name.endswith(".json"):
return normalized_file_name[:-5]
return normalized_file_name
def _queue_entries(state: str) -> list[dict]:
normalized_state = str(state or "").strip()
if normalized_state == "pending":
directory = _queue_pending_dir()
elif normalized_state == "dead_letter":
directory = _queue_dead_letter_dir()
else:
return []
entries: list[dict] = []
for file_name in _queue_file_names(directory):
file_path = os.path.join(directory, file_name)
try:
record = _read_json_file(file_path)
except Exception:
continue
entries.append(
{
"state": normalized_state,
"file_name": file_name,
"file_path": file_path,
"record_id": _queue_record_id(record, file_name=file_name),
"record": dict(record or {}),
}
)
return entries
def _matches_delivery_selector(entry: dict, selector: dict, *, allowed_states: set[str] | None = None) -> bool:
normalized_selector = dict(selector or {})
record = dict(entry.get("record") or {})
state = str(entry.get("state") or "").strip()
if allowed_states and state not in allowed_states:
return False
selected_state = str(normalized_selector.get("state") or "").strip()
if selected_state and state != selected_state:
return False
selected_record_id = str(
normalized_selector.get("record_id")
or normalized_selector.get("request_id")
or ""
).strip()
if selected_record_id and str(entry.get("record_id") or "").strip() != selected_record_id:
return False
selected_request_kind = str(normalized_selector.get("request_kind") or "").strip()
if selected_request_kind and str(record.get("kind") or "").strip() != selected_request_kind:
return False
selected_detail_code = str(normalized_selector.get("detail_code") or "").strip()
if selected_detail_code and str(record.get("last_detail_code") or "").strip() != selected_detail_code:
return False
return True
def _trim_queue_preview(records: list[dict], *, limit: int = 5) -> list[dict]:
safe_limit = max(1, int(limit or 5))
preview: list[dict] = []
for entry in list(records or [])[:safe_limit]:
record = dict(entry.get("record") or {})
preview.append(
{
"record_id": str(entry.get("record_id") or "").strip(),
"state": str(entry.get("state") or "").strip(),
"request_kind": str(record.get("kind") or "").strip(),
"request_id": str(record.get("request_id") or "").strip(),
"detail_code": str(record.get("last_detail_code") or "").strip(),
"created_at": str(record.get("created_at") or "").strip(),
"updated_at": str(record.get("updated_at") or "").strip(),
}
)
return preview
def _normalize_queue_limit(raw_value: object, *, default: int = 20, minimum: int = 1, maximum: int = 200) -> int:
try:
value = int(raw_value or default)
except Exception:
value = int(default)
return max(minimum, min(maximum, value))
def _replay_dead_letter_records(payload: dict | None = None) -> tuple[bool, str, dict]:
normalized_payload = dict(payload or {})
selector = dict(normalized_payload.get("selector") or {})
selector["state"] = str(selector.get("state") or "dead_letter").strip() or "dead_letter"
limit = _normalize_queue_limit(normalized_payload.get("limit"), default=20)
flush_after_replay = bool(normalized_payload.get("flush_after_replay", True))
reason = str(normalized_payload.get("reason") or "").strip()
matched_entries = [
entry
for entry in _queue_entries("dead_letter")
if _matches_delivery_selector(entry, selector, allowed_states={"dead_letter"})
][:limit]
if not matched_entries:
return False, "未找到符合条件的死信记录", {
"selector": selector,
"limit": limit,
"queue": _delivery_queue_snapshot(),
}
replayed_total = 0
for entry in matched_entries:
record = dict(entry.get("record") or {})
record["updated_at"] = _now_text()
record["replay_count"] = int(record.get("replay_count") or 0) + 1
record["last_replay_at"] = _now_text()
if reason:
record["last_replay_reason"] = reason
record.pop("dead_letter_at", None)
record.pop("dead_letter_reason", None)
_store_pending_delivery(record)
try:
os.remove(str(entry.get("file_path") or ""))
except FileNotFoundError:
pass
replayed_total += 1
flush_summary = {}
if flush_after_replay and replayed_total > 0:
flush_summary = _flush_delivery_queue(limit=replayed_total)
result = {
"selector": selector,
"limit": limit,
"replayed_total": replayed_total,
"flush_after_replay": flush_after_replay,
"flush_summary": flush_summary,
"matched_records": _trim_queue_preview(matched_entries),
"queue": _delivery_queue_snapshot(),
}
return True, f"已重放 {replayed_total} 条死信记录", result
def _store_discarded_delivery(record: dict, *, discarded_by: str = "", reason: str = "") -> str:
_ensure_queue_dirs()
discarded_record = dict(record or {})
discarded_record["discarded_at"] = _now_text()
discarded_record["discarded_by"] = str(discarded_by or "").strip()
discarded_record["discard_reason"] = str(reason or "").strip()
file_name = f"{int(time.time() * 1000)}-{uuid4().hex}.json"
file_path = os.path.join(_queue_discarded_dir(), file_name)
_write_json_file(file_path, discarded_record)
return file_path
def _discard_dead_letter_records(payload: dict | None = None) -> tuple[bool, str, dict]:
normalized_payload = dict(payload or {})
selector = dict(normalized_payload.get("selector") or {})
selector["state"] = str(selector.get("state") or "dead_letter").strip() or "dead_letter"
limit = _normalize_queue_limit(normalized_payload.get("limit"), default=20)
discarded_by = str(normalized_payload.get("discarded_by") or "").strip()
reason = str(normalized_payload.get("reason") or "").strip()
matched_entries = [
entry
for entry in _queue_entries("dead_letter")
if _matches_delivery_selector(entry, selector, allowed_states={"dead_letter"})
][:limit]
if not matched_entries:
return False, "未找到符合条件的死信记录", {
"selector": selector,
"limit": limit,
"queue": _delivery_queue_snapshot(),
}
discarded_total = 0
for entry in matched_entries:
_store_discarded_delivery(
dict(entry.get("record") or {}),
discarded_by=discarded_by,
reason=reason,
)
try:
os.remove(str(entry.get("file_path") or ""))
except FileNotFoundError:
pass
discarded_total += 1
result = {
"selector": selector,
"limit": limit,
"discarded_total": discarded_total,
"discarded_by": discarded_by,
"reason": reason,
"matched_records": _trim_queue_preview(matched_entries),
"queue": _delivery_queue_snapshot(),
}
return True, f"已丢弃 {discarded_total} 条死信记录", result
def _delivery_queue_snapshot() -> dict:
pending_dir = _queue_pending_dir()
dead_letter_dir = _queue_dead_letter_dir()
pending_files = _queue_file_names(pending_dir)
dead_letter_files = _queue_file_names(dead_letter_dir)
pending_head = _queue_head_record(pending_dir)
dead_letter_head = _queue_head_record(dead_letter_dir)
pending_count = len(pending_files)
dead_letter_count = len(dead_letter_files)
state = "healthy"
label = "正常"
reason = "当前没有待重试回执,也没有死信记录。"
if dead_letter_count > 0:
state = "dead_letter"
label = f"死信 {dead_letter_count}"
reason = "存在语义失败的回执/事件,自动重试已停止,建议人工查看。"
elif pending_count > 0:
state = "retrying"
label = f"待重试 {pending_count}"
reason = "存在待重试的回执/事件Node Agent 会在后续 heartbeat/poll 周期继续回放。"
elif not pending_head and not dead_letter_head:
label = "正常"
return {
"state": state,
"label": label,
"reason": reason,
"pending_count": pending_count,
"dead_letter_count": dead_letter_count,
"oldest_pending_at": str(pending_head.get("created_at") or "").strip(),
"oldest_pending_request_id": str(pending_head.get("request_id") or "").strip(),
"oldest_pending_kind": str(pending_head.get("kind") or "").strip(),
"oldest_dead_letter_at": str(
dead_letter_head.get("dead_letter_at") or dead_letter_head.get("created_at") or ""
).strip(),
"oldest_dead_letter_request_id": str(dead_letter_head.get("request_id") or "").strip(),
"oldest_dead_letter_kind": str(dead_letter_head.get("kind") or "").strip(),
"last_flush_at": str(_LAST_QUEUE_FLUSH_SUMMARY.get("last_flush_at") or "").strip(),
"last_flush_delivered": int(_LAST_QUEUE_FLUSH_SUMMARY.get("delivered") or 0),
"last_flush_deferred": int(_LAST_QUEUE_FLUSH_SUMMARY.get("deferred") or 0),
"last_flush_dead_letter": int(_LAST_QUEUE_FLUSH_SUMMARY.get("dead_letter") or 0),
}
def _response_detail_code(response: dict) -> str:
if not isinstance(response, dict):
return ""
top_level = str(response.get("detail_code") or "").strip()
if top_level:
return top_level
data = response.get("data")
if isinstance(data, dict):
return str(data.get("detail_code") or "").strip()
return ""
def _ensure_ok_response(response: dict, fallback_message: str) -> dict:
raw_code = response.get("code", 1)
try:
normalized_code = int(raw_code if raw_code not in (None, "") else 1)
except Exception:
normalized_code = 1
if normalized_code == 0:
return response
raise AgentResponseError(
str(response.get("message") or fallback_message),
detail_code=_response_detail_code(response),
payload=(response.get("data") if isinstance(response.get("data"), dict) else {}),
)
def _build_delivery_record(kind: str, path: str, payload: dict, request_id: str) -> dict:
return {
"kind": str(kind or "").strip() or "delivery",
"path": str(path or "").strip(),
"payload": dict(payload or {}),
"request_id": str(request_id or "").strip(),
"attempt_count": 0,
"created_at": _now_text(),
"updated_at": _now_text(),
"last_error": "",
"last_detail_code": "",
}
def _store_pending_delivery(record: dict) -> str:
_ensure_queue_dirs()
file_name = f"{int(time.time() * 1000)}-{uuid4().hex}.json"
file_path = os.path.join(_queue_pending_dir(), file_name)
_write_json_file(file_path, record)
return file_path
def _store_dead_letter_delivery(record: dict, *, reason: str, detail_code: str = "") -> str:
_ensure_queue_dirs()
dead_record = dict(record or {})
dead_record["dead_letter_reason"] = str(reason or "").strip()
dead_record["dead_letter_at"] = _now_text()
dead_record["last_error"] = str(reason or "").strip()
dead_record["last_detail_code"] = str(detail_code or "").strip()
file_name = f"{int(time.time() * 1000)}-{uuid4().hex}.json"
file_path = os.path.join(_queue_dead_letter_dir(), file_name)
_write_json_file(file_path, dead_record)
return file_path
def _dispatch_delivery_record(record: dict) -> dict:
kind = str(record.get("kind") or "delivery").strip() or "delivery"
path = str(record.get("path") or "").strip()
payload = record.get("payload") or {}
timeout = 15 if kind == "job_event" else 30
response = _post(path, payload, timeout=timeout)
return _ensure_ok_response(response, f"{kind} failed")
def _flush_delivery_queue(limit: int | None = None) -> dict:
global _LAST_QUEUE_FLUSH_SUMMARY
_ensure_queue_dirs()
safe_limit = max(1, int(limit or AGENT_QUEUE_FLUSH_LIMIT or 1))
summary = {
"scanned": 0,
"delivered": 0,
"deferred": 0,
"dead_letter": 0,
}
file_names = sorted(
file_name
for file_name in os.listdir(_queue_pending_dir())
if file_name.endswith(".json")
)[:safe_limit]
for file_name in file_names:
summary["scanned"] += 1
file_path = os.path.join(_queue_pending_dir(), file_name)
try:
record = _read_json_file(file_path)
except Exception as exc:
_store_dead_letter_delivery({"file_name": file_name}, reason=f"invalid queue record: {exc}")
try:
os.remove(file_path)
except FileNotFoundError:
pass
summary["dead_letter"] += 1
continue
try:
_dispatch_delivery_record(record)
try:
os.remove(file_path)
except FileNotFoundError:
pass
summary["delivered"] += 1
except AgentResponseError as exc:
record["attempt_count"] = int(record.get("attempt_count") or 0) + 1
record["updated_at"] = _now_text()
record["last_error"] = str(exc)
record["last_detail_code"] = exc.detail_code
if exc.detail_code in _PERMANENT_DELIVERY_DETAIL_CODES:
_store_dead_letter_delivery(record, reason=str(exc), detail_code=exc.detail_code)
try:
os.remove(file_path)
except FileNotFoundError:
pass
summary["dead_letter"] += 1
_log(
f"delivery moved to dead-letter: kind={record.get('kind')} request_id={record.get('request_id')} "
f"detail_code={exc.detail_code or '-'} message={exc}"
)
else:
_write_json_file(file_path, record)
summary["deferred"] += 1
except Exception as exc:
record["attempt_count"] = int(record.get("attempt_count") or 0) + 1
record["updated_at"] = _now_text()
record["last_error"] = str(exc)
_write_json_file(file_path, record)
summary["deferred"] += 1
summary["last_flush_at"] = _now_text()
_LAST_QUEUE_FLUSH_SUMMARY = dict(summary)
return summary
def _deliver_or_queue(
*,
kind: str,
path: str,
payload: dict,
request_id: str,
timeout: int = 30,
) -> dict:
record = _build_delivery_record(kind, path, payload, request_id)
try:
response = _post(path, payload, timeout=timeout)
_ensure_ok_response(response, f"{kind} failed")
return {"state": "delivered", "request_id": request_id, "detail_code": "", "response": response}
except AgentResponseError as exc:
if exc.detail_code in _PERMANENT_DELIVERY_DETAIL_CODES:
file_path = _store_dead_letter_delivery(record, reason=str(exc), detail_code=exc.detail_code)
_log(
f"{kind} dead-lettered: request_id={request_id} detail_code={exc.detail_code or '-'} "
f"message={exc} file={file_path}"
)
return {"state": "dead_letter", "request_id": request_id, "detail_code": exc.detail_code, "file_path": file_path}
file_path = _store_pending_delivery({**record, "last_error": str(exc), "last_detail_code": exc.detail_code})
_log(
f"{kind} queued for retry: request_id={request_id} detail_code={exc.detail_code or '-'} "
f"message={exc} file={file_path}"
)
return {"state": "queued", "request_id": request_id, "detail_code": exc.detail_code, "file_path": file_path}
except Exception as exc:
file_path = _store_pending_delivery({**record, "last_error": str(exc)})
_log(f"{kind} queued for retry: request_id={request_id} message={exc} file={file_path}")
return {"state": "queued", "request_id": request_id, "detail_code": "", "file_path": file_path}
def _headers() -> dict[str, str]:
return {
"Content-Type": "application/json",
"X-Domaincheck-Agent-Token": AGENT_TOKEN,
}
def _request(method: str, path: str, payload: dict | None = None, timeout: int = 30) -> dict:
if not CONTROL_PLANE_BASE_URL:
raise RuntimeError("OPS_CONTROL_PLANE_BASE_URL 未配置")
if not AGENT_TOKEN:
raise RuntimeError("OPS_AGENT_TOKEN 未配置")
url = f"{CONTROL_PLANE_BASE_URL}{path}"
data = json.dumps(payload or {}, ensure_ascii=False).encode("utf-8")
request = urllib.request.Request(url=url, data=data, headers=_headers(), method=method.upper())
with urllib.request.urlopen(request, timeout=timeout) as response:
body = response.read().decode("utf-8", errors="ignore")
return json.loads(body or "{}")
def _post(path: str, payload: dict, timeout: int = 30) -> dict:
return _request("POST", path, payload, timeout=timeout)
def _hostname() -> str:
try:
return socket.gethostname()
except Exception:
return ""
def _ip() -> str:
try:
return socket.gethostbyname(socket.gethostname())
except Exception:
return ""
def _base_payload() -> dict:
return {
"node_code": NODE_CODE,
"region": NODE_REGION,
"role": NODE_ROLE,
"title": NODE_CODE,
"hostname": _hostname(),
"ip": _ip(),
"agent_version": AGENT_VERSION,
"capabilities": AGENT_CAPABILITIES,
"labels": AGENT_LABELS,
"metadata": {
"service_names": {
"api": API_SERVICE_NAME,
"worker": WORKER_SERVICE_NAME,
"sync_agent": SYNC_AGENT_SERVICE_NAME,
"node_agent": NODE_AGENT_SERVICE_NAME,
},
"delivery_queue": _delivery_queue_snapshot(),
},
}
def _run(command: list[str], timeout: int = 60) -> tuple[int, str, str]:
completed = subprocess.run(command, capture_output=True, text=True, timeout=timeout)
return completed.returncode, completed.stdout.strip(), completed.stderr.strip()
def _execute_action(
action: str,
payload: dict,
*,
job_id: int | None = None,
job_context: dict | None = None,
) -> tuple[bool, str, dict]:
normalized_action = str(action or "").strip()
if normalized_action == "delivery.queue.flush":
limit = _normalize_queue_limit((payload or {}).get("limit"), default=20)
summary = _flush_delivery_queue(limit=limit)
return True, "Delivery Queue 已执行冲刷", {
"action": normalized_action,
"limit": limit,
"flush_summary": summary,
"queue": _delivery_queue_snapshot(),
}
if normalized_action == "delivery.queue.replay":
return _replay_dead_letter_records(payload)
if normalized_action == "delivery.queue.discard":
return _discard_dead_letter_records(payload)
if supports_structured_action(normalized_action):
return execute_structured_action(
normalized_action,
payload,
service_names=build_service_name_map(
api_service_name=API_SERVICE_NAME,
worker_service_name=WORKER_SERVICE_NAME,
sync_agent_service_name=SYNC_AGENT_SERVICE_NAME,
node_agent_service_name=NODE_AGENT_SERVICE_NAME,
),
runner=_run,
host_context={
"hostname": _hostname(),
"ip": _ip(),
},
)
if normalized_action == "deploy.release":
return execute_release_action(
dict(payload or {}),
run_command=_run,
default_api_service_name=API_SERVICE_NAME,
event_callback=(
(lambda event_type, message, level="info", payload=None: _job_event(
int(job_id),
event_type=event_type,
message=message,
level=level,
payload={
**dict(payload or {}),
"release_context": dict((job_context or {}).get("release_context") or {}),
"step_key": str((job_context or {}).get("step_key") or "").strip(),
},
summary_text=message,
focus_ref=dict((job_context or {}).get("focus_ref") or {}),
occurred_at=_now_text(),
))
if job_id
else None
),
urlopen_func=urllib.request.urlopen,
user_agent=f"domaincheck-node-agent/{AGENT_VERSION}",
)
return False, f"unsupported action: {normalized_action}", {"action": normalized_action}
def _register() -> None:
response = _post("/api/v1/ops/agent/register", _base_payload())
_ensure_ok_response(response, "agent register failed")
_log(f"registered: {response.get('message')}")
def _heartbeat() -> None:
response = _post("/api/v1/ops/agent/heartbeat", _base_payload())
_ensure_ok_response(response, "agent heartbeat failed")
def _pull_jobs() -> list[dict]:
response = _post(f"/api/v1/ops/agent/pull?limit=1", {"node_code": NODE_CODE})
_ensure_ok_response(response, "agent pull failed")
data = response.get("data") or {}
return list(data.get("jobs") or [])
def _normalize_agent_job(job: dict) -> dict:
normalized_job = dict(job or {})
focus_ref = normalized_job.get("focus_ref") if isinstance(normalized_job.get("focus_ref"), dict) else {}
release_context = (
normalized_job.get("release_context") if isinstance(normalized_job.get("release_context"), dict) else {}
)
step_ref = normalized_job.get("step_ref") if isinstance(normalized_job.get("step_ref"), dict) else {}
job_id = int(normalized_job.get("job_id") or normalized_job.get("id") or 0)
step_key = str(
normalized_job.get("step_key")
or step_ref.get("step_key")
or "dispatch"
).strip() or "dispatch"
step_title = str(
normalized_job.get("step_title")
or step_ref.get("step_title")
or normalized_job.get("summary")
or normalized_job.get("action")
or step_key
).strip() or step_key
if not focus_ref:
focus_ref = {
"kind": "ops_job",
"job_id": job_id,
"job_code": str(normalized_job.get("job_code") or "").strip(),
"action": str(normalized_job.get("action") or "").strip(),
"target_node_code": str(normalized_job.get("target_node_code") or "").strip(),
}
return {
**normalized_job,
"job_id": job_id,
"job_code": str(normalized_job.get("job_code") or "").strip(),
"job_type": str(normalized_job.get("job_type") or "ops_action").strip() or "ops_action",
"action": str(normalized_job.get("action") or "").strip(),
"payload": dict(normalized_job.get("payload") or {}),
"policy": dict(normalized_job.get("policy") or {}),
"focus_ref": focus_ref,
"release_context": dict(release_context or {}),
"step_key": step_key,
"step_title": step_title,
"step_ref": {
"step_id": int(step_ref.get("step_id") or 0),
"step_key": step_key,
"step_title": step_title,
},
}
def _job_start(job_id: int, *, job: dict | None = None) -> None:
normalized_job = _normalize_agent_job(job or {"job_id": job_id})
response = _post(
f"/api/v1/ops/agent/jobs/{job_id}/start",
{
"node_code": NODE_CODE,
"job_code": normalized_job.get("job_code") or "",
"action": normalized_job.get("action") or "",
"step_key": normalized_job.get("step_key") or "",
"focus_ref": normalized_job.get("focus_ref") or {},
},
)
_ensure_ok_response(response, "job start failed")
def _job_complete(
job_id: int,
*,
status: str,
stdout: str,
stderr: str,
result: dict,
error_message: str = "",
client_request_id: str | None = None,
duration_ms: int | None = None,
summary_text: str = "",
focus_ref: dict | None = None,
step_ref: dict | None = None,
release_context: dict | None = None,
) -> dict:
request_id = str(client_request_id or "").strip() or _new_request_id("complete", job_id)
payload = {
"node_code": NODE_CODE,
"status": status,
"stdout": stdout,
"stderr": stderr,
"result": result,
"error_message": error_message,
"client_request_id": request_id,
}
if duration_ms is not None:
payload["duration_ms"] = max(0, int(duration_ms or 0))
if str(summary_text or "").strip():
payload["summary_text"] = str(summary_text).strip()
if isinstance(focus_ref, dict) and focus_ref:
payload["focus_ref"] = dict(focus_ref)
if isinstance(step_ref, dict) and step_ref:
payload["step_ref"] = dict(step_ref)
if isinstance(release_context, dict) and release_context:
payload["release_context"] = dict(release_context)
return _deliver_or_queue(
kind="job_complete",
path=f"/api/v1/ops/agent/jobs/{job_id}/complete",
payload=payload,
request_id=request_id,
timeout=30,
)
def _job_event(
job_id: int,
*,
event_type: str,
message: str,
level: str = "info",
payload: dict | None = None,
client_event_id: str | None = None,
summary_text: str = "",
focus_ref: dict | None = None,
occurred_at: str = "",
) -> dict:
request_id = str(client_event_id or "").strip() or _new_request_id("event", job_id)
delivery_payload = {
"node_code": NODE_CODE,
"event_type": event_type,
"message": message,
"level": level,
"payload": payload or {},
"client_event_id": request_id,
}
if str(summary_text or "").strip():
delivery_payload["summary_text"] = str(summary_text).strip()
if isinstance(focus_ref, dict) and focus_ref:
delivery_payload["focus_ref"] = dict(focus_ref)
if str(occurred_at or "").strip():
delivery_payload["occurred_at"] = str(occurred_at).strip()
return _deliver_or_queue(
kind="job_event",
path=f"/api/v1/ops/agent/jobs/{job_id}/events",
payload=delivery_payload,
request_id=request_id,
timeout=15,
)
def _process_job(job: dict) -> None:
normalized_job = _normalize_agent_job(job)
job_id = int(normalized_job.get("job_id") or 0)
action = str(normalized_job.get("action") or "").strip()
payload = dict(normalized_job.get("payload") or {})
if job_id <= 0 or not action:
return
started_at = time.monotonic()
_log(
"job start: "
f"id={job_id} code={normalized_job.get('job_code') or '-'} "
f"type={normalized_job.get('job_type') or '-'} "
f"step={normalized_job.get('step_key') or '-'} action={action}"
)
start_delivery_state = "delivered"
start_delivery_error = ""
try:
_job_start(job_id, job=normalized_job)
except Exception as exc:
start_delivery_state = "failed_local"
start_delivery_error = str(exc)
_log(
"job start delivery failed, continue locally: "
f"id={job_id} code={normalized_job.get('job_code') or '-'} action={action} error={exc}"
)
_job_event(
job_id,
event_type="executor_received",
message=f"node agent accepted action {action}",
summary_text=f"已接单 {action}",
focus_ref=dict(normalized_job.get("focus_ref") or {}),
occurred_at=_now_text(),
payload={
"action": action,
"job_code": normalized_job.get("job_code") or "",
"job_type": normalized_job.get("job_type") or "",
"step_key": normalized_job.get("step_key") or "",
"step_title": normalized_job.get("step_title") or "",
"release_context": dict(normalized_job.get("release_context") or {}),
"start_delivery_state": start_delivery_state,
"start_delivery_error": start_delivery_error,
},
)
ok, message, result = _execute_action(action, payload, job_id=job_id, job_context=normalized_job)
stdout = str(result.get("stdout") or "")
stderr = str(result.get("stderr") or "")
duration_ms = max(0, int((time.monotonic() - started_at) * 1000))
summary_text = str(result.get("summary_text") or result.get("summary") or message).strip()
delivery = _job_complete(
job_id,
status="success" if ok else "failed",
stdout=stdout,
stderr=stderr,
result=result,
error_message="" if ok else message,
duration_ms=duration_ms,
summary_text=summary_text,
focus_ref=dict(normalized_job.get("focus_ref") or {}),
step_ref=dict(normalized_job.get("step_ref") or {}),
release_context=dict(normalized_job.get("release_context") or {}),
)
_log(
"job complete: "
f"id={job_id} code={normalized_job.get('job_code') or '-'} "
f"action={action} ok={ok} duration_ms={duration_ms} delivery={delivery.get('state')}"
)
def main() -> None:
if not NODE_CODE:
raise RuntimeError("NODE_CODE 未配置")
_log(f"starting node agent: node={NODE_CODE} role={NODE_ROLE} region={NODE_REGION}")
_ensure_queue_dirs()
_register()
last_heartbeat_at = 0.0
while True:
now = time.time()
try:
delivery_summary = _flush_delivery_queue(limit=AGENT_QUEUE_FLUSH_LIMIT)
if delivery_summary["delivered"] or delivery_summary["dead_letter"]:
_log(f"delivery queue flush: {delivery_summary}")
if now - last_heartbeat_at >= 15:
_heartbeat()
last_heartbeat_at = now
jobs = _pull_jobs()
if jobs:
for job in jobs:
_process_job(job)
else:
time.sleep(AGENT_POLL_INTERVAL_SECONDS)
except urllib.error.HTTPError as exc:
_log(f"http error: {exc.code}")
time.sleep(AGENT_POLL_INTERVAL_SECONDS)
except urllib.error.URLError as exc:
_log(f"url error: {exc}")
time.sleep(AGENT_POLL_INTERVAL_SECONDS)
except Exception as exc:
_log(f"loop error: {exc}")
time.sleep(AGENT_POLL_INTERVAL_SECONDS)
if __name__ == "__main__":
main()