test_cdc_monitor.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  1. # -*- coding:utf-8 -*-
  2. """
  3. bin/cdc-monitor.py 单测:复制槽判定 / Flink 作业存活 / 消息渲染 / main 串联。
  4. 不连真 PG(pgdb.connect + query 走替身)、不打真 JobManager(urlopen 走替身)、
  5. 不发真消息(Alerter 走替身)。
  6. 脚本路径含连字符,用 importlib.util 动态加载为模块。
  7. """
  8. import importlib.util
  9. import json
  10. import os
  11. import re
  12. import sys
  13. from datetime import datetime
  14. from unittest.mock import MagicMock
  15. import pytest
  16. PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
  17. SCRIPT_PATH = os.path.join(PROJECT_ROOT, 'bin', 'cdc-monitor.py')
  18. GB = 1024 ** 3
  19. LAG_CRIT = 5 * GB
  20. JM = 'cdhmaster02:8081'
  21. def _load_script():
  22. spec = importlib.util.spec_from_file_location('cdc_monitor', SCRIPT_PATH)
  23. mod = importlib.util.module_from_spec(spec)
  24. sys.modules['cdc_monitor'] = mod
  25. spec.loader.exec_module(mod)
  26. return mod
  27. MON = _load_script()
  28. def _stub_pg(monkeypatch, rows, conn=None):
  29. """替身 PG:connect 返回 conn,query 返回 rows。返回该 conn 供断言 close。"""
  30. conn = conn or MagicMock()
  31. monkeypatch.setattr(MON.pgdb, 'connect', lambda ds_ref: conn)
  32. monkeypatch.setattr(MON.pgdb, 'query', lambda c, sql: rows)
  33. return conn
  34. def _stub_flink(monkeypatch, jobs):
  35. """替身 JobManager:/jobs/overview 返回给定作业列表。返回捕获的请求参数。"""
  36. captured = {}
  37. def _open(url, timeout=None):
  38. captured['url'] = url
  39. captured['timeout'] = timeout
  40. resp = MagicMock()
  41. resp.read.return_value = json.dumps({'jobs': jobs}).encode('utf-8')
  42. resp.__enter__.return_value = resp
  43. return resp
  44. monkeypatch.setattr(MON.urllib.request, 'urlopen', _open)
  45. return captured
  46. # ---------- 复制槽 ----------
  47. def test_cdc_slot_inactive_alerts(monkeypatch):
  48. _stub_pg(monkeypatch, [('cdc_large', False, 0)])
  49. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  50. ('槽 cdc_large', '空闲,采集作业可能已停')]
  51. def test_non_cdc_slot_inactive_is_ignored(monkeypatch):
  52. """PG→PG 那条槽 apply 卡住时会反复重连,active 抖动不算告警。"""
  53. _stub_pg(monkeypatch, [('finance_pg2pg', False, 0)])
  54. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  55. def test_slot_lag_over_threshold_alerts(monkeypatch):
  56. _stub_pg(monkeypatch, [('cdc_small', True, 6 * GB)])
  57. alerts = MON.check_slots('postgresql/x', LAG_CRIT)
  58. assert alerts == [('槽 cdc_small', '滞后 6.0 GB')]
  59. def test_slot_lag_at_threshold_not_alerted(monkeypatch):
  60. """阈值是 >,正好等于不报。"""
  61. _stub_pg(monkeypatch, [('cdc_small', True, LAG_CRIT)])
  62. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  63. def test_lag_applies_to_non_cdc_slot_too(monkeypatch):
  64. """active 只对 cdc_ 判,滞后对所有槽都判。"""
  65. _stub_pg(monkeypatch, [('finance_pg2pg', True, 9 * GB)])
  66. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  67. ('槽 finance_pg2pg', '滞后 9.0 GB')]
  68. def test_null_lag_does_not_crash(monkeypatch):
  69. """confirmed_flush_lsn 为 NULL 时 lag 是 None,当 0 处理。"""
  70. _stub_pg(monkeypatch, [('cdc_large', True, None)])
  71. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  72. def test_inactive_and_lagging_yields_two_alerts(monkeypatch):
  73. _stub_pg(monkeypatch, [('cdc_large', False, 7 * GB)])
  74. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  75. ('槽 cdc_large', '空闲,采集作业可能已停'),
  76. ('槽 cdc_large', '滞后 7.0 GB'),
  77. ]
  78. def test_connection_closed_after_query(monkeypatch):
  79. conn = _stub_pg(monkeypatch, [])
  80. MON.check_slots('postgresql/x', LAG_CRIT)
  81. conn.close.assert_called_once()
  82. def test_db_failure_becomes_alert(monkeypatch):
  83. def _boom(ds_ref):
  84. raise RuntimeError('数据源 ini 不存在')
  85. monkeypatch.setattr(MON.pgdb, 'connect', _boom)
  86. alerts = MON.check_slots('postgresql/x', LAG_CRIT)
  87. assert len(alerts) == 1
  88. assert alerts[0][0] == '主库'
  89. assert '数据源 ini 不存在' in alerts[0][1]
  90. # ---------- Flink 作业存活 ----------
  91. def test_all_jobs_running_no_alert(monkeypatch):
  92. captured = _stub_flink(monkeypatch, [
  93. {'name': 'st-cdc-large', 'state': 'RUNNING'},
  94. {'name': 'st-cdc-small', 'state': 'RUNNING'},
  95. ])
  96. assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == []
  97. assert captured['url'] == 'http://cdhmaster02:8081/jobs/overview'
  98. def test_missing_job_alerts(monkeypatch):
  99. _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'RUNNING'}])
  100. assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
  101. ('作业 st-cdc-small', '不在 Flink 集群上')]
  102. def test_restarting_job_alerts(monkeypatch):
  103. """2026-09-07 的故障形态:作业名一直在列表里,只判存在会漏报 7 小时。"""
  104. _stub_flink(monkeypatch, [
  105. {'name': 'st-cdc-large', 'state': 'RESTARTING'},
  106. {'name': 'st-cdc-small', 'state': 'RESTARTING'},
  107. ])
  108. assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
  109. ('作业 st-cdc-large', '状态 RESTARTING'),
  110. ('作业 st-cdc-small', '状态 RESTARTING'),
  111. ]
  112. def test_canceled_job_alerts(monkeypatch):
  113. """取消的作业仍留在 /jobs/overview 里。"""
  114. _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'CANCELED'}])
  115. assert MON.check_flink(JM, ['st-cdc-large']) == [
  116. ('作业 st-cdc-large', '状态 CANCELED')]
  117. def test_stale_record_does_not_mask_running(monkeypatch):
  118. """同名作业的历史记录不能盖掉在跑的那条,RUNNING 优先。"""
  119. _stub_flink(monkeypatch, [
  120. {'name': 'st-cdc-large', 'state': 'CANCELED'},
  121. {'name': 'st-cdc-large', 'state': 'RUNNING'},
  122. ])
  123. assert MON.check_flink(JM, ['st-cdc-large']) == []
  124. def test_flink_failure_becomes_single_alert(monkeypatch):
  125. def _boom(url, timeout=None):
  126. raise OSError('Connection refused')
  127. monkeypatch.setattr(MON.urllib.request, 'urlopen', _boom)
  128. alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
  129. assert len(alerts) == 1 # 查不到就是一条,不按作业数量刷屏
  130. assert alerts[0][0] == 'Flink 集群'
  131. assert 'Connection refused' in alerts[0][1]
  132. # ---------- 渲染:走一圈全部六类告警 ----------
  133. FIXED_NOW = datetime(2026, 8, 26, 15, 21, 27)
  134. def test_render_single_alert():
  135. assert MON.render([('槽 cdc_large', '空闲,采集作业可能已停')], now=FIXED_NOW) == (
  136. '### <font color="warning">CDC 监控告警</font>\n'
  137. '> 检测时间:2026-08-26 15:21:27\n'
  138. '> 槽 cdc_large:<font color="warning">空闲,采集作业可能已停</font>')
  139. def test_render_defaults_to_current_time(monkeypatch):
  140. """不传 now 时取当前时间,格式 yyyy-MM-dd HH:mm:ss。"""
  141. out = MON.render([('主库', '连不上')])
  142. assert re.match(r'^> 检测时间:\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}$', out.split('\n')[1])
  143. def test_render_all_alert_kinds():
  144. alerts = [
  145. ('槽 cdc_large', '空闲,采集作业可能已停'),
  146. ('槽 cdc_small', '滞后 6.0 GB'),
  147. ('作业 st-cdc-large', '不在 Flink 集群上'),
  148. ('作业 st-cdc-small', '状态 RESTARTING'),
  149. ('主库', '复制槽查询失败:connection refused'),
  150. ('Flink 集群', '查询失败:timeout'),
  151. ]
  152. out = MON.render(alerts, now=FIXED_NOW)
  153. lines = out.split('\n')
  154. assert lines[0] == '### <font color="warning">CDC 监控告警</font>'
  155. assert lines[1] == '> 检测时间:2026-08-26 15:21:27'
  156. assert len(lines) == 8
  157. for (key, desc), line in zip(alerts, lines[2:]):
  158. assert line == '> {0}:<font color="warning">{1}</font>'.format(key, desc)
  159. # ---------- main 串联 ----------
  160. def _run_main(monkeypatch, slot_alerts, flink_alerts):
  161. """跑 main,返回替身 alerter 实例供断言。"""
  162. monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
  163. monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: list(slot_alerts))
  164. monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: list(flink_alerts))
  165. instance = MagicMock()
  166. monkeypatch.setattr(MON, 'Alerter', MagicMock(return_value=instance))
  167. MON.main()
  168. return instance
  169. def test_main_sends_nothing_when_healthy(monkeypatch):
  170. instance = _run_main(monkeypatch, [], [])
  171. instance.send_markdown.assert_not_called()
  172. def test_main_sends_once_with_all_alerts(monkeypatch):
  173. instance = _run_main(
  174. monkeypatch,
  175. [('槽 cdc_large', '空闲,采集作业可能已停')],
  176. [('作业 st-cdc-small', '状态 RESTARTING')])
  177. instance.send_markdown.assert_called_once()
  178. content = instance.send_markdown.call_args[0][0]
  179. assert '槽 cdc_large' in content
  180. assert '作业 st-cdc-small' in content
  181. def test_main_builds_alerter_before_checks(monkeypatch):
  182. """健康时也要构造 Alerter:配置坏了要立刻暴露,不能等真出事才发现发不出去。"""
  183. monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
  184. monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
  185. monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: [])
  186. fake_cls = MagicMock()
  187. monkeypatch.setattr(MON, 'Alerter', fake_cls)
  188. MON.main()
  189. fake_cls.assert_called_once_with(channel='default')
  190. def test_main_defaults_cover_both_jobs(monkeypatch):
  191. """默认值必须带上小表组,否则它挂了不报警。"""
  192. seen = {}
  193. monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
  194. monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
  195. monkeypatch.setattr(MON, 'check_flink',
  196. lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
  197. monkeypatch.setattr(MON, 'Alerter', MagicMock())
  198. MON.main()
  199. assert seen['jm'] == 'cdhmaster02:8081'
  200. assert seen['jobs'] == ['st-cdc-large', 'st-cdc-small']
  201. def test_main_passes_cli_overrides(monkeypatch):
  202. seen = {}
  203. monkeypatch.setattr(sys, 'argv', [
  204. 'cdc-monitor.py', '-ds', 'postgresql/other', '-channel', 'realtime',
  205. '-jm', 'other-host:9081', '-jobs', 'a, b ,', '-lag-gb', '2'])
  206. monkeypatch.setattr(MON, 'check_slots',
  207. lambda ds, lag: seen.update(ds=ds, lag=lag) or [])
  208. monkeypatch.setattr(MON, 'check_flink',
  209. lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
  210. fake_cls = MagicMock()
  211. monkeypatch.setattr(MON, 'Alerter', fake_cls)
  212. MON.main()
  213. assert seen['ds'] == 'postgresql/other'
  214. assert seen['lag'] == 2 * GB
  215. assert seen['jm'] == 'other-host:9081'
  216. assert seen['jobs'] == ['a', 'b'] # 逗号分隔去空白、丢空项
  217. fake_cls.assert_called_once_with(channel='realtime')