Compare commits

...

9 Commits

Author SHA1 Message Date
Misaka_Company
2ce03d21d9 Refine incremental sync audit observability
- Add per-cycle correlation ID ([cyc:xxxxxxxx]) threaded through Python
  logs and SQL audit tables for end-to-end traceability of any divergence.
- New ProductionDataBaseSync.SyncApplyRunLog table + @CycleID on usp_SyncApply
  for per-table apply auditing (pending/merged/deleted/applied, dead/error
  counts, duration). Audit write isolated in its own TRY/CATCH outside txn.
- Capture/cleanup phases now emit full detail: Insert->Delete downgrade
  warning with record_id/log_id, dedup_skipped as apply-stall signal,
  per-file summary, and access log-id ranges on cleanup.
- Queue health check surfaces error/dead rows with recent samples instead of
  silent accumulation (previously the top cause of data divergence).
- sql_writer uses INSERT...SELECT...WHERE NOT EXISTS for observable dedup;
  idle cycles lowered to DEBUG with periodic heartbeat.
- Backward compatible: old proc callers still work (CycleID nullable; legacy
  coarse-grained logging with one-time notice).

Excluded from this commit: CODE.md, build_code_doc.py (doc generation).
2026-07-31 17:10:23 +08:00
Misaka_Company
75fb6a3a01 feat(compare): 核对结果写入 logs/ 并为 114 增加开机计划任务
- compare.py: 新增 write_report(),每次运行把报告以日志形式写到
  logs/compare_<granularity>_<日期>.log(含生成时间 + 汇总头部),
  --report 仍兼容作为额外输出路径
- main.py: compare 子命令改用 write_report,默认落 logs/,stdout 仍打印
- run_compare_ids.cmd: 114 主机开机计划任务包装脚本(ID 级核对)

Co-Authored-By: WorkBuddy <workbuddy@tencent.com>
2026-07-17 11:06:21 +08:00
Misaka_Company
d71b7eab62 Migrate sync objects into ProductionDataBaseSync schema + add permanent audit archive
- Add sql/00_schema.sql: create dedicated ProductionDataBaseSync schema (idempotent)
- Move SyncQueue and usp_SyncApply from dbo into ProductionDataBaseSync
- Add sql/03_sync_log_archive.sql: permanent, append-only SyncLogArchive that
  records both OriginalOperateType and ProcessedOperateType plus the Access log
  OriginalTime, so pipeline divergences (e.g. Insert applied as Delete) stay
  reconstructible forever (SyncQueue is transient and only keeps processed type)
- config.py: inject sync_queue_table / archive_table / apply_proc (default to the
  new schema); SqlWriter takes these names instead of hardcoding dbo
- sql_writer.py: add ArchiveRow + insert_archive_row (dedup on source keys),
  parametrize queue/archive/proc names throughout
- capture.py: archive every consumed log row before enqueue (preserves evidence
  before cleanup deletes the Access log)
- service.py: pass the three names into SqlWriter
- tests: read queue/proc names from config instead of hardcoding dbo.SyncQueue
2026-07-16 12:33:01 +08:00
Misaka_Company
4179d3e232 feat(logging): archive historical logs daily into Archive/
At startup, previously produced logs are relocated into an Archive/
subfolder next to the active log: the project sync.log gets a
-YYYY-MM-DD suffix when archived, and NSSM's nssm_*.log captures are
moved as-is. The log root then only shows the current day's sync.log.
Idempotent handler setup is preserved.

Co-Authored-By: WorkBuddy <workbuddy@tencent.com>
2026-07-16 09:24:41 +08:00
Misaka_Company
d8a423c983 fix: make setup_logging idempotent to avoid duplicate log lines
setup_logging unconditionally addHandler'd on every call. Under the new
main.py entry the function was invoked twice for --loop mode (main.py and
service.run both called it), stacking two file + two console handlers so
every log line was written twice.

Clear any pre-existing root handlers before re-adding so repeated calls
always yield exactly one file handler + one console handler, regardless of
the caller. Preserves the standalone python -m sync.service entry point.
2026-07-15 14:29:00 +08:00
Misaka_Company
7d11eddc7a 📝 docs: document unified CLI, incremental sync, and compare 2026-07-15 14:03:27 +08:00
Misaka_Company
1b44ca28e4 feat(sync): add unified main.py CLI dispatcher 2026-07-15 14:03:24 +08:00
Misaka_Company
e9f012eab5 feat(sync): add data consistency compare (count + ID-set) 2026-07-15 14:03:21 +08:00
Misaka_Company
5ca0715801 ♻️ refactor(sync): extract shared target-table resolution 2026-07-15 14:03:19 +08:00
25 changed files with 1816 additions and 157 deletions

189
README.md
View File

@@ -1,50 +1,199 @@
# ProductionDataBaseSync_DataMacro # ProductionDataBaseSync_DataMacro
Access → SQL Server 增量同步(数据宏驱动)。设计详见 `docs/superpowers/specs/2026-07-14-access-datamacro-sync-design.md` Access → SQL Server 单向增量同步(数据宏驱动)。把各 Access 库的业务数据周期性同步到 SQL Server 镜像表,作为 Access → SQL Server 迁移期的过渡数据层
设计详见 `docs/superpowers/specs/2026-07-14-access-datamacro-sync-design.md`
## 背景
生产数据实际承载在网络共享下的多个 Access `.accdb`(按车间/年份分库)。旧机制靠客户端前端 VBA 写变更日志,但 VBA 只在特定表单事件触发,批量改表、直接改表等路径会绕过 → 漏数据。
改用 Access **数据宏**(表级引擎触发器,任何写路径必触发):每张业务表挂 `After Insert/Update/Delete`,变更写入各库本地 `TableChangeLog`。本程序就是「读各库日志 → 增量同步到 SQL」的搬运器理论上 100% 捕获、对客户端零侵入。
## 架构 ## 架构
每个 Access 库通过数据宏把变更写入本地 `TableChangeLog`;同步服务周期性把这些变更搬运到 SQL Server 镜像表。单库一轮分三段: 每个 Access 库通过数据宏把变更写入本地 `TableChangeLog`;同步服务周期性把这些变更搬运到 SQL Server 镜像表。单库一轮分三段:
| 阶段 | 动作 | 说明 | | 阶段 | 动作 | 说明 |
| --- | --- | --- | | --- | --- | --- |
| Capture | `SELECT` 读取 `TableChangeLog` 最旧的 N 条 | 只读 Access | | Capture | `SELECT` 读取 `TableChangeLog` 最旧的 N 条I/U 按 `RecordID` 回读整行 | 只读 Access |
| Apply | 调用 `dbo.usp_SyncApply` 写入 SQL 镜像表 | 按 `ID` 精确落库 | | Apply | 调用 `dbo.usp_SyncApply` 写入 SQL 镜像表 | 按 `ID` 精确落库,保序「最后操作胜」 |
| Cleanup | `DELETE` 已应用的 `TableChangeLog` 行 | 按 `ID` 列表精确删除,遇锁自动退避重试 | | Cleanup | `DELETE` 已应用的 `TableChangeLog` 行 | 按 `ID` 列表精确删除,遇锁自动退避重试 |
## 命令 关键设计点:
- **无水位线表**——Access 日志「应用成功即删」,日志本身就是待处理队列;`dbo.SyncQueue` 的唯一索引 `(SourceFile, SourceTable, SourceLogID)` 兜底去重,重复捕获幂等。
- **保序「最后操作胜」**——同一 `RecordID` 多次操作(先 Insert 后 Delete 等)按日志顺序取最后一条,保证最终态与 Access 一致。
- **每表一个事务**——单表失败只回滚该表,失败行标 `error` 重试,超限标 `dead` 待人工。
- **`SyncQueue` 长期保留**作审计/重试日志;`applied` 行清理后标 `cleaned`,超保留期再 purge控制表增长。
| 用途 | 命令 | ## 环境要求
| --- | --- |
| 增量同步服务(常驻,由 nssm 托管 `DataMacroSync` | `.venv/Scripts/python.exe -m sync.service` |
| 一次性全量同步TRUNCATE + 全量 INSERT所有库所有表 | `.venv/Scripts/python.exe -m sync.fullsync config.yaml` |
> 上述命令均使用项目自带的 `.venv`。同步类命令由 nssm 以服务方式运行,无需手动设置环境变量 - **Python 3.10+**(实测 3.13)。代码用 `X | None` 等新语法
- **ODBC 驱动**(系统级,非 pip 安装,需预先装好):
- `Microsoft Access Driver (*.accdb, *.mdb)`ACE Redist 2016
- `ODBC Driver 17 for SQL Server`
- **SQL Server ≥ 2017**(存储过程用 `STRING_AGG ... WITHIN GROUP`)。
- 执行账号需对目标表有 `ALTER` 权限(`SET IDENTITY_INSERT` 要求)。
## 安装
```bash
python -m venv .venv
.venv/Scripts/python.exe -m pip install --upgrade pip
.venv/Scripts/python.exe -m pip install -r requirements.txt
```
依赖(`requirements.txt``pyodbc``PyYAML``pydantic``pytest`
## SQL 端部署
首次需在目标库建好暂存表与 apply 存储过程(两个脚本都幂等,可重复执行):
```bash
sqlcmd -S <SERVER>,1433 -U <USER> -P <PASSWORD> -d <DB> -C -N o -i sql/01_sync_queue.sql
sqlcmd -S <SERVER>,1433 -U <USER> -P <PASSWORD> -d <DB> -C -N o -i sql/02_sync_apply.sql
```
- `sql/01_sync_queue.sql`:建 `dbo.SyncQueue` + 去重/清理索引 + `CleanedAt` 列。
- `sql/02_sync_apply.sql``dbo.usp_SyncApply` 集合化 apply 存储过程。
> 连接串/凭据以 `config.yaml` 为准README 不硬编码。
## 配置 ## 配置
编辑 `config.yaml`(从 `config.example.yaml` 复制并填入真实凭据)。关键段: 编辑 `config.yaml`(从 `config.example.yaml` 复制并填入真实凭据;该文件 gitignored)。关键段:
- `sql_server`SQL Server 连接串与 `SyncQueue` 表名 - **`sql_server`**`conn_str`ODBC 连接串`sync_queue_table`(默认 `dbo.SyncQueue`
- `access`ACE ODBC 驱动名与各根目录(`roots`)映射 - **`access`**`driver`ACE 驱动名)与 `roots`(年份→根目录映射,如 `2026: "\\\\srv\\生产进度表\\2026年数据"`)
- `runtime`:轮询间隔、批大小、重试与保留策略 - **`runtime`**`poll_interval_seconds`(轮询间隔)、`capture_batch_size`/`apply_batch_size`/`cleanup_batch_size`(各段批大小)、`max_retries`/`retry_backoff_seconds`(重试)、`cleanup_lock_retries`Access 锁重试次数)、`cleaned_retention_hours``cleaned` 行保留多久后 purge
- `files`:每个 Access 文件一条映射`file` / `root` / `schema` / `year_suffix` / `exclude_tables` / `include_tables`)。 - **`files`**:每个 Access 文件一条映射
- `file` / `root`(对应 `access.roots` 的 key/ `schema`SQL 目标 schema
- `year_suffix`:拼到表名后(`2026年数据``_YEAR2026``2025年数据`/合同表用 `""`)。
- `exclude_tables` / `include_tables`:排除/包含规则,**exclude 优先于 include**。`TableChangeLog` 必须排除。
## 命令行
统一入口 `main.py`(仓库根目录),三个功能块都用它调用。配置固定读取同目录的 `config.yaml`,不在命令中指定:
```bash
.venv/Scripts/python.exe main.py fullsync [--db FILE] [--table NAME] [--clear-change-log]
.venv/Scripts/python.exe main.py incremental [--loop] [--poll-interval N]
.venv/Scripts/python.exe main.py compare [--granularity count|ids] [--db FILE] [--table NAME] [--report PATH]
```
| 子命令 | 说明 | 退出码 |
| --- | --- | --- |
| `fullsync` | 一次性全量同步TRUNCATE + 批量 INSERT绕开增量队列。 | 0 |
| `incremental` | 增量同步一轮capture→apply→cleanup`--loop` 切持续轮询(服务模式)。 | 0 |
| `compare` | 数据一致性核对:默认行数总量,`--granularity ids` 精确到 ID 集合差异。 | 全一致 0 / 有不一致 1 |
> 三个块共用同一份 `config.yaml`,目标表集合完全一致(由 `sync.targets` 统一解析)。`main.py` 在根目录、自行把 `src/` 加入 `sys.path`,无需 `-m`、无需设环境变量;控制台强制 UTF-8中文表名不乱码。
>
> 旧入口 `-m sync.service` / `-m sync.fullsync` 保留为兼容,行为不变。
## 全量同步 ## 全量同步
用于从零重建镜像表或修复 Access 与 SQL 之间的漂移。会**清空目标表再全量写入**,绕开增量队列: 用于从零重建镜像表或修复 Access 与 SQL 之间的漂移。会**清空目标表再全量写入**,绕开增量队列:
```bash ```bash
.venv/Scripts/python.exe -m sync.fullsync config.yaml # 全部库、全部表 .venv/Scripts/python.exe main.py fullsync # 全部库、全部表
.venv/Scripts/python.exe -m sync.fullsync config.yaml --db OEM.accdb # 仅单个库 .venv/Scripts/python.exe main.py fullsync --db OEM.accdb # 仅单个库
.venv/Scripts/python.exe -m sync.fullsync config.yaml --table 表壳焊接记录 # 仅单表(作用于所有库) .venv/Scripts/python.exe main.py fullsync --table 表壳焊接记录 # 仅单表(作用于所有库)
.venv/Scripts/python.exe -m sync.fullsync config.yaml --clear-change-log # 同步后同时清空 TableChangeLog谨慎 .venv/Scripts/python.exe main.py fullsync --clear-change-log # 同步后同时清空 TableChangeLog谨慎
``` ```
- `year_suffix` 通过 `FileMapping` 拼到表名后(如 `表壳焊接记录``表壳焊接记录_YEAR2026`)。 - `year_suffix` 通过 `FileMapping` 拼到表名后(如 `表壳焊接记录``表壳焊接记录_YEAR2026`)。
- 写入时 `SET IDENTITY_INSERT ON`,保留 Access 原 ID保证后续增量的 `RecordID` 匹配不错位。
- 无镜像表的目标表按设计跳过(`target table missing`),不报错。 - 无镜像表的目标表按设计跳过(`target table missing`),不报错。
## 增量同步
以 Access 数据宏日志为唯一变更源,每轮跑一遍 capture → apply → cleanup三段说明见上面「架构」。这是**主用模式**,生产上常驻运行。两种调用方式:
```bash
.venv/Scripts/python.exe main.py incremental # 跑一轮就退出(手动/按需补跑)
.venv/Scripts/python.exe main.py incremental --loop # 持续轮询(服务模式,不退出)
.venv/Scripts/python.exe main.py incremental --loop --poll-interval 30 # 覆盖 runtime.poll_interval_seconds
```
- 生产环境以 nssm 服务 `DataMacroSync` 常驻(即 `--loop` 模式见下文「NSSM 服务」;手动单轮适合验证或临时补跑积压。
- `--loop` 持续轮询直到进程被停(`nssm stop` 或 Ctrl+C默认单轮跑完即退出。
- 单轮一次最多处理每库 `capture_batch_size` 条日志;积压多时连续跑几轮或用 `--loop` 直到清空。
- 每个文件的 capture/cleanup 独立隔离单文件失败不影响其它apply 失败不阻塞 cleanup失败行 `error` 下轮重试、超 `max_retries``dead` 待人工。
- 幂等:`SyncQueue` 唯一索引去重,重复 capture、中断续跑都不会重写或漏写。
- 与全量同步共用同一份 `config.yaml``sync.targets`,目标表集合完全一致;全量是「从零重建」的补充手段,不替代增量。
## 数据对比
核对 Access 源表与 SQL 镜像表是否一致。两种粒度:
- **行数总量**(默认):逐表比对 `COUNT(*)`
- **ID 集合**`--granularity ids`):逐表比对两边 `ID` 集合报告「Access 有 / SQL 无」与「SQL 有 / Access 无」的 ID每表前 50 个 + 总数)。
```bash
.venv/Scripts/python.exe main.py compare # 全部库、全部表,行数总量
.venv/Scripts/python.exe main.py compare --db 氩弧焊.accdb # 仅单个库
.venv/Scripts/python.exe main.py compare --granularity ids --table 表壳焊接记录 # 单表 ID 级
.venv/Scripts/python.exe main.py compare --report report.txt # 同时写入报告文件UTF-8
```
- 无镜像表按设计跳过(`[SKIPPED no mirror]`),不计为不一致——这类表多半是该排除却没排除(如 `*_停` 停用表、`USysApplicationLog`),可作为配置清理的线索。
- 任一表不一致时退出码 `1`(便于脚本化);全部一致为 `0`
- 实时增量同步存在秒级延迟窗口,刚写入 Access 的行可能尚未到 SQL属正常稍后再核或对照 `SyncQueue` 的 pending 行)。
## NSSM 服务114
增量同步在 host 114 上以 nssm 服务 `DataMacroSync` 常驻运行。常用操作(经 `ssh 114`
```bash
ssh 114 "nssm status DataMacroSync" # 查状态SERVICE_RUNNING / SERVICE_STOPPED
ssh 114 "nssm stop DataMacroSync" # 停
ssh 114 "nssm start DataMacroSync" # 起
ssh 114 "nssm restart DataMacroSync" # 重启
ssh 114 "nssm list" # 列出所有 nssm 服务
```
## 测试
```bash
.venv/Scripts/python.exe -m pytest # 仅单元测试(默认)
RUN_INTEGRATION=1 .venv/Scripts/python.exe -m pytest # 含集成测试(需能连真实 Access + SQL Server
```
- 单元测试用 mock不依赖数据库集成测试`@pytest.mark.integration`)连 `config.yaml` 里的真实库,且自带清理。
- `pyproject.toml` 仅用于配置 pytest`pythonpath = ["src", "."]``testpaths``integration` 标记)。
## 项目结构
```
main.py 统一命令行入口fullsync / incremental / compare
config.yaml 真实配置gitignoredconfig.example.yaml 是模板
requirements.txt 依赖
pyproject.toml pytest 配置
sql/
01_sync_queue.sql dbo.SyncQueue 建表 + 索引(幂等)
02_sync_apply.sql dbo.usp_SyncApply 存储过程
src/sync/
config.py Pydantic 配置模型 + load_config
targets.py 共享目标表解析exclude/include全量/增量/对比共用)
serialize.py Access 值 → JSON 可序列化
access_reader.py 读 Access日志/整行/计数/ID/删除日志)
sql_writer.py 写 SQLSyncQueue/apply/计数/ID/全量灌表)
capture.py 增量编排:读日志→回读整行→入队
cleanup.py 清理编排:回删已应用日志
service.py 主循环 cycle() / run()
fullsync.py 一次性全量同步
compare.py 数据一致性对比count / ids
logging_setup.py 日志配置(滚动文件 + 控制台)
tests/ 单元 + 集成测试
docs/superpowers/ 设计文档与实现计划
```
## 常见问题 ## 常见问题
- **cleanup 报 `-1102 无法更新;当前被锁定`**Access 是文件型数据库cleanup 反写 `DELETE` 与生产客户端数据宏写日志争用页级锁。服务已对锁冲突自动退避重试(`access_reader.delete_log_ids` 捕获 `pyodbc.Error` 并判断 `-1102`/被锁定)。偶发属正常,持续刷错再排查。 - **cleanup 报 `-1102 无法更新;当前被锁定`**Access 是文件型数据库cleanup 反写 `DELETE` 与生产客户端数据宏写日志争用页级锁。服务已对锁冲突自动退避重试(`access_reader.delete_log_ids` 捕获 `pyodbc.Error` 并判断 `-1102`/被锁定)。偶发属正常,持续刷错再排查。
- **`No module named 'pydantic_core'` / `pyodbc`**venv 解释器与轮子 ABI 不匹配(常见于 Python 3.13 装到 cp310 轮子)。修复:`.venv/Scripts/python.exe -m pip install --force-reinstall --no-cache-dir pyodbc pydantic` - **`No module named 'pydantic_core'` / `pyodbc`**venv 解释器与轮子 ABI 不匹配(常见于 Python 3.13 装到 cp310 轮子)。修复:`.venv/Scripts/python.exe -m pip install --force-reinstall --no-cache-dir pyodbc pydantic`
- **compare/fullsync 报 `[SKIPPED no mirror]` / `target table missing`**:该表在 Access 里但 SQL 端没有镜像表(多为 `*_停` 停用表、`USysApplicationLog` 等系统表,或尚未建镜像的新表)。若是该停用的表,加进对应 `exclude_tables`;若该同步,先在 SQL 建表再 fullsync。
- **`SyncQueue` 出现 `error`/`dead` 行**`error` 会在下轮自动重试(未超 `max_retries``dead` 是超限放弃,需人工看 `ErrorMsg` 排查后处理。
- **`UserWarning: Field name "schema" ... shadows ... BaseModel`**`FileMapping.schema` 字段名与 Pydantic 基类属性重名,仅告警、不影响功能。
- **控制台中文乱码**`main.py` 已强制 stdout/stderr 为 UTF-8若仍乱码设环境变量 `PYTHONIOENCODING=utf-8`,或用 `compare --report` 输出 UTF-8 文件。

118
main.py Normal file
View File

@@ -0,0 +1,118 @@
"""Unified command-line entry point for the Access -> SQL Server sync toolkit.
Run from the project root (no ``-m`` needed)::
python main.py fullsync [--db FILE] [--table NAME] [--clear-change-log]
python main.py incremental [--loop] [--poll-interval N]
python main.py compare [--granularity count|ids] [--db FILE] [--table NAME] [--report PATH]
Configuration is hard-coded to ``config.yaml`` next to this script -- it is not
a command-line argument, so all three blocks always use the same config (and
therefore the same target tables).
This file lives at the repo root, outside the ``src/`` package, so it puts
``src`` on ``sys.path`` itself to import ``sync.*`` regardless of how Python
was launched or whether the venv already has ``src`` on its path.
"""
from __future__ import annotations
import argparse
import logging
import os
import sys
# Make the src/ package importable when running this root script directly.
sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), "src"))
from sync.config import load_config
from sync.logging_setup import setup_logging
from sync.fullsync import full_sync
from sync import service
from sync.compare import compare, any_mismatch, format_report, write_report
CONFIG_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "config.yaml")
log = logging.getLogger("main")
def _parse_args(argv):
p = argparse.ArgumentParser(
prog="python main.py",
description="Access -> SQL Server sync toolkit (fullsync / incremental / compare).",
)
sub = p.add_subparsers(dest="command", required=True)
pf = sub.add_parser("fullsync", help="one-shot TRUNCATE + bulk INSERT")
pf.add_argument("--db", help="limit to one Access file (by FileMapping.file)")
pf.add_argument("--table", help="limit to one table (applies to all matched files)")
pf.add_argument("--clear-change-log", action="store_true",
help="after loading, clear TableChangeLog on the synced files")
pi = sub.add_parser("incremental", help="capture -> apply -> cleanup")
pi.add_argument("--loop", action="store_true",
help="run continuously (service mode); default is a single pass")
pi.add_argument("--poll-interval", type=int, dest="poll_interval",
help="override runtime.poll_interval_seconds (with --loop)")
pc = sub.add_parser("compare", help="compare Access vs SQL Server data")
pc.add_argument("--granularity", choices=["count", "ids"], default="count",
help="count = row totals (default); ids = ID-set membership diff")
pc.add_argument("--db", help="limit to one Access file")
pc.add_argument("--table", help="limit to one table")
pc.add_argument("--report", help="write the report to this file as well as stdout")
return p.parse_args(argv)
def _force_utf8_console():
"""Render Chinese table names correctly on a Windows GBK console.
compare prints to stdout by default; without this the default console
codepage mojibakes non-ASCII. No-op when stdout is already UTF-8 or when it
does not support reconfigure (e.g. some test-capture streams).
"""
for stream in (sys.stdout, sys.stderr):
try:
stream.reconfigure(encoding="utf-8", errors="replace")
except (AttributeError, ValueError):
pass
def main(argv=None) -> int:
"""Parse argv, load config, dispatch to the chosen block. Returns exit code."""
_force_utf8_console()
args = _parse_args(argv)
cfg = load_config(CONFIG_PATH)
setup_logging(cfg.logging)
if args.command == "fullsync":
full_sync(cfg, db_filter=args.db, table_filter=args.table,
clear_change_log=args.clear_change_log)
return 0
if args.command == "incremental":
if args.poll_interval is not None:
cfg.runtime.poll_interval_seconds = args.poll_interval
if args.loop:
service.run(cfg)
else:
service.cycle(cfg)
return 0
if args.command == "compare":
results = compare(cfg, granularity=args.granularity,
db_filter=args.db, table_filter=args.table)
# Persist the report as a dated log file under the logging directory
# (logs/ by default); --report still allows an extra custom path.
log_dir = os.path.dirname((cfg.logging or {}).get("path", "sync.log")) or "."
written = write_report(results, args.granularity, log_dir=log_dir,
extra_path=args.report)
print(format_report(results, args.granularity))
log.info("compare finished (granularity=%s) report=%s mismatch=%s",
args.granularity, written, any_mismatch(results))
return 1 if any_mismatch(results) else 0
return 2 # unreachable: argparse requires a subcommand
if __name__ == "__main__":
sys.exit(main())

