"""tmq_watch 模块测试(mock TMQ consumer)。""" from unittest.mock import MagicMock, patch import pytest from ts_anomaly_td.tmq_watch import WatchRunner def make_mock_conn(): conn = MagicMock() conn.execute.return_value = [(1000, 1.0), (2000, 2.0)] return conn def test_watch_runner_polling_fallback(monkeypatch): """TMQ 不可用时自动降级为轮询。""" conn = make_mock_conn() runner = WatchRunner(conn, "test", window_sec=0.1) # 强制 TMQ 不可用 runner._run_tmq = MagicMock(side_effect=ImportError("no TMQ")) # 跑一个轮询周期就停 poll_count = [0] def mock_poll(): poll_count[0] += 1 runner._mode = "polling" runner.stop() runner._run_polling = mock_poll runner.run() assert runner._mode == "polling" assert poll_count[0] == 1 def test_watch_runner_stop_on_signal(): """收到 stop 后退出循环。""" conn = make_mock_conn() runner = WatchRunner(conn, "test", window_sec=0.1) runner.stop() assert runner._stop is True def test_watch_runner_window_check(): """_check_window 调用 detection。""" conn = make_mock_conn() runner = WatchRunner(conn, "test") with patch("ts_anomaly_td.detection.detect_all_algos") as mock_det: mock_det.return_value = { "ksigma": {"windows": [(1000, 2000)], "error": None}, } runner._check_window([(1000, 1.0)]) mock_det.assert_called_once_with(conn, "test")