Răsfoiți Sursa

feat(monitor): cdc-monitor 加 -report 无条件推当前状态,两个 check 改返回事实+判定

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HAeFesmMXaZSrNTdFXE8Hu
tianyu.chu 1 săptămână în urmă
părinte
comite
d2992f8113
2 a modificat fișierele cu 219 adăugiri și 52 ștergeri
  1. 73 16
      bin/cdc-monitor.py
  2. 146 36
      tests/unit/alerter/test_cdc_monitor.py

+ 73 - 16
bin/cdc-monitor.py

@@ -3,7 +3,10 @@
 """
 CDC 链路监控:源库复制槽 + Flink 采集作业存活。
 
-DS 每 10 分钟调一次,无状态:每次独立判断。连续收到就是还没恢复,不再收到就是已恢复。
+两种模式,同一份检查逻辑与配置:
+  - 缺省(DS 每 10 分钟调一次):只在检测到异常时推企微。无状态,每次独立判断,
+    连续收到就是还没恢复,不再收到就是已恢复。
+  - -report(DS 手动触发,上线核对用):无条件推当前状态,正常也发。
 
 检查两类:
   1. 复制槽——源库**必须是主库**,从库上 pg_current_wal_lsn() 报
@@ -14,6 +17,9 @@ DS 每 10 分钟调一次,无状态:每次独立判断。连续收到就是
 
 连库失败 / Flink 查询失败本身也计入告警,不静默跳过——查不到状态和状态异常同样需要人介入。
 
+两个 check 都返回 (facts, alerts):facts 是当前值,供 -report 无条件列出;
+alerts 是判定结果,供缺省模式决定发不发。判定口径两种模式共用一份。
+
 Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留到真出事才发现告警链路是哑的。
 
 退出码:检测到异常仍返回 0(DS 任务不标红,异常靠企微通知);
@@ -22,7 +28,7 @@ Alerter 在开跑前构造:配置缺失或 key 未替换立刻暴露,不留
 CLI:
   python3 bin/cdc-monitor.py [-ds postgresql/prd-poyee-aliyun-cdc] [-channel default]
                              [-jm cdhmaster02:8081] [-jobs st-cdc-large,st-cdc-small]
-                             [-lag-gb 5]
+                             [-lag-gb 5] [-report]
 
 参数:
   -ds       主库 datasource ref,解析项目同级 ../datasource/{ds}.ini
@@ -30,6 +36,7 @@ CLI:
   -jm       Flink JobManager REST 地址
   -jobs     CDC 作业名,逗号分隔,逐个在 Flink 作业列表里找
   -lag-gb   槽滞后告警阈值(GB),正常应为 0
+  -report   无条件推当前状态,正常也发;不带则只在异常时告警
 """
 import argparse
 import json
@@ -53,16 +60,22 @@ _FLINK_TIMEOUT_SEC = 20
 
 
 def check_slots(ds_ref, lag_crit_bytes):
-    """查复制槽,返回 [(项, 说明)]。连库失败本身也算一项。"""
+    """查复制槽,返回 (facts, alerts)。
+
+    facts: [(槽名, active, 滞后字节)],连库失败时为空。
+    alerts: [(项, 说明)],连库失败本身也算一项。
+    """
+    facts = []
     alerts = []
     conn = None
     try:
         conn = pgdb.connect(ds_ref)
         for slot_name, active, lag_bytes in pgdb.query(conn, SLOT_SQL):
+            lag = int(lag_bytes or 0)
+            facts.append((slot_name, active, lag))
             # active 只对 cdc_ 槽判:PG→PG 那条 apply 卡住时会反复重连,会误报
             if slot_name.startswith('cdc_') and not active:
                 alerts.append(('槽 ' + slot_name, '空闲,采集作业可能已停'))
-            lag = int(lag_bytes or 0)
             if lag > lag_crit_bytes:
                 alerts.append(('槽 ' + slot_name, '滞后 %.1f GB' % (lag / 1024.0 ** 3)))
     except Exception as e:
@@ -70,11 +83,14 @@ def check_slots(ds_ref, lag_crit_bytes):
     finally:
         if conn is not None:
             conn.close()
-    return alerts
+    return facts, alerts
 
 
 def check_flink(jm, jobs):