View File

@@ -1,4 +1,4 @@
[tool.pytest.ini_options] [tool.pytest.ini_options]
pythonpath = ["src"] pythonpath = ["src", "."]
testpaths = ["tests"] testpaths = ["tests"]
markers = ["integration: marks tests requiring real Access/SQL Server"] markers = ["integration: marks tests requiring real Access/SQL Server"]

3
run_compare_ids.cmd Normal file
View File

@@ -0,0 +1,3 @@
@echo off
cd /d C:\Users\peng\Projects\ProductionDataBaseSync_DataMacro
.venv\Scripts\python.exe main.py compare --granularity ids >> logs\compare_task_stdout.log 2>&1

12
sql/00_schema.sql Normal file
View File

@@ -0,0 +1,12 @@
-- ProductionDataBaseSync: dedicated schema that groups every object owned by
-- the Access -> SQL Server sync (SyncQueue staging, SyncLogArchive audit store,
-- and the usp_SyncApply procedure). Keeping them out of dbo makes ownership and
-- housekeeping explicit.
-- Idempotent: CREATE SCHEMA must be the only statement in its batch, so it is
-- wrapped in EXEC() behind an existence check.
IF NOT EXISTS (SELECT 1 FROM sys.schemas WHERE name = 'ProductionDataBaseSync')
BEGIN
EXEC('CREATE SCHEMA ProductionDataBaseSync');
END
GO

View File

@@ -1,9 +1,10 @@
-- SyncQueue: staging table for the Access -> SQL Server one-way sync. -- SyncQueue: staging table for the Access -> SQL Server one-way sync.
-- Lives under the ProductionDataBaseSync schema (run 00_schema.sql first).
-- Idempotent: safe to re-run (table created only if absent; indexes only if absent). -- Idempotent: safe to re-run (table created only if absent; indexes only if absent).
IF OBJECT_ID('dbo.SyncQueue', 'U') IS NULL IF OBJECT_ID('ProductionDataBaseSync.SyncQueue', 'U') IS NULL
BEGIN BEGIN
CREATE TABLE dbo.SyncQueue ( CREATE TABLE ProductionDataBaseSync.SyncQueue (
QueueID bigint IDENTITY(1,1) NOT NULL, QueueID bigint IDENTITY(1,1) NOT NULL,
SourceFile nvarchar(255) NOT NULL, SourceFile nvarchar(255) NOT NULL,
SourceTable nvarchar(255) NOT NULL, SourceTable nvarchar(255) NOT NULL,
@@ -18,6 +19,7 @@ BEGIN
ErrorMsg nvarchar(max) NULL, ErrorMsg nvarchar(max) NULL,
CapturedAt datetime2 NOT NULL CONSTRAINT DF_SyncQueue_Captured DEFAULT sysdatetime(), CapturedAt datetime2 NOT NULL CONSTRAINT DF_SyncQueue_Captured DEFAULT sysdatetime(),
AppliedAt datetime2 NULL, AppliedAt datetime2 NULL,
CleanedAt datetime2 NULL,
CONSTRAINT PK_SyncQueue PRIMARY KEY CLUSTERED (QueueID) CONSTRAINT PK_SyncQueue PRIMARY KEY CLUSTERED (QueueID)
); );
END END
@@ -25,38 +27,27 @@ GO
IF NOT EXISTS (SELECT 1 FROM sys.indexes IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'UX_SyncQueue_Dedup' WHERE name = 'UX_SyncQueue_Dedup'
AND object_id = OBJECT_ID('dbo.SyncQueue')) AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
BEGIN BEGIN
CREATE UNIQUE INDEX UX_SyncQueue_Dedup CREATE UNIQUE INDEX UX_SyncQueue_Dedup
ON dbo.SyncQueue(SourceFile, SourceTable, SourceLogID); ON ProductionDataBaseSync.SyncQueue(SourceFile, SourceTable, SourceLogID);
END END
GO GO
IF NOT EXISTS (SELECT 1 FROM sys.indexes IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncQueue_Pending' WHERE name = 'IX_SyncQueue_Pending'
AND object_id = OBJECT_ID('dbo.SyncQueue')) AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
BEGIN BEGIN
CREATE INDEX IX_SyncQueue_Pending CREATE INDEX IX_SyncQueue_Pending
ON dbo.SyncQueue(Status, TargetSchema, TargetTable); ON ProductionDataBaseSync.SyncQueue(Status, TargetSchema, TargetTable);
END
GO
-- Cleanup bookkeeping: track when a queue row's Access log counterpart has
-- been physically deleted, so cleanup never re-deletes the same IDs and the
-- table can be purged to bound its growth.
IF NOT EXISTS (SELECT 1 FROM sys.columns
WHERE object_id = OBJECT_ID('dbo.SyncQueue')
AND name = 'CleanedAt')
BEGIN
ALTER TABLE dbo.SyncQueue ADD CleanedAt datetime2 NULL;
END END
GO GO
IF NOT EXISTS (SELECT 1 FROM sys.indexes IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncQueue_Cleaned' WHERE name = 'IX_SyncQueue_Cleaned'
AND object_id = OBJECT_ID('dbo.SyncQueue')) AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncQueue'))
BEGIN BEGIN
CREATE INDEX IX_SyncQueue_Cleaned CREATE INDEX IX_SyncQueue_Cleaned
ON dbo.SyncQueue(Status, CleanedAt); ON ProductionDataBaseSync.SyncQueue(Status, CleanedAt);
END END
GO GO

View File

