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。
142 lines
5.8 KiB
Python
142 lines
5.8 KiB
Python
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
|
||
# 入队标 dead,OperateType 保留原始 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"
|