-    """查 Flink session 集群的作业状态,返回 [(项, 说明)]。
+    """查 Flink session 集群的作业状态,返回 (facts, alerts)。
+
+    facts: [(作业名, state)],按 -jobs 顺序;集群上没有的 state 为 None。
+    只报 -jobs 点名的作业——这是 CDC 链路状态,不是集群状态。
 
     必须按 state 过滤:作业取消/失败后仍留在 /jobs/overview 列表里,
     只判名字存在会漏报——重启死循环时状态是 RESTARTING,名字一直在。
@@ -84,7 +100,7 @@ def check_flink(jm, jobs):
                                     timeout=_FLINK_TIMEOUT_SEC) as resp:
             overview = json.loads(resp.read().decode('utf-8'))
     except Exception as e:
-        return [('Flink 集群', '查询失败:%s' % e)]
+        return [], [('Flink 集群', '查询失败:%s' % e)]
 
     # 同名作业会留下历史记录,RUNNING 优先
     states = {}
@@ -92,25 +108,57 @@ def check_flink(jm, jobs):
         if j['name'] not in states or j['state'] == 'RUNNING':
             states[j['name']] = j['state']
 
+    facts = []
     alerts = []
     for job in jobs:
         state = states.get(job)
+        facts.append((job, state))
         if state is None:
             alerts.append(('作业 ' + job, '不在 Flink 集群上'))
         elif state != 'RUNNING':
             alerts.append(('作业 ' + job, '状态 %s' % state))
-    return alerts
+    return facts, alerts
+
+
+def _line(key, value, is_bad):
+    """一行正文:字段名默认色,只给值上色(对齐 kb/41 §3.5)。"""
+    return '> %s:<font color="%s">%s</font>' % (key, 'warning' if is_bad else 'info', value)
+
+
+def _head(title, color, now):
+    """标题 + 检测时间。时间用服务器本地时区(与 DS 告警消息同格式,同群里看着一致)。"""
+    return ['### <font color="%s">%s</font>' % (color, title),
+            '> 检测时间:%s' % (now or datetime.now()).strftime('%Y-%m-%d %H:%M:%S')]
 
 
 def render(alerts, now=None):
-    """拼企微 markdown。
+    """拼告警 markdown:只列异常项。now 留给单测注入固定值。"""
+    lines = _head('CDC 监控告警', 'warning', now)
+    lines.extend(_line(k, v, True) for k, v in alerts)
+    return '\n'.join(lines)
 
-    检测时间用服务器本地时区(与 DS 告警消息同格式,同群里看着一致)。
-    now 留给单测注入固定值,默认取当前时间。
+
+def render_status(slot_facts, flink_facts, alerts, now=None):
+    """拼状态报告 markdown:无条件列出每一项当前值。
+
+    异常项的值上 warning,正常项上 info;标题按有无异常切色。
+    查询失败时没有事实可列,把失败原因单独补在末尾,否则报告会是空的。
     """
-    lines = ['### <font color="warning">CDC 监控告警</font>',
-             '> 检测时间:%s' % (now or datetime.now()).strftime('%Y-%m-%d %H:%M:%S')]
-    lines.extend('> %s:<font color="warning">%s</font>' % (k, v) for k, v in alerts)
+    bad = set(k for k, _ in alerts)
+    lines = _head('CDC 链路状态', 'warning' if alerts else 'info', now)
+
+    fact_keys = set()
+    for slot_name, active, lag in slot_facts:
+        key = '槽 ' + slot_name
+        fact_keys.add(key)
+        lines.append(_line(key, '%s,滞后 %.1f GB' % (
+            '活跃' if active else '空闲', lag / 1024.0 ** 3), key in bad))
+    for job, state in flink_facts:
+        key = '作业 ' + job
+        fact_keys.add(key)
+        lines.append(_line(key, state or '不在集群上', key in bad))
+
+    lines.extend(_line(k, v, True) for k, v in alerts if k not in fact_keys)
     return '\n'.join(lines)
 
 
@@ -126,16 +174,25 @@ def main():
                         help='CDC 作业名,逗号分隔(默认 st-cdc-large,st-cdc-small)')
     parser.add_argument('-lag-gb', type=float, default=5, dest='lag_gb',
                         help='槽滞后告警阈值 GB(默认 5,正常应为 0)')
+    parser.add_argument('-report', action='store_true',
+                        help='无条件推当前状态,正常也发;不带则只在异常时告警')
     args = parser.parse_args()
 
     alerter = Alerter(channel=args.channel)
     jobs = [j.strip() for j in args.jobs.split(',') if j.strip()]
 
-    alerts = check_slots(args.ds, int(args.lag_gb * 1024 ** 3)) + check_flink(args.jm, jobs)
+    slot_facts, slot_alerts = check_slots(args.ds, int(args.lag_gb * 1024 ** 3))
+    flink_facts, flink_alerts = check_flink(args.jm, jobs)
+    alerts = slot_alerts + flink_alerts
 
     for key, desc in alerts:
         print('%s: %s' % (key, desc))  # 进 DS 任务日志
-    if alerts:
+
+    if args.report:
+        content = render_status(slot_facts, flink_facts, alerts)
+        print(content)                 # 发了什么,DS 任务日志里留一份
+        alerter.send_markdown(content)
+    elif alerts:
         alerter.send_markdown(render(alerts))
 
 

+ 146 - 36
tests/unit/alerter/test_cdc_monitor.py

@@ -1,6 +1,6 @@
 # -*- coding:utf-8 -*-
 """