@@ -10,24 +10,42 @@
-- Delete), which let a stale Delete outrank a newer Insert for the same RecordID -- Delete), which let a stale Delete outrank a newer Insert for the same RecordID
-- and silently drop a row whose true last op was an Insert. Routing off the one -- and silently drop a row whose true last op was an Insert. Routing off the one
-- winning row's OperateType guarantees only the genuine last op wins. -- winning row's OperateType guarantees only the genuine last op wins.
--
-- AUDIT NOTE: when ProductionDataBaseSync.SyncApplyRunLog exists (see
-- sql/04_sync_apply_runlog.sql), one audit row is written per target table and
-- per invocation: pending/distinct counts going in, the MERGE/DELETE rowcounts,
-- how many queue rows were flipped to applied (or to error/dead on failure),
-- the error message, the duration, and the caller-supplied @CycleID that ties
-- the row to the Python service log's [cyc:...] lines. The audit insert is
-- best-effort (wrapped in its own TRY/CATCH, outside the data transaction) and
-- can never fail the apply itself. @CycleID defaults to NULL so the legacy
-- single-parameter EXEC keeps working.
CREATE OR ALTER PROCEDURE dbo.usp_SyncApply CREATE OR ALTER PROCEDURE ProductionDataBaseSync.usp_SyncApply
@MaxRetries INT = 5 @MaxRetries INT = 5,
@CycleID NVARCHAR(40) = NULL
AS AS
BEGIN BEGIN
SET NOCOUNT ON; SET NOCOUNT ON;
DECLARE @sch NVARCHAR(128), @tbl NVARCHAR(255), @FullName NVARCHAR(514); DECLARE @sch NVARCHAR(128), @tbl NVARCHAR(255), @FullName NVARCHAR(514);
DECLARE @sql NVARCHAR(MAX), @cols NVARCHAR(MAX), @upd NVARCHAR(MAX), @ins NVARCHAR(MAX); DECLARE @sql NVARCHAR(MAX), @cols NVARCHAR(MAX), @upd NVARCHAR(MAX), @ins NVARCHAR(MAX);
DECLARE @hasRunLog bit =
CASE WHEN OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog', 'U') IS NOT NULL
THEN 1 ELSE 0 END;
DECLARE @t0 DATETIME2, @pending INT, @distinctRecs INT,
@merged INT, @deleted INT, @applied INT,
@errCnt INT, @deadCnt INT,
@errMsg NVARCHAR(MAX), @outcome VARCHAR(10);
-- Re-queue rows still under the retry budget. -- Re-queue rows still under the retry budget.
UPDATE dbo.SyncQueue UPDATE ProductionDataBaseSync.SyncQueue
SET Status = 'pending' SET Status = 'pending'
WHERE Status = 'error' AND RetryCount < @MaxRetries; WHERE Status = 'error' AND RetryCount < @MaxRetries;
DECLARE cur CURSOR LOCAL FAST_FORWARD FOR DECLARE cur CURSOR LOCAL FAST_FORWARD FOR
SELECT DISTINCT TargetSchema, TargetTable SELECT DISTINCT TargetSchema, TargetTable
FROM dbo.SyncQueue FROM ProductionDataBaseSync.SyncQueue
WHERE Status = 'pending'; WHERE Status = 'pending';
OPEN cur; OPEN cur;
@@ -37,6 +55,19 @@ BEGIN
BEGIN BEGIN
SET @FullName = QUOTENAME(@sch) + N'.' + QUOTENAME(@tbl); SET @FullName = QUOTENAME(@sch) + N'.' + QUOTENAME(@tbl);
-- Per-table audit state. Explicit reset every iteration: local
-- variables keep their previous value across cursor loops.
SELECT @t0 = SYSDATETIME(),
@merged = NULL, @deleted = NULL, @applied = NULL,
@errCnt = NULL, @deadCnt = NULL,
@errMsg = NULL, @outcome = 'ok',
@cols = NULL, @upd = NULL, @ins = NULL;
SELECT @pending = COUNT(*),
@distinctRecs = COUNT(DISTINCT RecordID)
FROM ProductionDataBaseSync.SyncQueue
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
BEGIN TRY BEGIN TRY
BEGIN TRAN; BEGIN TRAN;
-- @cols: comma-quoted column names (for INSERT target list) -- @cols: comma-quoted column names (for INSERT target list)
@@ -66,11 +97,14 @@ BEGIN
-- pending ops so the rn=1 row is the true last op for that -- pending ops so the rn=1 row is the true last op for that
-- RecordID; an earlier Delete can no longer outrank a newer -- RecordID; an earlier Delete can no longer outrank a newer
-- Insert for the same RecordID. -- Insert for the same RecordID.
-- AUDIT: @@ROWCOUNT is captured into @MergedOut IMMEDIATELY
-- after the MERGE -- SET IDENTITY_INSERT (like any SET option)
-- resets @@ROWCOUNT, so the order below is load-bearing.
SET @sql = N'SET IDENTITY_INSERT ' + @FullName + N' ON;' SET @sql = N'SET IDENTITY_INSERT ' + @FullName + N' ON;'
+ N';WITH ranked AS (' + N';WITH ranked AS ('
+ N' SELECT RecordID, OperateType, RowData,' + N' SELECT RecordID, OperateType, RowData,'
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn' + N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
+ N' FROM dbo.SyncQueue' + N' FROM ProductionDataBaseSync.SyncQueue'
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending''' + N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
+ N')' + N')'
+ N'MERGE ' + @FullName + N' WITH (HOLDLOCK) AS tgt' + N'MERGE ' + @FullName + N' WITH (HOLDLOCK) AS tgt'
@@ -81,28 +115,33 @@ BEGIN
+ N' WHEN MATCHED THEN UPDATE SET ' + @upd + N' WHEN MATCHED THEN UPDATE SET ' + @upd
+ N' WHEN NOT MATCHED THEN INSERT (ID,' + @cols + N')' + N' WHEN NOT MATCHED THEN INSERT (ID,' + @cols + N')'
+ N' VALUES (TRY_CAST(src.RecordID AS int),' + @ins + N');' + N' VALUES (TRY_CAST(src.RecordID AS int),' + @ins + N');'
+ N'SET @MergedOut = @@ROWCOUNT;'
+ N'SET IDENTITY_INSERT ' + @FullName + N' OFF;'; + N'SET IDENTITY_INSERT ' + @FullName + N' OFF;';
EXEC sp_executesql @sql, EXEC sp_executesql @sql,
N'@sch NVARCHAR(128),@tbl NVARCHAR(255)', @sch, @tbl; N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@MergedOut INT OUTPUT',
@sch, @tbl, @MergedOut = @merged OUTPUT;
-- Delete: winners (rn=1) whose winning OperateType is Delete. -- Delete: winners (rn=1) whose winning OperateType is Delete.
-- Same ranked CTE — only the genuine last op can be a delete. -- Same ranked CTE — only the genuine last op can be a delete.
SET @sql = N';WITH ranked AS (' SET @sql = N';WITH ranked AS ('
+ N' SELECT RecordID, OperateType,' + N' SELECT RecordID, OperateType,'
+ N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn' + N' ROW_NUMBER() OVER (PARTITION BY RecordID ORDER BY SourceLogID DESC) rn'
+ N' FROM dbo.SyncQueue' + N' FROM ProductionDataBaseSync.SyncQueue'
+ N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending''' + N' WHERE TargetSchema=@sch AND TargetTable=@tbl AND Status=''pending'''
+ N')' + N')'
+ N'DELETE t FROM ' + @FullName + N' t' + N'DELETE t FROM ' + @FullName + N' t'
+ N' JOIN (SELECT RecordID FROM ranked WHERE rn=1 AND OperateType=''Delete'') d' + N' JOIN (SELECT RecordID FROM ranked WHERE rn=1 AND OperateType=''Delete'') d'
+ N' ON t.ID = TRY_CAST(d.RecordID AS int);'; + N' ON t.ID = TRY_CAST(d.RecordID AS int);'
+ N'SET @DeletedOut = @@ROWCOUNT;';
EXEC sp_executesql @sql, EXEC sp_executesql @sql,
N'@sch NVARCHAR(128),@tbl NVARCHAR(255)', @sch, @tbl; N'@sch NVARCHAR(128),@tbl NVARCHAR(255),@DeletedOut INT OUTPUT',
@sch, @tbl, @DeletedOut = @deleted OUTPUT;
END END
UPDATE dbo.SyncQueue UPDATE ProductionDataBaseSync.SyncQueue
SET Status = 'applied', AppliedAt = SYSDATETIME() SET Status = 'applied', AppliedAt = SYSDATETIME()
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending'; WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
SET @applied = @@ROWCOUNT;
COMMIT; COMMIT;
END TRY END TRY
@@ -113,13 +152,48 @@ BEGIN
-- issues its own rollback. We roll back here so the partial per-table -- issues its own rollback. We roll back here so the partial per-table
-- work is discarded before marking rows as 'error'/'dead'. -- work is discarded before marking rows as 'error'/'dead'.
IF @@TRANCOUNT > 0 ROLLBACK; IF @@TRANCOUNT > 0 ROLLBACK;
UPDATE dbo.SyncQueue SELECT @errMsg = ERROR_MESSAGE(), @outcome = 'error';
SET Status = CASE WHEN RetryCount + 1 >= @MaxRetries THEN 'dead' ELSE 'error' END,
RetryCount = RetryCount + 1, -- Split of the previous single CASE update so the audit row can
ErrorMsg = ERROR_MESSAGE() -- report exactly how many rows died vs. how many will be retried.
-- Net effect on SyncQueue is identical.
UPDATE ProductionDataBaseSync.SyncQueue
SET Status = 'dead', RetryCount = RetryCount + 1, ErrorMsg = @errMsg
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending'
AND RetryCount + 1 >= @MaxRetries;
SET @deadCnt = @@ROWCOUNT;
UPDATE ProductionDataBaseSync.SyncQueue
SET Status = 'error', RetryCount = RetryCount + 1, ErrorMsg = @errMsg
WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending'; WHERE TargetSchema = @sch AND TargetTable = @tbl AND Status = 'pending';
SET @errCnt = @@ROWCOUNT;
END CATCH END CATCH
-- Best-effort audit row: sits outside the data transaction (after
-- COMMIT/ROLLBACK) so it survives either outcome, and its own failure
-- can never break the apply loop.
IF @hasRunLog = 1
BEGIN
BEGIN TRY
INSERT ProductionDataBaseSync.SyncApplyRunLog
(CycleID, TargetSchema, TargetTable,
PendingCount, DistinctRecords,
MergedCount, DeletedCount, AppliedCount,
ErrorCount, DeadCount,
Outcome, ErrorMsg, StartedAt, DurationMs)
VALUES
(@CycleID, @sch, @tbl,
@pending, @distinctRecs,
@merged, @deleted, @applied,
@errCnt, @deadCnt,
@outcome, @errMsg, @t0,
DATEDIFF(millisecond, @t0, SYSDATETIME()));
END TRY
BEGIN CATCH
PRINT 'SyncApplyRunLog insert failed: ' + ERROR_MESSAGE();
END CATCH
END
FETCH NEXT FROM cur INTO @sch, @tbl; FETCH NEXT FROM cur INTO @sch, @tbl;
END END

View File

@@ -0,0 +1,71 @@
-- SyncLogArchive: permanent, append-only audit store for every Access change-log
-- row consumed by the incremental sync. This is the durable evidence layer that
-- SyncQueue is NOT: SyncQueue is a transient work queue (purged 24h after a row
-- is cleaned) and it only stores the *processed* OperateType, so a downgrade
-- (e.g. capture turning an Insert into a Delete when the row is momentarily
-- unreadable) erases the original intent. This table preserves both the ORIGINAL
-- operate type recorded by the Access data macro and the PROCESSED type actually
-- sent to SQL Server, plus the original Access log timestamp and the row payload.
--
-- With this in place, cases like the 14287 incident (a real Insert applied to SQL
-- as a Delete) stay fully reconstructible: OriginalOperateType != ProcessedOperateType
-- flags exactly where the pipeline diverged from the source.
--
-- Retention: PERMANENT. No purge job touches this table. SyncQueue keeps its
-- short-lived queue role; this table keeps history. Lives under the
-- ProductionDataBaseSync schema (run 00_schema.sql first).
-- Idempotent: safe to re-run.
IF OBJECT_ID('ProductionDataBaseSync.SyncLogArchive', 'U') IS NULL
BEGIN
CREATE TABLE ProductionDataBaseSync.SyncLogArchive (
ArchiveID bigint IDENTITY(1,1) NOT NULL,
SourceFile nvarchar(255) NOT NULL,
SourceTable nvarchar(255) NOT NULL,
SourceLogID bigint NOT NULL,
RecordID nvarchar(50) NOT NULL,
TargetSchema nvarchar(128) NOT NULL,
TargetTable nvarchar(255) NOT NULL,
OriginalOperateType varchar(10) NOT NULL, -- as recorded by the Access data macro
ProcessedOperateType varchar(10) NOT NULL, -- as actually sent to SyncQueue / SQL
RowData nvarchar(max) NULL, -- captured row payload (NULL for Delete)
OriginalTime datetime2 NULL, -- Access TableChangeLog.Time (previously discarded)
CapturedAt datetime2 NOT NULL
CONSTRAINT DF_SyncLogArchive_Captured DEFAULT sysdatetime(),
CONSTRAINT PK_SyncLogArchive PRIMARY KEY CLUSTERED (ArchiveID)
);
END
GO
-- Dedup: one archive row per source log entry. Capture may re-run the same log
-- row if a prior cycle's apply failed (the Access log is only deleted after a
-- successful apply), so the write path uses IF NOT EXISTS on these keys.
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'UX_SyncLogArchive_Dedup'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
BEGIN
CREATE UNIQUE INDEX UX_SyncLogArchive_Dedup
ON ProductionDataBaseSync.SyncLogArchive(SourceFile, SourceTable, SourceLogID);
END
GO
-- Evidence lookup by table + record (e.g. "show every log ever seen for ID 14287").
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncLogArchive_Record'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
BEGIN
CREATE INDEX IX_SyncLogArchive_Record
ON ProductionDataBaseSync.SyncLogArchive(SourceTable, RecordID);
END
GO
-- Fast filter for the anomaly the archive exists to catch: original != processed.
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncLogArchive_Divergence'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncLogArchive'))
BEGIN
CREATE INDEX IX_SyncLogArchive_Divergence
ON ProductionDataBaseSync.SyncLogArchive(OriginalOperateType, ProcessedOperateType)
INCLUDE (SourceFile, SourceTable, RecordID, CapturedAt);
END
GO

View File

