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" 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 _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) -> 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", "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", "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;" "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) _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") _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: _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