-bin/cdc-monitor.py 单测:复制槽判定 / Flink 作业存活 / 消息渲染 / main 串联。
+bin/cdc-monitor.py 单测:复制槽判定 / Flink 作业存活 / 告警与状态渲染 / main 串联。
 
 不连真 PG(pgdb.connect + query 走替身)、不打真 JobManager(urlopen 走替身)、
 不发真消息(Alerter 走替身)。
@@ -43,6 +43,12 @@ def _stub_pg(monkeypatch, rows, conn=None):
     return conn
 
 
+def _slot_alerts(monkeypatch, rows):
+    """只取判定结果,facts 另有用例覆盖。"""
+    _stub_pg(monkeypatch, rows)
+    return MON.check_slots('postgresql/x', LAG_CRIT)[1]
+
+
 def _stub_flink(monkeypatch, jobs):
     """替身 JobManager:/jobs/overview 返回给定作业列表。返回捕获的请求参数。"""
     captured = {}
@@ -59,48 +65,41 @@ def _stub_flink(monkeypatch, jobs):
     return captured
 
 
-# ---------- 复制槽 ----------
+# ---------- 复制槽:判定 ----------
 
 def test_cdc_slot_inactive_alerts(monkeypatch):
-    _stub_pg(monkeypatch, [('cdc_large', False, 0)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == [
+    assert _slot_alerts(monkeypatch, [('cdc_large', False, 0)]) == [
         ('槽 cdc_large', '空闲,采集作业可能已停')]
 
 
 def test_non_cdc_slot_inactive_is_ignored(monkeypatch):
     """PG→PG 那条槽 apply 卡住时会反复重连,active 抖动不算告警。"""
-    _stub_pg(monkeypatch, [('finance_pg2pg', False, 0)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == []
+    assert _slot_alerts(monkeypatch, [('finance_pg2pg', False, 0)]) == []
 
 
 def test_slot_lag_over_threshold_alerts(monkeypatch):
-    _stub_pg(monkeypatch, [('cdc_small', True, 6 * GB)])
-    alerts = MON.check_slots('postgresql/x', LAG_CRIT)
-    assert alerts == [('槽 cdc_small', '滞后 6.0 GB')]
+    assert _slot_alerts(monkeypatch, [('cdc_small', True, 6 * GB)]) == [
+        ('槽 cdc_small', '滞后 6.0 GB')]
 
 
 def test_slot_lag_at_threshold_not_alerted(monkeypatch):
     """阈值是 >,正好等于不报。"""
-    _stub_pg(monkeypatch, [('cdc_small', True, LAG_CRIT)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == []
+    assert _slot_alerts(monkeypatch, [('cdc_small', True, LAG_CRIT)]) == []
 
 
 def test_lag_applies_to_non_cdc_slot_too(monkeypatch):
     """active 只对 cdc_ 判,滞后对所有槽都判。"""
-    _stub_pg(monkeypatch, [('finance_pg2pg', True, 9 * GB)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == [
+    assert _slot_alerts(monkeypatch, [('finance_pg2pg', True, 9 * GB)]) == [
         ('槽 finance_pg2pg', '滞后 9.0 GB')]
 
 
 def test_null_lag_does_not_crash(monkeypatch):
     """confirmed_flush_lsn 为 NULL 时 lag 是 None,当 0 处理。"""
-    _stub_pg(monkeypatch, [('cdc_large', True, None)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == []
+    assert _slot_alerts(monkeypatch, [('cdc_large', True, None)]) == []
 
 
 def test_inactive_and_lagging_yields_two_alerts(monkeypatch):
-    _stub_pg(monkeypatch, [('cdc_large', False, 7 * GB)])
-    assert MON.check_slots('postgresql/x', LAG_CRIT) == [
+    assert _slot_alerts(monkeypatch, [('cdc_large', False, 7 * GB)]) == [
         ('槽 cdc_large', '空闲,采集作业可能已停'),
         ('槽 cdc_large', '滞后 7.0 GB'),
     ]
@@ -116,12 +115,29 @@ def test_db_failure_becomes_alert(monkeypatch):
     def _boom(ds_ref):
         raise RuntimeError('数据源 ini 不存在')
     monkeypatch.setattr(MON.pgdb, 'connect', _boom)
-    alerts = MON.check_slots('postgresql/x', LAG_CRIT)
+    facts, alerts = MON.check_slots('postgresql/x', LAG_CRIT)
+    assert facts == []
     assert len(alerts) == 1
     assert alerts[0][0] == '主库'
     assert '数据源 ini 不存在' in alerts[0][1]
 
 
+# ---------- 复制槽:事实 ----------
+
+def test_slot_facts_cover_every_slot(monkeypatch):
+    """健康的槽也要出现在 facts 里,-report 才有东西可列。"""
+    _stub_pg(monkeypatch, [('cdc_large', True, 0),
+                           ('finance_pg2pg', True, 3 * GB)])
+    facts, alerts = MON.check_slots('postgresql/x', LAG_CRIT)
+    assert facts == [('cdc_large', True, 0), ('finance_pg2pg', True, 3 * GB)]
+    assert alerts == []
+
+
+def test_slot_facts_normalize_null_lag(monkeypatch):
+    _stub_pg(monkeypatch, [('cdc_large', True, None)])
+    assert MON.check_slots('postgresql/x', LAG_CRIT)[0] == [('cdc_large', True, 0)]
+
+
 # ---------- Flink 作业存活 ----------
 
 def test_all_jobs_running_no_alert(monkeypatch):
@@ -129,14 +145,17 @@ def test_all_jobs_running_no_alert(monkeypatch):
         {'name': 'st-cdc-large', 'state': 'RUNNING'},
         {'name': 'st-cdc-small', 'state': 'RUNNING'},
     ])
-    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == []
+    facts, alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
+    assert alerts == []
+    assert facts == [('st-cdc-large', 'RUNNING'), ('st-cdc-small', 'RUNNING')]
     assert captured['url'] == 'http://cdhmaster02:8081/jobs/overview'
 
 
 def test_missing_job_alerts(monkeypatch):
     _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'RUNNING'}])
-    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
-        ('作业 st-cdc-small', '不在 Flink 集群上')]
+    facts, alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
+    assert alerts == [('作业 st-cdc-small', '不在 Flink 集群上')]
+    assert facts == [('st-cdc-large', 'RUNNING'), ('st-cdc-small', None)]
 
 
 def test_restarting_job_alerts(monkeypatch):
@@ -145,7 +164,7 @@ def test_restarting_job_alerts(monkeypatch):
         {'name': 'st-cdc-large', 'state': 'RESTARTING'},
         {'name': 'st-cdc-small', 'state': 'RESTARTING'},
     ])
-    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small']) == [
+    assert MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])[1] == [
         ('作业 st-cdc-large', '状态 RESTARTING'),
         ('作业 st-cdc-small', '状态 RESTARTING'),
     ]
@@ -154,7 +173,7 @@ def test_restarting_job_alerts(monkeypatch):
 def test_canceled_job_alerts(monkeypatch):
     """取消的作业仍留在 /jobs/overview 里。"""
     _stub_flink(monkeypatch, [{'name': 'st-cdc-large', 'state': 'CANCELED'}])
-    assert MON.check_flink(JM, ['st-cdc-large']) == [
+    assert MON.check_flink(JM, ['st-cdc-large'])[1] == [
         ('作业 st-cdc-large', '状态 CANCELED')]
 
 
@@ -164,20 +183,30 @@ def test_stale_record_does_not_mask_running(monkeypatch):
         {'name': 'st-cdc-large', 'state': 'CANCELED'},
         {'name': 'st-cdc-large', 'state': 'RUNNING'},
     ])
-    assert MON.check_flink(JM, ['st-cdc-large']) == []
+    assert MON.check_flink(JM, ['st-cdc-large']) == ([('st-cdc-large', 'RUNNING')], [])
+
+
+def test_untracked_job_not_reported(monkeypatch):
+    """集群上的其他作业不进 facts——这是 CDC 链路状态,不是集群状态。"""
+    _stub_flink(monkeypatch, [
+        {'name': 'st-finaltest', 'state': 'CANCELED'},
+        {'name': 'st-cdc-large', 'state': 'RUNNING'},
+    ])
+    assert MON.check_flink(JM, ['st-cdc-large']) == ([('st-cdc-large', 'RUNNING')], [])
 
 
 def test_flink_failure_becomes_single_alert(monkeypatch):
     def _boom(url, timeout=None):
         raise OSError('Connection refused')
     monkeypatch.setattr(MON.urllib.request, 'urlopen', _boom)
-    alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
+    facts, alerts = MON.check_flink(JM, ['st-cdc-large', 'st-cdc-small'])
+    assert facts == []
     assert len(alerts) == 1          # 查不到就是一条,不按作业数量刷屏
     assert alerts[0][0] == 'Flink 集群'
     assert 'Connection refused' in alerts[0][1]
 
 
-# ---------- 渲染:走一圈全部六类告警 ----------
+# ---------- 告警渲染:走一圈全部六类告警 ----------
 
 FIXED_NOW = datetime(2026, 8, 26, 15, 21, 27)
 
@@ -213,13 +242,73 @@ def test_render_all_alert_kinds():
         assert line == '> {0}:<font color="warning">{1}</font>'.format(key, desc)
 
 
+# ---------- 状态报告渲染 ----------
+
+HEALTHY_SLOTS = [('cdc_large', True, 0), ('cdc_small', True, 0),
+                 ('finance_pg2pg', True, 107374182)]
+HEALTHY_JOBS = [('st-cdc-large', 'RUNNING'), ('st-cdc-small', 'RUNNING')]
+
+
+def test_status_all_healthy():
+    assert MON.render_status(HEALTHY_SLOTS, HEALTHY_JOBS, [], now=FIXED_NOW) == (
+        '### <font color="info">CDC 链路状态</font>\n'
+        '> 检测时间:2026-08-26 15:21:27\n'
+        '> 槽 cdc_large:<font color="info">活跃,滞后 0.0 GB</font>\n'
+        '> 槽 cdc_small:<font color="info">活跃,滞后 0.0 GB</font>\n'
+        '> 槽 finance_pg2pg:<font color="info">活跃,滞后 0.1 GB</font>\n'
+        '> 作业 st-cdc-large:<font color="info">RUNNING</font>\n'
+        '> 作业 st-cdc-small:<font color="info">RUNNING</font>')
+
+
+def test_status_marks_only_bad_items():
+    """异常项转 warning,其余保持 info,标题跟着转 warning。"""
+    slots = [('cdc_large', False, 0), ('cdc_small', True, 0)]
+    jobs = [('st-cdc-large', 'RESTARTING'), ('st-cdc-small', 'RUNNING')]
+    alerts = [('槽 cdc_large', '空闲,采集作业可能已停'),
+              ('作业 st-cdc-large', '状态 RESTARTING')]
+    lines = MON.render_status(slots, jobs, alerts, now=FIXED_NOW).split('\n')
+    assert lines[0] == '### <font color="warning">CDC 链路状态</font>'
+    assert lines[2] == '> 槽 cdc_large:<font color="warning">空闲,滞后 0.0 GB</font>'
+    assert lines[3] == '> 槽 cdc_small:<font color="info">活跃,滞后 0.0 GB</font>'
+    assert lines[4] == '> 作业 st-cdc-large:<font color="warning">RESTARTING</font>'
+    assert lines[5] == '> 作业 st-cdc-small:<font color="info">RUNNING</font>'
+
+
+def test_status_shows_missing_job():
+    lines = MON.render_status([], [('st-cdc-small', None)],
+                              [('作业 st-cdc-small', '不在 Flink 集群上')],
+                              now=FIXED_NOW).split('\n')
+    assert lines[2] == '> 作业 st-cdc-small:<font color="warning">不在集群上</font>'
+
+
+def test_status_falls_back_to_failure_reason():
+    """查询失败时没有事实可列,报告不能是空的。"""
+    alerts = [('主库', '复制槽查询失败:connection refused'),
+              ('Flink 集群', '查询失败:timeout')]
+    lines = MON.render_status([], [], alerts, now=FIXED_NOW).split('\n')
+    assert lines[0] == '### <font color="warning">CDC 链路状态</font>'
+    assert lines[2] == '> 主库:<font color="warning">复制槽查询失败:connection refused</font>'
+    assert lines[3] == '> Flink 集群:<font color="warning">查询失败:timeout</font>'
+    assert len(lines) == 4
+
+
+def test_status_does_not_repeat_items_already_listed():
+    """槽/作业的告警已经体现在它自己那行上,末尾不再重复一遍。"""
+    slots = [('cdc_large', False, 0)]
+    alerts = [('槽 cdc_large', '空闲,采集作业可能已停')]
+    assert len(MON.render_status(slots, [], alerts, now=FIXED_NOW).split('\n')) == 3
+
+
 # ---------- main 串联 ----------
 
-def _run_main(monkeypatch, slot_alerts, flink_alerts):
+def _run_main(monkeypatch, slot_alerts, flink_alerts, argv=None,
+              slot_facts=(), flink_facts=()):
     """跑 main,返回替身 alerter 实例供断言。"""
-    monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
-    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: list(slot_alerts))
-    monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: list(flink_alerts))
+    monkeypatch.setattr(sys, 'argv', argv or ['cdc-monitor.py'])
+    monkeypatch.setattr(MON, 'check_slots',
+                        lambda ds, lag: (list(slot_facts), list(slot_alerts)))
+    monkeypatch.setattr(MON, 'check_flink',
+                        lambda jm, jobs: (list(flink_facts), list(flink_alerts)))
     instance = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', MagicMock(return_value=instance))
     MON.main()
