From 05a6b42cce2e159bea910ef6cf4fb5c3c5ad0f74 Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Tue, 18 Aug 2026 22:22:47 +0530 Subject: [PATCH] Add Local Agent dashboard SQLite history and manual updates --- app/modules/documents/agent_package.py | 15 +- .../README_ERP_LOCAL_AGENT.txt | 18 ++ .../erp_local_agent/__init__.py | 2 +- .../erp_local_agent/config.py | 11 +- .../erp_local_agent/dashboard.py | 167 +++++++++++ .../local_agent_runtime/erp_local_agent/db.py | 267 ++++++++++++++++-- .../erp_local_agent/main.py | 29 +- .../erp_local_agent/sync.py | 5 +- .../erp_local_agent/tunnel.py | 22 +- .../erp_local_agent/updater.py | 181 +++++++++--- .../local_agent_runtime/open_dashboard.bat | 2 + 11 files changed, 640 insertions(+), 79 deletions(-) create mode 100644 app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt create mode 100644 app/modules/documents/local_agent_runtime/erp_local_agent/dashboard.py create mode 100644 app/modules/documents/local_agent_runtime/open_dashboard.bat diff --git a/app/modules/documents/agent_package.py b/app/modules/documents/agent_package.py index 84bc9f5..b45a8ee 100644 --- a/app/modules/documents/agent_package.py +++ b/app/modules/documents/agent_package.py @@ -4,7 +4,7 @@ import io from pathlib import Path import zipfile -ERP_LOCAL_AGENT_VERSION = "1.1.0" +ERP_LOCAL_AGENT_VERSION = "1.2.0" ERP_LOCAL_AGENT_NAME = "ERP Local Agent" RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime" @@ -30,7 +30,11 @@ def build_agent_env(*, erp_base_url: str, node_code: str, node_secret: str, stor f"TUNNEL_ENABLED={str(bool(tunnel_enabled)).lower()}\n" f"TUNNEL_RECONNECT_SECONDS={int(tunnel_reconnect_seconds or 10)}\n" "AUTO_UPDATE=true\n" + "AUTO_INSTALL_UPDATES=false\n" "UPDATE_CHECK_INTERVAL_SECONDS=300\n" + "DASHBOARD_ENABLED=true\n" + "DASHBOARD_HOST=127.0.0.1\n" + "DASHBOARD_PORT=8788\n" ) @@ -47,7 +51,14 @@ def _build_zip(*, env_text: str | None, include_env: bool, include_admin_readme: if include_env and env_text is not None: dst.writestr(".env", env_text) if include_admin_readme: - dst.writestr("README_ERP_LOCAL_AGENT.txt", f"ERP Local Agent {ERP_LOCAL_AGENT_VERSION}\nStorage behaviour and D:\\AuditFirmStorage are preserved. Future compatible releases self-update securely from ERP.\n") + dst.writestr( + "README_ERP_LOCAL_AGENT.txt", + f"ERP Local Agent {ERP_LOCAL_AGENT_VERSION}\n" + "Existing storage, WebSocket tunnel, Tally and accounting .act functionality are preserved.\n" + "Local dashboard: http://127.0.0.1:8788\n" + "The agent checks for updates automatically but installation is always initiated by the local user.\n" + "The existing .env, data, logs, .venv and client storage are preserved during updates.\n", + ) return buffer.getvalue() diff --git a/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt b/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt new file mode 100644 index 0000000..ef4f37a --- /dev/null +++ b/app/modules/documents/local_agent_runtime/README_ERP_LOCAL_AGENT.txt @@ -0,0 +1,18 @@ +ERP Local Agent 1.2.0 + +Existing storage, WebSocket tunnel, Tally and client .act functionality are preserved. + +Local dashboard: + http://127.0.0.1:8788 + or run open_dashboard.bat + +Updates: + - The agent checks the ERP for new versions automatically. + - Updates are NOT installed automatically. + - Open the local dashboard, Download Update, then Install Downloaded Update. + - .env, .venv, data, logs and client storage are preserved during update. + +Local operational database: + data\agent.db + +Client accounting .act databases remain separate under the configured STORAGE_ROOT. diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py index c94f386..90bfaca 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/__init__.py @@ -1,2 +1,2 @@ -__version__ = "1.1.0" +__version__ = "1.2.0" AGENT_NAME = "ERP Local Agent" diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/config.py b/app/modules/documents/local_agent_runtime/erp_local_agent/config.py index 1fac679..3ada930 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/config.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/config.py @@ -2,8 +2,9 @@ from __future__ import annotations import os from dataclasses import dataclass -from urllib.parse import urlencode from pathlib import Path +from urllib.parse import urlencode + from dotenv import load_dotenv from . import __version__ @@ -26,6 +27,10 @@ class AgentConfig: tunnel_reconnect_seconds: int = 10 auto_update: bool = True update_check_interval_seconds: int = 300 + auto_install_updates: bool = False + dashboard_enabled: bool = True + dashboard_host: str = "127.0.0.1" + dashboard_port: int = 8788 @property def headers(self) -> dict[str, str]: @@ -109,4 +114,8 @@ def load_config(env_file: str | None = None) -> AgentConfig: tunnel_reconnect_seconds=_get_int("TUNNEL_RECONNECT_SECONDS", 10), auto_update=_get_bool("AUTO_UPDATE", True), update_check_interval_seconds=max(60, _get_int("UPDATE_CHECK_INTERVAL_SECONDS", 300)), + auto_install_updates=_get_bool("AUTO_INSTALL_UPDATES", False), + dashboard_enabled=_get_bool("DASHBOARD_ENABLED", True), + dashboard_host=os.getenv("DASHBOARD_HOST", "127.0.0.1").strip() or "127.0.0.1", + dashboard_port=max(1, min(65535, _get_int("DASHBOARD_PORT", 8788))), ) diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/dashboard.py b/app/modules/documents/local_agent_runtime/erp_local_agent/dashboard.py new file mode 100644 index 0000000..9bf2ec7 --- /dev/null +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/dashboard.py @@ -0,0 +1,167 @@ +from __future__ import annotations + +import json +import shutil +import threading +import time +from http import HTTPStatus +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any + +from . import AGENT_NAME, __version__ +from .tally import TallyLiveConnector + + +class AgentDashboard: + def __init__(self, config, db, updater, logger, install_dir: Path): + self.config = config + self.db = db + self.updater = updater + self.logger = logger + self.install_dir = Path(install_dir).resolve() + self.server: ThreadingHTTPServer | None = None + + def start(self) -> None: + if not self.config.dashboard_enabled: + return + dashboard = self + + class Handler(BaseHTTPRequestHandler): + def log_message(self, fmt, *args): + dashboard.logger.debug("Dashboard: " + fmt, *args) + + def _json(self, payload: Any, status: int = 200): + data = json.dumps(payload, ensure_ascii=False, default=str).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Cache-Control", "no-store") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def _html(self, text: str): + data = text.encode("utf-8") + self.send_response(HTTPStatus.OK) + self.send_header("Content-Type", "text/html; charset=utf-8") + self.send_header("Cache-Control", "no-store") + self.send_header("Content-Length", str(len(data))) + self.end_headers() + self.wfile.write(data) + + def do_GET(self): + if self.path == "/" or self.path.startswith("/?"): + return self._html(dashboard._page()) + if self.path.startswith("/api/status"): + return self._json(dashboard.status()) + if self.path.startswith("/api/history"): + return self._json(dashboard.history()) + self.send_error(404) + + def do_POST(self): + try: + if self.path == "/api/update/check": + return self._json({"ok": True, "update": dashboard.updater.check_for_update(force=True)}) + if self.path == "/api/update/download": + return self._json({"ok": True, "update": dashboard.updater.download_latest()}) + if self.path == "/api/update/install": + state = dashboard.updater.status() + if not state.get("downloaded"): + raise RuntimeError("Download and verify the update before installing it.") + threading.Thread(target=dashboard._delayed_install, daemon=True).start() + return self._json({"ok": True, "message": "Update installation started. The agent will restart automatically."}) + if self.path == "/api/tally/check": + status = TallyLiveConnector().status() + dashboard.db.record_tally_status(status) + return self._json({"ok": True, "tally": status}) + self.send_error(404) + except Exception as exc: + dashboard.logger.exception("Dashboard action failed: %s", exc) + dashboard.db.record_event("ERROR", "dashboard_action", str(exc)) + return self._json({"ok": False, "error": str(exc)}, 500) + + self.server = ThreadingHTTPServer((self.config.dashboard_host, self.config.dashboard_port), Handler) + threading.Thread(target=self.server.serve_forever, name="erp-local-agent-dashboard", daemon=True).start() + self.db.set_meta("dashboard_url", f"http://{self.config.dashboard_host}:{self.config.dashboard_port}") + self.db.record_event("INFO", "dashboard_started", f"Local dashboard started on {self.config.dashboard_host}:{self.config.dashboard_port}") + self.logger.info("ERP Local Agent dashboard: http://%s:%s", self.config.dashboard_host, self.config.dashboard_port) + + def _delayed_install(self): + time.sleep(1.0) + self.updater.install_downloaded() + + def status(self) -> dict[str, Any]: + total, used, free = shutil.disk_usage(self.config.storage_root) + meta = self.db.meta_snapshot() + latest_tally = self.db.latest_row("tally_status_history") + tally = {"connected": None, "company_count": 0, "companies": [], "error": None} + if latest_tally: + try: + tally = json.loads(latest_tally.get("payload_json") or "{}") or tally + except Exception: + pass + return { + "agent": { + "name": AGENT_NAME, + "version": __version__, + "node_code": self.config.node_code, + "erp_base_url": self.config.erp_base_url, + "tunnel_enabled": self.config.tunnel_enabled, + "connection_state": meta.get("connection_state", "unknown"), + "last_connection_event_utc": meta.get("last_connection_event_utc"), + "last_heartbeat_utc": meta.get("last_heartbeat_utc"), + }, + "storage": { + "root": str(self.config.storage_root), + "total_bytes": total, + "used_bytes": used, + "free_bytes": free, + }, + "tally": tally, + "update": self.updater.status(), + "dashboard": {"host": self.config.dashboard_host, "port": self.config.dashboard_port}, + } + + def history(self) -> dict[str, Any]: + return { + "connections": self.db.recent_rows("connection_history", 20), + "updates": self.db.recent_rows("update_history", 20), + "commands": self.db.recent_rows("command_history", 20), + "events": self.db.recent_rows("event_log", 30), + } + + @staticmethod + def _page() -> str: + return r''' + +ERP Local Agent + +

