Add Local Agent dashboard SQLite history and manual updates
This commit is contained in:
@@ -4,7 +4,7 @@ import io
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
import zipfile
|
import zipfile
|
||||||
|
|
||||||
ERP_LOCAL_AGENT_VERSION = "1.1.0"
|
ERP_LOCAL_AGENT_VERSION = "1.2.0"
|
||||||
ERP_LOCAL_AGENT_NAME = "ERP Local Agent"
|
ERP_LOCAL_AGENT_NAME = "ERP Local Agent"
|
||||||
RUNTIME_ROOT = Path(__file__).resolve().parent / "local_agent_runtime"
|
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_ENABLED={str(bool(tunnel_enabled)).lower()}\n"
|
||||||
f"TUNNEL_RECONNECT_SECONDS={int(tunnel_reconnect_seconds or 10)}\n"
|
f"TUNNEL_RECONNECT_SECONDS={int(tunnel_reconnect_seconds or 10)}\n"
|
||||||
"AUTO_UPDATE=true\n"
|
"AUTO_UPDATE=true\n"
|
||||||
|
"AUTO_INSTALL_UPDATES=false\n"
|
||||||
"UPDATE_CHECK_INTERVAL_SECONDS=300\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:
|
if include_env and env_text is not None:
|
||||||
dst.writestr(".env", env_text)
|
dst.writestr(".env", env_text)
|
||||||
if include_admin_readme:
|
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()
|
return buffer.getvalue()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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.
|
||||||
@@ -1,2 +1,2 @@
|
|||||||
__version__ = "1.1.0"
|
__version__ = "1.2.0"
|
||||||
AGENT_NAME = "ERP Local Agent"
|
AGENT_NAME = "ERP Local Agent"
|
||||||
|
|||||||
@@ -2,8 +2,9 @@ from __future__ import annotations
|
|||||||
|
|
||||||
import os
|
import os
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from urllib.parse import urlencode
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
from urllib.parse import urlencode
|
||||||
|
|
||||||
from dotenv import load_dotenv
|
from dotenv import load_dotenv
|
||||||
|
|
||||||
from . import __version__
|
from . import __version__
|
||||||
@@ -26,6 +27,10 @@ class AgentConfig:
|
|||||||
tunnel_reconnect_seconds: int = 10
|
tunnel_reconnect_seconds: int = 10
|
||||||
auto_update: bool = True
|
auto_update: bool = True
|
||||||
update_check_interval_seconds: int = 300
|
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
|
@property
|
||||||
def headers(self) -> dict[str, str]:
|
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),
|
tunnel_reconnect_seconds=_get_int("TUNNEL_RECONNECT_SECONDS", 10),
|
||||||
auto_update=_get_bool("AUTO_UPDATE", True),
|
auto_update=_get_bool("AUTO_UPDATE", True),
|
||||||
update_check_interval_seconds=max(60, _get_int("UPDATE_CHECK_INTERVAL_SECONDS", 300)),
|
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))),
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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'''<!doctype html>
|
||||||
|
<html><head><meta charset="utf-8"><meta name="viewport" content="width=device-width,initial-scale=1">
|
||||||
|
<title>ERP Local Agent</title>
|
||||||
|
<style>
|
||||||
|
body{font-family:Segoe UI,Arial,sans-serif;background:#f4f6f8;color:#17202a;margin:0}.wrap{max-width:1100px;margin:28px auto;padding:0 18px}h1{margin:0 0 4px}.muted{color:#6b7280}.grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(230px,1fr));gap:14px;margin:22px 0}.card{background:#fff;border:1px solid #e5e7eb;border-radius:12px;padding:18px;box-shadow:0 1px 2px rgba(0,0,0,.04)}.label{font-size:12px;text-transform:uppercase;letter-spacing:.06em;color:#6b7280}.value{font-size:22px;font-weight:650;margin-top:8px}.ok{color:#087f5b}.bad{color:#c92a2a}.warn{color:#b26a00}button{border:0;border-radius:8px;padding:10px 14px;margin:4px 6px 4px 0;font-weight:600;cursor:pointer;background:#111827;color:#fff}button.secondary{background:#e5e7eb;color:#111827}button:disabled{opacity:.5;cursor:not-allowed}.section{background:#fff;border:1px solid #e5e7eb;border-radius:12px;padding:18px;margin:14px 0}table{width:100%;border-collapse:collapse;font-size:13px}th,td{text-align:left;padding:8px;border-bottom:1px solid #eee;vertical-align:top}#message{margin-top:10px;padding:10px;border-radius:8px;display:none}.msgok{display:block!important;background:#e6fcf5;color:#087f5b}.msgbad{display:block!important;background:#fff5f5;color:#c92a2a}code{background:#f3f4f6;padding:2px 5px;border-radius:4px}</style></head>
|
||||||
|
<body><div class="wrap"><h1>ERP Local Agent</h1><div class="muted">Local dashboard — available only on this computer</div>
|
||||||
|
<div class="grid">
|
||||||
|
<div class="card"><div class="label">ERP Connection</div><div id="connection" class="value">Loading…</div><div id="heartbeat" class="muted"></div></div>
|
||||||
|
<div class="card"><div class="label">Agent Version</div><div id="version" class="value">—</div><div id="latest" class="muted"></div></div>
|
||||||
|
<div class="card"><div class="label">TallyPrime</div><div id="tally" class="value">—</div><div id="companies" class="muted"></div></div>
|
||||||
|
<div class="card"><div class="label">Storage Free</div><div id="storage" class="value">—</div><div id="storageRoot" class="muted"></div></div>
|
||||||
|
</div>
|
||||||
|
<div class="section"><h2>Updates</h2><p class="muted">The agent checks for new versions automatically, but it will never install an update without a local user clicking Install.</p>
|
||||||
|
<button onclick="action('/api/update/check')">Check for Update</button><button class="secondary" id="download" onclick="action('/api/update/download')">Download Update</button><button id="install" onclick="installUpdate()">Install Downloaded Update</button><div id="message"></div></div>
|
||||||
|
<div class="section"><h2>Tally</h2><button onclick="action('/api/tally/check')">Test Tally Connection</button><p class="muted">Tally is checked locally at <code>127.0.0.1:9000</code>. Port 9000 is not exposed to the internet.</p></div>
|
||||||
|
<div class="section"><h2>Recent History</h2><div id="history" class="muted">Loading…</div></div>
|
||||||
|
</div>
|
||||||
|
<script>
|
||||||
|
function gb(n){return (n/1073741824).toFixed(1)+' GB'}
|
||||||
|
async function refresh(){try{const r=await fetch('/api/status',{cache:'no-store'});const s=await r.json();
|
||||||
|
let c=s.agent.connection_state||'unknown';let ce=document.getElementById('connection');ce.textContent=c.charAt(0).toUpperCase()+c.slice(1);ce.className='value '+(c==='connected'?'ok':'bad');
|
||||||
|
document.getElementById('heartbeat').textContent=s.agent.last_heartbeat_utc?'Last heartbeat: '+s.agent.last_heartbeat_utc:'No heartbeat recorded';
|
||||||
|
document.getElementById('version').textContent=s.agent.version;document.getElementById('latest').textContent='Latest: '+s.update.latest_version+(s.update.update_available?' • update available':' • current');
|
||||||
|
let te=document.getElementById('tally');let known=s.tally.connected!==null;te.textContent=!known?'Not checked':(s.tally.connected?'Connected':'Offline');te.className='value '+(!known?'warn':(s.tally.connected?'ok':'bad'));document.getElementById('companies').textContent=!known?'Use Test Tally Connection':(s.tally.connected?(s.tally.company_count+' loaded compan'+(s.tally.company_count===1?'y':'ies')):(s.tally.error||''));
|
||||||
|
document.getElementById('storage').textContent=gb(s.storage.free_bytes);document.getElementById('storageRoot').textContent=s.storage.root;
|
||||||
|
document.getElementById('download').disabled=!s.update.update_available;document.getElementById('install').disabled=!s.update.downloaded;await refreshHistory();}catch(e){console.log(e)}}
|
||||||
|
async function refreshHistory(){const r=await fetch('/api/history',{cache:'no-store'});const h=await r.json();let rows=h.events.slice(0,12).map(x=>'<tr><td>'+x.occurred_at_utc+'</td><td>'+x.level+'</td><td>'+x.event_type+'</td><td>'+escapeHtml(x.message)+'</td></tr>').join('');document.getElementById('history').innerHTML='<table><thead><tr><th>Time UTC</th><th>Level</th><th>Event</th><th>Message</th></tr></thead><tbody>'+rows+'</tbody></table>'}
|
||||||
|
function escapeHtml(v){return String(v||'').replace(/[&<>'"]/g,m=>({'&':'&','<':'<','>':'>',"'":''','"':'"'}[m]))}
|
||||||
|
async function action(url){msg('Working…',true);try{const r=await fetch(url,{method:'POST'});const j=await r.json();if(!r.ok||!j.ok)throw new Error(j.error||'Operation failed');msg('Completed successfully.',true);await refresh()}catch(e){msg(e.message,false)}}
|
||||||
|
async function installUpdate(){if(!confirm('Install the downloaded update now? The ERP Local Agent will restart.'))return;msg('Starting update installation…',true);try{const r=await fetch('/api/update/install',{method:'POST'});const j=await r.json();if(!r.ok||!j.ok)throw new Error(j.error||'Install failed');msg(j.message||'Update started. Refresh this page after the agent restarts.',true)}catch(e){msg(e.message,false)}}
|
||||||
|
function msg(t,ok){let e=document.getElementById('message');e.textContent=t;e.className=ok?'msgok':'msgbad'}
|
||||||
|
refresh();setInterval(refresh,5000);
|
||||||
|
</script></body></html>'''
|
||||||
@@ -1,8 +1,14 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
import sqlite3
|
import sqlite3
|
||||||
from pathlib import Path
|
|
||||||
from datetime import datetime, timezone
|
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:
|
class LocalDB:
|
||||||
@@ -12,56 +18,275 @@ class LocalDB:
|
|||||||
self._init()
|
self._init()
|
||||||
|
|
||||||
def connect(self):
|
def connect(self):
|
||||||
conn = sqlite3.connect(self.db_path)
|
conn = sqlite3.connect(self.db_path, timeout=30)
|
||||||
conn.row_factory = sqlite3.Row
|
conn.row_factory = sqlite3.Row
|
||||||
|
conn.execute("PRAGMA busy_timeout=30000")
|
||||||
return conn
|
return conn
|
||||||
|
|
||||||
def _init(self) -> None:
|
def _init(self) -> None:
|
||||||
with self.connect() as conn:
|
with self.connect() as conn:
|
||||||
conn.execute(
|
conn.executescript(
|
||||||
"""
|
"""
|
||||||
|
PRAGMA journal_mode=WAL;
|
||||||
CREATE TABLE IF NOT EXISTS processed_storage_jobs (
|
CREATE TABLE IF NOT EXISTS processed_storage_jobs (
|
||||||
job_id TEXT PRIMARY KEY,
|
job_id TEXT PRIMARY KEY,
|
||||||
local_path TEXT NOT NULL,
|
local_path TEXT NOT NULL,
|
||||||
sha256 TEXT NOT NULL,
|
sha256 TEXT NOT NULL,
|
||||||
file_size INTEGER NOT NULL,
|
file_size INTEGER NOT NULL,
|
||||||
processed_at TEXT NOT NULL
|
processed_at TEXT NOT NULL
|
||||||
)
|
);
|
||||||
"""
|
|
||||||
)
|
|
||||||
conn.execute(
|
|
||||||
"""
|
|
||||||
CREATE TABLE IF NOT EXISTS processed_download_requests (
|
CREATE TABLE IF NOT EXISTS processed_download_requests (
|
||||||
request_id TEXT PRIMARY KEY,
|
request_id TEXT PRIMARY KEY,
|
||||||
local_path TEXT NOT NULL,
|
local_path TEXT NOT NULL,
|
||||||
sha256 TEXT NOT NULL,
|
sha256 TEXT NOT NULL,
|
||||||
file_size INTEGER NOT NULL,
|
file_size INTEGER NOT NULL,
|
||||||
processed_at TEXT 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()
|
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:
|
def record_storage_job(self, job_id: str, local_path: str, sha256: str, file_size: int) -> None:
|
||||||
with self.connect() as conn:
|
with self.connect() as conn:
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"""
|
"""INSERT OR REPLACE INTO processed_storage_jobs
|
||||||
INSERT OR REPLACE INTO processed_storage_jobs
|
(job_id, local_path, sha256, file_size, processed_at) VALUES (?, ?, ?, ?, ?)""",
|
||||||
(job_id, local_path, sha256, file_size, processed_at)
|
(str(job_id), local_path, sha256, int(file_size), _utc_now()),
|
||||||
VALUES (?, ?, ?, ?, ?)
|
|
||||||
""",
|
|
||||||
(str(job_id), local_path, sha256, int(file_size), datetime.now(timezone.utc).isoformat()),
|
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
|
||||||
def record_download_request(self, request_id: str, local_path: str, sha256: str, file_size: int) -> None:
|
def record_download_request(self, request_id: str, local_path: str, sha256: str, file_size: int) -> None:
|
||||||
with self.connect() as conn:
|
with self.connect() as conn:
|
||||||
conn.execute(
|
conn.execute(
|
||||||
"""
|
"""INSERT OR REPLACE INTO processed_download_requests
|
||||||
INSERT OR REPLACE INTO processed_download_requests
|
(request_id, local_path, sha256, file_size, processed_at) VALUES (?, ?, ?, ?, ?)""",
|
||||||
(request_id, local_path, sha256, file_size, processed_at)
|
(str(request_id), local_path, sha256, int(file_size), _utc_now()),
|
||||||
VALUES (?, ?, ?, ?, ?)
|
|
||||||
""",
|
|
||||||
(str(request_id), local_path, sha256, int(file_size), datetime.now(timezone.utc).isoformat()),
|
|
||||||
)
|
)
|
||||||
conn.commit()
|
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}
|
||||||
|
|||||||
@@ -6,8 +6,10 @@ from pathlib import Path
|
|||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
|
|
||||||
|
from . import __version__
|
||||||
from .client import ERPClient
|
from .client import ERPClient
|
||||||
from .config import load_config
|
from .config import load_config
|
||||||
|
from .dashboard import AgentDashboard
|
||||||
from .db import LocalDB
|
from .db import LocalDB
|
||||||
from .logger import setup_logger
|
from .logger import setup_logger
|
||||||
from .sync import StorageAgent
|
from .sync import StorageAgent
|
||||||
@@ -26,21 +28,33 @@ def main() -> int:
|
|||||||
args = build_parser().parse_args()
|
args = build_parser().parse_args()
|
||||||
root = Path.cwd()
|
root = Path.cwd()
|
||||||
logger = setup_logger(root)
|
logger = setup_logger(root)
|
||||||
|
db = None
|
||||||
try:
|
try:
|
||||||
config = load_config(args.env)
|
config = load_config(args.env)
|
||||||
db = LocalDB(root / "data" / "agent.db")
|
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)
|
client = ERPClient(config)
|
||||||
agent = StorageAgent(config, client, db, logger)
|
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:
|
if args.once:
|
||||||
updater.maybe_update(force=True)
|
updater.maybe_check(force=True)
|
||||||
agent.run_once()
|
agent.run_once()
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
def update_worker():
|
def update_worker():
|
||||||
while True:
|
while True:
|
||||||
updater.maybe_update()
|
updater.maybe_check()
|
||||||
time.sleep(max(60, config.update_check_interval_seconds))
|
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:
|
if config.tunnel_enabled:
|
||||||
asyncio.run(StorageAgentTunnel(agent).run_forever())
|
asyncio.run(StorageAgentTunnel(agent).run_forever())
|
||||||
else:
|
else:
|
||||||
@@ -48,9 +62,16 @@ def main() -> int:
|
|||||||
return 0
|
return 0
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
logger.info("ERP Local Agent stopped by user")
|
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
|
return 0
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.exception("ERP Local Agent failed: %s", 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
|
return 1
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -49,13 +49,16 @@ class StorageAgent:
|
|||||||
"free_bytes": free,
|
"free_bytes": free,
|
||||||
"agent_version": __version__,
|
"agent_version": __version__,
|
||||||
"agent_name": "ERP Local Agent",
|
"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(),
|
"agent_time_utc": datetime.now(timezone.utc).isoformat(),
|
||||||
}
|
}
|
||||||
self.client.heartbeat(payload)
|
self.client.heartbeat(payload)
|
||||||
self.last_heartbeat = now
|
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))
|
self.logger.info("Heartbeat sent. free_gb=%.2f", free / (1024 ** 3))
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
self.db.record_heartbeat("failed", detail=str(exc))
|
||||||
self.logger.exception("Heartbeat failed: %s", exc)
|
self.logger.exception("Heartbeat failed: %s", exc)
|
||||||
|
|
||||||
def process_storage_jobs(self) -> None:
|
def process_storage_jobs(self) -> None:
|
||||||
|
|||||||
@@ -6,8 +6,8 @@ from typing import Any
|
|||||||
|
|
||||||
import websockets
|
import websockets
|
||||||
|
|
||||||
from .sync import StorageAgent
|
|
||||||
from .commands import AgentCommandProcessor
|
from .commands import AgentCommandProcessor
|
||||||
|
from .sync import StorageAgent
|
||||||
|
|
||||||
|
|
||||||
class StorageAgentTunnel:
|
class StorageAgentTunnel:
|
||||||
@@ -15,20 +15,30 @@ class StorageAgentTunnel:
|
|||||||
self.agent = agent
|
self.agent = agent
|
||||||
self.config = agent.config
|
self.config = agent.config
|
||||||
self.logger = agent.logger
|
self.logger = agent.logger
|
||||||
|
self.db = agent.db
|
||||||
self.command_processor = AgentCommandProcessor(self.config, self.logger)
|
self.command_processor = AgentCommandProcessor(self.config, self.logger)
|
||||||
|
|
||||||
async def run_forever(self) -> None:
|
async def run_forever(self) -> None:
|
||||||
self.logger.info("Starting tunnel mode for node=%s", self.config.node_code)
|
self.logger.info("Starting tunnel mode for node=%s", self.config.node_code)
|
||||||
|
delay = max(5, self.config.tunnel_reconnect_seconds)
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
await self._connect_once()
|
await self._connect_once()
|
||||||
|
delay = max(5, self.config.tunnel_reconnect_seconds)
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
raise
|
raise
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
self.logger.exception("Tunnel disconnected/error: %s", 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:
|
async def _connect_once(self) -> None:
|
||||||
|
endpoint = self.config.tunnel_url.split("?", 1)[0]
|
||||||
async with websockets.connect(
|
async with websockets.connect(
|
||||||
self.config.tunnel_url,
|
self.config.tunnel_url,
|
||||||
ping_interval=30,
|
ping_interval=30,
|
||||||
@@ -36,6 +46,8 @@ class StorageAgentTunnel:
|
|||||||
close_timeout=10,
|
close_timeout=10,
|
||||||
max_size=1024 * 1024,
|
max_size=1024 * 1024,
|
||||||
) as websocket:
|
) 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()}))
|
await websocket.send(self._json_text({"type": "ready", "agent_time_utc": self._now()}))
|
||||||
self.agent._maybe_heartbeat(force=True)
|
self.agent._maybe_heartbeat(force=True)
|
||||||
async for raw in websocket:
|
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_storage_jobs_from_payload, jobs)
|
||||||
await asyncio.to_thread(self.agent.process_download_requests_from_payload, requests_)
|
await asyncio.to_thread(self.agent.process_download_requests_from_payload, requests_)
|
||||||
for command in commands:
|
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)
|
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(result))
|
||||||
await websocket.send(self._json_text({
|
await websocket.send(self._json_text({
|
||||||
"type": "agent_status",
|
"type": "agent_status",
|
||||||
"jobs_seen": len(jobs),
|
"jobs_seen": len(jobs),
|
||||||
"requests_seen": len(requests_),
|
"requests_seen": len(requests_),
|
||||||
"commands_seen": len(commands),
|
"commands_seen": len(commands),
|
||||||
"capabilities": ["storage", "tally", "accounting_act"],
|
"capabilities": ["storage", "tally", "accounting_act", "local_dashboard", "manual_updates"],
|
||||||
"agent_time_utc": self._now(),
|
"agent_time_utc": self._now(),
|
||||||
}))
|
}))
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ import os
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
import shutil
|
import shutil
|
||||||
import subprocess
|
import subprocess
|
||||||
|
import threading
|
||||||
import time
|
import time
|
||||||
import zipfile
|
import zipfile
|
||||||
|
|
||||||
@@ -12,39 +13,78 @@ from . import __version__
|
|||||||
|
|
||||||
|
|
||||||
class AgentUpdater:
|
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.config = config
|
||||||
self.client = client
|
self.client = client
|
||||||
self.logger = logger
|
self.logger = logger
|
||||||
self.install_dir = install_dir.resolve()
|
self.install_dir = install_dir.resolve()
|
||||||
|
self.db = db
|
||||||
self.last_check = 0.0
|
self.last_check = 0.0
|
||||||
self.busy = False
|
self.busy = False
|
||||||
|
self.lock = threading.RLock()
|
||||||
|
self.latest_manifest: dict = {}
|
||||||
|
|
||||||
def maybe_update(self, force: bool = False) -> bool:
|
def _record(self, event_type: str, to_version: str = "", status: str = "", **kwargs) -> None:
|
||||||
if not self.config.auto_update or self.busy:
|
if self.db is not None:
|
||||||
return False
|
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()
|
now = time.time()
|
||||||
if not force and now - self.last_check < self.config.update_check_interval_seconds:
|
if not force and now - self.last_check < self.config.update_check_interval_seconds and self.latest_manifest:
|
||||||
return False
|
return self.status()
|
||||||
self.last_check = now
|
self.last_check = now
|
||||||
try:
|
try:
|
||||||
manifest = self.client.update_manifest()
|
manifest = self.client.update_manifest()
|
||||||
latest = str(manifest.get("latest_version") or "").strip()
|
latest = str(manifest.get("latest_version") or "").strip()
|
||||||
if not latest or latest == __version__:
|
|
||||||
return False
|
|
||||||
expected = str(manifest.get("sha256") or "").lower().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
|
||||||
|
try:
|
||||||
|
state = self.check_for_update(force=force)
|
||||||
|
return bool(state.get("update_available"))
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
|
|
||||||
|
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()
|
package = self.client.download_update_package()
|
||||||
actual = hashlib.sha256(package).hexdigest().lower()
|
actual = hashlib.sha256(package).hexdigest().lower()
|
||||||
if len(expected) != 64 or actual != expected:
|
if len(expected) != 64 or actual != expected:
|
||||||
raise RuntimeError("ERP Local Agent update package SHA256 verification failed.")
|
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)
|
|
||||||
return False
|
|
||||||
|
|
||||||
def _stage_and_restart(self, latest: str, package: bytes) -> None:
|
|
||||||
self.busy = True
|
|
||||||
updates = self.install_dir / "updates"
|
updates = self.install_dir / "updates"
|
||||||
updates.mkdir(parents=True, exist_ok=True)
|
updates.mkdir(parents=True, exist_ok=True)
|
||||||
safe = "".join(ch for ch in latest if ch.isalnum() or ch in ".-_") or "update"
|
safe = "".join(ch for ch in latest if ch.isalnum() or ch in ".-_") or "update"
|
||||||
@@ -57,38 +97,87 @@ class AgentUpdater:
|
|||||||
archive.extractall(staged)
|
archive.extractall(staged)
|
||||||
if not (staged / "erp_local_agent" / "__init__.py").exists():
|
if not (staged / "erp_local_agent" / "__init__.py").exists():
|
||||||
raise RuntimeError("ERP Local Agent update package is incomplete.")
|
raise RuntimeError("ERP Local Agent update package is incomplete.")
|
||||||
script_path = updates / f"apply_{safe}.ps1"
|
self._record("download", latest, "downloaded", sha256=actual, package_path=str(package_path))
|
||||||
script_path.write_text(self._powershell_update_script(staged), encoding="utf-8")
|
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 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)
|
flags = getattr(subprocess, "CREATE_NEW_PROCESS_GROUP", 0) | getattr(subprocess, "DETACHED_PROCESS", 0)
|
||||||
subprocess.Popen(
|
subprocess.Popen(
|
||||||
["powershell.exe", "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(script_path), "-ParentPid", str(os.getpid())],
|
["powershell.exe", "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(script_path), "-ParentPid", str(os.getpid())],
|
||||||
cwd=str(self.install_dir), creationflags=flags, close_fds=True,
|
cwd=str(self.install_dir), creationflags=flags, close_fds=True,
|
||||||
)
|
)
|
||||||
self.logger.warning("ERP Local Agent update %s staged; restarting.", latest)
|
self.logger.warning("ERP Local Agent update %s installation requested by local user; restarting.", latest)
|
||||||
|
time.sleep(0.3)
|
||||||
os._exit(0)
|
os._exit(0)
|
||||||
|
|
||||||
def _powershell_update_script(self, staged: Path) -> str:
|
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("'", "''")
|
install = str(self.install_dir).replace("'", "''")
|
||||||
stage = str(staged.resolve()).replace("'", "''")
|
stage = str(staged.resolve()).replace("'", "''")
|
||||||
|
latest_safe = str(latest).replace("'", "''")
|
||||||
lines = [
|
lines = [
|
||||||
"param([int]$ParentPid)",
|
"param([int]$ParentPid)",
|
||||||
'$ErrorActionPreference = "Stop"',
|
'$ErrorActionPreference = "Stop"',
|
||||||
"$InstallDir = '" + install + "'",
|
"$InstallDir = '" + install + "'",
|
||||||
"$StagedDir = '" + stage + "'",
|
"$StagedDir = '" + stage + "'",
|
||||||
|
"$TargetVersion = '" + latest_safe + "'",
|
||||||
"$TaskName = 'ERP Local Agent'",
|
"$TaskName = 'ERP Local Agent'",
|
||||||
"$BackupDir = Join-Path $InstallDir ('updates\\backup_' + (Get-Date -Format 'yyyyMMdd_HHmmss'))",
|
"$BackupDir = Join-Path $InstallDir ('updates\\backup_' + (Get-Date -Format 'yyyyMMdd_HHmmss'))",
|
||||||
"try {",
|
"try {",
|
||||||
" if ($ParentPid -gt 0) { try { Wait-Process -Id $ParentPid -Timeout 120 -ErrorAction SilentlyContinue } catch {} }",
|
" if ($ParentPid -gt 0) { try { Wait-Process -Id $ParentPid -Timeout 120 -ErrorAction SilentlyContinue } catch {} }",
|
||||||
" try { Stop-ScheduledTask -TaskName $TaskName -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",
|
" 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 }",
|
" 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",
|
" 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'",
|
" $Python = Join-Path $InstallDir '.venv\\Scripts\\python.exe'",
|
||||||
" if (-not (Test-Path $Python)) { throw 'ERP Local Agent virtual environment is missing.' }",
|
" if (-not (Test-Path $Python)) { throw 'ERP Local Agent virtual environment is missing.' }",
|
||||||
" & $Python -m pip install -r (Join-Path $InstallDir 'requirements.txt')",
|
" & $Python -m pip install -r (Join-Path $InstallDir 'requirements.txt')",
|
||||||
" if ($LASTEXITCODE -ne 0) { throw 'Dependency update failed.' }",
|
" 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",
|
" Start-ScheduledTask -TaskName $TaskName",
|
||||||
"} catch {",
|
"} 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 {}",
|
" 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 {}",
|
||||||
|
|||||||
@@ -0,0 +1,2 @@
|
|||||||
|
@echo off
|
||||||
|
start "" "http://127.0.0.1:8788"
|
||||||
Reference in New Issue
Block a user