Files
arrr-erp/app/modules/documents/local_agent_runtime/ERPAgentSupervisor.pyw
T
2026-09-30 22:28:39 +05:30

564 lines
21 KiB
Python

from __future__ import annotations
import argparse
import ctypes
from ctypes import wintypes
import hashlib
import json
import os
from pathlib import Path
import shutil
import signal
import subprocess
import time
import urllib.request
INSTALL_ROOT = Path(__file__).resolve().parent
DATA_DIR = INSTALL_ROOT / "data"
LOG_DIR = INSTALL_ROOT / "logs"
UPDATES_DIR = INSTALL_ROOT / "updates"
LOCK_FILE = DATA_DIR / "supervisor.lock"
STATE_FILE = DATA_DIR / "supervisor_state.json"
OWNER_FILE = DATA_DIR / "supervisor_owner.json"
REQUEST_FILE = UPDATES_DIR / "supervisor_request.json"
UPDATE_PROGRESS_FILE = DATA_DIR / "update_progress.json"
FAILED_UPDATES_FILE = DATA_DIR / "failed_updates.json"
WORKER_LOG = LOG_DIR / "worker-supervisor.log"
SUPERVISOR_LOG = LOG_DIR / "supervisor.log"
DASHBOARD_URL = "http://127.0.0.1:8788"
WORKER_MODULE = "erp_local_agent.main"
RESTART_DELAY_SECONDS = 3
HEALTH_TIMEOUT_SECONDS = 90
ERROR_ALREADY_EXISTS = 183
def _supervisor_mutex_name() -> str:
normalized = str(INSTALL_ROOT.resolve()).casefold().encode("utf-8", errors="ignore")
suffix = hashlib.sha256(normalized).hexdigest()[:20]
return f"Global\\ARRR_ERP_Local_Agent_SUPERVISOR_{suffix}"
def _acquire_os_mutex():
if os.name != "nt":
return object()
kernel32 = ctypes.WinDLL("kernel32", use_last_error=True)
kernel32.CreateMutexW.argtypes = (ctypes.c_void_p, wintypes.BOOL, wintypes.LPCWSTR)
kernel32.CreateMutexW.restype = wintypes.HANDLE
ctypes.set_last_error(0)
handle = kernel32.CreateMutexW(None, False, _supervisor_mutex_name())
if not handle:
raise ctypes.WinError(ctypes.get_last_error())
if ctypes.get_last_error() == ERROR_ALREADY_EXISTS:
kernel32.CloseHandle(handle)
return None
return handle
def _close_os_mutex(handle) -> None:
if os.name != "nt" or handle is None or not isinstance(handle, int):
return
try:
ctypes.WinDLL("kernel32", use_last_error=True).CloseHandle(handle)
except Exception:
pass
def _create_no_window() -> int:
return getattr(subprocess, "CREATE_NO_WINDOW", 0)
def _log(message: str) -> None:
LOG_DIR.mkdir(parents=True, exist_ok=True)
line = f"{time.strftime('%Y-%m-%d %H:%M:%S')} {message}\n"
try:
with SUPERVISOR_LOG.open("a", encoding="utf-8") as handle:
handle.write(line)
except Exception:
pass
def _write_json_atomic(path: Path, payload: dict) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
tmp = path.with_suffix(path.suffix + ".tmp")
tmp.write_text(json.dumps(payload, indent=2, sort_keys=True), encoding="utf-8")
tmp.replace(path)
def _state(status: str, message: str, **extra) -> None:
payload = {
"status": status,
"message": message,
"supervisor_pid": os.getpid(),
"updated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
}
payload.update(extra)
_write_json_atomic(STATE_FILE, payload)
_log(f"{status}: {message}")
def _update_progress(phase: str, percent: int, message: str, *, version: str = "", status: str = "installing") -> None:
payload = {
"operation": "install",
"status": status,
"phase": phase,
"percent": max(0, min(100, int(percent))),
"message": message,
"target_version": version,
"updated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
"supervisor_pid": os.getpid(),
}
_write_json_atomic(UPDATE_PROGRESS_FILE, payload)
def _read_failed_updates() -> dict:
try:
if FAILED_UPDATES_FILE.exists():
value = json.loads(FAILED_UPDATES_FILE.read_text(encoding="utf-8"))
return value if isinstance(value, dict) else {}
except Exception:
pass
return {}
def _mark_failed_update(version: str, error: str) -> None:
if not version:
return
payload = _read_failed_updates()
previous = payload.get(version) if isinstance(payload.get(version), dict) else {}
payload[version] = {
"failed_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
"error": str(error)[:2000],
"attempts": int(previous.get("attempts") or 0) + 1,
}
_write_json_atomic(FAILED_UPDATES_FILE, payload)
def _clear_failed_update(version: str) -> None:
payload = _read_failed_updates()
if version in payload:
payload.pop(version, None)
_write_json_atomic(FAILED_UPDATES_FILE, payload)
def _dashboard_ready(timeout: float = 1.5) -> bool:
try:
with urllib.request.urlopen(DASHBOARD_URL + "/api/status", timeout=timeout) as response:
return int(getattr(response, "status", 200)) < 500
except Exception:
return False
def _python(console: bool = True) -> Path:
name = "python.exe" if console else "pythonw.exe"
candidate = INSTALL_ROOT / ".venv" / "Scripts" / name
if candidate.exists():
return candidate
fallback = INSTALL_ROOT / ".venv" / "Scripts" / "python.exe"
if fallback.exists():
return fallback
raise RuntimeError(f"ERP Local Agent Python environment not found under {INSTALL_ROOT / '.venv'}")
def _terminate_process_tree(process: subprocess.Popen | None) -> None:
if process is None or process.poll() is not None:
return
try:
subprocess.run(
["taskkill", "/PID", str(process.pid), "/T", "/F"],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
creationflags=_create_no_window(),
check=False,
)
except Exception:
try:
process.terminate()
except Exception:
pass
try:
process.wait(timeout=20)
except Exception:
pass
class Supervisor:
def __init__(self):
self.worker: subprocess.Popen | None = None
self.running = True
self._os_mutex = None
self._worker_log_handle = None
def acquire_lock(self) -> bool:
"""Acquire the only authoritative supervisor singleton guard.
The Global Windows mutex is atomic across SYSTEM and user sessions. PID/owner
files are diagnostics only and must never veto a successfully acquired mutex.
"""
_log("startup: acquiring Global supervisor mutex")
self._os_mutex = _acquire_os_mutex()
if self._os_mutex is None:
_log("startup: another supervisor already owns the Global mutex; exiting")
return False
DATA_DIR.mkdir(parents=True, exist_ok=True)
LOCK_FILE.write_text(str(os.getpid()), encoding="utf-8")
_write_json_atomic(
OWNER_FILE,
{
"supervisor_pid": os.getpid(),
"install_root": str(INSTALL_ROOT),
"started_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
"singleton": "global_mutex",
},
)
_log("startup: Global supervisor mutex acquired")
return True
def release_lock(self) -> None:
try:
if LOCK_FILE.exists() and LOCK_FILE.read_text(encoding="utf-8").strip() == str(os.getpid()):
LOCK_FILE.unlink()
except Exception:
pass
try:
if OWNER_FILE.exists():
owner = json.loads(OWNER_FILE.read_text(encoding="utf-8"))
if int(owner.get("supervisor_pid") or 0) == os.getpid():
OWNER_FILE.unlink()
except Exception:
pass
_close_os_mutex(self._os_mutex)
self._os_mutex = None
def start_worker(self) -> None:
if self.worker is not None and self.worker.poll() is None:
return
LOG_DIR.mkdir(parents=True, exist_ok=True)
py = _python(console=True)
self._worker_log_handle = WORKER_LOG.open("a", encoding="utf-8")
self._worker_log_handle.write(f"\n--- worker start {time.strftime('%Y-%m-%d %H:%M:%S')} ---\n")
self._worker_log_handle.flush()
self.worker = subprocess.Popen(
[str(py), "-m", WORKER_MODULE],
cwd=str(INSTALL_ROOT),
stdout=self._worker_log_handle,
stderr=self._worker_log_handle,
creationflags=_create_no_window(),
)
_state("running", "ERP Local Agent worker started.", worker_pid=self.worker.pid)
def stop_worker(self) -> None:
if self.worker is None:
return
_state("stopping", "Stopping ERP Local Agent worker.")
_terminate_process_tree(self.worker)
self.worker = None
if self._worker_log_handle is not None:
try:
self._worker_log_handle.close()
except Exception:
pass
self._worker_log_handle = None
@staticmethod
def _read_request() -> dict | None:
if not REQUEST_FILE.exists():
return None
try:
payload = json.loads(REQUEST_FILE.read_text(encoding="utf-8"))
return payload if isinstance(payload, dict) else None
except Exception as exc:
_state("error", f"Invalid supervisor update request: {exc}")
try:
REQUEST_FILE.unlink()
except Exception:
pass
return None
@staticmethod
def _clear_request() -> None:
try:
REQUEST_FILE.unlink()
except FileNotFoundError:
pass
@staticmethod
def _validate_staged(staged: Path, expected_version: str = "") -> None:
if not staged.exists() or not staged.is_dir():
raise RuntimeError(f"Staged update directory not found: {staged}")
package_dir = staged / "erp_local_agent"
init_path = package_dir / "__init__.py"
dashboard_path = package_dir / "dashboard.py"
main_path = package_dir / "main.py"
if not init_path.is_file():
raise RuntimeError("Staged update does not contain erp_local_agent/__init__.py.")
if not dashboard_path.is_file() or not main_path.is_file():
raise RuntimeError("Staged update is missing dashboard.py or main.py.")
if not (staged / "requirements.txt").is_file():
raise RuntimeError("Staged update does not contain requirements.txt.")
# Validate the candidate in-place BEFORE the working worker is stopped.
# This specifically protects the dashboard identity contract that caused
# older staged builds to fail on: from . import AGENT_NAME, __version__.
py = _python(console=True)
code = (
"import sys;"
f"sys.path.insert(0, {str(staged)!r});"
"import erp_local_agent;"
"assert getattr(erp_local_agent, 'AGENT_NAME', '') == 'ERP Local Agent', 'AGENT_NAME missing/invalid';"
"assert getattr(erp_local_agent, '__version__', ''), '__version__ missing';"
"from erp_local_agent import dashboard,main,commands,updater;"
"print(erp_local_agent.__version__)"
)
result = subprocess.run(
[str(py), "-c", code],
cwd=str(staged),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
creationflags=_create_no_window(),
timeout=90,
)
if result.returncode != 0:
raise RuntimeError("Staged update preflight import validation failed:\n" + result.stdout[-4000:])
staged_version = (result.stdout or "").strip().splitlines()[-1].strip()
if expected_version and staged_version != expected_version:
raise RuntimeError(
f"Staged update version mismatch: requested={expected_version}, package={staged_version or 'missing'}."
)
@staticmethod
def _backup_runtime(version: str) -> Path:
backup = UPDATES_DIR / f"supervisor_backup_{time.strftime('%Y%m%d_%H%M%S')}_{version}"
backup.mkdir(parents=True, exist_ok=True)
agent = INSTALL_ROOT / "erp_local_agent"
if agent.exists():
shutil.copytree(agent, backup / "erp_local_agent")
for name in (
"requirements.txt",
"README_ERP_LOCAL_AGENT.txt",
"run_agent.bat",
"run_once.bat",
"check_config.py",
"check_config.bat",
"open_dashboard.bat",
"Open ERP Local Agent Dashboard.vbs",
"install_task_scheduler.bat",
"start_task_scheduler.bat",
"stop_task_scheduler.bat",
"status_task_scheduler.bat",
"uninstall_task_scheduler.bat",
"ERPAgentSupervisor.pyw",
"ERPAgentDashboard.pyw",
):
src = INSTALL_ROOT / name
if src.exists():
shutil.copy2(src, backup / name)
return backup
@staticmethod
def _replace_worker_runtime(staged: Path) -> None:
target = INSTALL_ROOT / "erp_local_agent"
new_target = INSTALL_ROOT / "erp_local_agent.__new__"
old_target = INSTALL_ROOT / "erp_local_agent.__old__"
shutil.rmtree(new_target, ignore_errors=True)
shutil.rmtree(old_target, ignore_errors=True)
shutil.copytree(staged / "erp_local_agent", new_target)
if target.exists():
target.rename(old_target)
try:
new_target.rename(target)
except Exception:
if old_target.exists() and not target.exists():
old_target.rename(target)
raise
shutil.rmtree(old_target, ignore_errors=True)
for name in (
"requirements.txt",
"README_ERP_LOCAL_AGENT.txt",
"run_agent.bat",
"run_once.bat",
"check_config.py",
"check_config.bat",
"open_dashboard.bat",
"Open ERP Local Agent Dashboard.vbs",
"install_task_scheduler.bat",
"start_task_scheduler.bat",
"stop_task_scheduler.bat",
"status_task_scheduler.bat",
"uninstall_task_scheduler.bat",
"ERPAgentSupervisor.pyw",
"ERPAgentDashboard.pyw",
):
src = staged / name
if src.exists():
shutil.copy2(src, INSTALL_ROOT / name)
@staticmethod
def _install_requirements_and_verify() -> None:
py = _python(console=True)
req = INSTALL_ROOT / "requirements.txt"
result = subprocess.run(
[str(py), "-m", "pip", "install", "-r", str(req)],
cwd=str(INSTALL_ROOT),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
creationflags=_create_no_window(),
timeout=300,
)
if result.returncode != 0:
raise RuntimeError("Dependency installation failed:\n" + result.stdout[-4000:])
code = (
"import erp_local_agent;"
"assert getattr(erp_local_agent, 'AGENT_NAME', '') == 'ERP Local Agent', 'AGENT_NAME missing/invalid';"
"from erp_local_agent import accounting_store,commands,tally,updater,dashboard,main;"
"print(erp_local_agent.__version__)"
)
result = subprocess.run(
[str(py), "-c", code],
cwd=str(INSTALL_ROOT),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
creationflags=_create_no_window(),
timeout=60,
)
if result.returncode != 0:
raise RuntimeError("Updated worker import verification failed:\n" + result.stdout[-4000:])
@staticmethod
def _restore_backup(backup: Path) -> None:
target = INSTALL_ROOT / "erp_local_agent"
backup_agent = backup / "erp_local_agent"
if backup_agent.exists():
shutil.rmtree(target, ignore_errors=True)
shutil.copytree(backup_agent, target)
for item in backup.iterdir():
if item.name == "erp_local_agent":
continue
if item.is_file():
shutil.copy2(item, INSTALL_ROOT / item.name)
def apply_update(self, request: dict) -> None:
version = str(request.get("version") or "").strip()
staged_text = str(request.get("staged_dir") or "").strip()
if not version or not staged_text:
raise RuntimeError("Supervisor update request is missing version/staged_dir.")
staged = Path(staged_text).resolve()
self._validate_staged(staged, expected_version=version)
_update_progress("preparing", 8, f"Preparing ERP Local Agent {version} update.", version=version)
_state("installing", f"Installing ERP Local Agent {version}.", target_version=version, progress_pct=8)
backup = None
try:
_update_progress("stopping", 18, "Stopping the current Local Agent worker safely.", version=version)
self.stop_worker()
_update_progress("backup", 30, "Creating rollback backup.", version=version)
backup = self._backup_runtime(version)
_update_progress("replacing", 48, "Replacing Local Agent runtime files.", version=version)
self._replace_worker_runtime(staged)
_update_progress("dependencies", 66, "Verifying runtime dependencies and imports.", version=version)
self._install_requirements_and_verify()
_update_progress("starting", 82, "Starting the updated Local Agent.", version=version)
self.start_worker()
deadline = time.time() + HEALTH_TIMEOUT_SECONDS
while time.time() < deadline:
if self.worker is not None and self.worker.poll() is not None:
raise RuntimeError(f"Updated worker exited with code {self.worker.returncode}.")
if _dashboard_ready():
_update_progress("complete", 100, f"ERP Local Agent {version} updated successfully.", version=version, status="updated")
_clear_failed_update(version)
_state(
"updated",
f"ERP Local Agent {version} installed and restarted successfully.",
target_version=version,
worker_pid=(self.worker.pid if self.worker else None),
backup_path=str(backup),
progress_pct=100,
)
return
_update_progress("health_check", 90, "Waiting for the updated dashboard to become ready.", version=version)
time.sleep(1)
raise RuntimeError("Updated worker did not become healthy within the restart timeout.")
except Exception as exc:
_mark_failed_update(version, str(exc))
_update_progress("rollback", 92, f"Update failed; restoring the previous runtime: {exc}", version=version, status="rollback")
_state("rollback", f"Update {version} failed. Restoring previous worker.", target_version=version)
self.stop_worker()
if backup is not None:
self._restore_backup(backup)
self.start_worker()
_update_progress("error", 100, f"Update failed and the previous version was restored: {exc}", version=version, status="error")
raise
def run(self) -> int:
if not self.acquire_lock():
return 0
try:
_state("starting", "ERP Local Agent supervisor starting.")
self.start_worker()
while self.running:
request = self._read_request()
if request and str(request.get("action") or "") == "install_update":
self._clear_request()
try:
self.apply_update(request)
except Exception as exc:
_state("error", f"Update failed and worker was restored/restarted: {exc}")
if self.worker is None or self.worker.poll() is not None:
code = None if self.worker is None else self.worker.returncode
_state("restarting", f"ERP Local Agent worker stopped (exit={code}); restarting.")
time.sleep(RESTART_DELAY_SECONDS)
self.start_worker()
time.sleep(1)
return 0
except BaseException as exc:
_log(f"fatal: {type(exc).__name__}: {exc}")
try:
_state("fatal", f"Supervisor stopped unexpectedly: {exc}")
except Exception:
pass
raise
finally:
try:
self.stop_worker()
finally:
self.release_lock()
_state("stopped", "ERP Local Agent supervisor stopped.")
def main() -> int:
parser = argparse.ArgumentParser(description="ERP Local Agent desktop/background supervisor")
parser.add_argument("--background", action="store_true")
parser.add_argument("--replace", action="store_true", help="Compatibility flag; the Global mutex remains authoritative")
parser.parse_args()
_log(f"startup: supervisor process entered main pid={os.getpid()} root={INSTALL_ROOT}")
supervisor = Supervisor()
def stop_handler(*_args):
supervisor.running = False
try:
signal.signal(signal.SIGTERM, stop_handler)
signal.signal(signal.SIGINT, stop_handler)
except Exception:
pass
return supervisor.run()
if __name__ == "__main__":
try:
raise SystemExit(main())
except SystemExit:
raise
except BaseException as exc:
_log(f"startup-fatal: {type(exc).__name__}: {exc}")
raise