video_monitor/tests/test_analysis_runtime.py
2026-09-04 18:16:14 +08:00

337 lines
15 KiB
Python

"""Regression coverage for camera restart, response routing and inference health."""
import multiprocessing as mp
import queue
import threading
import time
import unittest
from types import SimpleNamespace
from unittest.mock import patch
import numpy as np
from app.analysis.manager import AnalysisManager
from app.analysis.pipeline import CameraPipeline
from app.analysis.process_worker import pipeline_process_main
from app.analysis.remote_detector import RemoteDetector
from app.analysis.tracker import IoUTracker
from app.analysis.preview_sessions import PreviewSessionRegistry
class CommandPipeline:
"""Delay run() so control threads must survive the pre-start window."""
def __init__(self, **kwargs):
self._running = False
self._started_event = threading.Event()
self.zones = []
def run(self):
time.sleep(0.2)
self._running = True
self._started_event.set()
deadline = time.monotonic() + 3
while self._running and time.monotonic() < deadline:
time.sleep(0.01)
assert not self._running, 'stop command was lost before run()'
assert self.zones == [{'id': 12}], 'reload command was lost'
def status(self):
return {'running': self._running}
def set_zone_polygons(self, zones):
self.zones = zones
def stop(self):
self._running = False
def command_process(commands):
with patch('app.analysis.pipeline.CameraPipeline', CommandPipeline):
pipeline_process_main({'stream_id': 1, 'rtsp_url': 'test'},
queue.Queue(), commands, {})
def wait_response(responses, entered):
entered.set()
responses.get(timeout=60)
class AnalysisRuntimeTests(unittest.TestCase):
def make_router(self):
manager = object.__new__(AnalysisManager)
manager._infer_routes = {}
manager._infer_routes_lock = threading.Lock()
return manager
def test_stop_and_reload_commands_survive_startup(self):
ctx = mp.get_context('spawn')
commands = ctx.Queue()
commands.put({'cmd': 'reload_zones', 'zones': [{'id': 12}]})
commands.put({'cmd': 'stop'})
proc = ctx.Process(target=command_process, args=(commands,))
proc.start()
try:
proc.join(10)
self.assertEqual(proc.exitcode, 0)
finally:
if proc.is_alive():
proc.terminate()
proc.join(2)
commands.cancel_join_thread()
commands.close()
def test_restart_uses_new_channel_after_consumer_is_killed(self):
ctx = mp.get_context('spawn')
router = self.make_router()
old, fresh, other = ctx.Queue(), ctx.Queue(), ctx.Queue()
router._infer_routes.update(old=old, fresh=fresh, other=other)
entered = ctx.Event()
proc = ctx.Process(target=wait_response, args=(old, entered))
proc.start()
try:
self.assertTrue(entered.wait(10))
deadline = time.monotonic() + 3
while time.monotonic() < deadline:
if not old._rlock.acquire(False):
break
old._rlock.release()
time.sleep(0.01)
else:
self.fail('consumer did not acquire response queue lock')
proc.terminate()
proc.join(3)
old.put({'req_id': 'unreadable'})
with self.assertRaises(queue.Empty):
old.get(timeout=0.1)
router._close_inference_channel({'response_channel': 'old'})
router._send_inference_response('old', {'req_id': 'late'})
router._send_inference_response('fresh', {'req_id': 'new'})
router._send_inference_response('other', {'req_id': 'camera-2'})
self.assertEqual(fresh.get(timeout=2)['req_id'], 'new')
self.assertEqual(other.get(timeout=2)['req_id'], 'camera-2')
finally:
if proc.is_alive():
proc.terminate()
proc.join(2)
for key in list(router._infer_routes):
router._close_inference_channel({'response_channel': key})
def test_remote_timeout_is_not_an_empty_success(self):
detector = RemoteDetector({'id': 5}, queue.Queue(), queue.Queue(), timeout=0.03)
with self.assertRaises(TimeoutError):
detector.detect(np.zeros((32, 32, 3), dtype=np.uint8))
def test_downscaled_inference_coordinates_are_restored(self):
restored = RemoteDetector._restore_coordinates([{
'box': [10, 20, 100, 200],
'keypoints': [[15, 25, .9]],
'polygon': [[10, 20], [100, 200]],
}], 2.0, 2.0)[0]
self.assertEqual(restored['box'], [20.0, 40.0, 200.0, 400.0])
self.assertEqual(restored['keypoints'][0][:2], [30.0, 50.0])
self.assertEqual(restored['polygon'][1], [200.0, 400.0])
def test_tracker_is_one_to_one_and_confirms_medium_score_twice(self):
tracker = IoUTracker()
detections = [
{'label': 'person', 'score': .45, 'box': [0, 0, 20, 40]},
{'label': 'person', 'score': .8, 'box': [30, 0, 50, 40]},
]
first, _, _, _ = tracker.update(detections, 1, timestamp=1.0)
self.assertEqual(len({t['track_id'] for t in first}), 2)
self.assertFalse(next(t for t in first if t['score'] == .45)['confirmed'])
second, _, new_ids, _ = tracker.update([
{'label': 'person', 'score': .5, 'box': [2, 0, 22, 40]},
{'label': 'person', 'score': .82, 'box': [28, 0, 48, 40]},
], 2, timestamp=4.0)
self.assertEqual(new_ids, [])
self.assertEqual(len({t['track_id'] for t in second}), 2)
self.assertTrue(all(t['confirmed'] for t in second))
def test_low_fps_track_survives_short_miss(self):
tracker = IoUTracker()
active, _, _, _ = tracker.update([
{'label': 'person', 'score': .8, 'box': [0, 0, 20, 40]},
], 1, timestamp=1.0)
tid = active[0]['track_id']
tracker.update([], 2, timestamp=4.0)
active, ended, new_ids, _ = tracker.update([
{'label': 'person', 'score': .8, 'box': [2, 0, 22, 40]},
], 3, timestamp=7.0)
self.assertNotIn(tid, ended)
self.assertEqual(new_ids, [])
self.assertEqual(active[0]['track_id'], tid)
def test_preview_uses_capture_timestamp_and_normalized_velocity(self):
snapshots = []
pipe = CameraPipeline(1, 'test', 'test', on_preview=snapshots.append)
pipe._last_w, pipe._last_h = 100, 200
pipe._frame_index = 2
pipe._tracker.update([
{'label': 'person', 'score': .8, 'box': [10, 20, 30, 80]},
], 1, timestamp=10.0)
pipe._tracker.update([
{'label': 'person', 'score': .8, 'box': [20, 30, 40, 90]},
], 2, timestamp=11.0)
pipe._publish_preview(11.0, processed_ts=11.25)
snapshot = snapshots[-1]
self.assertEqual(snapshot['timestamp'], 11.0)
self.assertEqual(snapshot['processed_timestamp'], 11.25)
self.assertAlmostEqual(snapshot['inference_latency'], .25)
self.assertEqual(snapshot['tracks'][0]['box'], [.2, .15, .4, .45])
self.assertEqual(snapshot['tracks'][0]['velocity'], [.05, .025, .05, .025])
def test_formal_preview_temporarily_boosts_and_restores_analysis_fps(self):
class Handle:
def __init__(self): self.calls = []
def set_analyze_fps(self, fps): self.calls.append(float(fps))
class Process:
@staticmethod
def is_alive(): return True
handle = Handle()
manager = object.__new__(AnalysisManager)
manager._lock = threading.RLock()
manager._pipelines = {87: {
'running': True, 'mode': 'process', 'process': Process(),
'handle': handle, 'preview_only': False, 'analyze_fps': .67,
}}
zone = SimpleNamespace(stream=SimpleNamespace(id=87))
with patch('monitor_runtime.licensing.require_license'):
self.assertEqual(manager.start_preview(zone), (True, 'formal'))
self.assertEqual(handle.calls, [5.0])
self.assertEqual(manager._pipelines[87]['preview_restore_fps'], .67)
self.assertFalse(manager.stop_preview(87))
self.assertEqual(handle.calls, [5.0, .67])
self.assertNotIn('preview_restore_fps', manager._pipelines[87])
def test_medium_score_confirmation_requires_consecutive_hits(self):
tracker = IoUTracker()
tracker.update([{'label': 'person', 'score': .4, 'box': [0, 0, 20, 40]}], 1, timestamp=1)
tracker.update([], 2, timestamp=2)
active, _, _, _ = tracker.update([
{'label': 'person', 'score': .5, 'box': [1, 0, 21, 40]},
], 3, timestamp=3)
self.assertFalse(active[0]['confirmed'])
def test_person_foot_point_and_two_outside_hits(self):
events = []
zone = {'id': 12, 'name': 'office', 'coords': [[0, .8], [1, .8], [1, 1], [0, 1]],
'alarm_repeat_sec': 30, 'biz_algorithms': [{
'id': 1, 'name': 'person intrusion', 'flow_type': 1,
'small_model_id': 5, 'target_labels': ['person'], 'post_process': 'AREA'}]}
pipe = CameraPipeline(1, 'test', 'test', zone_polygons=[zone], on_event=events.append)
pipe._last_w, pipe._last_h = 100, 100
inside = {'track_id': 7, 'label': 'person', 'score': .8, 'box': [40, 40, 60, 90],
'algorithm_id': 5, 'confirmed': True, 'observed': True}
with patch.object(pipe, '_emit_biz_alarm', return_value=True) as alarm:
pipe._check_zones([inside], 1, np.zeros((100, 100, 3), dtype=np.uint8))
self.assertTrue(alarm.called) # 中心 y=65 在外,脚点 y=87.5 在区域内
outside = dict(inside, box=[40, 20, 60, 50])
pipe._check_zones([outside], 2, np.zeros((100, 100, 3), dtype=np.uint8))
self.assertEqual(pipe._last_zone_state[7], {12})
pipe._check_zones([outside], 3, np.zeros((100, 100, 3), dtype=np.uint8))
self.assertEqual(pipe._last_zone_state[7], set())
def test_unobserved_retained_track_does_not_leave_or_reenter(self):
pipe = CameraPipeline(1, 'test', 'test')
active, _, _, _ = pipe._tracker.update([
{'label': 'person', 'score': .8, 'box': [10, 10, 20, 30]},
], 1, timestamp=time.time())
tid = active[0]['track_id']
pipe._last_zone_state = {tid: {12}}
pipe._track_enter_ts = {(tid, 12): 10.0}
pipe._check_zones([], 2, np.zeros((20, 20, 3), dtype=np.uint8))
self.assertEqual(pipe._last_zone_state[tid], {12})
self.assertEqual(pipe._track_enter_ts[(tid, 12)], 10.0)
def test_area_alarm_repeat_is_independent_from_detection_rate(self):
pipe = CameraPipeline(1, 'test', 'test')
rule = {'id': 1, 'post_process': 'AREA', 'target_labels': ['person'],
'flow_type': 1, 'small_model_id': 5}
zone = {'id': 12, 'detect_interval_sec': 1, 'alarm_repeat_sec': 30,
'biz_algorithms': [rule]}
track = {'track_id': 7, 'label': 'person', 'score': .8, 'algorithm_id': 5,
'box': [0, 0, 10, 20]}
with patch.object(pipe, '_emit_biz_alarm', return_value=True) as alarm:
pipe._fire_area_alarms(track, 12, zone, set(), None, track['box'], 100.0, 7)
pipe._fire_area_alarms(track, 12, zone, {12}, None, track['box'], 101.0, 7)
pipe._fire_area_alarms(track, 12, zone, {12}, None, track['box'], 130.0, 7)
self.assertEqual(alarm.call_count, 2)
def test_preview_session_owner_and_last_viewer_cleanup(self):
class Related(list):
def filter(self, **_kwargs): return self
def select_related(self, *_args): return self
class Model:
id, name, state = 5, 'YOLO26x', 1
class Rule:
flow_type, detector_model, small_model = 1, None, Model()
target_labels = '["person"]'
class Stream:
id, app, name = 87, 'live', 'camera01'
class Zone:
id, stream_id, stream, name, color = 12, 87, Stream(), 'office', '#169F85'
coordinates = '[[0,0],[1,0],[1,1]]'
algorithms = Related([Rule()])
registry = PreviewSessionRegistry()
with patch('app.analysis.preview_sessions.AnalysisManager') as manager_cls:
manager = manager_cls.return_value
manager.start_preview.return_value = (True, 'preview')
manager.preview_mode.return_value = 'preview'
first = registry.start('owner-a', Zone())
second = registry.start('owner-b', Zone())
with self.assertRaises(PermissionError):
registry.data('owner-b', first['session_id'])
self.assertTrue(registry.stop('owner-a', first['session_id']))
manager.stop_preview.assert_not_called()
self.assertTrue(registry.stop('owner-b', second['session_id']))
manager.stop_preview.assert_called_once_with(87)
registry._running = False
def test_worker_errors_are_preserved(self):
requests, responses = queue.Queue(), queue.Queue()
detector = RemoteDetector({'id': 5}, requests, responses, timeout=2,
response_channel='camera-1')
def reply():
msg = requests.get(timeout=2)
self.assertEqual(msg['response_channel'], 'camera-1')
responses.put({'req_id': msg['req_id'], 'ok': False, 'error': 'engine load failed'})
worker = threading.Thread(target=reply)
worker.start()
try:
with self.assertRaisesRegex(RuntimeError, 'engine load failed'):
detector.detect(np.zeros((32, 32, 3), dtype=np.uint8))
finally:
worker.join(3)
def test_failed_frame_is_not_counted_and_success_recovers(self):
pipe = CameraPipeline(1, 'test', 'test')
frame = np.zeros((32, 32, 3), dtype=np.uint8)
def process(_):
pipe.stop()
raise TimeoutError('response timeout')
with patch.object(pipe, '_decode_loop'), \
patch.object(pipe, '_pop_latest_frame', return_value=frame), \
patch.object(pipe, '_process_frame', side_effect=process):
pipe.run()
self.assertEqual(pipe.status()['analyzed_count'], 0)
self.assertEqual(pipe.status()['analysis_health'], 'error')
pipe._last_analyze_ts = 0
with patch.object(pipe, '_decode_loop'), \
patch.object(pipe, '_pop_latest_frame', return_value=frame), \
patch.object(pipe, '_process_frame', side_effect=lambda _: pipe.stop()):
pipe.run()
self.assertEqual(pipe.status()['analyzed_count'], 1)
self.assertEqual(pipe.status()['analysis_health'], 'running')
self.assertEqual(pipe.status()['analysis_error'], '')
def test_low_frequency_success_keeps_nonzero_fps(self):
pipe = CameraPipeline(1, 'test', 'test', analyze_fps=1 / 60)
pipe._last_analysis_ts = time.time()
pipe._analysis_fps = 1 / 60
self.assertEqual(pipe.status()['analysis_health'], 'running')
self.assertGreater(pipe.status()['analysis_fps'], 0)
if __name__ == '__main__':
unittest.main()