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
+
+
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.
+
+
+'''
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"