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

553 lines
19 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 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"
OWNER_FILE = DATA_DIR / "supervisor_owner.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
ERROR_ALREADY_EXISTS = 183
SINGLETON_WATCHDOG_SECONDS = 10
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 _kill_duplicate_supervisors() -> None:
if os.name != "nt":
return
root = str(INSTALL_ROOT).replace("'", "''")
script = (
"$root='" + root + "';$self=" + str(os.getpid()) + ";"
"Get-CimInstance Win32_Process -ErrorAction SilentlyContinue | "
"Where-Object { $_.ProcessId -ne $self -and $_.CommandLine -and "
"$_.CommandLine -match 'ERPAgentSupervisor\\.pyw' -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,
)
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,
)
def _enforce_singleton_processes(*, keep_worker_pid: int | None = None) -> None:
"""Kill stale/duplicate supervisors and workers for this install only.
The Global named mutex is the primary atomic guard. This watchdog is a
recovery layer for legacy processes started before the hardened runtime.
"""
if os.name != "nt":
return
root = str(INSTALL_ROOT.resolve()).replace("'", "''")
keep_worker = int(keep_worker_pid or 0)
script = rf"""
$root='{root}'
$selfPid={os.getpid()}
$keepWorker={keep_worker}
Get-CimInstance Win32_Process -ErrorAction SilentlyContinue | ForEach-Object {{
$cmd=[string]$_.CommandLine
if (-not $cmd -or $cmd -notlike ('*'+$root+'*')) {{ return }}
if (($cmd -match 'ERPAgentSupervisor\.pyw') -and ($_.ProcessId -ne $selfPid)) {{
Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue
return
}}
if (($cmd -match 'erp_local_agent\.main') -and ($_.ProcessId -ne $keepWorker)) {{
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
self._os_mutex = None
def acquire_lock(self) -> bool:
# The named mutex is atomic at the Windows kernel level and prevents
# the lock-file race that allowed multiple supervisors before 1.22.3.
self._os_mutex = _acquire_os_mutex()
if self._os_mutex is None:
return False
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")
_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()),
})
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
_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:
_kill_duplicate_supervisors()
_state("starting", "ERP Local Agent supervisor starting.")
self.start_worker()
last_singleton_watchdog = 0.0
while self.running:
now = time.time()
if now - last_singleton_watchdog >= SINGLETON_WATCHDOG_SECONDS:
_enforce_singleton_processes(
keep_worker_pid=(self.worker.pid if self.worker is not None and self.worker.poll() is None else None)
)
last_singleton_watchdog = now
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.add_argument("--replace", action="store_true", help="Explicit controlled replacement of an older supervisor instance")
args = parser.parse_args()
if args.replace and os.name == "nt":
# Explicit restart/update only. Normal duplicate launches never disturb
# the healthy owner; they simply fail the Global mutex and exit.
_kill_duplicate_supervisors()
deadline = time.time() + 15
while time.time() < deadline:
time.sleep(0.25)
break
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())