Files
arrr-erp/app/modules/accounting/agent_bridge.py
T
2026-08-18 14:39:19 +05:30

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)