@@ -0,0 +1,75 @@
-- SyncApplyRunLog: per-invocation, per-target-table audit of what usp_SyncApply
-- actually did. This closes the biggest observability gap in the pipeline: the
-- apply phase used to be a black box ("apply done") -- queue rows could flip
-- to 'error' or 'dead' with no trace in the service log, and there was no
-- record of how many rows a MERGE/DELETE touched at any point in time. With
-- this table, "what did apply do to table X around time T, and did it fail?"
-- is a single indexed query, and CycleID joins each row back to the exact
-- [cyc:xxxxxxxx] lines in the Python service log.
--
-- Written by usp_SyncApply (sql/02_sync_apply.sql) as a best-effort insert per
-- (invocation, target table). Idle cycles write nothing (the proc's cursor
-- only visits tables that have pending rows), so growth tracks real change
-- traffic, not poll frequency. Retention: unmanaged by default; if it ever
-- grows large, purge by StartedAt, e.g.
-- DELETE FROM ProductionDataBaseSync.SyncApplyRunLog
-- WHERE StartedAt < DATEADD(day, -90, SYSDATETIME());
--
-- Lives under the ProductionDataBaseSync schema (run 00_schema.sql first).
-- Idempotent: safe to re-run. Deploy alongside the updated 02_sync_apply.sql;
-- ordering is forgiving either way (the proc checks for this table and skips
-- the audit insert when it is absent).
IF OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog', 'U') IS NULL
BEGIN
CREATE TABLE ProductionDataBaseSync.SyncApplyRunLog (
RunLogID bigint IDENTITY(1,1) NOT NULL,
CycleID nvarchar(40) NULL, -- correlation id from the Python service ([cyc:...])
TargetSchema nvarchar(128) NOT NULL,
TargetTable nvarchar(255) NOT NULL,
PendingCount int NOT NULL, -- pending queue rows seen for this table
DistinctRecords int NULL, -- distinct RecordIDs among them
MergedCount int NULL, -- rows affected by the MERGE (insert + update)
DeletedCount int NULL, -- rows affected by the DELETE
AppliedCount int NULL, -- queue rows flipped to 'applied'
ErrorCount int NULL, -- queue rows flipped to 'error' (will retry)
DeadCount int NULL, -- queue rows flipped to 'dead' (retries exhausted)
Outcome varchar(10) NOT NULL, -- 'ok' | 'error'
ErrorMsg nvarchar(max) NULL,
StartedAt datetime2 NOT NULL
CONSTRAINT DF_SyncApplyRunLog_Started DEFAULT sysdatetime(),
DurationMs int NULL,
CONSTRAINT PK_SyncApplyRunLog PRIMARY KEY CLUSTERED (RunLogID)
);
END
GO
-- Join back to the Python service log of one cycle.
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncApplyRunLog_Cycle'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
BEGIN
CREATE INDEX IX_SyncApplyRunLog_Cycle
ON ProductionDataBaseSync.SyncApplyRunLog(CycleID);
END
GO
-- "What happened to this table around time T?"
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncApplyRunLog_Table'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
BEGIN
CREATE INDEX IX_SyncApplyRunLog_Table
ON ProductionDataBaseSync.SyncApplyRunLog(TargetSchema, TargetTable, StartedAt);
END
GO
-- Fast scan for failed applies.
IF NOT EXISTS (SELECT 1 FROM sys.indexes
WHERE name = 'IX_SyncApplyRunLog_Outcome'
AND object_id = OBJECT_ID('ProductionDataBaseSync.SyncApplyRunLog'))
BEGIN
CREATE INDEX IX_SyncApplyRunLog_Outcome
ON ProductionDataBaseSync.SyncApplyRunLog(Outcome, StartedAt);
END
GO

View File

@@ -94,7 +94,9 @@ class AccessReader:
``cursor.rowcount``), so the caller can report honest counts. A no-op ``cursor.rowcount``), so the caller can report honest counts. A no-op
(returns 0) when ``ids`` is empty. Retries with linear backoff because (returns 0) when ``ids`` is empty. Retries with linear backoff because
the live client may briefly hold a page lock on ``TableChangeLog``. the live client may briefly hold a page lock on ``TableChangeLog``.
Raises on final lock failure. Each retry is logged at WARNING (previously a silent sleep) and the
final lock failure at ERROR before raising, so lock churn on the
production files is visible in the audit trail.
""" """
if not ids: if not ids:
return 0 return 0
@@ -121,8 +123,20 @@ class AccessReader:
msg = str(e) msg = str(e)
is_lock = "被锁定" in msg or "-1102" in msg is_lock = "被锁定" in msg or "-1102" in msg
if is_lock and attempt < retries - 1: if is_lock and attempt < retries - 1:
log.warning(
"TableChangeLog delete lock contention "
"(attempt %d/%d) db=%s ids=%s..%s -- retrying: %s",
attempt + 1, retries, self.db_path,
chunk[0], chunk[-1], msg[:200],
)
time.sleep(0.2 * (attempt + 1)) time.sleep(0.2 * (attempt + 1))
else: else:
if is_lock:
log.error(
"TableChangeLog delete still locked after %d "
"attempts db=%s ids=%s..%s -- giving up",
retries, self.db_path, chunk[0], chunk[-1],
)
raise raise
return total return total
@@ -163,6 +177,18 @@ class AccessReader:
cur.execute("SELECT ID FROM TableChangeLog ORDER BY ID") cur.execute("SELECT ID FROM TableChangeLog ORDER BY ID")
return [r[0] for r in cur.fetchall()] return [r[0] for r in cur.fetchall()]
def count_rows(self, table: str) -> int:
"""Return ``COUNT(*)`` for ``table`` (compare count check)."""
cur = self._connect().cursor()
cur.execute(f'SELECT COUNT(*) FROM "{table}"')
return cur.fetchone()[0]
def read_ids(self, table: str) -> list:
"""Return every ``ID`` from ``table``, ascending (compare ID-set check)."""
cur = self._connect().cursor()
cur.execute(f'SELECT ID FROM "{table}" ORDER BY ID')
return [r[0] for r in cur.fetchall()]
def close(self): def close(self):
if self._conn: if self._conn:
self._conn.close() self._conn.close()

View File

@@ -1,42 +1,175 @@
"""Capture phase: read Access change-log rows and stage them into SyncQueue.
Every consumed ``TableChangeLog`` row is appended to the permanent audit store
(``SyncLogArchive``) BEFORE it is enqueued, so the original evidence survives
cleanup. On top of that, this module now emits a detailed audit trail to the
service log:
- WARNING for every operate-type downgrade (Insert/Update -> Delete because
the source row was unreadable at capture time), with record id, log id and
the Access log timestamp -- the exact event class behind the 14287 incident,
previously invisible in the text log;
- WARNING for unknown operate types, with enough identity (log id / record id
/ time) to locate and repair the offending log row manually;
- a per-file INFO summary: rows read, newly enqueued, dedup-skipped
(re-capture after a failed apply -- a symptom worth noticing), downgraded,
out-of-scope, unknown ops, the processed log-ID range and a per-operation
breakdown;
- DEBUG detail for individual out-of-scope and dedup skips.
"""
from __future__ import annotations from __future__ import annotations
import json, logging
import json
import logging
from dataclasses import dataclass, field
from .access_reader import AccessReader from .access_reader import AccessReader
from .sql_writer import SqlWriter, QueueRow from .sql_writer import SqlWriter, QueueRow, ArchiveRow
from .config import FileMapping, SyncConfig from .config import FileMapping, SyncConfig
from .targets import is_synced_table
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
def capture_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter, cfg: SyncConfig) -> int:
exclude = set(fm.exclude_tables or []) @dataclass
include = set(fm.include_tables) if fm.include_tables else None class CaptureStats:
"""Counters for one capture pass (per file, or aggregated per cycle)."""
read: int = 0 # change-log rows read from Access
enqueued: int = 0 # rows newly inserted into SyncQueue
dedup_skipped: int = 0 # already queued (re-capture after a failed apply)
downgraded: int = 0 # Insert/Update downgraded to Delete
out_of_scope: int = 0 # log rows for tables outside the sync scope
unknown_op: int = 0 # log rows with an unrecognised OperateType
min_log_id: int | None = None
max_log_id: int | None = None
ops: dict[str, int] = field(default_factory=dict) # processed op -> count
def note_op(self, op: str) -> None:
self.ops[op] = self.ops.get(op, 0) + 1
def merge(self, other: "CaptureStats") -> None:
"""Fold *other* into this instance (cycle-level aggregation)."""
self.read += other.read
self.enqueued += other.enqueued
self.dedup_skipped += other.dedup_skipped
self.downgraded += other.downgraded
self.out_of_scope += other.out_of_scope
self.unknown_op += other.unknown_op
for k, v in other.ops.items():
self.ops[k] = self.ops.get(k, 0) + v
for attr, pick in (("min_log_id", min), ("max_log_id", max)):
a, b = getattr(self, attr), getattr(other, attr)
if a is None:
setattr(self, attr, b)
elif b is not None:
setattr(self, attr, pick(a, b))
def ops_str(self) -> str:
return ",".join(f"{k}={v}" for k, v in sorted(self.ops.items())) or "-"
def capture_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter,
cfg: SyncConfig) -> CaptureStats:
"""Capture one file's pending change-log rows. Returns detailed stats.
NOTE: previously returned a bare int (rows enqueued); that count is now
``stats.enqueued``. ``sync.service.cycle`` is the only in-repo caller and
has been updated accordingly.
"""
rows = reader.read_log(cfg.runtime.capture_batch_size) rows = reader.read_log(cfg.runtime.capture_batch_size)
n = 0 st = CaptureStats(read=len(rows))
if rows: # read_log orders by ID ascending
st.min_log_id, st.max_log_id = rows[0].id, rows[-1].id
for lr in rows: for lr in rows:
if lr.table_name in exclude: if not is_synced_table(fm, lr.table_name):
continue st.out_of_scope += 1
if include is not None and lr.table_name not in include: log.debug("capture skip (out of scope) file=%s table=%s log_id=%s",
fm.file, lr.table_name, lr.id)
continue continue
op = lr.operate_type op = lr.operate_type
row_data = None row_data = None
if op in ("Insert", "Update"): if op in ("Insert", "Update"):
d = reader.read_row(lr.table_name, lr.record_id) d = reader.read_row(lr.table_name, lr.record_id)
if d is None: if d is None:
op = "Delete" # 行已删,降级 # 行已删(或此刻不可读),降级为 Delete。这是数据差异排查的
# 头号嫌疑事件(参见 14287 事件),必须在文本日志显式留痕,
# 而不只是写入 SyncLogArchive。
op = "Delete"
st.downgraded += 1
log.warning(
"capture DOWNGRADE %s->Delete file=%s table=%s "
"record_id=%s log_id=%s log_time=%s "
"(source row unreadable at capture time; original intent "
"preserved in SyncLogArchive.OriginalOperateType)",
lr.operate_type, fm.file, lr.table_name,
lr.record_id, lr.id, lr.time,
)
else: else:
row_data = json.dumps(d, ensure_ascii=False) row_data = json.dumps(d, ensure_ascii=False)
elif op != "Delete": elif op != "Delete":
log.warning("unknown OperateType %r in %s log %s", op, fm.file, lr.id) st.unknown_op += 1
log.warning(
"capture UNKNOWN OperateType %r file=%s table=%s "
"record_id=%s log_id=%s log_time=%s -- row skipped; it will "
"be re-read every cycle until removed from TableChangeLog",
op, fm.file, lr.table_name, lr.record_id, lr.id, lr.time,
)
continue continue
target_schema = fm.schema
target_table = fm.target_table(lr.table_name)
# Persist the ORIGINAL log entry to the permanent audit store BEFORE the
# queue insert (and long before cleanup deletes the Access log). This
# keeps both the source operate type (lr.operate_type) and the processed
# one (op) so a downgrade like Insert->Delete stays reconstructible.
writer.insert_archive_row(ArchiveRow(
source_file=fm.file,
source_table=lr.table_name,
source_log_id=lr.id,
record_id=lr.record_id,
target_schema=target_schema,
target_table=target_table,
original_operate_type=lr.operate_type,
processed_operate_type=op,
row_data=row_data,
original_time=lr.time,
))
qr = QueueRow( qr = QueueRow(
source_file=fm.file, source_file=fm.file,
source_table=lr.table_name, source_table=lr.table_name,
record_id=lr.record_id, record_id=lr.record_id,
target_schema=fm.schema, target_schema=target_schema,
target_table=fm.target_table(lr.table_name), target_table=target_table,
source_log_id=lr.id, source_log_id=lr.id,
operate_type=op, operate_type=op,
row_data=row_data, row_data=row_data,
) )
writer.insert_queue_row(qr) if writer.insert_queue_row(qr):
n += 1 st.enqueued += 1
return n st.note_op(op)
else:
# Dedup hit: this log row was already staged by an earlier cycle
# whose apply failed (the Access log row is only deleted after a
# successful apply). A persistently non-zero dedup count therefore
# points straight at a stuck apply -- see the queue-health WARNINGs
# emitted by sync.service.cycle.
st.dedup_skipped += 1
log.debug(
"capture dedup-skip (already queued) file=%s table=%s "
"log_id=%s record_id=%s op=%s",
fm.file, lr.table_name, lr.id, lr.record_id, op,
)
if st.read:
log.info(
"capture file=%s read=%d enqueued=%d dedup_skipped=%d "
"downgraded=%d out_of_scope=%d unknown_op=%d "
"log_ids=%s..%s ops={%s}",
fm.file, st.read, st.enqueued, st.dedup_skipped, st.downgraded,
st.out_of_scope, st.unknown_op, st.min_log_id, st.max_log_id,
st.ops_str(),
)
else:
log.debug("capture file=%s: change log empty", fm.file)
return st

View File

@@ -1,10 +1,16 @@
"""Cleanup phase: delete applied Access log rows. """Cleanup phase: delete applied Access log rows.
After ``dbo.usp_SyncApply`` flips queue rows to ``applied``, those rows' After ``usp_SyncApply`` flips queue rows to ``applied``, those rows'
``SourceLogID`` values are no longer needed on the Access side. This module ``SourceLogID`` values are no longer needed on the Access side. This module
asks the writer which log IDs have been applied for a given source file and asks the writer which log IDs have been applied for a given source file and
deletes them from ``TableChangeLog`` via the reader, in batches with lock deletes them from ``TableChangeLog`` via the reader, in batches with lock
retry. The delete is the only mutation that touches the Access side. retry. The delete is the only mutation that touches the Access side.
Audit trail: the INFO line records the exact log-ID range removed from each
file, and a WARNING is raised when fewer rows were deleted than expected --
i.e. some applied log IDs were already absent from the Access log (removed
externally, or by a previously interrupted run), which is worth knowing when
reconstructing what happened around a divergence.
""" """
from __future__ import annotations from __future__ import annotations
import logging import logging
@@ -28,9 +34,25 @@ def cleanup_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter, cfg:
""" """
ids = writer.applied_log_ids(fm.file) ids = writer.applied_log_ids(fm.file)
if not ids: if not ids:
log.debug("cleanup file=%s: no applied rows to clean", fm.file)
return 0 return 0
log.debug("cleanup file=%s: deleting %d applied log rows ids=%s..%s",
fm.file, len(ids), ids[0], ids[-1])
deleted = reader.delete_log_ids( deleted = reader.delete_log_ids(
ids, cfg.runtime.cleanup_batch_size, cfg.runtime.cleanup_lock_retries ids, cfg.runtime.cleanup_batch_size, cfg.runtime.cleanup_lock_retries
) )
if deleted != len(ids):
log.warning(
"cleanup file=%s: deleted %d of %d applied log rows "
"(some log IDs were already absent from Access -- removed "
"externally or by an earlier interrupted run)",
fm.file, deleted, len(ids),
)
writer.mark_cleaned(fm.file, ids) writer.mark_cleaned(fm.file, ids)
if deleted:
log.info(
"cleanup file=%s: removed %d access log rows ids=%s..%s; "
"queue rows marked cleaned",
fm.file, deleted, ids[0], ids[-1],
)
return deleted return deleted

192
src/sync/compare.py Normal file
View File