@@ -242,11 +331,32 @@ def test_main_sends_once_with_all_alerts(monkeypatch):
     assert '作业 st-cdc-small' in content
 
 
+def test_main_report_sends_even_when_healthy(monkeypatch):
+    """-report 是上线核对用的,正常也要发。"""
+    instance = _run_main(monkeypatch, [], [],
+                         argv=['cdc-monitor.py', '-report'],
+                         slot_facts=HEALTHY_SLOTS, flink_facts=HEALTHY_JOBS)
+    instance.send_markdown.assert_called_once()
+    content = instance.send_markdown.call_args[0][0]
+    assert content.startswith('### <font color="info">CDC 链路状态</font>')
+    assert '> 作业 st-cdc-large:<font color="info">RUNNING</font>' in content
+
+
+def test_main_report_shows_abnormal_state(monkeypatch):
+    """-report 遇到异常发的仍是状态报告,不是告警消息。"""
+    instance = _run_main(monkeypatch, [], [('作业 st-cdc-large', '状态 RESTARTING')],
+                         argv=['cdc-monitor.py', '-report'],
+                         flink_facts=[('st-cdc-large', 'RESTARTING')])
+    content = instance.send_markdown.call_args[0][0]
+    assert content.startswith('### <font color="warning">CDC 链路状态</font>')
+    assert '> 作业 st-cdc-large:<font color="warning">RESTARTING</font>' in content
+
+
 def test_main_builds_alerter_before_checks(monkeypatch):
     """健康时也要构造 Alerter:配置坏了要立刻暴露,不能等真出事才发现发不出去。"""
     monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
