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)