ERP Local Agent

Local dashboard — available only on this computer
+
+
ERP Connection
Loading…
+
Agent Version
—
+
TallyPrime
—
+
Storage Free
—
+
+

Updates

The agent checks for new versions automatically, but it will never install an update without a local user clicking Install.

+
+

Tally

Tally is checked locally at 127.0.0.1:9000. Port 9000 is not exposed to the internet.

+

Recent History

Loading…
+
+''' diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/db.py b/app/modules/documents/local_agent_runtime/erp_local_agent/db.py index 25188e1..8280964 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/db.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/db.py @@ -1,8 +1,14 @@ from __future__ import annotations +import json import sqlite3 -from pathlib import Path from datetime import datetime, timezone +from pathlib import Path +from typing import Any + + +def _utc_now() -> str: + return datetime.now(timezone.utc).isoformat() class LocalDB: @@ -12,56 +18,275 @@ class LocalDB: self._init() def connect(self): - conn = sqlite3.connect(self.db_path) + conn = sqlite3.connect(self.db_path, timeout=30) conn.row_factory = sqlite3.Row + conn.execute("PRAGMA busy_timeout=30000") return conn def _init(self) -> None: with self.connect() as conn: - conn.execute( + conn.executescript( """ + PRAGMA journal_mode=WAL; CREATE TABLE IF NOT EXISTS processed_storage_jobs ( job_id TEXT PRIMARY KEY, local_path TEXT NOT NULL, sha256 TEXT NOT NULL, file_size INTEGER NOT NULL, processed_at TEXT NOT NULL - ) - """ - ) - conn.execute( - """ + ); CREATE TABLE IF NOT EXISTS processed_download_requests ( request_id TEXT PRIMARY KEY, local_path TEXT NOT NULL, sha256 TEXT NOT NULL, file_size INTEGER NOT NULL, processed_at TEXT NOT NULL - ) + ); + CREATE TABLE IF NOT EXISTS agent_meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL, + updated_at_utc TEXT NOT NULL + ); + CREATE TABLE IF NOT EXISTS connection_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_type TEXT NOT NULL, + occurred_at_utc TEXT NOT NULL, + transport TEXT NOT NULL DEFAULT 'websocket', + endpoint TEXT NOT NULL DEFAULT '', + detail TEXT NULL + ); + CREATE TABLE IF NOT EXISTS heartbeat_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + occurred_at_utc TEXT NOT NULL, + status TEXT NOT NULL, + free_bytes INTEGER NULL, + detail TEXT NULL + ); + CREATE TABLE IF NOT EXISTS command_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + command_id TEXT NOT NULL DEFAULT '', + action TEXT NOT NULL DEFAULT '', + received_at_utc TEXT NOT NULL, + completed_at_utc TEXT NULL, + status TEXT NOT NULL DEFAULT 'received', + error_message TEXT NULL + ); + CREATE TABLE IF NOT EXISTS update_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_type TEXT NOT NULL, + from_version TEXT NOT NULL DEFAULT '', + to_version TEXT NOT NULL DEFAULT '', + occurred_at_utc TEXT NOT NULL, + status TEXT NOT NULL DEFAULT '', + sha256 TEXT NOT NULL DEFAULT '', + package_path TEXT NOT NULL DEFAULT '', + error_message TEXT NULL, + details_json TEXT NULL + ); + CREATE TABLE IF NOT EXISTS diagnostic_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + diagnostic_type TEXT NOT NULL, + occurred_at_utc TEXT NOT NULL, + status TEXT NOT NULL, + detail TEXT NULL + ); + CREATE TABLE IF NOT EXISTS tally_status_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + checked_at_utc TEXT NOT NULL, + connected INTEGER NOT NULL, + company_count INTEGER NOT NULL DEFAULT 0, + error_message TEXT NULL, + payload_json TEXT NULL + ); + CREATE TABLE IF NOT EXISTS storage_status_history ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + checked_at_utc TEXT NOT NULL, + storage_root TEXT NOT NULL, + total_bytes INTEGER NULL, + used_bytes INTEGER NULL, + free_bytes INTEGER NULL + ); + CREATE TABLE IF NOT EXISTS event_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + level TEXT NOT NULL, + event_type TEXT NOT NULL, + occurred_at_utc TEXT NOT NULL, + message TEXT NOT NULL, + details_json TEXT NULL + ); + CREATE INDEX IF NOT EXISTS ix_connection_history_time ON connection_history(occurred_at_utc); + CREATE INDEX IF NOT EXISTS ix_heartbeat_history_time ON heartbeat_history(occurred_at_utc); + CREATE INDEX IF NOT EXISTS ix_command_history_time ON command_history(received_at_utc); + CREATE INDEX IF NOT EXISTS ix_update_history_time ON update_history(occurred_at_utc); + CREATE INDEX IF NOT EXISTS ix_event_log_time ON event_log(occurred_at_utc); """ ) conn.commit() + def set_meta(self, key: str, value: Any) -> None: + with self.connect() as conn: + conn.execute( + """INSERT INTO agent_meta(key, value, updated_at_utc) VALUES (?, ?, ?) + ON CONFLICT(key) DO UPDATE SET value=excluded.value, updated_at_utc=excluded.updated_at_utc""", + (str(key), str(value), _utc_now()), + ) + conn.commit() + + def get_meta(self, key: str, default: str | None = None) -> str | None: + with self.connect() as conn: + row = conn.execute("SELECT value FROM agent_meta WHERE key=?", (str(key),)).fetchone() + return str(row["value"]) if row else default + + def record_event(self, level: str, event_type: str, message: str, details: Any = None) -> None: + details_json = None if details is None else json.dumps(details, ensure_ascii=False, default=str) + with self.connect() as conn: + conn.execute( + "INSERT INTO event_log(level, event_type, occurred_at_utc, message, details_json) VALUES (?, ?, ?, ?, ?)", + (str(level).upper(), str(event_type), _utc_now(), str(message), details_json), + ) + conn.commit() + + def record_connection(self, event_type: str, endpoint: str, detail: str | None = None) -> None: + with self.connect() as conn: + conn.execute( + "INSERT INTO connection_history(event_type, occurred_at_utc, transport, endpoint, detail) VALUES (?, ?, 'websocket', ?, ?)", + (str(event_type), _utc_now(), str(endpoint or ""), detail), + ) + conn.commit() + self.set_meta("connection_state", "connected" if event_type == "connected" else "disconnected") + self.set_meta("last_connection_event_utc", _utc_now()) + + def record_heartbeat(self, status: str, free_bytes: int | None = None, detail: str | None = None) -> None: + now = _utc_now() + with self.connect() as conn: + conn.execute( + "INSERT INTO heartbeat_history(occurred_at_utc, status, free_bytes, detail) VALUES (?, ?, ?, ?)", + (now, str(status), free_bytes, detail), + ) + conn.commit() + self.set_meta("last_heartbeat_status", status) + self.set_meta("last_heartbeat_utc", now) + + def record_storage_status(self, storage_root: str, total_bytes: int, used_bytes: int, free_bytes: int) -> None: + with self.connect() as conn: + conn.execute( + "INSERT INTO storage_status_history(checked_at_utc, storage_root, total_bytes, used_bytes, free_bytes) VALUES (?, ?, ?, ?, ?)", + (_utc_now(), str(storage_root), int(total_bytes), int(used_bytes), int(free_bytes)), + ) + conn.commit() + + def command_received(self, command_id: str, action: str) -> int: + with self.connect() as conn: + cursor = conn.execute( + "INSERT INTO command_history(command_id, action, received_at_utc, status) VALUES (?, ?, ?, 'received')", + (str(command_id or ""), str(action or ""), _utc_now()), + ) + conn.commit() + return int(cursor.lastrowid) + + def command_completed(self, row_id: int, ok: bool, error: str | None = None) -> None: + with self.connect() as conn: + conn.execute( + "UPDATE command_history SET completed_at_utc=?, status=?, error_message=? WHERE id=?", + (_utc_now(), "completed" if ok else "failed", error, int(row_id)), + ) + conn.commit() + + def record_update(self, event_type: str, from_version: str, to_version: str, status: str = "", sha256: str = "", package_path: str = "", error_message: str | None = None, details: Any = None) -> None: + details_json = None if details is None else json.dumps(details, ensure_ascii=False, default=str) + with self.connect() as conn: + conn.execute( + """INSERT INTO update_history(event_type, from_version, to_version, occurred_at_utc, status, sha256, package_path, error_message, details_json) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""", + (event_type, from_version, to_version, _utc_now(), status, sha256, package_path, error_message, details_json), + ) + conn.commit() + + def record_diagnostic(self, diagnostic_type: str, status: str, detail: str | None = None) -> None: + with self.connect() as conn: + conn.execute( + "INSERT INTO diagnostic_history(diagnostic_type, occurred_at_utc, status, detail) VALUES (?, ?, ?, ?)", + (str(diagnostic_type), _utc_now(), str(status), detail), + ) + conn.commit() + + def record_tally_status(self, status: dict[str, Any]) -> None: + with self.connect() as conn: + conn.execute( + "INSERT INTO tally_status_history(checked_at_utc, connected, company_count, error_message, payload_json) VALUES (?, ?, ?, ?, ?)", + ( + _utc_now(), + 1 if status.get("connected") else 0, + int(status.get("company_count") or 0), + str(status.get("error") or "") or None, + json.dumps(status, ensure_ascii=False, default=str), + ), + ) + conn.commit() + def record_storage_job(self, job_id: str, local_path: str, sha256: str, file_size: int) -> None: with self.connect() as conn: conn.execute( - """ - INSERT OR REPLACE INTO processed_storage_jobs - (job_id, local_path, sha256, file_size, processed_at) - VALUES (?, ?, ?, ?, ?) - """, - (str(job_id), local_path, sha256, int(file_size), datetime.now(timezone.utc).isoformat()), + """INSERT OR REPLACE INTO processed_storage_jobs + (job_id, local_path, sha256, file_size, processed_at) VALUES (?, ?, ?, ?, ?)""", + (str(job_id), local_path, sha256, int(file_size), _utc_now()), ) conn.commit() def record_download_request(self, request_id: str, local_path: str, sha256: str, file_size: int) -> None: with self.connect() as conn: conn.execute( - """ - INSERT OR REPLACE INTO processed_download_requests - (request_id, local_path, sha256, file_size, processed_at) - VALUES (?, ?, ?, ?, ?) - """, - (str(request_id), local_path, sha256, int(file_size), datetime.now(timezone.utc).isoformat()), + """INSERT OR REPLACE INTO processed_download_requests + (request_id, local_path, sha256, file_size, processed_at) VALUES (?, ?, ?, ?, ?)""", + (str(request_id), local_path, sha256, int(file_size), _utc_now()), ) conn.commit() + + def recent_rows(self, table: str, limit: int = 20) -> list[dict[str, Any]]: + allowed = { + "connection_history", "heartbeat_history", "command_history", "update_history", + "diagnostic_history", "tally_status_history", "storage_status_history", "event_log", + } + if table not in allowed: + raise ValueError("Unsupported history table") + limit = max(1, min(int(limit), 200)) + with self.connect() as conn: + rows = conn.execute(f"SELECT * FROM {table} ORDER BY id DESC LIMIT ?", (limit,)).fetchall() + return [dict(row) for row in rows] + + + def latest_row(self, table: str) -> dict[str, Any] | None: + rows = self.recent_rows(table, 1) + return rows[0] if rows else None + + def prune_history(self) -> None: + rules = { + "heartbeat_history": 14, + "diagnostic_history": 30, + "tally_status_history": 30, + "storage_status_history": 30, + "connection_history": 90, + "command_history": 90, + "event_log": 90, + } + with self.connect() as conn: + for table, days in rules.items(): + time_column = { + "heartbeat_history": "occurred_at_utc", + "diagnostic_history": "occurred_at_utc", + "tally_status_history": "checked_at_utc", + "storage_status_history": "checked_at_utc", + "connection_history": "occurred_at_utc", + "command_history": "received_at_utc", + "event_log": "occurred_at_utc", + }[table] + conn.execute( + f"DELETE FROM {table} WHERE datetime({time_column}) < datetime('now', ?)", + (f"-{days} days",), + ) + conn.commit() + + def meta_snapshot(self) -> dict[str, str]: + with self.connect() as conn: + rows = conn.execute("SELECT key, value FROM agent_meta ORDER BY key").fetchall() + return {str(row["key"]): str(row["value"]) for row in rows} diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/main.py b/app/modules/documents/local_agent_runtime/erp_local_agent/main.py index 74b26a5..2111eb9 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/main.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/main.py @@ -6,8 +6,10 @@ from pathlib import Path import threading import time +from . import __version__ from .client import ERPClient from .config import load_config +from .dashboard import AgentDashboard from .db import LocalDB from .logger import setup_logger from .sync import StorageAgent @@ -26,21 +28,33 @@ def main() -> int: args = build_parser().parse_args() root = Path.cwd() logger = setup_logger(root) + db = None try: config = load_config(args.env) db = LocalDB(root / "data" / "agent.db") + db.prune_history() + db.set_meta("agent_version", __version__) + db.set_meta("node_code", config.node_code) + db.set_meta("erp_base_url", config.erp_base_url) + db.record_event("INFO", "agent_started", f"ERP Local Agent {__version__} started") client = ERPClient(config) agent = StorageAgent(config, client, db, logger) - updater = AgentUpdater(config, client, logger, root) + updater = AgentUpdater(config, client, logger, root, db=db) + dashboard = AgentDashboard(config, db, updater, logger, root) + dashboard.start() + if args.once: - updater.maybe_update(force=True) + updater.maybe_check(force=True) agent.run_once() return 0 + def update_worker(): while True: - updater.maybe_update() + updater.maybe_check() time.sleep(max(60, config.update_check_interval_seconds)) - threading.Thread(target=update_worker, name="erp-local-agent-updater", daemon=True).start() + + threading.Thread(target=update_worker, name="erp-local-agent-update-checker", daemon=True).start() + if config.tunnel_enabled: asyncio.run(StorageAgentTunnel(agent).run_forever()) else: @@ -48,9 +62,16 @@ def main() -> int: return 0 except KeyboardInterrupt: logger.info("ERP Local Agent stopped by user") + if db is not None: + db.record_event("INFO", "agent_stopped", "ERP Local Agent stopped by user") return 0 except Exception as exc: logger.exception("ERP Local Agent failed: %s", exc) + if db is not None: + try: + db.record_event("ERROR", "agent_failed", str(exc)) + except Exception: + pass return 1 diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py b/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py index 07f0a84..572eda6 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/sync.py @@ -49,13 +49,16 @@ class StorageAgent: "free_bytes": free, "agent_version": __version__, "agent_name": "ERP Local Agent", - "capabilities": ["storage", "tally", "accounting_act"], + "capabilities": ["storage", "tally", "accounting_act", "local_dashboard", "manual_updates"], "agent_time_utc": datetime.now(timezone.utc).isoformat(), } self.client.heartbeat(payload) self.last_heartbeat = now + self.db.record_heartbeat("sent", free_bytes=free) + self.db.record_storage_status(str(self.config.storage_root), total, used, free) self.logger.info("Heartbeat sent. free_gb=%.2f", free / (1024 ** 3)) except Exception as exc: + self.db.record_heartbeat("failed", detail=str(exc)) self.logger.exception("Heartbeat failed: %s", exc) def process_storage_jobs(self) -> None: diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py index c3bc2ba..73f4000 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/tunnel.py @@ -6,8 +6,8 @@ from typing import Any import websockets -from .sync import StorageAgent from .commands import AgentCommandProcessor +from .sync import StorageAgent class StorageAgentTunnel: @@ -15,20 +15,30 @@ class StorageAgentTunnel: self.agent = agent self.config = agent.config self.logger = agent.logger + self.db = agent.db self.command_processor = AgentCommandProcessor(self.config, self.logger) async def run_forever(self) -> None: self.logger.info("Starting tunnel mode for node=%s", self.config.node_code) + delay = max(5, self.config.tunnel_reconnect_seconds) while True: try: await self._connect_once() + delay = max(5, self.config.tunnel_reconnect_seconds) except KeyboardInterrupt: raise except Exception as exc: self.logger.exception("Tunnel disconnected/error: %s", exc) - await asyncio.sleep(max(5, self.config.tunnel_reconnect_seconds)) + try: + self.db.record_connection("disconnected", self.config.tunnel_url.split("?", 1)[0], str(exc)) + self.db.record_event("WARNING", "tunnel_disconnected", str(exc)) + except Exception: + pass + await asyncio.sleep(delay) + delay = min(60, max(delay + 5, delay * 2)) async def _connect_once(self) -> None: + endpoint = self.config.tunnel_url.split("?", 1)[0] async with websockets.connect( self.config.tunnel_url, ping_interval=30, @@ -36,6 +46,8 @@ class StorageAgentTunnel: close_timeout=10, max_size=1024 * 1024, ) as websocket: + self.db.record_connection("connected", endpoint) + self.db.record_event("INFO", "tunnel_connected", f"WebSocket connected to {endpoint}") await websocket.send(self._json_text({"type": "ready", "agent_time_utc": self._now()})) self.agent._maybe_heartbeat(force=True) async for raw in websocket: @@ -53,14 +65,18 @@ class StorageAgentTunnel: await asyncio.to_thread(self.agent.process_storage_jobs_from_payload, jobs) await asyncio.to_thread(self.agent.process_download_requests_from_payload, requests_) for command in commands: + command_id = str(command.get("command_id") or command.get("id") or "") + action = str(command.get("action") or "") + history_id = self.db.command_received(command_id, action) result = await asyncio.to_thread(self.command_processor.process, command) + self.db.command_completed(history_id, bool(result.get("ok")), result.get("error")) await websocket.send(self._json_text(result)) await websocket.send(self._json_text({ "type": "agent_status", "jobs_seen": len(jobs), "requests_seen": len(requests_), "commands_seen": len(commands), - "capabilities": ["storage", "tally", "accounting_act"], + "capabilities": ["storage", "tally", "accounting_act", "local_dashboard", "manual_updates"], "agent_time_utc": self._now(), })) continue diff --git a/app/modules/documents/local_agent_runtime/erp_local_agent/updater.py b/app/modules/documents/local_agent_runtime/erp_local_agent/updater.py index d9386d1..64fb6b1 100644 --- a/app/modules/documents/local_agent_runtime/erp_local_agent/updater.py +++ b/app/modules/documents/local_agent_runtime/erp_local_agent/updater.py @@ -5,6 +5,7 @@ import os from pathlib import Path import shutil import subprocess +import threading import time import zipfile @@ -12,83 +13,171 @@ from . import __version__ class AgentUpdater: - def __init__(self, config, client, logger, install_dir: Path): + """Manual-install updater. + + Background work checks the ERP manifest only. Download and installation happen only + when the local user requests them from the dashboard. + """ + + def __init__(self, config, client, logger, install_dir: Path, db=None): self.config = config self.client = client self.logger = logger self.install_dir = install_dir.resolve() + self.db = db self.last_check = 0.0 self.busy = False + self.lock = threading.RLock() + self.latest_manifest: dict = {} - def maybe_update(self, force: bool = False) -> bool: - if not self.config.auto_update or self.busy: + def _record(self, event_type: str, to_version: str = "", status: str = "", **kwargs) -> None: + if self.db is not None: + try: + self.db.record_update(event_type, __version__, to_version, status=status, **kwargs) + except Exception: + self.logger.exception("Could not record updater history") + + def check_for_update(self, force: bool = False) -> dict: + with self.lock: + now = time.time() + if not force and now - self.last_check < self.config.update_check_interval_seconds and self.latest_manifest: + return self.status() + self.last_check = now + try: + manifest = self.client.update_manifest() + latest = str(manifest.get("latest_version") or "").strip() + expected = str(manifest.get("sha256") or "").lower().strip() + self.latest_manifest = dict(manifest or {}) + available = bool(latest and latest != __version__) + if self.db is not None: + self.db.set_meta("latest_agent_version", latest or __version__) + self.db.set_meta("update_available", "1" if available else "0") + self.db.set_meta("last_update_check_utc", time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())) + self._record("check", latest, "available" if available else "current", sha256=expected) + return self.status() + except Exception as exc: + self._record("check", "", "failed", error_message=str(exc)) + self.logger.exception("ERP Local Agent update check failed: %s", exc) + raise + + def maybe_check(self, force: bool = False) -> bool: + if not self.config.auto_update: return False - now = time.time() - if not force and now - self.last_check < self.config.update_check_interval_seconds: - return False - self.last_check = now try: - manifest = self.client.update_manifest() - latest = str(manifest.get("latest_version") or "").strip() - if not latest or latest == __version__: - return False - expected = str(manifest.get("sha256") or "").lower().strip() - package = self.client.download_update_package() - actual = hashlib.sha256(package).hexdigest().lower() - if len(expected) != 64 or actual != expected: - raise RuntimeError("ERP Local Agent update package SHA256 verification failed.") - self._stage_and_restart(latest, package) - return True - except Exception as exc: - self.logger.exception("ERP Local Agent update check failed: %s", exc) + state = self.check_for_update(force=force) + return bool(state.get("update_available")) + except Exception: return False - def _stage_and_restart(self, latest: str, package: bytes) -> None: - self.busy = True - updates = self.install_dir / "updates" - updates.mkdir(parents=True, exist_ok=True) - safe = "".join(ch for ch in latest if ch.isalnum() or ch in ".-_") or "update" - package_path = updates / f"ERP_Local_Agent_{safe}.zip" - package_path.write_bytes(package) - staged = updates / f"staged_{safe}" - shutil.rmtree(staged, ignore_errors=True) - staged.mkdir(parents=True, exist_ok=True) - with zipfile.ZipFile(package_path, "r") as archive: - archive.extractall(staged) - if not (staged / "erp_local_agent" / "__init__.py").exists(): - raise RuntimeError("ERP Local Agent update package is incomplete.") - script_path = updates / f"apply_{safe}.ps1" - script_path.write_text(self._powershell_update_script(staged), encoding="utf-8") - flags = getattr(subprocess, "CREATE_NEW_PROCESS_GROUP", 0) | getattr(subprocess, "DETACHED_PROCESS", 0) - subprocess.Popen( - ["powershell.exe", "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(script_path), "-ParentPid", str(os.getpid())], - cwd=str(self.install_dir), creationflags=flags, close_fds=True, - ) - self.logger.warning("ERP Local Agent update %s staged; restarting.", latest) - os._exit(0) + def download_latest(self) -> dict: + with self.lock: + if self.busy: + raise RuntimeError("Another update operation is already in progress.") + self.busy = True + try: + if not self.latest_manifest: + self.check_for_update(force=True) + latest = str(self.latest_manifest.get("latest_version") or "").strip() + expected = str(self.latest_manifest.get("sha256") or "").lower().strip() + if not latest or latest == __version__: + return self.status() + package = self.client.download_update_package() + actual = hashlib.sha256(package).hexdigest().lower() + if len(expected) != 64 or actual != expected: + raise RuntimeError("ERP Local Agent update package SHA256 verification failed.") + updates = self.install_dir / "updates" + updates.mkdir(parents=True, exist_ok=True) + safe = "".join(ch for ch in latest if ch.isalnum() or ch in ".-_") or "update" + package_path = updates / f"ERP_Local_Agent_{safe}.zip" + package_path.write_bytes(package) + staged = updates / f"staged_{safe}" + shutil.rmtree(staged, ignore_errors=True) + staged.mkdir(parents=True, exist_ok=True) + with zipfile.ZipFile(package_path, "r") as archive: + archive.extractall(staged) + if not (staged / "erp_local_agent" / "__init__.py").exists(): + raise RuntimeError("ERP Local Agent update package is incomplete.") + self._record("download", latest, "downloaded", sha256=actual, package_path=str(package_path)) + if self.db is not None: + self.db.set_meta("staged_update_version", latest) + self.db.set_meta("staged_update_path", str(staged)) + return self.status() + except Exception as exc: + latest = str(self.latest_manifest.get("latest_version") or "") + self._record("download", latest, "failed", error_message=str(exc)) + raise + finally: + self.busy = False - def _powershell_update_script(self, staged: Path) -> str: + def install_downloaded(self) -> None: + with self.lock: + latest = str((self.db.get_meta("staged_update_version") if self.db else "") or self.latest_manifest.get("latest_version") or "").strip() + if not latest: + raise RuntimeError("No downloaded ERP Local Agent update is available to install.") + safe = "".join(ch for ch in latest if ch.isalnum() or ch in ".-_") or "update" + staged = self.install_dir / "updates" / f"staged_{safe}" + if not (staged / "erp_local_agent" / "__init__.py").exists(): + raise RuntimeError("The downloaded update is not staged correctly. Download it again.") + script_path = self.install_dir / "updates" / f"apply_{safe}.ps1" + script_path.write_text(self._powershell_update_script(staged, latest), encoding="utf-8") + self._record("install_requested", latest, "pending", package_path=str(staged)) + flags = getattr(subprocess, "CREATE_NEW_PROCESS_GROUP", 0) | getattr(subprocess, "DETACHED_PROCESS", 0) + subprocess.Popen( + ["powershell.exe", "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(script_path), "-ParentPid", str(os.getpid())], + cwd=str(self.install_dir), creationflags=flags, close_fds=True, + ) + self.logger.warning("ERP Local Agent update %s installation requested by local user; restarting.", latest) + time.sleep(0.3) + os._exit(0) + + def status(self) -> dict: + latest = str(self.latest_manifest.get("latest_version") or (self.db.get_meta("latest_agent_version") if self.db else "") or __version__).strip() + staged_version = str((self.db.get_meta("staged_update_version") if self.db else "") or "").strip() + staged_ok = False + if staged_version: + safe = "".join(ch for ch in staged_version if ch.isalnum() or ch in ".-_") or "update" + staged_ok = (self.install_dir / "updates" / f"staged_{safe}" / "erp_local_agent" / "__init__.py").exists() + return { + "current_version": __version__, + "latest_version": latest, + "update_available": bool(latest and latest != __version__), + "downloaded": bool(staged_ok), + "staged_version": staged_version if staged_ok else "", + "busy": bool(self.busy), + "automatic_check_enabled": bool(self.config.auto_update), + "automatic_install_enabled": False, + } + + def _powershell_update_script(self, staged: Path, latest: str) -> str: install = str(self.install_dir).replace("'", "''") stage = str(staged.resolve()).replace("'", "''") + latest_safe = str(latest).replace("'", "''") lines = [ "param([int]$ParentPid)", '$ErrorActionPreference = "Stop"', "$InstallDir = '" + install + "'", "$StagedDir = '" + stage + "'", + "$TargetVersion = '" + latest_safe + "'", "$TaskName = 'ERP Local Agent'", "$BackupDir = Join-Path $InstallDir ('updates\\backup_' + (Get-Date -Format 'yyyyMMdd_HHmmss'))", "try {", " if ($ParentPid -gt 0) { try { Wait-Process -Id $ParentPid -Timeout 120 -ErrorAction SilentlyContinue } catch {} }", " try { Stop-ScheduledTask -TaskName $TaskName -ErrorAction SilentlyContinue } catch {}", + " Start-Sleep -Seconds 2", + " Get-CimInstance Win32_Process | Where-Object { $_.CommandLine -match 'erp_local_agent\\.main' -and $_.CommandLine -like ('*' + $InstallDir + '*') } | ForEach-Object { try { Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue } catch {} }", + " Start-Sleep -Seconds 2", " New-Item -ItemType Directory -Force -Path $BackupDir | Out-Null", " if (Test-Path (Join-Path $InstallDir 'erp_local_agent')) { Copy-Item (Join-Path $InstallDir 'erp_local_agent') -Destination $BackupDir -Recurse -Force }", - " Remove-Item (Join-Path $InstallDir 'erp_local_agent') -Recurse -Force -ErrorAction SilentlyContinue", + " Remove-Item (Join-Path $InstallDir 'erp_local_agent') -Recurse -Force -ErrorAction Stop", " Copy-Item (Join-Path $StagedDir 'erp_local_agent') -Destination $InstallDir -Recurse -Force", - " foreach ($Item in @('requirements.txt','run_agent.bat','run_once.bat','check_config.py','check_config.bat','start_task_scheduler.bat','stop_task_scheduler.bat','status_task_scheduler.bat','install_task_scheduler.bat','uninstall_task_scheduler.bat')) { $Source=Join-Path $StagedDir $Item; if (Test-Path $Source) { Copy-Item $Source -Destination (Join-Path $InstallDir $Item) -Force } }", + " foreach ($Item in @('requirements.txt','run_agent.bat','run_once.bat','check_config.py','check_config.bat','open_dashboard.bat','start_task_scheduler.bat','stop_task_scheduler.bat','status_task_scheduler.bat','install_task_scheduler.bat','uninstall_task_scheduler.bat')) { $Source=Join-Path $StagedDir $Item; if (Test-Path $Source) { Copy-Item $Source -Destination (Join-Path $InstallDir $Item) -Force } }", " $Python = Join-Path $InstallDir '.venv\\Scripts\\python.exe'", " if (-not (Test-Path $Python)) { throw 'ERP Local Agent virtual environment is missing.' }", " & $Python -m pip install -r (Join-Path $InstallDir 'requirements.txt')", " if ($LASTEXITCODE -ne 0) { throw 'Dependency update failed.' }", + " Push-Location $InstallDir", + " try { $InstalledVersion = & $Python -c \"import erp_local_agent; print(getattr(erp_local_agent,'__version__',''))\" } finally { Pop-Location }", + " if (($InstalledVersion | Out-String).Trim() -ne $TargetVersion) { throw ('Installed version verification failed. Expected ' + $TargetVersion + ', got ' + (($InstalledVersion | Out-String).Trim())) }", " Start-ScheduledTask -TaskName $TaskName", "} catch {", " try { if (Test-Path (Join-Path $BackupDir 'erp_local_agent')) { Remove-Item (Join-Path $InstallDir 'erp_local_agent') -Recurse -Force -ErrorAction SilentlyContinue; Copy-Item (Join-Path $BackupDir 'erp_local_agent') -Destination $InstallDir -Recurse -Force }; Start-ScheduledTask -TaskName $TaskName -ErrorAction SilentlyContinue } catch {}", diff --git a/app/modules/documents/local_agent_runtime/open_dashboard.bat b/app/modules/documents/local_agent_runtime/open_dashboard.bat new file mode 100644 index 0000000..ff05dfd --- /dev/null +++ b/app/modules/documents/local_agent_runtime/open_dashboard.bat @@ -0,0 +1,2 @@ +@echo off +start "" "http://127.0.0.1:8788"