@@ -0,0 +1,192 @@
"""Data consistency check: Access source tables vs their SQL Server mirrors.
Two granularities:
- ``count`` (default): row-count totals per table.
- ``ids``: ID-set membership -- which IDs exist only in Access or only in SQL.
Tables whose SQL mirror does not exist are skipped (same rule fullsync uses)
and reported as ``skipped``; they do not count as mismatches. Uses
``targets.resolve_synced_tables`` so compare visits exactly the same table set
as full sync -- empirically confirming the two stay aligned.
"""
from __future__ import annotations
import datetime as _dt
import logging
import os
from dataclasses import dataclass, field
from .config import FileMapping, SyncConfig
from .access_reader import AccessReader
from .sql_writer import SqlWriter
from .targets import resolve_synced_tables
log = logging.getLogger("sync.compare")
@dataclass
class TableResult:
"""One table's comparison outcome."""
file: str
access_table: str
target_schema: str
target_table: str
status: str # "match" | "mismatch" | "skipped" | "error"
access_count: int | None = None
sql_count: int | None = None
missing_in_sql: list = field(default_factory=list) # IDs in Access, not SQL
extra_in_sql: list = field(default_factory=list) # IDs in SQL, not Access
error: str | None = None
def compare_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter,
granularity: str = "count") -> list[TableResult]:
"""Compare every in-scope table in one Access file against its SQL mirror.
``granularity`` is ``"count"`` (row totals, default) or ``"ids"`` (ID-set
membership). Skips tables with no SQL mirror. Per-table errors are caught
so one bad table does not abort the file.
"""
results: list[TableResult] = []
for access_table in resolve_synced_tables(fm, reader):
target = fm.target_table(access_table)
if not writer.table_exists(fm.schema, target):
results.append(TableResult(fm.file, access_table, fm.schema, target, "skipped"))
log.warning("skip %s -> %s.%s (target table missing)",
access_table, fm.schema, target)
continue
try:
if granularity == "ids":
a_ids = set(reader.read_ids(access_table))
s_ids = set(writer.read_target_ids(fm.schema, target))
missing = sorted(a_ids - s_ids)
extra = sorted(s_ids - a_ids)
status = "match" if not missing and not extra else "mismatch"
results.append(TableResult(
fm.file, access_table, fm.schema, target, status,
access_count=len(a_ids), sql_count=len(s_ids),
missing_in_sql=missing, extra_in_sql=extra,
))
else:
a = reader.count_rows(access_table)
s = writer.count_target(fm.schema, target)
status = "match" if a == s else "mismatch"
results.append(TableResult(
fm.file, access_table, fm.schema, target, status,
access_count=a, sql_count=s,
))
except Exception as e:
results.append(TableResult(
fm.file, access_table, fm.schema, target, "error", error=str(e)
))
log.exception("compare failed for %s -> %s.%s", access_table, fm.schema, target)
return results
def compare(cfg: SyncConfig, granularity: str = "count",
db_filter: str | None = None,
table_filter: str | None = None) -> list[TableResult]:
"""Compare all (optionally filtered) configured files.
``db_filter`` limits to one Access file; ``table_filter`` restricts every
file to that one table (overrides include_tables), mirroring fullsync.
"""
files = cfg.files
if db_filter:
files = [f for f in files if f.file == db_filter]
if not files:
log.warning("no file matches --db %r", db_filter)
return []
if table_filter:
files = [f.model_copy(update={"include_tables": [table_filter]}) for f in files]
writer = SqlWriter(cfg.sql_server.conn_str, cfg.sql_server.sync_queue_table)
results: list[TableResult] = []
try:
for fm in files:
reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
try:
results.extend(compare_file(fm, reader, writer, granularity))
finally:
reader.close()
finally:
writer.close()
return results
def any_mismatch(results: list[TableResult]) -> bool:
"""True if any compared table diverged (skipped/error do not count)."""
return any(r.status == "mismatch" for r in results)
def format_report(results: list[TableResult], granularity: str = "count") -> str:
"""Render a human-readable per-table report."""
lines = []
for r in results:
base = f"{r.file}: {r.access_table} -> {r.target_schema}.{r.target_table}"
if r.status == "skipped":
lines.append(f"{base} [SKIPPED no mirror]")
elif r.status == "error":
lines.append(f"{base} [ERROR {r.error}]")
elif granularity == "ids":
lines.append(
f"{base} access={r.access_count} sql={r.sql_count} "
f"missing_in_sql={len(r.missing_in_sql)} extra_in_sql={len(r.extra_in_sql)} "
f"[{r.status.upper()}]"
)
if r.missing_in_sql:
lines.append(f" missing_in_sql (first 50): {r.missing_in_sql[:50]}")
if r.extra_in_sql:
lines.append(f" extra_in_sql (first 50): {r.extra_in_sql[:50]}")
else:
lines.append(f"{base} access={r.access_count} sql={r.sql_count} [{r.status.upper()}]")
return "\n".join(lines)
def summarize(results: list[TableResult]) -> dict:
"""Tally the outcome buckets: match / mismatch / skipped / error."""
s = {"match": 0, "mismatch": 0, "skipped": 0, "error": 0}
for r in results:
s[r.status] = s.get(r.status, 0) + 1
return s
def write_report(results: list[TableResult], granularity: str = "count",
log_dir: str = ".", run_dt: _dt.datetime | None = None,
extra_path: str | None = None) -> str:
"""Render the full report and persist it as a dated log file under *log_dir*.
Always writes ``<log_dir>/compare_<granularity>_<YYYY-MM-DD>.log``,
overwriting the day's previous run (one report per day; the logging system's
per-day archival later moves yesterday's file into ``logs/Archive/``). When
*extra_path* is given (the ``--report`` CLI option) the identical content is
also written there for backward compatibility. Returns the primary log path.
"""
run_dt = run_dt or _dt.datetime.now()
stats = summarize(results)
body = format_report(results, granularity)
header = (
"============================================================\n"
" 数据一致性核对报告 / Compare Report\n"
f" 生成时间 : {run_dt.strftime('%Y-%m-%d %H:%M:%S')}\n"
f" 粒度 : {granularity}\n"
f" 比对表合计 : {stats['match'] + stats['mismatch']}\n"
f" 一致 MATCH : {stats['match']}\n"
f" 不一致 MISMATCH : {stats['mismatch']}\n"
f" 跳过 SKIPPED : {stats['skipped']}\n"
f" 错误 ERROR : {stats['error']}\n"
"============================================================\n"
)
report = header + body + "\n"
log_dir = log_dir or "."
os.makedirs(log_dir, exist_ok=True)
primary = os.path.join(
log_dir, f"compare_{granularity}_{run_dt.strftime('%Y-%m-%d')}.log"
)
with open(primary, "w", encoding="utf-8") as f:
f.write(report)
if extra_path:
with open(extra_path, "w", encoding="utf-8") as f:
f.write(report)
return primary

View File

@@ -5,7 +5,13 @@ from pydantic import BaseModel, ConfigDict, Field
class SqlServerConfig(BaseModel): class SqlServerConfig(BaseModel):
conn_str: str conn_str: str
sync_queue_table: str = "dbo.SyncQueue" sync_queue_table: str = "ProductionDataBaseSync.SyncQueue"
archive_table: str = "ProductionDataBaseSync.SyncLogArchive"
apply_proc: str = "ProductionDataBaseSync.usp_SyncApply"
# Per-invocation, per-table apply audit written by usp_SyncApply
# (see sql/04_sync_apply_runlog.sql). Default matches the shipped script;
# existing config.yaml files need no change.
apply_runlog_table: str = "ProductionDataBaseSync.SyncApplyRunLog"
class AccessConfig(BaseModel): class AccessConfig(BaseModel):
model_config = ConfigDict(coerce_numbers_to_str=True) model_config = ConfigDict(coerce_numbers_to_str=True)
@@ -21,6 +27,9 @@ class RuntimeConfig(BaseModel):
cleanup_batch_size: int = 200 cleanup_batch_size: int = 200
cleanup_lock_retries: int = 3 cleanup_lock_retries: int = 3
cleaned_retention_hours: int = 24 cleaned_retention_hours: int = 24
# Idle cycles now log at DEBUG; the service emits an INFO heartbeat at this
# interval while idle so a quiet log still proves the service is alive.
idle_heartbeat_seconds: int = 600
class FileMapping(BaseModel): class FileMapping(BaseModel):
model_config = ConfigDict(coerce_numbers_to_str=True) model_config = ConfigDict(coerce_numbers_to_str=True)

View File

@@ -22,6 +22,7 @@ from .config import load_config, FileMapping, SyncConfig
from .access_reader import AccessReader from .access_reader import AccessReader
from .sql_writer import SqlWriter from .sql_writer import SqlWriter
from .logging_setup import setup_logging from .logging_setup import setup_logging
from .targets import resolve_synced_tables
log = logging.getLogger("sync.fullsync") log = logging.getLogger("sync.fullsync")
@@ -29,21 +30,11 @@ log = logging.getLogger("sync.fullsync")
def resolve_tables(reader: AccessReader, fm: FileMapping) -> list[str]: def resolve_tables(reader: AccessReader, fm: FileMapping) -> list[str]:
"""Tables to fully sync for one file, after exclude/include rules. """Tables to fully sync for one file, after exclude/include rules.
Mirrors the precedence used by ``capture.capture_file``: a table in Thin wrapper over ``targets.resolve_synced_tables`` so fullsync, capture
``exclude_tables`` is dropped even if it also appears in ``include_tables``. and compare share one resolution path. System tables (``MSys*`` / ``~*``)
System tables (``MSys*`` / ``~*``) are already filtered by are already filtered by ``AccessReader.list_user_tables``.
``AccessReader.list_user_tables``.
""" """
exclude = set(fm.exclude_tables or []) return resolve_synced_tables(fm, reader)
include = set(fm.include_tables) if fm.include_tables else None
out = []
for t in reader.list_user_tables():
if t in exclude:
continue
if include is not None and t not in include:
continue
out.append(t)
return out
def full_sync_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter) -> dict: def full_sync_file(fm: FileMapping, reader: AccessReader, writer: SqlWriter) -> dict:

View File

@@ -2,10 +2,124 @@
Configures the root logger with a RotatingFileHandler (10 MB x 5, UTF-8) plus Configures the root logger with a RotatingFileHandler (10 MB x 5, UTF-8) plus
a console StreamHandler. The level and log path come from ``cfg.logging``. a console StreamHandler. The level and log path come from ``cfg.logging``.
Every record written to the log file carries a cycle correlation id
(``[cyc:xxxxxxxx]``): ``sync.service.cycle`` allocates one per pass via
``set_cycle_id`` and ``_CycleIdFilter`` injects it into each record, so every
capture/apply/cleanup line of one pass -- and the matching
``SyncApplyRunLog.CycleID`` rows on SQL Server -- can be correlated with a
single grep. Outside a cycle (fullsync, compare, startup) the field is ``-``.
Log files are managed per-day: at startup any previously produced log
(including the project's own ``sync.log`` and NSSM's ``nssm_*.log`` captures)
is relocated into an ``Archive/`` subfolder next to the active log. Where a log
lacks a timestamp, a ``-YYYY-MM-DD`` suffix is added (derived from its first
log line, falling back to mtime) so historical files carry a date. The log
root therefore only ever shows the current day's ``sync.log``.
""" """
import contextvars
import datetime
import logging import logging
import logging.handlers import logging.handlers
import os import os
import re
import shutil
_LOG_DATE_RE = re.compile(r"^(\d{4}-\d{2}-\d{2})")
# Current cycle correlation id ("-" when not inside a service cycle).
_cycle_id: contextvars.ContextVar[str] = contextvars.ContextVar(
"sync_cycle_id", default="-"
)
def set_cycle_id(cycle_id: str | None) -> None:
"""Set (or clear, with ``None``) the id stamped on every log record.
Called by ``sync.service.cycle`` at the start/end of each pass. The same
id is passed to ``usp_SyncApply`` so SQL-side ``SyncApplyRunLog`` rows can
be joined back to the exact log lines of the cycle that produced them.
"""
_cycle_id.set(cycle_id or "-")
class _CycleIdFilter(logging.Filter):
"""Inject the current cycle id into every record as ``record.cycle``.
Attached to the handlers (not the logger) so records emitted through any
module logger -- capture, cleanup, access_reader, ... -- are covered.
"""
def filter(self, record: logging.LogRecord) -> bool: # noqa: A003
record.cycle = _cycle_id.get()
return True
def _first_line_date(path: str) -> str | None:
"""Best-effort extraction of the first log line's YYYY-MM-DD date."""
try:
with open(path, "r", encoding="utf-8", errors="replace") as f:
for line in f:
m = _LOG_DATE_RE.match(line)
if m:
return m.group(1)
except OSError:
pass
return None
def _archive_completed_logs(log_dir: str, active_name: str):
"""Move already-produced logs out of *log_dir* into *log_dir*/Archive.
The active log (*active_name*) is left in place only when it belongs to the
current day; an older ``sync.log`` is archived (with a ``-YYYY-MM-DD``
suffix) so a fresh one can be opened. NSSM's own ``nssm_*.log`` captures are
timestamped already and are moved as-is. Moves are best-effort: files locked
by another process (e.g. NSSM's live handles) are skipped.
"""
archive_dir = os.path.join(log_dir, "Archive")
os.makedirs(archive_dir, exist_ok=True)
today = datetime.date.today()
for name in os.listdir(log_dir):
src = os.path.join(log_dir, name)
if not os.path.isfile(src):
continue
if name == "Archive":
continue
if not (name.endswith(".log") or ".log." in name):
continue
# Keep NSSM's live, currently-open handles in place.
if name in ("nssm_stderr.log", "nssm_stdout.log"):
continue
if name == active_name:
# Only archive the active log if it is from a previous day.
log_date = _first_line_date(src)
log_date = (
datetime.date.fromisoformat(log_date)
if log_date
else datetime.date.fromtimestamp(os.path.getmtime(src))
)
if log_date >= today:
continue # today's log: keep appending
new_name = f"sync-{log_date.strftime('%Y-%m-%d')}.log"
else:
new_name = name
dst = os.path.join(archive_dir, new_name)
if os.path.exists(dst):
stem, ext = os.path.splitext(new_name)
i = 2
while os.path.exists(os.path.join(archive_dir, f"{stem}({i}){ext}")):
i += 1
dst = os.path.join(archive_dir, f"{stem}({i}){ext}")
try:
shutil.move(src, dst)
except OSError:
# Locked by another process (e.g. NSSM holding the file open).
pass
def setup_logging(cfg_dict: dict | None): def setup_logging(cfg_dict: dict | None):
@@ -13,18 +127,45 @@ def setup_logging(cfg_dict: dict | None):
``cfg_dict`` is ``SyncConfig.logging`` (a dict or None). ``level`` is a ``cfg_dict`` is ``SyncConfig.logging`` (a dict or None). ``level`` is a
logging-level name string (default ``"INFO"``); ``path`` is the log file logging-level name string (default ``"INFO"``); ``path`` is the log file
path (default ``"sync.log"``). The parent directory is created if missing. path (default ``"sync.log"``). The parent directory is created if missing,
and any previously produced logs are archived before the new handler opens.
""" """
level = getattr(logging, (cfg_dict or {}).get("level", "INFO")) if cfg_dict else logging.INFO level = getattr(logging, (cfg_dict or {}).get("level", "INFO")) if cfg_dict else logging.INFO
path = (cfg_dict or {}).get("path", "sync.log") path = (cfg_dict or {}).get("path", "sync.log")
os.makedirs(os.path.dirname(path) or ".", exist_ok=True) log_dir = os.path.dirname(path) or "."
os.makedirs(log_dir, exist_ok=True)
# Archive everything from prior runs so the root only shows today's log.
_archive_completed_logs(log_dir, os.path.basename(path))
# Idempotent: drop any handlers already attached to the root logger before
# re-adding. setup_logging can be called from more than one entry point
# (e.g. main.py and service.run), and the old code unconditionally
# addHandler'd each time, stacking duplicate handlers so every log line was
# written twice. Clearing first means repeated calls always yield exactly
# one file handler + one console handler, regardless of caller.
root = logging.getLogger()
for old in list(root.handlers):
root.removeHandler(old)
try:
old.close()
except Exception:
pass
root.setLevel(level)
cycle_filter = _CycleIdFilter()
h = logging.handlers.RotatingFileHandler( h = logging.handlers.RotatingFileHandler(
path, maxBytes=10 * 1024 * 1024, backupCount=5, encoding="utf-8" path, maxBytes=10 * 1024 * 1024, backupCount=5, encoding="utf-8"
) )
h.setFormatter(logging.Formatter("%(asctime)s %(levelname)s [%(name)s] %(message)s")) h.addFilter(cycle_filter)
root = logging.getLogger() h.setFormatter(logging.Formatter(
root.setLevel(level) "%(asctime)s %(levelname)s [%(name)s] [cyc:%(cycle)s] %(message)s"
))
root.addHandler(h) root.addHandler(h)
sh = logging.StreamHandler() sh = logging.StreamHandler()
sh.addFilter(cycle_filter)
# Console keeps the short format (fullsync/compare are interactive there).
sh.setFormatter(logging.Formatter("%(levelname)s %(message)s")) sh.setFormatter(logging.Formatter("%(levelname)s %(message)s"))
root.addHandler(sh) root.addHandler(sh)

