2026-08-30 22:22:11 +08:00
|
|
|
"""Explicit lifecycle for process-owning background services.
|
|
|
|
|
|
|
|
|
|
Only the process holding ``ServiceLeaderLock`` may bind SIP, manage ZLM,
|
|
|
|
|
spawn recording processes, auto-proxy streams, or emit telemetry.
|
|
|
|
|
"""
|
|
|
|
|
import atexit
|
|
|
|
|
import json
|
|
|
|
|
import logging
|
|
|
|
|
import os
|
|
|
|
|
import threading
|
|
|
|
|
import time
|
|
|
|
|
from pathlib import Path
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger("app.services.lifecycle")
|
2026-09-04 18:16:14 +08:00
|
|
|
from monitor_runtime.paths import DATA_ROOT
|
|
|
|
|
PROJECT_ROOT = DATA_ROOT
|
2026-08-30 22:22:11 +08:00
|
|
|
DEFAULT_LOCK_PATH = PROJECT_ROOT / ".runtime" / "service-leader.lock"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class ServiceLeaderLock:
|
|
|
|
|
def __init__(self, path=None):
|
|
|
|
|
self.path = Path(path or os.environ.get("MONITOR_SERVICE_LOCK", DEFAULT_LOCK_PATH)).resolve()
|
|
|
|
|
self._handle = None
|
|
|
|
|
|
|
|
|
|
@property
|
|
|
|
|
def acquired(self):
|
|
|
|
|
return self._handle is not None
|
|
|
|
|
|
|
|
|
|
def acquire(self):
|
|
|
|
|
if self.acquired:
|
|
|
|
|
return True
|
|
|
|
|
self.path.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
|
handle = open(self.path, "a+b")
|
|
|
|
|
try:
|
|
|
|
|
if os.name == "nt":
|
|
|
|
|
import msvcrt
|
|
|
|
|
handle.seek(0, os.SEEK_END)
|
|
|
|
|
if handle.tell() == 0:
|
|
|
|
|
handle.write(b"\0")
|
|
|
|
|
handle.flush()
|
|
|
|
|
handle.seek(0)
|
|
|
|
|
msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1)
|
|
|
|
|
else:
|
|
|
|
|
import fcntl
|
|
|
|
|
fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
|
|
|
except (OSError, IOError):
|
|
|
|
|
handle.close()
|
|
|
|
|
return False
|
|
|
|
|
self._handle = handle
|
|
|
|
|
metadata = json.dumps({"pid": os.getpid(), "started_at": int(time.time())}).encode("utf-8")
|
|
|
|
|
handle.seek(0)
|
|
|
|
|
handle.truncate()
|
|
|
|
|
handle.write(metadata)
|
|
|
|
|
handle.flush()
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
def release(self):
|
|
|
|
|
handle = self._handle
|
|
|
|
|
if handle is None:
|
|
|
|
|
return
|
|
|
|
|
try:
|
|
|
|
|
handle.seek(0)
|
|
|
|
|
if os.name == "nt":
|
|
|
|
|
import msvcrt
|
|
|
|
|
msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1)
|
|
|
|
|
else:
|
|
|
|
|
import fcntl
|
|
|
|
|
fcntl.flock(handle.fileno(), fcntl.LOCK_UN)
|
|
|
|
|
finally:
|
|
|
|
|
handle.close()
|
|
|
|
|
self._handle = None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class BackgroundCoordinator:
|
|
|
|
|
def __init__(self):
|
|
|
|
|
self._stop = threading.Event()
|
|
|
|
|
self._thread = None
|
|
|
|
|
|
|
|
|
|
def start(self):
|
|
|
|
|
if self._thread and self._thread.is_alive():
|
|
|
|
|
return
|
|
|
|
|
self._stop.clear()
|
|
|
|
|
self._thread = threading.Thread(target=self._run, name="service-coordinator", daemon=True)
|
|
|
|
|
self._thread.start()
|
|
|
|
|
|
|
|
|
|
def stop(self):
|
|
|
|
|
self._stop.set()
|
|
|
|
|
if self._thread and self._thread.is_alive():
|
|
|
|
|
self._thread.join(timeout=5)
|
|
|
|
|
|
|
|
|
|
def _run(self):
|
|
|
|
|
from app.utils.GlobalUtils import GlobalUtils, g_config, g_logger
|
|
|
|
|
|
|
|
|
|
if getattr(g_config, "autoAddStreamProxy", False):
|
|
|
|
|
delay = max(0, int(getattr(g_config, "autoAddStreamProxySleep", 15)))
|
|
|
|
|
if not self._stop.wait(delay):
|
|
|
|
|
try:
|
|
|
|
|
ok, msg = GlobalUtils.addAllStreamProxy()
|
|
|
|
|
g_logger.info("autoAddStreamProxy ok=%s msg=%s", ok, msg)
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
g_logger.warning("autoAddStreamProxy failed: %s", exc)
|
|
|
|
|
|
|
|
|
|
sequence = 0
|
|
|
|
|
interval = max(60, int(os.environ.get("MONITOR_TELEMETRY_INTERVAL", "4800")))
|
|
|
|
|
while not self._stop.wait(interval):
|
|
|
|
|
if not getattr(g_config, "telemetryEnabled", False):
|
|
|
|
|
continue
|
|
|
|
|
sequence += 1
|
|
|
|
|
try:
|
|
|
|
|
from app.services.telemetry import send_heartbeat
|
|
|
|
|
send_heartbeat(sequence)
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
g_logger.warning("telemetry heartbeat failed: %s", exc)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class ServiceManager:
|
|
|
|
|
def __init__(self, lock=None):
|
|
|
|
|
self.lock = lock or ServiceLeaderLock()
|
|
|
|
|
self.coordinator = BackgroundCoordinator()
|
|
|
|
|
self._started = False
|
|
|
|
|
self._guard = threading.RLock()
|
|
|
|
|
|
|
|
|
|
@property
|
|
|
|
|
def is_leader(self):
|
|
|
|
|
return self.lock.acquired
|
|
|
|
|
|
|
|
|
|
def start(self):
|
2026-09-04 18:16:14 +08:00
|
|
|
from monitor_runtime.licensing import require_license
|
|
|
|
|
require_license()
|
2026-08-30 22:22:11 +08:00
|
|
|
with self._guard:
|
|
|
|
|
if self._started:
|
|
|
|
|
return True
|
|
|
|
|
if not self.lock.acquire():
|
|
|
|
|
logger.info("background services skipped: another process is leader")
|
|
|
|
|
return False
|
|
|
|
|
try:
|
|
|
|
|
from app.utils.GlobalUtils import g_config, g_gb28181SipServer
|
|
|
|
|
if getattr(g_config, "autoStartMedia", False):
|
|
|
|
|
from app.utils.MediaServerManager import get_media_server_manager
|
|
|
|
|
ok, info = get_media_server_manager().start()
|
|
|
|
|
logger.info("ZLM explicit start: ok=%s %s", ok, info)
|
2026-09-04 18:16:14 +08:00
|
|
|
if not ok:
|
|
|
|
|
raise RuntimeError("ZLM startup failed: " + str(info))
|
2026-08-30 22:22:11 +08:00
|
|
|
g_gb28181SipServer.start()
|
|
|
|
|
if getattr(g_config, "recordingEnabled", False):
|
|
|
|
|
from app.recording.manager import get_recording_manager
|
|
|
|
|
get_recording_manager().start()
|
|
|
|
|
self.coordinator.start()
|
|
|
|
|
self._started = True
|
|
|
|
|
logger.info("background services started as leader pid=%s", os.getpid())
|
|
|
|
|
return True
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.exception("background service startup failed")
|
|
|
|
|
self.stop()
|
|
|
|
|
raise
|
|
|
|
|
|
|
|
|
|
def reconcile(self):
|
|
|
|
|
if not self.is_leader:
|
|
|
|
|
return False
|
|
|
|
|
from app.utils.GlobalUtils import g_config
|
|
|
|
|
from app.recording.manager import get_recording_manager
|
|
|
|
|
recording = get_recording_manager()
|
|
|
|
|
if getattr(g_config, "recordingEnabled", False):
|
|
|
|
|
recording.start()
|
|
|
|
|
else:
|
|
|
|
|
recording.stop()
|
|
|
|
|
return True
|
|
|
|
|
|
|
|
|
|
def stop(self):
|
|
|
|
|
with self._guard:
|
|
|
|
|
if not self.is_leader:
|
|
|
|
|
return
|
|
|
|
|
self.coordinator.stop()
|
|
|
|
|
try:
|
2026-09-04 18:16:14 +08:00
|
|
|
from app.analysis.manager import shutdown_analysis
|
|
|
|
|
shutdown_analysis()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.exception("analysis shutdown failed")
|
|
|
|
|
try:
|
2026-08-30 22:22:11 +08:00
|
|
|
from app.recording.manager import get_recording_manager
|
|
|
|
|
get_recording_manager().stop()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.exception("recording shutdown failed")
|
|
|
|
|
try:
|
|
|
|
|
from app.utils.GlobalUtils import g_gb28181SipServer
|
|
|
|
|
g_gb28181SipServer.stop()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.exception("SIP shutdown failed")
|
|
|
|
|
try:
|
|
|
|
|
from app.utils.MediaServerManager import get_media_server_manager
|
|
|
|
|
media = get_media_server_manager()
|
|
|
|
|
if media.managed_pid():
|
|
|
|
|
media.stop()
|
|
|
|
|
except Exception:
|
|
|
|
|
logger.exception("ZLM shutdown failed")
|
|
|
|
|
self._started = False
|
|
|
|
|
self.lock.release()
|
|
|
|
|
logger.info("background services stopped")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
_MANAGER = ServiceManager()
|
|
|
|
|
atexit.register(_MANAGER.stop)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def get_service_manager():
|
|
|
|
|
return _MANAGER
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def is_service_leader():
|
|
|
|
|
return _MANAGER.is_leader
|