Files
ProductionDataBaseSync_Data…/tests/test_capture.py
Misaka_Company 7cef5153e6 fix(capture): 读不到源行时按日志年龄延迟重试,根治降级Delete致数据丢失
Insert/Update 日志回读不到源行时,旧逻辑无条件降级为 Delete,在 ACE
引擎可见性延迟(批量插入约11s窗口)下误伤,导致该 Insert 的行不仅没进
SQL 反而被空打 Delete 并清理 Access 日志,数据永久丢失(8/4 事件根因:
氩弧焊.接收 16255-16259 缺失)。

改为按 Access 日志年龄(capture_defer_seconds=60)决定动作:
- 年龄 < 阈值:DEFER,跳过不入队,Access 日志保留,下一轮 cycle 重读;
- 年龄 >= 阈值:AGED-OUT,入队直接标 dead 待人工(usp_SyncApply 只选
  pending,不会执行任何破坏性 SQL),归档表记 ProcessedOperateType=AgedOut。

零 schema/SQL 改动:cleanup 只删 Status=applied 的日志,DEFER 行不入队
则无 applied 行、AGED-OUT 标 dead 非 applied,两者 Access 日志均保留,
与既有 cleanup/apply/service 逻辑天然自洽。

详见 docs/plan-capture-defer-by-age.md 与 docs/incremental-sync-flow.md。
2026-08-05 11:00:58 +08:00

142 lines
5.8 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import datetime as _dt
from unittest.mock import MagicMock
from sync.config import FileMapping, SyncConfig, AccessConfig, RuntimeConfig, SqlServerConfig
from sync.access_reader import LogRow
from sync.capture import capture_file
def _cfg(defer=60):
return SyncConfig(sql_server=SqlServerConfig(conn_str="x"),
access=AccessConfig(driver="d", roots={"2026": "r"}),
runtime=RuntimeConfig(capture_defer_seconds=defer),
files=[])
def _fm():
return FileMapping(file="x.accdb", root="2026", schema="s",
year_suffix="_YEAR2026", exclude_tables=["TableChangeLog"])
def test_capture_insert_reads_row_and_queues():
cfg = _cfg()
fm = FileMapping(file="氩弧焊.accdb", root="2026", schema="TIGWelding",
year_suffix="_YEAR2026", exclude_tables=["TableChangeLog"])
reader = MagicMock()
reader.read_log.return_value = [LogRow(10, "表壳焊接记录", "34041", "Insert", _dt.datetime.now())]
reader.read_row.return_value = {"ID": 34041, "订单号": "X1"}
writer = MagicMock()
st = capture_file(fm, reader, writer, cfg)
assert st.enqueued == 1
assert st.deferred == 0 and st.aged_out == 0
args = writer.insert_queue_row.call_args[0][0]
assert args.target_schema == "TIGWelding"
assert args.target_table == "表壳焊接记录_YEAR2026"
assert args.operate_type == "Insert"
assert '"订单号": "X1"' in args.row_data
assert args.source_log_id == 10
def test_capture_unreadable_young_row_is_deferred():
# Insert 日志年龄 < capture_defer_seconds → 跳过,不入队、不写 archive
cfg = _cfg(defer=60)
reader = MagicMock()
reader.read_log.return_value = [LogRow(11, "T", "5", "Insert", _dt.datetime.now())]
reader.read_row.return_value = None # 回读不到
writer = MagicMock()
st = capture_file(_fm(), reader, writer, cfg)
assert st.deferred == 1
assert st.enqueued == 0 and st.aged_out == 0
writer.insert_queue_row.assert_not_called()
writer.insert_dead_row.assert_not_called()
writer.insert_archive_row.assert_not_called()
def test_capture_unreadable_young_update_also_deferred():
# Update 同样走 defer 路径(不降级 Delete
cfg = _cfg(defer=60)
reader = MagicMock()
reader.read_log.return_value = [LogRow(12, "T", "6", "Update", _dt.datetime.now())]
reader.read_row.return_value = None
writer = MagicMock()
st = capture_file(_fm(), reader, writer, cfg)
assert st.deferred == 1
writer.insert_queue_row.assert_not_called()
writer.insert_dead_row.assert_not_called()
def test_capture_unreadable_aged_row_marked_dead():
# Insert 日志年龄 >= capture_defer_seconds → 入队标 dead不执行破坏性 SQL
cfg = _cfg(defer=60)
old_time = _dt.datetime.now() - _dt.timedelta(seconds=200)
reader = MagicMock()
reader.read_log.return_value = [LogRow(13, "T", "7", "Insert", old_time)]
reader.read_row.return_value = None # 仍读不到
writer = MagicMock()
st = capture_file(_fm(), reader, writer, cfg)
assert st.aged_out == 1
assert st.enqueued == 0 and st.deferred == 0
# archive 留痕ProcessedOperateType=AgedOut
arch = writer.insert_archive_row.call_args[0][0]
assert arch.original_operate_type == "Insert"
assert arch.processed_operate_type == "AgedOut"
assert arch.row_data is None
# 入队标 deadOperateType 保留原始 Insert
qr, err = writer.insert_dead_row.call_args[0]
assert qr.operate_type == "Insert" # 不降级
assert qr.row_data is None
assert "200s" in err
def test_capture_defer_zero_disables_retry():
# capture_defer_seconds=0 → 任何读不到都立即判死(边界)
cfg = _cfg(defer=0)
reader = MagicMock()
reader.read_log.return_value = [LogRow(14, "T", "8", "Insert", _dt.datetime.now())]
reader.read_row.return_value = None
writer = MagicMock()
st = capture_file(_fm(), reader, writer, cfg)
assert st.aged_out == 1
assert st.deferred == 0
def test_capture_delete_never_reads_row():
# Delete 日志无需回读,直接入队
cfg = _cfg()
reader = MagicMock()
reader.read_log.return_value = [LogRow(15, "T", "9", "Delete", _dt.datetime.now())]
writer = MagicMock()
st = capture_file(_fm(), reader, writer, cfg)
assert st.enqueued == 1
reader.read_row.assert_not_called()
qr = writer.insert_queue_row.call_args[0][0]
assert qr.operate_type == "Delete"
assert qr.row_data is None
def test_capture_skips_excluded_tables():
cfg = _cfg()
fm = FileMapping(file="x.accdb", root="2026", schema="s", year_suffix="_YEAR2026",
exclude_tables=["TableChangeLog", "氩弧焊每日催货落实记录_停"])
reader = MagicMock()
reader.read_log.return_value = [LogRow(1, "TableChangeLog", "1", "Insert", _dt.datetime.now()),
LogRow(2, "氩弧焊每日催货落实记录_停", "1", "Insert", _dt.datetime.now())]
writer = MagicMock()
st = capture_file(fm, reader, writer, cfg)
assert st.enqueued == 0 and st.out_of_scope == 2
writer.insert_queue_row.assert_not_called()
def test_capture_include_tables_filter():
cfg = _cfg()
fm = FileMapping(file="x.accdb", root="2026", schema="inspectionRecords", year_suffix="_YEAR2026",
exclude_tables=["TableChangeLog"], include_tables=["检验合格记录表"])
reader = MagicMock()
reader.read_log.return_value = [LogRow(1, "检验合格记录表", "1", "Insert", _dt.datetime.now()),
LogRow(2, "其它表", "1", "Insert", _dt.datetime.now())]
reader.read_row.return_value = {"ID": 1}
writer = MagicMock()
st = capture_file(fm, reader, writer, cfg)
assert st.enqueued == 1
assert writer.insert_queue_row.call_args[0][0].target_table == "检验合格记录表_YEAR2026"