View File

@@ -2,9 +2,24 @@
``cycle(cfg)`` runs one full pass over every configured file: ``cycle(cfg)`` runs one full pass over every configured file:
1. capture — read each file's change log and stage rows into SyncQueue; 1. capture — read each file's change log and stage rows into SyncQueue;
2. apply — drain the queue via ``dbo.usp_SyncApply``; 2. apply — drain the queue via ``usp_SyncApply``;
3. cleanup — delete applied log rows from each file's Access log. 3. cleanup — delete applied log rows from each file's Access log.
Observability (rebuilt so data divergence is traceable from the log alone):
- every cycle gets a short correlation id; ``logging_setup`` stamps it on each
log line as ``[cyc:xxxxxxxx]`` and the same id is passed to ``usp_SyncApply``
so ``SyncApplyRunLog`` rows on SQL Server join back to the exact log lines of
the cycle that produced them;
- after apply, the per-table run-log rows (pending/merged/deleted/applied/
error/dead counts, duration, error message) are read back and logged --
the old single "apply done" line hid all of this;
- queue health is checked every cycle: ``error``/``dead`` rows, which
previously accumulated in complete silence, now emit WARNINGs with per-row
samples (table, record id, retry count, error message) -- these are exactly
the changes that exist in Access but never reached SQL Server;
- idle cycles log at DEBUG so INFO stays high-signal; ``run`` emits a periodic
idle heartbeat so a quiet log still proves the service is alive.
Each file's capture and cleanup is wrapped in its own try/except so one Each file's capture and cleanup is wrapped in its own try/except so one
file's failure is logged and the cycle continues; the writer is always closed file's failure is logged and the cycle continues; the writer is always closed
in a ``finally``. ``run(cfg)`` loops ``cycle`` with a sleep; ``main()`` in a ``finally``. ``run(cfg)`` loops ``cycle`` with a sleep; ``main()``
@@ -13,14 +28,15 @@ loads the config from ``argv[1]`` (default ``config.yaml``).
from __future__ import annotations from __future__ import annotations
import sys import sys
import time import time
import uuid
import logging import logging
from .config import load_config from .config import load_config
from .access_reader import AccessReader from .access_reader import AccessReader
from .sql_writer import SqlWriter from .sql_writer import SqlWriter
from .capture import capture_file from .capture import capture_file, CaptureStats
from .cleanup import cleanup_file from .cleanup import cleanup_file
from .logging_setup import setup_logging from .logging_setup import setup_logging, set_cycle_id
log = logging.getLogger("sync.service") log = logging.getLogger("sync.service")
@@ -30,47 +46,156 @@ def run(cfg):
Configures logging once on entry. Intended to be started by ``main()`` Configures logging once on entry. Intended to be started by ``main()``
under the service host (e.g. NSSM). Not unit-tested (infinite loop); under the service host (e.g. NSSM). Not unit-tested (infinite loop);
``cycle()`` is the testable unit. ``cycle()`` is the testable unit. Emits an idle heartbeat every
``runtime.idle_heartbeat_seconds`` so a quiet log still proves liveness
now that idle cycles log at DEBUG.
""" """
setup_logging(cfg.logging) setup_logging(cfg.logging)
log.info(
"service started: files=%d poll_interval=%ds idle_heartbeat=%ds",
len(cfg.files), cfg.runtime.poll_interval_seconds,
cfg.runtime.idle_heartbeat_seconds,
)
idle_since = None
idle_cycles = 0
while True: while True:
cycle(cfg) active = cycle(cfg)
now = time.monotonic()
if active:
idle_since, idle_cycles = None, 0
else:
idle_cycles += 1
if idle_since is None:
idle_since = now
elif now - idle_since >= cfg.runtime.idle_heartbeat_seconds:
log.info("idle heartbeat: %d cycles with no changes in the "
"last %ds", idle_cycles, int(now - idle_since))
idle_since, idle_cycles = now, 0
time.sleep(cfg.runtime.poll_interval_seconds) time.sleep(cfg.runtime.poll_interval_seconds)
def cycle(cfg): def cycle(cfg) -> bool:
"""One capture -> apply -> cleanup pass over all files. """One capture -> apply -> cleanup pass over all files.
Per-file capture/cleanup failures are logged and do not abort the cycle. Returns True when the cycle did any work (captured / applied / cleaned /
Apply failure does not block cleanup. The writer is always closed in a purged anything); ``run`` uses this for idle-heartbeat pacing. Per-file
capture/cleanup failures are logged and do not abort the cycle. Apply
failure does not block cleanup. The writer is always closed in a
``finally``. Safe to call directly from tests (does not sleep or loop). ``finally``. Safe to call directly from tests (does not sleep or loop).
""" """
writer = SqlWriter(cfg.sql_server.conn_str, cfg.sql_server.sync_queue_table) cycle_id = uuid.uuid4().hex[:8]
set_cycle_id(cycle_id)
t0 = time.monotonic()
activity = False
writer = SqlWriter(
cfg.sql_server.conn_str,
cfg.sql_server.sync_queue_table,
cfg.sql_server.archive_table,
cfg.sql_server.apply_proc,
cfg.sql_server.apply_runlog_table,
)
try: try:
total_captured = 0 # ---- capture ------------------------------------------------------
total = CaptureStats()
for fm in cfg.files: for fm in cfg.files:
reader = AccessReader(fm.source_path(cfg), cfg.access.driver) reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
try: try:
n = capture_file(fm, reader, writer, cfg) total.merge(capture_file(fm, reader, writer, cfg))
total_captured += n
except Exception: except Exception:
log.exception("capture failed for %s", fm.file) log.exception("capture failed for %s", fm.file)
finally: finally:
reader.close() reader.close()
log.info("captured %d rows", total_captured) if total.read:
activity = True
log.info(
"capture summary: read=%d enqueued=%d dedup_skipped=%d "
"downgraded=%d out_of_scope=%d unknown_op=%d ops={%s}",
total.read, total.enqueued, total.dedup_skipped,
total.downgraded, total.out_of_scope, total.unknown_op,
total.ops_str(),
)
else:
log.debug("capture summary: no new change-log rows in any file")
# ---- apply ----------------------------------------------------------
try: try:
writer.call_apply(cfg.runtime.max_retries) writer.call_apply(cfg.runtime.max_retries, cycle_id)
runs = writer.apply_run_results(cycle_id)
if runs is None:
# Audit infra (SyncApplyRunLog / updated proc) not deployed
# yet: keep the legacy coarse line. sql_writer already warned
# once about what to deploy.
log.info("apply done") log.info("apply done")
elif not runs:
if total.enqueued:
# Rows were staged but the proc wrote no audit rows: the
# table exists but the deployed proc is probably the old
# version that does not write into it.
log.info("apply done (proc wrote no run-log rows -- "
"re-run the updated sql/02_sync_apply.sql?)")
else:
log.debug("apply done: queue was empty")
else:
activity = True
for r in runs:
if r["Outcome"] == "ok":
log.info(
"apply ok table=%s.%s pending=%s records=%s "
"merged=%s deleted=%s applied=%s dur_ms=%s",
r["TargetSchema"], r["TargetTable"],
r["PendingCount"], r["DistinctRecords"],
r["MergedCount"], r["DeletedCount"],
r["AppliedCount"], r["DurationMs"],
)
else:
log.warning(
"apply FAILED table=%s.%s pending=%s records=%s "
"-> error=%s dead=%s dur_ms=%s msg=%s",
r["TargetSchema"], r["TargetTable"],
r["PendingCount"], r["DistinctRecords"],
r["ErrorCount"], r["DeadCount"],
r["DurationMs"], r["ErrorMsg"],
)
except Exception: except Exception:
log.exception("apply failed") log.exception("apply failed")
# ---- queue health ---------------------------------------------------
# error/dead rows are changes that exist in Access but never reached
# the SQL mirror. Before this check they accumulated with zero trace
# in this log -- the classic "data diverged, no idea why" scenario.
try:
status = writer.queue_status_summary()
stuck = status.get("error", 0) + status.get("dead", 0)
if stuck:
log.warning(
"queue health: %d stuck row(s) [%s] -- these changes are "
"NOT in SQL Server and will show up as data divergence",
stuck,
",".join(f"{k}={v}" for k, v in sorted(status.items())),
)
for s in writer.queue_error_samples(10):
log.warning(
" stuck row: table=%s.%s record_id=%s "
"source_log_id=%s op=%s status=%s retries=%s "
"captured_at=%s err=%s",
s["TargetSchema"], s["TargetTable"], s["RecordID"],
s["SourceLogID"], s["OperateType"], s["Status"],
s["RetryCount"], s["CapturedAt"], s["ErrorMsg"],
)
elif status.get("pending", 0):
log.warning(
"queue health: %d row(s) still pending after apply "
"(apply may have failed this cycle)", status["pending"],
)
except Exception:
log.exception("queue health check failed")
# ---- cleanup --------------------------------------------------------
for fm in cfg.files: for fm in cfg.files:
reader = AccessReader(fm.source_path(cfg), cfg.access.driver) reader = AccessReader(fm.source_path(cfg), cfg.access.driver)
try: try:
c = cleanup_file(fm, reader, writer, cfg) if cleanup_file(fm, reader, writer, cfg):
if c: activity = True
log.info("cleaned %d log rows from %s", c, fm.file)
except Exception: except Exception:
log.exception("cleanup failed for %s", fm.file) log.exception("cleanup failed for %s", fm.file)
finally: finally:
@@ -81,12 +206,18 @@ def cycle(cfg):
try: try:
purged = writer.purge_cleaned(cfg.runtime.cleaned_retention_hours) purged = writer.purge_cleaned(cfg.runtime.cleaned_retention_hours)
if purged: if purged:
log.info("purged %d cleaned queue rows", purged) activity = True
log.info("purged %d cleaned queue rows (older than %dh)",
purged, cfg.runtime.cleaned_retention_hours)
except Exception: except Exception:
log.exception("purge failed") log.exception("purge failed")
(log.info if activity else log.debug)(
"cycle finished in %.2fs", time.monotonic() - t0)
return activity
finally: finally:
writer.close() writer.close()
set_cycle_id(None)
def main(): def main():

View File

