test_cdc_monitor.py 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. # -*- coding:utf-8 -*-
  2. """
  3. bin/cdc-monitor.py 单测:复制槽判定 / YARN 存活 / 消息渲染 / main 串联。
  4. 不连真 PG(pgdb.connect + query 走替身)、不跑真 yarn(subprocess 走替身)、
  5. 不发真消息(Alerter 走替身)。
  6. 脚本路径含连字符,用 importlib.util 动态加载为模块。
  7. """
  8. import importlib.util
  9. import os
  10. import re
  11. import sys
  12. from datetime import datetime
  13. from unittest.mock import MagicMock
  14. import pytest
  15. PROJECT_ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
  16. SCRIPT_PATH = os.path.join(PROJECT_ROOT, 'bin', 'cdc-monitor.py')
  17. GB = 1024 ** 3
  18. LAG_CRIT = 5 * GB
  19. def _load_script():
  20. spec = importlib.util.spec_from_file_location('cdc_monitor', SCRIPT_PATH)
  21. mod = importlib.util.module_from_spec(spec)
  22. sys.modules['cdc_monitor'] = mod
  23. spec.loader.exec_module(mod)
  24. return mod
  25. MON = _load_script()
  26. def _stub_pg(monkeypatch, rows, conn=None):
  27. """替身 PG:connect 返回 conn,query 返回 rows。返回该 conn 供断言 close。"""
  28. conn = conn or MagicMock()
  29. monkeypatch.setattr(MON.pgdb, 'connect', lambda ds_ref: conn)
  30. monkeypatch.setattr(MON.pgdb, 'query', lambda c, sql: rows)
  31. return conn
  32. # ---------- 复制槽 ----------
  33. def test_cdc_slot_inactive_alerts(monkeypatch):
  34. _stub_pg(monkeypatch, [('cdc_large', False, 0)])
  35. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  36. ('槽 cdc_large', '空闲,采集作业可能已停')]
  37. def test_non_cdc_slot_inactive_is_ignored(monkeypatch):
  38. """PG→PG 那条槽 apply 卡住时会反复重连,active 抖动不算告警。"""
  39. _stub_pg(monkeypatch, [('finance_pg2pg', False, 0)])
  40. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  41. def test_slot_lag_over_threshold_alerts(monkeypatch):
  42. _stub_pg(monkeypatch, [('cdc_small', True, 6 * GB)])
  43. alerts = MON.check_slots('postgresql/x', LAG_CRIT)
  44. assert alerts == [('槽 cdc_small', '滞后 6.0 GB')]
  45. def test_slot_lag_at_threshold_not_alerted(monkeypatch):
  46. """阈值是 >,正好等于不报。"""
  47. _stub_pg(monkeypatch, [('cdc_small', True, LAG_CRIT)])
  48. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  49. def test_lag_applies_to_non_cdc_slot_too(monkeypatch):
  50. """active 只对 cdc_ 判,滞后对所有槽都判。"""
  51. _stub_pg(monkeypatch, [('finance_pg2pg', True, 9 * GB)])
  52. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  53. ('槽 finance_pg2pg', '滞后 9.0 GB')]
  54. def test_null_lag_does_not_crash(monkeypatch):
  55. """confirmed_flush_lsn 为 NULL 时 lag 是 None,当 0 处理。"""
  56. _stub_pg(monkeypatch, [('cdc_large', True, None)])
  57. assert MON.check_slots('postgresql/x', LAG_CRIT) == []
  58. def test_inactive_and_lagging_yields_two_alerts(monkeypatch):
  59. _stub_pg(monkeypatch, [('cdc_large', False, 7 * GB)])
  60. assert MON.check_slots('postgresql/x', LAG_CRIT) == [
  61. ('槽 cdc_large', '空闲,采集作业可能已停'),
  62. ('槽 cdc_large', '滞后 7.0 GB'),
  63. ]
  64. def test_connection_closed_after_query(monkeypatch):
  65. conn = _stub_pg(monkeypatch, [])
  66. MON.check_slots('postgresql/x', LAG_CRIT)
  67. conn.close.assert_called_once()
  68. def test_db_failure_becomes_alert(monkeypatch):
  69. def _boom(ds_ref):
  70. raise RuntimeError('数据源 ini 不存在')
  71. monkeypatch.setattr(MON.pgdb, 'connect', _boom)
  72. alerts = MON.check_slots('postgresql/x', LAG_CRIT)
  73. assert len(alerts) == 1
  74. assert alerts[0][0] == '主库'
  75. assert '数据源 ini 不存在' in alerts[0][1]
  76. # ---------- YARN ----------
  77. def test_all_jobs_running_no_alert(monkeypatch):
  78. monkeypatch.setattr(MON.subprocess, 'check_output',
  79. lambda *a, **kw: b'app-1 st-cdc-large RUNNING\napp-2 st-cdc-small RUNNING')
  80. assert MON.check_yarn(['st-cdc-large', 'st-cdc-small']) == []
  81. def test_missing_job_alerts(monkeypatch):
  82. monkeypatch.setattr(MON.subprocess, 'check_output',
  83. lambda *a, **kw: b'app-1 st-cdc-large RUNNING')
  84. assert MON.check_yarn(['st-cdc-large', 'st-cdc-small']) == [
  85. ('作业 st-cdc-small', '不在 YARN RUNNING 列表')]
  86. def test_yarn_failure_becomes_single_alert(monkeypatch):
  87. def _boom(*a, **kw):
  88. raise OSError('yarn: command not found')
  89. monkeypatch.setattr(MON.subprocess, 'check_output', _boom)
  90. alerts = MON.check_yarn(['st-cdc-large', 'st-cdc-small'])
  91. assert len(alerts) == 1 # 查不到就是一条,不按作业数量刷屏
  92. assert alerts[0][0] == 'YARN'
  93. assert 'command not found' in alerts[0][1]
  94. # ---------- 渲染:走一圈全部五类告警 ----------
  95. FIXED_NOW = datetime(2026, 8, 26, 15, 21, 27)
  96. def test_render_single_alert():
  97. assert MON.render([('槽 cdc_large', '空闲,采集作业可能已停')], now=FIXED_NOW) == (
  98. '### <font color="warning">CDC 监控告警</font>\n'
  99. '> 检测时间:2026-08-26 15:21:27\n'
  100. '> 槽 cdc_large:<font color="warning">空闲,采集作业可能已停</font>')
  101. def test_render_defaults_to_current_time(monkeypatch):
  102. """不传 now 时取当前时间,格式 yyyy-MM-dd HH:mm:ss。"""
  103. out = MON.render([('主库', '连不上')])
  104. assert re.match(r'^> 检测时间:\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}$', out.split('\n')[1])
  105. def test_render_all_alert_kinds():
  106. alerts = [
  107. ('槽 cdc_large', '空闲,采集作业可能已停'),
  108. ('槽 cdc_small', '滞后 6.0 GB'),
  109. ('作业 st-cdc-large', '不在 YARN RUNNING 列表'),
  110. ('主库', '复制槽查询失败:connection refused'),
  111. ('YARN', '查询失败:timeout'),
  112. ]
  113. out = MON.render(alerts, now=FIXED_NOW)
  114. lines = out.split('\n')
  115. assert lines[0] == '### <font color="warning">CDC 监控告警</font>'
  116. assert lines[1] == '> 检测时间:2026-08-26 15:21:27'
  117. assert len(lines) == 7
  118. for (key, desc), line in zip(alerts, lines[2:]):
  119. assert line == '> {0}:<font color="warning">{1}</font>'.format(key, desc)
  120. # ---------- main 串联 ----------
  121. def _run_main(monkeypatch, slot_alerts, yarn_alerts):
  122. """跑 main,返回替身 alerter 实例供断言。"""
  123. monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
  124. monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: list(slot_alerts))
  125. monkeypatch.setattr(MON, 'check_yarn', lambda jobs: list(yarn_alerts))
  126. instance = MagicMock()
  127. monkeypatch.setattr(MON, 'Alerter', MagicMock(return_value=instance))
  128. MON.main()
  129. return instance
  130. def test_main_sends_nothing_when_healthy(monkeypatch):
  131. instance = _run_main(monkeypatch, [], [])
  132. instance.send_markdown.assert_not_called()
  133. def test_main_sends_once_with_all_alerts(monkeypatch):
  134. instance = _run_main(
  135. monkeypatch,
  136. [('槽 cdc_large', '空闲,采集作业可能已停')],
  137. [('作业 st-cdc-small', '不在 YARN RUNNING 列表')])
  138. instance.send_markdown.assert_called_once()
  139. content = instance.send_markdown.call_args[0][0]
  140. assert '槽 cdc_large' in content
  141. assert '作业 st-cdc-small' in content
  142. def test_main_builds_alerter_before_checks(monkeypatch):
  143. """健康时也要构造 Alerter:配置坏了要立刻暴露,不能等真出事才发现发不出去。"""
  144. monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
  145. monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
  146. monkeypatch.setattr(MON, 'check_yarn', lambda jobs: [])
  147. fake_cls = MagicMock()
  148. monkeypatch.setattr(MON, 'Alerter', fake_cls)
  149. MON.main()
  150. fake_cls.assert_called_once_with(channel='default')
  151. def test_main_passes_cli_overrides(monkeypatch):
  152. seen = {}
  153. monkeypatch.setattr(sys, 'argv', [
  154. 'cdc-monitor.py', '-ds', 'postgresql/other', '-channel', 'realtime',
  155. '-jobs', 'a, b ,', '-lag-gb', '2'])
  156. monkeypatch.setattr(MON, 'check_slots',
  157. lambda ds, lag: seen.update(ds=ds, lag=lag) or [])
  158. monkeypatch.setattr(MON, 'check_yarn', lambda jobs: seen.update(jobs=jobs) or [])
  159. fake_cls = MagicMock()
  160. monkeypatch.setattr(MON, 'Alerter', fake_cls)
  161. MON.main()
  162. assert seen['ds'] == 'postgresql/other'
  163. assert seen['lag'] == 2 * GB
  164. assert seen['jobs'] == ['a', 'b'] # 逗号分隔去空白、丢空项
  165. fake_cls.assert_called_once_with(channel='realtime')