video_monitor/workshop_monitor/runtime.py
2026-09-06 22:26:35 +08:00

264 lines
12 KiB
Python

"""独立于布控分析流水线的车间实时定位进程。"""
import multiprocessing as mp
import queue
import threading
import time
from .calibration import foot_point_world
from .fusion import GlobalFusionTracker
class _CameraReader(threading.Thread):
def __init__(self, config, target_fps):
super().__init__(name="workshop-camera-%s" % config["camera_id"], daemon=True)
self.config = config
self.target_fps = target_fps
self._lock = threading.Lock()
self._running = True
self._frame = None
self._timestamp = 0.0
self._sequence = 0
self._health = {"stream_health": "connecting", "stalled_sec": 0.0}
def run(self):
from app.analysis.frames import FrameSource
source = FrameSource(self.config["rtsp_url"], target_fps=max(2, int(self.target_fps * 2)))
try:
while self._running:
ok, frame = source.read()
with self._lock:
self._health = source.health_snapshot()
if ok and frame is not None:
self._frame = frame
self._timestamp = time.time()
self._sequence += 1
if not ok:
time.sleep(0.1)
finally:
source.close()
def latest(self):
with self._lock:
return self._sequence, self._timestamp, self._frame, dict(self._health)
def close(self):
self._running = False
def _build_engine(spec):
from app.analysis.engines.factory import EngineFactory
engine = EngineFactory.create(
spec["inference_engine"], model_file=spec["model_file"], labels=spec.get("labels") or [],
input_size=tuple(spec.get("input_size") or (640, 640)),
conf_threshold=float(spec.get("conf_threshold", 0.4)),
iou_threshold=float(spec.get("iou_threshold", 0.5)),
algorithm_type=spec.get("algorithm_type", "yolo"),
task_type=spec.get("task_type", "detect"), device=spec.get("device", "cpu"),
target_labels=spec.get("target_labels") or [],
)
if not engine.load():
raise RuntimeError("模型加载失败: %s" % spec.get("name", spec.get("model_file", "")))
return engine
def workshop_worker_main(config, state_queue, command_queue):
readers = []
try:
detector = _build_engine(config["detector"])
reid = _build_engine(config["reid_model"]) if config.get("reid_model") else None
tracker_cls = __import__("app.analysis.tracker", fromlist=["IoUTracker"]).IoUTracker
trackers = {int(c["camera_id"]): tracker_cls() for c in config["cameras"]}
readers = [_CameraReader(c, config["analysis_fps"]) for c in config["cameras"]]
for reader in readers:
reader.start()
fusion = GlobalFusionTracker(
config["fusion_radius_m"], config["observation_window_sec"],
config["max_speed_mps"], config["lost_ttl_sec"], config["trail_sec"],
)
last_seq = {int(c["camera_id"]): 0 for c in config["cameras"]}
last_process = {int(c["camera_id"]): 0.0 for c in config["cameras"]}
latest_observations = {int(c["camera_id"]): [] for c in config["cameras"]}
frame_index = {int(c["camera_id"]): 0 for c in config["cameras"]}
state_sequence = 0
period = 1.0 / max(0.1, float(config["analysis_fps"]))
running = True
while running:
try:
while True:
command = command_queue.get_nowait()
if command.get("cmd") == "stop":
running = False
except queue.Empty:
pass
if not running:
break
changed = False
camera_states = []
now = time.time()
for camera, reader in zip(config["cameras"], readers):
cid = int(camera["camera_id"])
sequence, captured_at, frame, health = reader.latest()
if frame is not None and sequence != last_seq[cid] and now - last_process[cid] >= period:
last_seq[cid] = sequence
last_process[cid] = now
frame_index[cid] += 1
try:
detections = [d for d in detector.detect(frame)
if d.get("label") in config["target_labels"]]
active, _ended, _new, _idx = trackers[cid].update(
detections, frame_index[cid], timestamp=captured_at)
confirmed = [t for t in active if t.get("confirmed") and t.get("observed")]
embeddings = {}
if reid and confirmed:
valid, values = reid.extract_embeddings(frame, [t["box"] for t in confirmed])
for output_index, track_index in enumerate(valid):
embeddings[confirmed[track_index]["track_id"]] = values[output_index]
h, w = frame.shape[:2]
observations = []
calibration = camera["calibration"]
for track in confirmed:
world = foot_point_world(track["box"], calibration["homography"],
config["width_m"], config["height_m"])
if world is None:
continue
box = [float(track["box"][0]) / w, float(track["box"][1]) / h,
float(track["box"][2]) / w, float(track["box"][3]) / h]
observations.append({
"camera_id": cid, "stream_id": camera["stream_id"],
"local_track_id": int(track["track_id"]), "class": track.get("label", "person"),
"score": float(track.get("score", 0)), "bbox": box,
"x": world[0], "y": world[1], "timestamp": captured_at,
"calibration_error_m": calibration.get("validation_mean_m") or 1.0,
"embedding": embeddings.get(track["track_id"]),
})
latest_observations[cid] = observations
health["analysis_health"] = "running"
health["analysis_error"] = ""
changed = True
except Exception as exc:
health["analysis_health"] = "error"
health["analysis_error"] = str(exc)
fresh = [o for o in latest_observations[cid]
if now - float(o.get("timestamp", 0)) <= config["observation_window_sec"]]
camera_states.append({
"camera_id": cid, "stream_id": camera["stream_id"], "slot": camera["slot"],
"display_name": camera["display_name"], **health, "observations": fresh,
})
if changed or state_sequence == 0:
all_observations = [o for rows in latest_observations.values() for o in rows]
targets = fusion.update(all_observations, now)
for obs in all_observations:
choices = [t for t in targets if t["class"] == obs["class"] and
obs["camera_id"] in t["source_camera_ids"]]
if choices:
obs["global_id"] = min(choices, key=lambda t: (t["x"] - obs["x"]) ** 2 +
(t["y"] - obs["y"]) ** 2)["global_id"]
obs.pop("embedding", None)
state_sequence += 1
payload = {"kind": "state", "sequence": state_sequence, "timestamp": now,
"running": True, "cameras": camera_states, "targets": targets}
try:
while True:
state_queue.get_nowait()
except queue.Empty:
pass
state_queue.put(payload)
time.sleep(0.02)
except Exception as exc:
state_queue.put({"kind": "error", "running": False, "error": str(exc), "timestamp": time.time()})
finally:
for reader in readers:
reader.close()
for reader in readers:
reader.join(timeout=2.0)
state_queue.put({"kind": "stopped", "running": False, "timestamp": time.time()})
class WorkshopRuntimeManager:
_instance = None
_instance_lock = threading.Lock()
def __new__(cls):
if cls._instance is None:
with cls._instance_lock:
if cls._instance is None:
cls._instance = super().__new__(cls)
cls._instance._initialized = False
return cls._instance
def __init__(self):
if self._initialized:
return
self._initialized = True
self._lock = threading.RLock()
self._process = None
self._state_queue = None
self._command_queue = None
self._listener = None
self._state = {"running": False, "sequence": 0, "cameras": [], "targets": []}
def _listen(self):
while self._process is not None:
try:
message = self._state_queue.get(timeout=0.5)
with self._lock:
self._state = message
except queue.Empty:
if self._process is not None and not self._process.is_alive():
with self._lock:
self._state = {**self._state, "running": False}
if self._state.get("kind") != "stopped" and not self._state.get("error"):
self._state["error"] = "车间定位进程已退出"
break
def start(self, config):
with self._lock:
if self._process is not None and self._process.is_alive():
return True, "already running"
context = mp.get_context("spawn")
self._state_queue = context.Queue(maxsize=4)
self._command_queue = context.Queue(maxsize=8)
self._process = context.Process(
target=workshop_worker_main, args=(config, self._state_queue, self._command_queue),
name="workshop-monitor", daemon=True,
)
self._state = {"running": True, "sequence": 0, "cameras": [], "targets": [],
"timestamp": time.time()}
self._process.start()
self._listener = threading.Thread(target=self._listen, name="workshop-state-listener", daemon=True)
self._listener.start()
return True, "started"
def stop(self):
with self._lock:
process = self._process
if process is None or not process.is_alive():
self._state = {**self._state, "running": False}
return True, "already stopped"
try:
self._command_queue.put({"cmd": "stop"}, timeout=1.0)
except Exception:
pass
process.join(timeout=8.0)
if process.is_alive():
process.terminate()
process.join(timeout=3.0)
with self._lock:
self._process = None
self._state = {**self._state, "running": False}
return True, "stopped"
def snapshot(self, since=None):
with self._lock:
state = dict(self._state)
sequence = int(state.get("sequence") or 0)
if since is not None and sequence <= int(since or 0):
return {"changed": False, "sequence": sequence, "running": bool(state.get("running")),
"timestamp": state.get("timestamp"), "error": state.get("error", "")}
state["changed"] = True
return state
def get_runtime_manager():
return WorkshopRuntimeManager()