@@ -1,26 +1,50 @@
"""SQL-side writer for the Access -> SQL Server sync. """SQL-side writer for the Access -> SQL Server sync.
SqlWriter owns the pyodbc connection used to (a) dedup-insert rows into SqlWriter owns the pyodbc connection used to (a) dedup-insert rows into the
``dbo.SyncQueue``, (b) invoke ``dbo.usp_SyncApply`` to drain the queue, and staging queue (``ProductionDataBaseSync.SyncQueue`` by default), (b) invoke the
(c) report which SourceLogIDs have been applied. apply proc (``ProductionDataBaseSync.usp_SyncApply``) to drain the queue,
(c) report which SourceLogIDs have been applied, and (d) append every consumed
Access change-log row to the permanent audit store
(``ProductionDataBaseSync.SyncLogArchive``). The queue/archive/proc names are
injected from config so the whole sync can live under a dedicated schema.
Observability additions:
- the dedup inserts now report whether a row was actually inserted (True) or
already present (False), so capture can log dedup hits -- the tell-tale of a
re-capture after a failed apply;
- ``call_apply`` passes the service's cycle correlation id to the proc, which
records one audit row per target table into
``ProductionDataBaseSync.SyncApplyRunLog`` (see sql/04_sync_apply_runlog.sql);
- ``apply_run_results`` reads those rows back so the service log shows
per-table merged/deleted/applied/error/dead counts instead of "apply done";
- ``queue_status_summary`` / ``queue_error_samples`` surface stuck
(``error``/``dead``) queue rows, which previously accumulated silently.
The connection is opened with ``autocommit=True`` on purpose. ``usp_SyncApply`` The connection is opened with ``autocommit=True`` on purpose. ``usp_SyncApply``
manages its own transaction internally (BEGIN TRAN ... ROLLBACK on error); if manages its own transaction internally (BEGIN TRAN ... ROLLBACK on error); if
the caller held an outer implicit transaction, the proc's ROLLBACK would the caller held an outer implicit transaction, the proc's ROLLBACK would
cascade and raise SQL error 266. The dedup ``IF NOT EXISTS ... INSERT`` is a cascade and raise SQL error 266. The dedup ``INSERT .. SELECT .. WHERE NOT
single statement that is atomic under autocommit, so no explicit transaction is EXISTS`` is a single statement that is atomic under autocommit, so no explicit
needed on the write path either. transaction is needed on the write path either.
""" """
from __future__ import annotations from __future__ import annotations
import logging
from dataclasses import dataclass from dataclasses import dataclass
import pyodbc import pyodbc
log = logging.getLogger(__name__)
# Process-wide deployment-state flags (SqlWriter instances are recreated every
# cycle, so per-instance flags would re-warn each cycle).
_proc_lacks_cycle_id = False # deployed usp_SyncApply predates @CycleID
_runlog_missing_noted = False # SyncApplyRunLog table not deployed yet
@dataclass @dataclass
class QueueRow: class QueueRow:
"""A single staged change to enqueue into ``dbo.SyncQueue``.""" """A single staged change to enqueue into the SyncQueue table."""
source_file: str source_file: str
source_table: str source_table: str
@@ -32,35 +56,67 @@ class QueueRow:
row_data: str | None row_data: str | None
class SqlWriter: @dataclass
"""Writes to SyncQueue and drives the apply proc over a pyodbc connection.""" class ArchiveRow:
"""A single Access change-log row to append to ``SyncLogArchive``.
def __init__(self, conn_str: str, queue_table: str = "dbo.SyncQueue"): Captures both the ORIGINAL operate type recorded by the Access data macro
and the PROCESSED operate type actually sent to the queue, so a downgrade
(e.g. Insert -> Delete when the row is momentarily unreadable) stays visible
forever. ``original_time`` preserves the Access log's own timestamp, which
the queue path discards.
"""
source_file: str
source_table: str
source_log_id: int
record_id: str
target_schema: str
target_table: str
original_operate_type: str
processed_operate_type: str
row_data: str | None
original_time: object
class SqlWriter:
"""Writes to SyncQueue/SyncLogArchive and drives the apply proc."""
def __init__(
self,
conn_str: str,
queue_table: str = "ProductionDataBaseSync.SyncQueue",
archive_table: str = "ProductionDataBaseSync.SyncLogArchive",
apply_proc: str = "ProductionDataBaseSync.usp_SyncApply",
runlog_table: str = "ProductionDataBaseSync.SyncApplyRunLog",
):
self.conn_str = conn_str self.conn_str = conn_str
self.queue_table = queue_table self.queue_table = queue_table
self.archive_table = archive_table
self.apply_proc = apply_proc
self.runlog_table = runlog_table
# autocommit=True: usp_SyncApply manages its own transaction internally. # autocommit=True: usp_SyncApply manages its own transaction internally.
# An outer pyodbc transaction would conflict on ROLLBACK (SQL error 266). # An outer pyodbc transaction would conflict on ROLLBACK (SQL error 266).
self._conn = pyodbc.connect(conn_str, autocommit=True) self._conn = pyodbc.connect(conn_str, autocommit=True)
def insert_queue_row(self, row: QueueRow) -> None: def insert_queue_row(self, row: QueueRow) -> bool:
"""Idempotently enqueue ``row`` (dedup on SourceFile/Table/LogID). """Idempotently enqueue ``row`` (dedup on SourceFile/Table/LogID).
``IF NOT EXISTS ... INSERT`` is a single statement, atomic under Returns True when a new queue row was inserted, False on a dedup hit
autocommit. The unique index UX_SyncQueue_Dedup is the DB backstop. (the row was already staged -- typically a re-capture after a failed
apply left the Access log row in place). ``INSERT .. SELECT .. WHERE
NOT EXISTS`` is a single statement, atomic under autocommit, and its
deterministic rowcount (0/1) is what makes the dedup outcome
observable for the capture audit log. The unique index
UX_SyncQueue_Dedup is the DB backstop.
""" """
# IF NOT EXISTS and the VALUES list each carry their own ? markers;
# pyodbc binds them positionally, so the 3 dedup keys are supplied
# twice (once for the EXISTS check, once for the INSERT).
cur = self._conn.cursor() cur = self._conn.cursor()
cur.execute( cur.execute(
"IF NOT EXISTS (SELECT 1 FROM dbo.SyncQueue " f"INSERT INTO {self.queue_table}(SourceFile,SourceTable,SourceLogID,"
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?) "
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,"
"TargetSchema,TargetTable,RecordID,OperateType,RowData,Status) " "TargetSchema,TargetTable,RecordID,OperateType,RowData,Status) "
"VALUES (?,?,?,?,?,?,?,?, 'pending')", "SELECT ?,?,?,?,?,?,?,?,'pending' "
row.source_file, f"WHERE NOT EXISTS (SELECT 1 FROM {self.queue_table} "
row.source_table, "WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?)",
row.source_log_id,
row.source_file, row.source_file,
row.source_table, row.source_table,
row.source_log_id, row.source_log_id,
@@ -69,22 +125,147 @@ class SqlWriter:
row.record_id, row.record_id,
row.operate_type, row.operate_type,
row.row_data, row.row_data,
row.source_file,
row.source_table,
row.source_log_id,
) )
# autocommit: statement already committed. # autocommit: statement already committed.
return cur.rowcount > 0
def call_apply(self, max_retries: int) -> None: def insert_archive_row(self, row: ArchiveRow) -> bool:
"""Drain the pending queue via the stored procedure. """Append ``row`` to the permanent audit store (dedup on source keys).
``usp_SyncApply`` flips rows to ``applied`` (or ``error`` after retries). Written during capture, BEFORE cleanup deletes the Access log, so the
original evidence survives even after the queue row is purged. Stores
both the original and processed operate types plus the Access log time.
The ``WHERE NOT EXISTS`` guard keeps the first archive record if
capture re-runs the same log id (a prior cycle's apply failed and the
Access log persisted). Returns True when a new archive row was
written, False on a dedup hit.
""" """
cur = self._conn.cursor() cur = self._conn.cursor()
cur.execute("EXEC dbo.usp_SyncApply ?", max_retries) cur.execute(
f"INSERT INTO {self.archive_table}(SourceFile,SourceTable,SourceLogID,"
"RecordID,TargetSchema,TargetTable,OriginalOperateType,"
"ProcessedOperateType,RowData,OriginalTime) "
"SELECT ?,?,?,?,?,?,?,?,?,? "
f"WHERE NOT EXISTS (SELECT 1 FROM {self.archive_table} "
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?)",
row.source_file,
row.source_table,
row.source_log_id,
row.record_id,
row.target_schema,
row.target_table,
row.original_operate_type,
row.processed_operate_type,
row.row_data,
row.original_time,
row.source_file,
row.source_table,
row.source_log_id,
)
# autocommit: statement already committed.
return cur.rowcount > 0
def call_apply(self, max_retries: int, cycle_id: str | None = None) -> None:
"""Drain the pending queue via the stored procedure.
``usp_SyncApply`` flips rows to ``applied`` (or ``error``/``dead``
after retries) and -- once the updated proc plus the SyncApplyRunLog
table are deployed -- writes one audit row per target table tagged
with ``cycle_id``, so SQL-side apply stats join back to the service
log's ``[cyc:...]`` lines. Falls back to the legacy single-parameter
signature when the deployed proc predates ``@CycleID`` (SQL errors
8144/8145), so rolling out the Python side first keeps working.
"""
global _proc_lacks_cycle_id
cur = self._conn.cursor()
if cycle_id is not None and not _proc_lacks_cycle_id:
try:
cur.execute(
f"EXEC {self.apply_proc} @MaxRetries=?, @CycleID=?",
max_retries, cycle_id,
)
return
except pyodbc.Error as e:
msg = str(e)
if "8144" in msg or "8145" in msg:
_proc_lacks_cycle_id = True
log.warning(
"apply proc %s does not accept @CycleID yet -- run the "
"updated sql/02_sync_apply.sql to enable per-table "
"apply auditing; falling back to the legacy signature",
self.apply_proc,
)
else:
raise
cur.execute(f"EXEC {self.apply_proc} ?", max_retries)
def apply_run_results(self, cycle_id: str) -> list[dict] | None:
"""Per-table apply outcomes recorded by the proc for ``cycle_id``.
Reads ``SyncApplyRunLog`` (pending/distinct counts, MERGE/DELETE
rowcounts, applied/error/dead queue-row counts, duration and error
message per target table). Returns ``None`` when the audit
infrastructure is not deployed yet (legacy proc, or table missing) so
the caller can fall back to coarse logging; returns ``[]`` when the
queue was simply empty.
"""
global _runlog_missing_noted
if _proc_lacks_cycle_id:
return None # legacy proc never writes run-log rows
cur = self._conn.cursor()
try:
cur.execute(
"SELECT TargetSchema,TargetTable,PendingCount,DistinctRecords,"
"MergedCount,DeletedCount,AppliedCount,ErrorCount,DeadCount,"
"Outcome,ErrorMsg,StartedAt,DurationMs "
f"FROM {self.runlog_table} WHERE CycleID=? ORDER BY RunLogID",
cycle_id,
)
except pyodbc.Error:
if not _runlog_missing_noted:
_runlog_missing_noted = True
log.warning(
"run-log table %s not available -- run "
"sql/04_sync_apply_runlog.sql to enable per-table apply "
"stats in this log", self.runlog_table,
)
return None
cols = [c[0] for c in cur.description]
return [dict(zip(cols, r)) for r in cur.fetchall()]
def queue_status_summary(self) -> dict[str, int]:
"""Row counts per Status (pending/applied/error/dead/cleaned)."""
cur = self._conn.cursor()
cur.execute(
f"SELECT Status, COUNT(*) FROM {self.queue_table} GROUP BY Status"
)
return {r[0]: r[1] for r in cur.fetchall()}
def queue_error_samples(self, limit: int = 10) -> list[dict]:
"""Most recent ``error``/``dead`` queue rows, for WARNING-level triage.
These are exactly the changes that exist in Access but never reached
the SQL mirror -- the prime suspects for any data divergence.
"""
cur = self._conn.cursor()
cur.execute(
"SELECT TOP (?) TargetSchema,TargetTable,RecordID,SourceLogID,"
"OperateType,Status,RetryCount,ErrorMsg,CapturedAt "
f"FROM {self.queue_table} WHERE Status IN ('error','dead') "
"ORDER BY QueueID DESC",
limit,
)
cols = [c[0] for c in cur.description]
return [dict(zip(cols, r)) for r in cur.fetchall()]
def applied_log_ids(self, source_file: str) -> list[int]: def applied_log_ids(self, source_file: str) -> list[int]:
"""Return applied SourceLogIDs for ``source_file`` in ascending order.""" """Return applied SourceLogIDs for ``source_file`` in ascending order."""
cur = self._conn.cursor() cur = self._conn.cursor()
cur.execute( cur.execute(
"SELECT SourceLogID FROM dbo.SyncQueue " f"SELECT SourceLogID FROM {self.queue_table} "
"WHERE SourceFile=? AND Status='applied' ORDER BY SourceLogID", "WHERE SourceFile=? AND Status='applied' ORDER BY SourceLogID",
source_file, source_file,
) )
@@ -107,7 +288,7 @@ class SqlWriter:
chunk = source_log_ids[i:i + 1000] chunk = source_log_ids[i:i + 1000]
placeholders = ",".join("?" * len(chunk)) placeholders = ",".join("?" * len(chunk))
cur.execute( cur.execute(
f"UPDATE dbo.SyncQueue SET Status='cleaned', " f"UPDATE {self.queue_table} SET Status='cleaned', "
f"CleanedAt=sysdatetime() " f"CleanedAt=sysdatetime() "
f"WHERE SourceFile=? AND Status='applied' " f"WHERE SourceFile=? AND Status='applied' "
f"AND SourceLogID IN ({placeholders})", f"AND SourceLogID IN ({placeholders})",
@@ -120,12 +301,12 @@ class SqlWriter:
Keeps the table bounded: cleanup marks rows ``cleaned`` every cycle, Keeps the table bounded: cleanup marks rows ``cleaned`` every cycle,
and this removes the old ones after a short audit/debug window so and this removes the old ones after a short audit/debug window so
``dbo.SyncQueue`` stops growing without bound. Returns the number of the queue table stops growing without bound. Returns the number of
rows removed. rows removed.
""" """
cur = self._conn.cursor() cur = self._conn.cursor()
cur.execute( cur.execute(
"DELETE FROM dbo.SyncQueue " f"DELETE FROM {self.queue_table} "
"WHERE Status='cleaned' " "WHERE Status='cleaned' "
"AND CleanedAt < DATEADD(hour, -?, GETDATE())", "AND CleanedAt < DATEADD(hour, -?, GETDATE())",
retention_hours, retention_hours,
@@ -143,6 +324,18 @@ class SqlWriter:
) )
return cur.fetchone() is not None return cur.fetchone() is not None
def count_target(self, schema: str, table: str) -> int:
"""Return ``COUNT(*)`` for ``[schema].[table]`` (compare count check)."""
cur = self._conn.cursor()
cur.execute(f"SELECT COUNT(*) FROM [{schema}].[{table}]")
return cur.fetchone()[0]
def read_target_ids(self, schema: str, table: str) -> list:
"""Return every ``ID`` from ``[schema].[table]``, ascending (compare IDs)."""
cur = self._conn.cursor()
cur.execute(f"SELECT ID FROM [{schema}].[{table}] ORDER BY ID")
return [r[0] for r in cur.fetchall()]
def _has_identity(self, schema: str, table: str) -> bool: def _has_identity(self, schema: str, table: str) -> bool:
"""True if ``[schema].[table]`` has an IDENTITY column (the ``ID`` PK).""" """True if ``[schema].[table]`` has an IDENTITY column (the ``ID`` PK)."""
cur = self._conn.cursor() cur = self._conn.cursor()

37
src/sync/targets.py Normal file
View File

@@ -0,0 +1,37 @@
"""Shared target-table resolution for fullsync, incremental capture, and compare.
All three pipelines must agree on which Access tables are in scope and how each
maps to its SQL Server target. Centralising the exclude/include rule here makes
that guarantee structural instead of duplicated across three files.
- ``is_synced_table`` applies the per-file exclude/include rule (exclude wins).
- ``resolve_synced_tables`` returns the in-scope user tables in
``list_user_tables`` order; system tables (``MSys*`` / ``~*``) are already
filtered by ``AccessReader.list_user_tables``.
Target *naming* is shared via ``FileMapping.target_table`` (name + year_suffix)
and ``FileMapping.schema``, so a given Access table resolves to the same
``(schema, table)`` everywhere.
"""
from __future__ import annotations
from .config import FileMapping
def is_synced_table(fm: FileMapping, table_name: str) -> bool:
"""True if ``table_name`` is in sync scope for ``fm``.
``exclude_tables`` wins over ``include_tables``: a table listed in both is
excluded. When ``include_tables`` is None, every non-excluded table is in
scope.
"""
if table_name in (fm.exclude_tables or []):
return False
if fm.include_tables is not None:
return table_name in fm.include_tables
return True
def resolve_synced_tables(fm: FileMapping, reader) -> list[str]:
"""In-scope user tables for ``fm``, in ``list_user_tables`` order."""
return [t for t in reader.list_user_tables() if is_synced_table(fm, t)]

View File

@@ -6,6 +6,7 @@ No real log rows are deleted.
""" """
import os import os
import pytest import pytest
from unittest.mock import MagicMock
from sync.access_reader import AccessReader from sync.access_reader import AccessReader
from sync.config import load_config from sync.config import load_config
@@ -35,3 +36,23 @@ def test_read_log_and_row_and_delete():
r.delete_log_ids([], 100, 3) # empty list -> no-op, must not raise r.delete_log_ids([], 100, 3) # empty list -> no-op, must not raise
finally: finally:
r.close() r.close()
def test_count_rows_executes_count_sql_and_returns_value():
r = AccessReader("dummy.accdb", "{Microsoft Access Driver (*.accdb, *.mdb)}")
cur = MagicMock()
cur.fetchone.return_value = (42,)
r._conn = MagicMock() # bypass lazy connect
r._conn.cursor.return_value = cur
assert r.count_rows("表壳焊接记录") == 42
cur.execute.assert_called_once_with('SELECT COUNT(*) FROM "表壳焊接记录"')
def test_read_ids_returns_ordered_id_list():
r = AccessReader("dummy.accdb", "{Microsoft Access Driver (*.accdb, *.mdb)}")
cur = MagicMock()
cur.fetchall.return_value = [(1,), (3,), (5,)]
r._conn = MagicMock()
r._conn.cursor.return_value = cur
assert r.read_ids("T") == [1, 3, 5]
cur.execute.assert_called_once_with('SELECT ID FROM "T" ORDER BY ID')

View File

@@ -1,15 +1,23 @@
"""Integration test for dbo.usp_SyncApply. """Integration test for usp_SyncApply.
Validates: IDENTITY-preserving INSERT, last-write-wins UPDATE, BIT conversion Validates: IDENTITY-preserving INSERT, last-write-wins UPDATE, BIT conversion
from JSON ``true``/``false``, and DELETE of the last op. Creates a throwaway from JSON ``true``/``false``, and DELETE of the last op. Creates a throwaway
schema ``sync_test`` and table ``ApplyDemo_YEAR2026`` and cleans them up at the schema ``sync_test`` and table ``ApplyDemo_YEAR2026`` and cleans them up at the
end so no residue is left on CompanyDB. end so no residue is left on CompanyDB.
Queue table and apply proc names are read from ``config.yaml`` so the test
follows whichever schema the deployment targets (currently ProductionDataBaseSync).
""" """
import pytest import pytest
from sync.config import load_config
@pytest.mark.integration @pytest.mark.integration
def test_upsert_insert_update_delete_last_write_wins(sql_conn): def test_upsert_insert_update_delete_last_write_wins(sql_conn):
cfg = load_config("config.yaml")
qt = cfg.sql_server.sync_queue_table
proc = cfg.sql_server.apply_proc
cur = sql_conn.cursor() cur = sql_conn.cursor()
sch, tbl = "sync_test", "ApplyDemo_YEAR2026" sch, tbl = "sync_test", "ApplyDemo_YEAR2026"
@@ -27,17 +35,17 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
"时间 DATETIME2 NULL, " "时间 DATETIME2 NULL, "
"标记 BIT NULL)" "标记 BIT NULL)"
) )
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'") cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
# Insert then a later Update for the same RecordID=1 -> last write wins. # Insert then a later Update for the same RecordID=1 -> last write wins.
cur.execute( cur.execute(
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema," "INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
"TargetTable,RecordID,OperateType,RowData,Status) " "TargetTable,RecordID,OperateType,RowData,Status) "
"VALUES('t.accdb','ApplyDemo',1,'sync_test','ApplyDemo_YEAR2026','1'," "VALUES('t.accdb','ApplyDemo',1,'sync_test','ApplyDemo_YEAR2026','1',"
"'Insert','{\"ID\":1,\"名字\":\"A\",\"数量\":3,\"时间\":\"2026-01-01T00:00:00\",\"标记\":true}','pending')" "'Insert','{\"ID\":1,\"名字\":\"A\",\"数量\":3,\"时间\":\"2026-01-01T00:00:00\",\"标记\":true}','pending')"
) )
cur.execute( cur.execute(
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema," "INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
"TargetTable,RecordID,OperateType,RowData,Status) " "TargetTable,RecordID,OperateType,RowData,Status) "
"VALUES('t.accdb','ApplyDemo',2,'sync_test','ApplyDemo_YEAR2026','1'," "VALUES('t.accdb','ApplyDemo',2,'sync_test','ApplyDemo_YEAR2026','1',"
"'Update','{\"ID\":1,\"名字\":\"A2\",\"数量\":5,\"时间\":\"2026-01-02T00:00:00\",\"标记\":false}','pending')" "'Update','{\"ID\":1,\"名字\":\"A2\",\"数量\":5,\"时间\":\"2026-01-02T00:00:00\",\"标记\":false}','pending')"
@@ -45,7 +53,7 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
sql_conn.commit() sql_conn.commit()
# --- Act 1: apply upsert ---------------------------------------------- # --- Act 1: apply upsert ----------------------------------------------
cur.execute("EXEC dbo.usp_SyncApply @MaxRetries=5") cur.execute("EXEC " + proc + " @MaxRetries=5")
sql_conn.commit() sql_conn.commit()
# --- Assert 1: the later Update wins; BIT false -> 0 ------------------ # --- Assert 1: the later Update wins; BIT false -> 0 ------------------
@@ -57,15 +65,15 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
assert row.标记 == 0 # BIT false assert row.标记 == 0 # BIT false
# --- Act 2: a later Delete wins --------------------------------------- # --- Act 2: a later Delete wins ---------------------------------------
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'") cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
cur.execute( cur.execute(
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,TargetSchema," "INSERT " + qt + "(SourceFile,SourceTable,SourceLogID,TargetSchema,"
"TargetTable,RecordID,OperateType,RowData,Status) " "TargetTable,RecordID,OperateType,RowData,Status) "
"VALUES('t.accdb','ApplyDemo',3,'sync_test','ApplyDemo_YEAR2026','1'," "VALUES('t.accdb','ApplyDemo',3,'sync_test','ApplyDemo_YEAR2026','1',"
"'Delete',NULL,'pending')" "'Delete',NULL,'pending')"
) )
sql_conn.commit() sql_conn.commit()
cur.execute("EXEC dbo.usp_SyncApply @MaxRetries=5") cur.execute("EXEC " + proc + " @MaxRetries=5")
sql_conn.commit() sql_conn.commit()
# --- Assert 2: row removed -------------------------------------------- # --- Assert 2: row removed --------------------------------------------
@@ -77,5 +85,5 @@ def test_upsert_insert_update_delete_last_write_wins(sql_conn):
"IF OBJECT_ID('sync_test.ApplyDemo_YEAR2026') IS NOT NULL " "IF OBJECT_ID('sync_test.ApplyDemo_YEAR2026') IS NOT NULL "
"DROP TABLE sync_test.ApplyDemo_YEAR2026" "DROP TABLE sync_test.ApplyDemo_YEAR2026"
) )
cur.execute("DELETE dbo.SyncQueue WHERE TargetSchema='sync_test'") cur.execute("DELETE " + qt + " WHERE TargetSchema='sync_test'")
sql_conn.commit() sql_conn.commit()

105
tests/test_compare.py Normal file
View File

@@ -0,0 +1,105 @@
from unittest.mock import MagicMock
from sync.config import FileMapping
from sync.compare import compare_file, any_mismatch, format_report, TableResult
def _fm(**kw):
base = dict(file="x.accdb", root="2026", schema="s", year_suffix="_YEAR2026")
base.update(kw)
return FileMapping(**base)
def _reader_with_tables(tables):
r = MagicMock()
r.list_user_tables.return_value = tables
return r
def test_count_match():
fm = _fm(exclude_tables=["TableChangeLog"])
reader = _reader_with_tables(["T1", "TableChangeLog"])
reader.count_rows.return_value = 5
writer = MagicMock()
writer.table_exists.return_value = True
writer.count_target.return_value = 5
res = compare_file(fm, reader, writer, "count")
assert len(res) == 1
assert res[0].access_table == "T1"
assert res[0].status == "match"
assert res[0].access_count == 5 and res[0].sql_count == 5
def test_count_mismatch():
fm = _fm()
reader = _reader_with_tables(["T1"])
reader.count_rows.return_value = 5
writer = MagicMock()
writer.table_exists.return_value = True
writer.count_target.return_value = 7
res = compare_file(fm, reader, writer, "count")
assert res[0].status == "mismatch"
def test_skip_when_mirror_missing():
fm = _fm()
reader = _reader_with_tables(["T1"])
writer = MagicMock()
writer.table_exists.return_value = False
res = compare_file(fm, reader, writer, "count")
assert res[0].status == "skipped"
writer.count_target.assert_not_called()
writer.read_target_ids.assert_not_called()
def test_ids_match():
fm = _fm()
reader = _reader_with_tables(["T1"])
reader.read_ids.return_value = [1, 2, 3]
writer = MagicMock()
writer.table_exists.return_value = True
writer.read_target_ids.return_value = [1, 2, 3]
res = compare_file(fm, reader, writer, "ids")
assert res[0].status == "match"
assert res[0].missing_in_sql == []
assert res[0].extra_in_sql == []
def test_ids_reports_missing_and_extra():
fm = _fm()
reader = _reader_with_tables(["T1"])
reader.read_ids.return_value = [1, 2, 3]
writer = MagicMock()
writer.table_exists.return_value = True
writer.read_target_ids.return_value = [2, 3, 4]
res = compare_file(fm, reader, writer, "ids")
assert res[0].status == "mismatch"
assert res[0].missing_in_sql == [1] # in Access, not in SQL
assert res[0].extra_in_sql == [4] # in SQL, not in Access
def test_excluded_tables_not_compared():
fm = _fm(exclude_tables=["TableChangeLog"])
reader = _reader_with_tables(["T1", "TableChangeLog"])
reader.count_rows.return_value = 1
writer = MagicMock()
writer.table_exists.return_value = True
writer.count_target.return_value = 1
res = compare_file(fm, reader, writer, "count")
assert [r.access_table for r in res] == ["T1"]
def test_any_mismatch_detects_mismatch_only():
r_match = TableResult("f", "T", "s", "T_YEAR2026", "match", 1, 1)
r_skip = TableResult("f", "T2", "s", "T2_YEAR2026", "skipped")
r_mis = TableResult("f", "T3", "s", "T3_YEAR2026", "mismatch", 1, 2)
assert any_mismatch([r_match, r_skip]) is False
assert any_mismatch([r_match, r_mis]) is True
def test_format_report_count():
r = TableResult("x.accdb", "T1", "s", "T1_YEAR2026", "mismatch", 5, 7)
rep = format_report([r], "count")
assert "x.accdb: T1 -> s.T1_YEAR2026" in rep
assert "access=5 sql=7" in rep
assert "[MISMATCH]" in rep

70
tests/test_main.py Normal file
View File

@@ -0,0 +1,70 @@
"""Routing/argparse tests for the root main.py dispatcher.
The backend functions (full_sync, service.cycle/run, compare) are mocked so
these tests verify command dispatch, arg parsing, report output and exit codes
without touching Access or SQL Server.
"""
from unittest.mock import patch
import main
from sync.compare import TableResult
def test_fullsync_routes_with_filters():
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync") as fs, patch("main.service") as svc:
rc = main.main(["fullsync", "--db", "OEM.accdb", "--table", "表壳焊接记录",
"--clear-change-log"])
fs.assert_called_once()
assert fs.call_args.kwargs["db_filter"] == "OEM.accdb"
assert fs.call_args.kwargs["table_filter"] == "表壳焊接记录"
assert fs.call_args.kwargs["clear_change_log"] is True
svc.cycle.assert_not_called()
assert rc == 0
def test_incremental_default_runs_single_cycle():
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync"), patch("main.service") as svc:
rc = main.main(["incremental"])
svc.cycle.assert_called_once()
svc.run.assert_not_called()
assert rc == 0
def test_incremental_loop_runs_service():
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync"), patch("main.service") as svc:
rc = main.main(["incremental", "--loop", "--poll-interval", "5"])
svc.run.assert_called_once()
svc.cycle.assert_not_called()
assert rc == 0
def test_compare_default_granularity_is_count(capsys):
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync"), patch("main.service"), patch("main.compare") as cmp:
cmp.return_value = [TableResult("x.accdb", "T", "s", "T_YEAR2026", "match", 5, 5)]
rc = main.main(["compare"])
assert cmp.call_args.kwargs["granularity"] == "count"
assert rc == 0
assert "x.accdb: T -> s.T_YEAR2026" in capsys.readouterr().out
def test_compare_ids_granularity_and_mismatch_exit_code():
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync"), patch("main.service"), patch("main.compare") as cmp:
cmp.return_value = [TableResult("x.accdb", "T", "s", "T_YEAR2026", "mismatch", 5, 7)]
rc = main.main(["compare", "--granularity", "ids"])
assert cmp.call_args.kwargs["granularity"] == "ids"
assert rc == 1
def test_compare_writes_report_file(tmp_path):
with patch("main.load_config"), patch("main.setup_logging"), \
patch("main.full_sync"), patch("main.service"), patch("main.compare") as cmp:
cmp.return_value = [TableResult("x.accdb", "T", "s", "T_YEAR2026", "match", 1, 1)]
rep = tmp_path / "report.txt"
rc = main.main(["compare", "--report", str(rep)])
assert rc == 0
assert "x.accdb: T -> s.T_YEAR2026" in rep.read_text(encoding="utf-8")

View File

@@ -10,6 +10,7 @@ credentials are hardcoded here. The test self-cleans using a throwaway
""" """
import os import os
import pytest import pytest
from unittest.mock import MagicMock
from sync.sql_writer import SqlWriter, QueueRow from sync.sql_writer import SqlWriter, QueueRow
from sync.config import load_config from sync.config import load_config
@@ -20,10 +21,16 @@ def test_insert_dedup_and_applied_ids():
if not os.environ.get("RUN_INTEGRATION"): if not os.environ.get("RUN_INTEGRATION"):
pytest.skip("integration") pytest.skip("integration")
cfg = load_config("config.yaml") cfg = load_config("config.yaml")
w = SqlWriter(cfg.sql_server.conn_str, "dbo.SyncQueue") qt = cfg.sql_server.sync_queue_table
w = SqlWriter(
cfg.sql_server.conn_str,
qt,
cfg.sql_server.archive_table,
cfg.sql_server.apply_proc,
)
try: try:
cur = w._conn.cursor() cur = w._conn.cursor()
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'") cur.execute(f"DELETE {qt} WHERE SourceFile='sqlw_test.accdb'")
qr = QueueRow( qr = QueueRow(
source_file="sqlw_test.accdb", source_file="sqlw_test.accdb",
source_table="T", source_table="T",
@@ -37,7 +44,7 @@ def test_insert_dedup_and_applied_ids():
w.insert_queue_row(qr) w.insert_queue_row(qr)
w.insert_queue_row(qr) # duplicate must be deduped (ignored) w.insert_queue_row(qr) # duplicate must be deduped (ignored)
cur.execute( cur.execute(
"SELECT COUNT(*) FROM dbo.SyncQueue " f"SELECT COUNT(*) FROM {qt} "
"WHERE SourceFile='sqlw_test.accdb' AND SourceLogID=100" "WHERE SourceFile='sqlw_test.accdb' AND SourceLogID=100"
) )
assert cur.fetchone()[0] == 1 assert cur.fetchone()[0] == 1
@@ -46,10 +53,37 @@ def test_insert_dedup_and_applied_ids():
w.call_apply(max_retries=5) w.call_apply(max_retries=5)
cur.execute( cur.execute(
"UPDATE dbo.SyncQueue SET Status='applied' " f"UPDATE {qt} SET Status='applied' "
"WHERE SourceFile='sqlw_test.accdb'" "WHERE SourceFile='sqlw_test.accdb'"
) )
assert w.applied_log_ids("sqlw_test.accdb") == [100] assert w.applied_log_ids("sqlw_test.accdb") == [100]
finally: finally:
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'") cur.execute(f"DELETE {qt} WHERE SourceFile='sqlw_test.accdb'")
w.close() w.close()
def _writer_with_cursor(fetchone=None, fetchall=None):
"""A SqlWriter whose pyodbc connection is a mock (no real connect)."""
w = SqlWriter.__new__(SqlWriter)
w.conn_str = "dummy"
w.queue_table = "dbo.SyncQueue"
cur = MagicMock()
if fetchone is not None:
cur.fetchone.return_value = fetchone
if fetchall is not None:
cur.fetchall.return_value = fetchall
w._conn = MagicMock()
w._conn.cursor.return_value = cur
return w, cur
def test_count_target_executes_count_sql_and_returns_value():
w, cur = _writer_with_cursor(fetchone=(7,))
assert w.count_target("s", "T_YEAR2026") == 7
cur.execute.assert_called_once_with("SELECT COUNT(*) FROM [s].[T_YEAR2026]")
def test_read_target_ids_returns_ordered_id_list():
w, cur = _writer_with_cursor(fetchall=[(2,), (4,), (6,)])
assert w.read_target_ids("s", "T_YEAR2026") == [2, 4, 6]
cur.execute.assert_called_once_with("SELECT ID FROM [s].[T_YEAR2026] ORDER BY ID")

53
tests/test_targets.py Normal file
View File

@@ -0,0 +1,53 @@
from unittest.mock import MagicMock
from sync.config import FileMapping
from sync.targets import is_synced_table, resolve_synced_tables
def _fm(**kw):
base = dict(file="x.accdb", root="2026", schema="s", year_suffix="_YEAR2026")
base.update(kw)
return FileMapping(**base)
def test_is_synced_table_excludes_listed():
fm = _fm(exclude_tables=["TableChangeLog", "一车间每日催货落实记录_停"])
assert is_synced_table(fm, "TableChangeLog") is False
assert is_synced_table(fm, "一车间每日催货落实记录_停") is False
assert is_synced_table(fm, "表壳焊接记录") is True
def test_is_synced_table_include_restricts():
fm = _fm(exclude_tables=["TableChangeLog"], include_tables=["检验合格记录表"])
assert is_synced_table(fm, "检验合格记录表") is True
assert is_synced_table(fm, "其它表") is False
def test_is_synced_table_exclude_beats_include():
fm = _fm(exclude_tables=["TableChangeLog"],
include_tables=["TableChangeLog", "检验合格记录表"])
assert is_synced_table(fm, "TableChangeLog") is False
assert is_synced_table(fm, "检验合格记录表") is True
def test_is_synced_table_no_filters_includes_all():
fm = _fm()
assert is_synced_table(fm, "任意表") is True
def test_resolve_synced_tables_filters_user_tables():
fm = _fm(exclude_tables=["TableChangeLog"])
reader = MagicMock()
reader.list_user_tables.return_value = [
"TableChangeLog", "表壳焊接记录", "超压", "氩弧焊每日催货落实记录_停",
]
assert resolve_synced_tables(fm, reader) == [
"表壳焊接记录", "超压", "氩弧焊每日催货落实记录_停",
]
def test_resolve_synced_tables_with_include():
fm = _fm(exclude_tables=["TableChangeLog"], include_tables=["检验合格记录表"])
reader = MagicMock()
reader.list_user_tables.return_value = ["TableChangeLog", "检验合格记录表", "其它表"]
assert resolve_synced_tables(fm, reader) == ["检验合格记录表"]