-    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
-    monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: [])
+    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: ([], []))
+    monkeypatch.setattr(MON, 'check_flink', lambda jm, jobs: ([], []))
     fake_cls = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', fake_cls)
     MON.main()
@@ -257,9 +367,9 @@ def test_main_defaults_cover_both_jobs(monkeypatch):
     """默认值必须带上小表组,否则它挂了不报警。"""
     seen = {}
     monkeypatch.setattr(sys, 'argv', ['cdc-monitor.py'])
-    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: [])
+    monkeypatch.setattr(MON, 'check_slots', lambda ds, lag: ([], []))
     monkeypatch.setattr(MON, 'check_flink',
-                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
+                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or ([], []))
     monkeypatch.setattr(MON, 'Alerter', MagicMock())
 
     MON.main()
@@ -274,9 +384,9 @@ def test_main_passes_cli_overrides(monkeypatch):
         'cdc-monitor.py', '-ds', 'postgresql/other', '-channel', 'realtime',
         '-jm', 'other-host:9081', '-jobs', 'a, b ,', '-lag-gb', '2'])
     monkeypatch.setattr(MON, 'check_slots',
-                        lambda ds, lag: seen.update(ds=ds, lag=lag) or [])
+                        lambda ds, lag: seen.update(ds=ds, lag=lag) or ([], []))
     monkeypatch.setattr(MON, 'check_flink',
-                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or [])
+                        lambda jm, jobs: seen.update(jm=jm, jobs=jobs) or ([], []))
     fake_cls = MagicMock()
     monkeypatch.setattr(MON, 'Alerter', fake_cls)