Add persistent desktop supervisor for ERP Local Agent

This commit is contained in:
A R R R Associates
2026-08-20 20:01:59 +05:30
parent f7ebd32597
commit de8b32b06d
14 changed files with 707 additions and 72 deletions
@@ -0,0 +1,418 @@
from __future__ import annotations
import argparse
import json
import os
from pathlib import Path
import shutil
import signal
import subprocess
import sys
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"
REQUEST_FILE = UPDATES_DIR / "supervisor_request.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 = 60
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 _pid_alive(pid: int) -> bool:
if pid <= 0:
return False
try:
os.kill(pid, 0)
return True
except OSError:
return False
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 _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:
return
if 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=15)
except Exception:
pass
def _kill_unmanaged_workers() -> None:
script = (
"$root='" + str(INSTALL_ROOT).replace("'", "''") + "';"
"Get-CimInstance Win32_Process -ErrorAction SilentlyContinue | "
"Where-Object { $_.CommandLine -and $_.CommandLine -match 'erp_local_agent\\.main' "
"-and $_.CommandLine -like ('*'+$root+'*') } | "
"ForEach-Object { Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue }"
)
subprocess.run(
["powershell.exe", "-NoProfile", "-ExecutionPolicy", "Bypass", "-Command", script],
cwd=str(INSTALL_ROOT),
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
creationflags=_create_no_window(),
check=False,
)
class Supervisor:
def __init__(self):
self.worker: subprocess.Popen | None = None
self.running = True
def acquire_lock(self) -> bool:
DATA_DIR.mkdir(parents=True, exist_ok=True)
if LOCK_FILE.exists():
try:
existing = int(LOCK_FILE.read_text(encoding="utf-8").strip() or "0")
except Exception:
existing = 0
if existing and existing != os.getpid() and _pid_alive(existing):
return False
LOCK_FILE.write_text(str(os.getpid()), encoding="utf-8")
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
def start_worker(self) -> None:
if self.worker is not None and self.worker.poll() is None:
return
_kill_unmanaged_workers()
LOG_DIR.mkdir(parents=True, exist_ok=True)
py = _python(console=True)
handle = WORKER_LOG.open("a", encoding="utf-8")
handle.write(f"\n--- worker start {time.strftime('%Y-%m-%d %H:%M:%S')} ---\n")
handle.flush()
self.worker = subprocess.Popen(
[str(py), "-m", WORKER_MODULE],
cwd=str(INSTALL_ROOT),
stdout=handle,
stderr=handle,
creationflags=_create_no_window(),
)
_state("running", "ERP Local Agent worker started.", worker_pid=self.worker.pid)
def stop_worker(self) -> None:
_state("stopping", "Stopping ERP Local Agent worker.")
_terminate_process_tree(self.worker)
self.worker = None
_kill_unmanaged_workers()
@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) -> None:
if not staged.exists() or not staged.is_dir():
raise RuntimeError(f"Staged update directory not found: {staged}")
if not (staged / "erp_local_agent" / "__init__.py").is_file():
raise RuntimeError("Staged update does not contain erp_local_agent/__init__.py.")
if not (staged / "requirements.txt").is_file():
raise RuntimeError("Staged update does not contain requirements.txt.")
@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",
"install_task_scheduler.bat",
"start_task_scheduler.bat",
"stop_task_scheduler.bat",
"status_task_scheduler.bat",
"uninstall_task_scheduler.bat",
):
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)
# Root helper files are safe to refresh while the stable supervisor is running.
# Supervisor/dashboard launchers themselves are intentionally excluded here;
# they are updated only by the explicit direct installer.
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",
"install_task_scheduler.bat",
"start_task_scheduler.bat",
"stop_task_scheduler.bat",
"status_task_scheduler.bat",
"uninstall_task_scheduler.bat",
):
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;"
"from erp_local_agent import accounting_store,commands,tally,updater;"
"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)
_state("installing", f"Installing ERP Local Agent {version}.", target_version=version)
self.stop_worker()
backup = self._backup_runtime(version)
try:
self._replace_worker_runtime(staged)
self._install_requirements_and_verify()
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():
_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),
)
return
time.sleep(1)
raise RuntimeError("Updated worker did not become healthy within the restart timeout.")
except Exception:
_state("rollback", f"Update {version} failed. Restoring previous worker.", target_version=version)
self.stop_worker()
self._restore_backup(backup)
self.start_worker()
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
finally:
self.stop_worker()
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.parse_args()
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__":
raise SystemExit(main())