feat: capture orchestration with include/exclude and delete-downgrade
Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
42
src/sync/capture.py
Normal file
42
src/sync/capture.py
Normal file
@@ -0,0 +1,42 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
import json, logging
|
||||||
|
from .access_reader import AccessReader
|
||||||
|
from .sql_writer import SqlWriter, QueueRow
|
||||||
|
from .config import FileMapping, SyncConfig
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
def capture_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter, cfg: SyncConfig) -> int:
|
||||||
|
exclude = set(fm.exclude_tables or [])
|
||||||
|
include = set(fm.include_tables) if fm.include_tables else None
|
||||||
|
rows = reader.read_log(cfg.runtime.capture_batch_size)
|
||||||
|
n = 0
|
||||||
|
for lr in rows:
|
||||||
|
if lr.table_name in exclude:
|
||||||
|
continue
|
||||||
|
if include is not None and lr.table_name not in include:
|
||||||
|
continue
|
||||||
|
op = lr.operate_type
|
||||||
|
row_data = None
|
||||||
|
if op in ("Insert", "Update"):
|
||||||
|
d = reader.read_row(lr.table_name, lr.record_id)
|
||||||
|
if d is None:
|
||||||
|
op = "Delete" # 行已删,降级
|
||||||
|
else:
|
||||||
|
row_data = json.dumps(d, ensure_ascii=False)
|
||||||
|
elif op != "Delete":
|
||||||
|
log.warning("unknown OperateType %r in %s log %s", op, fm.file, lr.id)
|
||||||
|
continue
|
||||||
|
qr = QueueRow(
|
||||||
|
source_file=fm.file,
|
||||||
|
source_table=lr.table_name,
|
||||||
|
record_id=lr.record_id,
|
||||||
|
target_schema=fm.schema,
|
||||||
|
target_table=fm.target_table(lr.table_name),
|
||||||
|
source_log_id=lr.id,
|
||||||
|
operate_type=op,
|
||||||
|
row_data=row_data,
|
||||||
|
)
|
||||||
|
writer.insert_queue_row(qr)
|
||||||
|
n += 1
|
||||||
|
return n
|
||||||
61
tests/test_capture.py
Normal file
61
tests/test_capture.py
Normal file
@@ -0,0 +1,61 @@
|
|||||||
|
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():
|
||||||
|
return SyncConfig(sql_server=SqlServerConfig(conn_str="x"),
|
||||||
|
access=AccessConfig(driver="d", roots={"2026":"r"}),
|
||||||
|
runtime=RuntimeConfig(),
|
||||||
|
files=[])
|
||||||
|
|
||||||
|
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", None)]
|
||||||
|
reader.read_row.return_value = {"ID": 34041, "订单号": "X1"}
|
||||||
|
writer = MagicMock()
|
||||||
|
n = capture_file(fm, reader, writer, cfg)
|
||||||
|
assert n == 1
|
||||||
|
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_update_missing_row_downgrades_to_delete():
|
||||||
|
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(11, "T", "5", "Update", None)]
|
||||||
|
reader.read_row.return_value = None # 行已删
|
||||||
|
writer = MagicMock()
|
||||||
|
n = capture_file(fm, reader, writer, cfg)
|
||||||
|
assert n == 1
|
||||||
|
args = writer.insert_queue_row.call_args[0][0]
|
||||||
|
assert args.operate_type == "Delete"
|
||||||
|
assert args.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", None),
|
||||||
|
LogRow(2, "氩弧焊每日催货落实记录_停", "1", "Insert", None)]
|
||||||
|
writer = MagicMock()
|
||||||
|
assert capture_file(fm, reader, writer, cfg) == 0
|
||||||
|
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", None),
|
||||||
|
LogRow(2, "其它表", "1", "Insert", None)]
|
||||||
|
reader.read_row.return_value = {"ID": 1}
|
||||||
|
writer = MagicMock()
|
||||||
|
assert capture_file(fm, reader, writer, cfg) == 1
|
||||||
|
assert writer.insert_queue_row.call_args[0][0].target_table == "检验合格记录表_YEAR2026"
|
||||||
Reference in New Issue
Block a user