2026-08-30 22:22:11 +08:00
|
|
|
|
"""Monitor 分析子进程入口
|
|
|
|
|
|
|
|
|
|
|
|
每路摄像头在独立进程中运行 CameraPipeline,绕过 GIL,与主进程通过 Queue 通信:
|
|
|
|
|
|
- event_queue: 子 → 主,上报检测/区域/追踪事件
|
|
|
|
|
|
- cmd_queue: 主 → 子,stop / reload_zones
|
|
|
|
|
|
- status_dict: 共享状态(Manager.dict),供 openStatus 读取 FPS
|
|
|
|
|
|
"""
|
|
|
|
|
|
import json
|
|
|
|
|
|
import logging
|
|
|
|
|
|
import multiprocessing as mp
|
|
|
|
|
|
import os
|
|
|
|
|
|
import queue
|
|
|
|
|
|
import time
|
|
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger("analysis.process_worker")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _algorithm_spec_from_dict(d):
|
|
|
|
|
|
return d
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-09-04 18:16:14 +08:00
|
|
|
|
def _build_detectors_in_process(algorithm_specs, infer_req_q=None, infer_resp_q=None,
|
|
|
|
|
|
response_channel=None):
|
2026-08-30 22:22:11 +08:00
|
|
|
|
"""在子进程中构造检测器;若提供 infer_req_q/infer_resp_q 则走主进程共享推理池"""
|
|
|
|
|
|
detectors = []
|
|
|
|
|
|
if infer_req_q is not None and infer_resp_q is not None:
|
|
|
|
|
|
from app.analysis.remote_detector import RemoteDetector
|
|
|
|
|
|
for spec in algorithm_specs:
|
|
|
|
|
|
detectors.append({
|
|
|
|
|
|
"algorithm_id": spec.get("id", 0),
|
|
|
|
|
|
"algorithm_name": spec.get("name", ""),
|
2026-09-04 18:16:14 +08:00
|
|
|
|
"engine": RemoteDetector(spec, infer_req_q, infer_resp_q,
|
|
|
|
|
|
response_channel=response_channel),
|
|
|
|
|
|
"target_labels": spec.get("target_labels") or [],
|
|
|
|
|
|
"device": spec.get("device") or "cpu",
|
2026-08-30 22:22:11 +08:00
|
|
|
|
})
|
|
|
|
|
|
else:
|
|
|
|
|
|
from app.analysis.worker_pool import DetectorWorkerPool
|
|
|
|
|
|
pool = DetectorWorkerPool()
|
|
|
|
|
|
|
|
|
|
|
|
class _AlgoObj(object):
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
for spec in algorithm_specs:
|
|
|
|
|
|
o = _AlgoObj()
|
|
|
|
|
|
for k, v in spec.items():
|
|
|
|
|
|
setattr(o, k, v)
|
|
|
|
|
|
eng = pool.get_detector(o)
|
|
|
|
|
|
if eng:
|
|
|
|
|
|
detectors.append({
|
|
|
|
|
|
"algorithm_id": spec.get("id", 0),
|
|
|
|
|
|
"algorithm_name": spec.get("name", ""),
|
|
|
|
|
|
"engine": eng,
|
2026-09-04 18:16:14 +08:00
|
|
|
|
"target_labels": spec.get("target_labels") or [],
|
|
|
|
|
|
"device": spec.get("device") or "cpu",
|
2026-08-30 22:22:11 +08:00
|
|
|
|
})
|
|
|
|
|
|
return detectors
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def pipeline_process_main(config, event_queue, cmd_queue, status_dict,
|
|
|
|
|
|
infer_req_q=None, infer_resp_q=None):
|
|
|
|
|
|
"""子进程主函数(spawn 入口,config 必须为纯 dict)"""
|
|
|
|
|
|
from app.utils.Logger import LOG_FORMAT
|
|
|
|
|
|
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)
|
|
|
|
|
|
log = logging.getLogger("analysis.process_worker")
|
|
|
|
|
|
|
|
|
|
|
|
from app.analysis.pipeline import CameraPipeline
|
|
|
|
|
|
from app.analysis.motion import MotionDetector
|
|
|
|
|
|
|
|
|
|
|
|
stream_id = config["stream_id"]
|
|
|
|
|
|
stream_code = config.get("stream_code", str(stream_id))
|
|
|
|
|
|
rtsp_url = config["rtsp_url"]
|
|
|
|
|
|
target_fps = config.get("target_fps", 5)
|
|
|
|
|
|
analyze_fps = config.get("analyze_fps", target_fps)
|
|
|
|
|
|
zones = config.get("zones") or []
|
|
|
|
|
|
algorithm_specs = config.get("algorithms") or []
|
|
|
|
|
|
use_shared_inference = bool(config.get("use_shared_inference", True))
|
|
|
|
|
|
storage_alarm_dir = config.get("storage_alarm_dir") or ""
|
|
|
|
|
|
static_dir = config.get("static_dir") or ""
|
|
|
|
|
|
|
|
|
|
|
|
def on_event(ev):
|
|
|
|
|
|
try:
|
|
|
|
|
|
event_queue.put(("event", ev), timeout=2.0)
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
log.warning("事件入队失败: %s", e)
|
|
|
|
|
|
|
|
|
|
|
|
def on_track_snapshot(sid, frame_index, active, has_motion):
|
|
|
|
|
|
try:
|
|
|
|
|
|
event_queue.put(("touch", {
|
|
|
|
|
|
"stream_id": sid,
|
|
|
|
|
|
"active": active,
|
|
|
|
|
|
"has_motion": has_motion,
|
|
|
|
|
|
}), timeout=1.0)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
2026-09-04 18:16:14 +08:00
|
|
|
|
def on_preview(payload):
|
|
|
|
|
|
try:
|
|
|
|
|
|
event_queue.put(("preview", payload), timeout=0.2)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
2026-08-30 22:22:11 +08:00
|
|
|
|
detectors = _build_detectors_in_process(
|
|
|
|
|
|
algorithm_specs,
|
|
|
|
|
|
infer_req_q if use_shared_inference else None,
|
|
|
|
|
|
infer_resp_q if use_shared_inference else None,
|
2026-09-04 18:16:14 +08:00
|
|
|
|
config.get("response_channel"),
|
2026-08-30 22:22:11 +08:00
|
|
|
|
)
|
|
|
|
|
|
if not detectors and not algorithm_specs:
|
|
|
|
|
|
log.info("pipeline[%s] 无算法,仅运动检测", stream_code)
|
|
|
|
|
|
|
|
|
|
|
|
motion = MotionDetector()
|
|
|
|
|
|
pipeline = CameraPipeline(
|
|
|
|
|
|
stream_id=stream_id,
|
|
|
|
|
|
stream_code=stream_code,
|
|
|
|
|
|
rtsp_url=rtsp_url,
|
|
|
|
|
|
detectors=detectors,
|
|
|
|
|
|
motion=motion,
|
|
|
|
|
|
target_fps=target_fps,
|
|
|
|
|
|
analyze_fps=analyze_fps,
|
|
|
|
|
|
on_event=on_event,
|
|
|
|
|
|
on_track_snapshot=on_track_snapshot,
|
2026-09-04 18:16:14 +08:00
|
|
|
|
on_preview=on_preview,
|
|
|
|
|
|
alarm_enabled=not bool(config.get("preview_only", False)),
|
2026-08-30 22:22:11 +08:00
|
|
|
|
zone_polygons=zones,
|
|
|
|
|
|
storage_alarm_dir=storage_alarm_dir,
|
|
|
|
|
|
static_dir=static_dir,
|
|
|
|
|
|
)
|
|
|
|
|
|
pipeline._algorithm_name = ", ".join(d.get("name", "") for d in algorithm_specs) or "motion-only"
|
|
|
|
|
|
|
|
|
|
|
|
import threading
|
2026-09-04 18:16:14 +08:00
|
|
|
|
finished = threading.Event()
|
|
|
|
|
|
|
|
|
|
|
|
def wait_started():
|
|
|
|
|
|
while not finished.is_set():
|
|
|
|
|
|
if pipeline._started_event.wait(0.1):
|
|
|
|
|
|
return True
|
|
|
|
|
|
return False
|
2026-08-30 22:22:11 +08:00
|
|
|
|
|
|
|
|
|
|
def status_reporter():
|
|
|
|
|
|
# 等待 pipeline.run() 启动(_running 在 run() 里才置 True,
|
|
|
|
|
|
# 否则 while pipeline._running 条件不满足会立即退出,导致 status_dict 永远为空)
|
2026-09-04 18:16:14 +08:00
|
|
|
|
if not wait_started():
|
|
|
|
|
|
return
|
|
|
|
|
|
while not finished.is_set():
|
2026-08-30 22:22:11 +08:00
|
|
|
|
try:
|
|
|
|
|
|
st = pipeline.status()
|
|
|
|
|
|
status_dict[str(stream_id)] = st
|
|
|
|
|
|
except Exception:
|
2026-09-04 18:16:14 +08:00
|
|
|
|
log.exception("pipeline[%s] 状态上报失败", stream_code)
|
|
|
|
|
|
finished.wait(1.0)
|
2026-08-30 22:22:11 +08:00
|
|
|
|
|
|
|
|
|
|
reporter = threading.Thread(target=status_reporter, name="status-%s" % stream_id, daemon=True)
|
|
|
|
|
|
reporter.start()
|
|
|
|
|
|
|
|
|
|
|
|
def cmd_listener():
|
2026-09-04 18:16:14 +08:00
|
|
|
|
if not wait_started():
|
|
|
|
|
|
return
|
|
|
|
|
|
while not finished.is_set():
|
2026-08-30 22:22:11 +08:00
|
|
|
|
try:
|
|
|
|
|
|
cmd = cmd_queue.get(timeout=0.5)
|
|
|
|
|
|
except queue.Empty:
|
|
|
|
|
|
continue
|
|
|
|
|
|
if not cmd:
|
|
|
|
|
|
continue
|
|
|
|
|
|
try:
|
|
|
|
|
|
if cmd.get("cmd") == "stop":
|
|
|
|
|
|
pipeline.stop()
|
|
|
|
|
|
break
|
|
|
|
|
|
if cmd.get("cmd") == "reload_zones":
|
|
|
|
|
|
pipeline.set_zone_polygons(cmd.get("zones") or [])
|
|
|
|
|
|
if cmd.get("analyze_fps") is not None:
|
|
|
|
|
|
pipeline.set_analyze_fps(cmd.get("analyze_fps"))
|
2026-09-04 18:16:14 +08:00
|
|
|
|
if cmd.get("cmd") == "set_analyze_fps":
|
|
|
|
|
|
pipeline.set_analyze_fps(cmd.get("analyze_fps") or 1)
|
2026-08-30 22:22:11 +08:00
|
|
|
|
except Exception as e:
|
|
|
|
|
|
log.exception("pipeline[%s] 命令处理失败: %s", stream_code, e)
|
|
|
|
|
|
|
|
|
|
|
|
cmd_thread = threading.Thread(target=cmd_listener, name="cmd-%s" % stream_id, daemon=True)
|
|
|
|
|
|
cmd_thread.start()
|
|
|
|
|
|
|
|
|
|
|
|
log.info("pipeline 子进程启动 stream=%s url=%s", stream_code, rtsp_url)
|
|
|
|
|
|
try:
|
|
|
|
|
|
pipeline.run()
|
|
|
|
|
|
finally:
|
2026-09-04 18:16:14 +08:00
|
|
|
|
finished.set()
|
|
|
|
|
|
reporter.join(timeout=2)
|
|
|
|
|
|
cmd_thread.join(timeout=1)
|
2026-08-30 22:22:11 +08:00
|
|
|
|
try:
|
|
|
|
|
|
status_dict.pop(str(stream_id), None)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
log.info("pipeline 子进程退出 stream=%s", stream_code)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class PipelineProcessHandle(object):
|
|
|
|
|
|
"""主进程侧对子进程的封装"""
|
|
|
|
|
|
|
|
|
|
|
|
def __init__(self, stream_id, process, event_queue, cmd_queue, status_dict):
|
|
|
|
|
|
self.stream_id = stream_id
|
|
|
|
|
|
self.process = process
|
|
|
|
|
|
self.event_queue = event_queue
|
|
|
|
|
|
self.cmd_queue = cmd_queue
|
|
|
|
|
|
self.status_dict = status_dict
|
|
|
|
|
|
self.running = True
|
|
|
|
|
|
|
|
|
|
|
|
def stop(self, timeout=5):
|
|
|
|
|
|
self.running = False
|
|
|
|
|
|
try:
|
|
|
|
|
|
self.cmd_queue.put({"cmd": "stop"}, timeout=1.0)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
try:
|
|
|
|
|
|
self.process.join(timeout=timeout)
|
|
|
|
|
|
if self.process.is_alive():
|
2026-09-04 18:16:14 +08:00
|
|
|
|
logger.warning("pipeline[%s] 未在 %ss 内停止,强制结束子进程", self.stream_id, timeout)
|
2026-08-30 22:22:11 +08:00
|
|
|
|
self.process.terminate()
|
|
|
|
|
|
self.process.join(timeout=2)
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
pass
|
|
|
|
|
|
|
|
|
|
|
|
def reload_zones(self, zones, analyze_fps=None):
|
|
|
|
|
|
try:
|
|
|
|
|
|
payload = {"cmd": "reload_zones", "zones": zones}
|
|
|
|
|
|
if analyze_fps is not None:
|
|
|
|
|
|
payload["analyze_fps"] = analyze_fps
|
|
|
|
|
|
self.cmd_queue.put(payload, timeout=1.0)
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.warning("reload_zones 发送失败: %s", e)
|
|
|
|
|
|
|
2026-09-04 18:16:14 +08:00
|
|
|
|
def set_analyze_fps(self, analyze_fps):
|
|
|
|
|
|
try:
|
|
|
|
|
|
self.cmd_queue.put({"cmd": "set_analyze_fps", "analyze_fps": float(analyze_fps)}, timeout=1.0)
|
|
|
|
|
|
except Exception as e:
|
|
|
|
|
|
logger.warning("set_analyze_fps 发送失败: %s", e)
|
|
|
|
|
|
|
2026-08-30 22:22:11 +08:00
|
|
|
|
def status(self):
|
|
|
|
|
|
try:
|
|
|
|
|
|
return self.status_dict.get(str(self.stream_id))
|
|
|
|
|
|
except Exception:
|
|
|
|
|
|
return None
|