133 lines
4.2 KiB
Python
133 lines
4.2 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import tempfile
|
|
import time
|
|
from typing import Any
|
|
from uuid import uuid4
|
|
|
|
|
|
COMMAND_TIMEOUT_SECONDS = 25
|
|
COMMAND_MAX_AGE_SECONDS = 120
|
|
|
|
|
|
def _bus_root() -> Path:
|
|
configured = (os.getenv("ERP_AGENT_COMMAND_BUS_ROOT") or "").strip()
|
|
if configured:
|
|
root = Path(configured).expanduser().resolve()
|
|
else:
|
|
root = Path(tempfile.gettempdir()) / "arrr_erp_agent_command_bus"
|
|
root.mkdir(parents=True, exist_ok=True)
|
|
return root
|
|
|
|
|
|
def _safe_node_code(node_code: str) -> str:
|
|
return "".join(ch for ch in str(node_code or "") if ch.isalnum() or ch in "-_.")[:100]
|
|
|
|
|
|
def _node_dir(node_code: str) -> Path:
|
|
safe = _safe_node_code(node_code)
|
|
if not safe:
|
|
raise ValueError("Storage node code is required.")
|
|
path = _bus_root() / safe
|
|
path.mkdir(parents=True, exist_ok=True)
|
|
return path
|
|
|
|
|
|
def _atomic_json_write(path: Path, payload: dict[str, Any]) -> None:
|
|
temp = path.with_suffix(path.suffix + ".tmp")
|
|
temp.write_text(json.dumps(payload, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
|
|
temp.replace(path)
|
|
|
|
|
|
def _read_json(path: Path) -> dict[str, Any] | None:
|
|
try:
|
|
data = json.loads(path.read_text(encoding="utf-8"))
|
|
return data if isinstance(data, dict) else None
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def _cleanup(node_code: str) -> None:
|
|
now = time.time()
|
|
for path in _node_dir(node_code).glob("*.json"):
|
|
try:
|
|
if now - path.stat().st_mtime > COMMAND_MAX_AGE_SECONDS:
|
|
path.unlink(missing_ok=True)
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
def enqueue_agent_command(node_code: str, action: str, payload: dict[str, Any] | None = None) -> str:
|
|
_cleanup(node_code)
|
|
command_id = uuid4().hex
|
|
body = {
|
|
"command_id": command_id,
|
|
"action": str(action or "").strip(),
|
|
"payload": payload or {},
|
|
"created_at_epoch": time.time(),
|
|
}
|
|
_atomic_json_write(_node_dir(node_code) / f"command_{command_id}.json", body)
|
|
return command_id
|
|
|
|
|
|
def list_pending_agent_commands(node_code: str, limit: int = 10) -> list[dict[str, Any]]:
|
|
_cleanup(node_code)
|
|
base = _node_dir(node_code)
|
|
items: list[dict[str, Any]] = []
|
|
for path in sorted(base.glob("command_*.json"), key=lambda p: p.stat().st_mtime):
|
|
command_id = path.stem.removeprefix("command_")
|
|
if (base / f"response_{command_id}.json").exists():
|
|
continue
|
|
payload = _read_json(path)
|
|
if payload:
|
|
items.append(payload)
|
|
if len(items) >= max(1, int(limit)):
|
|
break
|
|
return items
|
|
|
|
|
|
def record_agent_command_result(node_code: str, result: dict[str, Any]) -> None:
|
|
command_id = str(result.get("command_id") or "").strip()
|
|
if not command_id:
|
|
return
|
|
response = {
|
|
"command_id": command_id,
|
|
"ok": bool(result.get("ok")),
|
|
"result": result.get("result"),
|
|
"error": result.get("error"),
|
|
"agent_time_utc": result.get("agent_time_utc"),
|
|
"received_at_epoch": time.time(),
|
|
}
|
|
_atomic_json_write(_node_dir(node_code) / f"response_{command_id}.json", response)
|
|
|
|
|
|
def wait_for_agent_result(node_code: str, command_id: str, timeout_seconds: int = COMMAND_TIMEOUT_SECONDS) -> dict[str, Any]:
|
|
base = _node_dir(node_code)
|
|
response_path = base / f"response_{command_id}.json"
|
|
command_path = base / f"command_{command_id}.json"
|
|
deadline = time.monotonic() + max(1, int(timeout_seconds))
|
|
while time.monotonic() < deadline:
|
|
result = _read_json(response_path)
|
|
if result:
|
|
try:
|
|
response_path.unlink(missing_ok=True)
|
|
command_path.unlink(missing_ok=True)
|
|
except Exception:
|
|
pass
|
|
return result
|
|
time.sleep(0.2)
|
|
raise TimeoutError("ERP Local Agent did not respond before the command timeout.")
|
|
|
|
|
|
def request_agent_command(
|
|
node_code: str,
|
|
action: str,
|
|
payload: dict[str, Any] | None = None,
|
|
timeout_seconds: int = COMMAND_TIMEOUT_SECONDS,
|
|
) -> dict[str, Any]:
|
|
command_id = enqueue_agent_command(node_code, action, payload)
|
|
return wait_for_agent_result(node_code, command_id, timeout_seconds=timeout_seconds)
|