feat: sql writer with dedup insert and apply call
Adds SqlWriter: a pyodbc-backed writer that dedup-inserts into dbo.SyncQueue (IF NOT EXISTS guarded by UX_SyncQueue_Dedup), invokes dbo.usp_SyncApply, and reports applied SourceLogIDs. Connection is opened with autocommit=True per the controller revision: usp_SyncApply manages its own transaction internally (BEGIN/ROLLBACK), and an outer pyodbc transaction would conflict on ROLLBACK (SQL error 266). The dedup IF NOT EXISTS...INSERT is a single atomic statement. Integration test self-cleans via SourceFile='sqlw_test.accdb' marker; conn_str comes from the gitignored config.yaml (no hardcoded creds). Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
95
src/sync/sql_writer.py
Normal file
95
src/sync/sql_writer.py
Normal file
@@ -0,0 +1,95 @@
|
|||||||
|
"""SQL-side writer for the Access -> SQL Server sync.
|
||||||
|
|
||||||
|
SqlWriter owns the pyodbc connection used to (a) dedup-insert rows into
|
||||||
|
``dbo.SyncQueue``, (b) invoke ``dbo.usp_SyncApply`` to drain the queue, and
|
||||||
|
(c) report which SourceLogIDs have been applied.
|
||||||
|
|
||||||
|
The connection is opened with ``autocommit=True`` on purpose. ``usp_SyncApply``
|
||||||
|
manages its own transaction internally (BEGIN TRAN ... ROLLBACK on error); if
|
||||||
|
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
|
||||||
|
single statement that is atomic under autocommit, so no explicit transaction is
|
||||||
|
needed on the write path either.
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
import pyodbc
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class QueueRow:
|
||||||
|
"""A single staged change to enqueue into ``dbo.SyncQueue``."""
|
||||||
|
|
||||||
|
source_file: str
|
||||||
|
source_table: str
|
||||||
|
record_id: str
|
||||||
|
target_schema: str
|
||||||
|
target_table: str
|
||||||
|
source_log_id: int
|
||||||
|
operate_type: str
|
||||||
|
row_data: str | None
|
||||||
|
|
||||||
|
|
||||||
|
class SqlWriter:
|
||||||
|
"""Writes to SyncQueue and drives the apply proc over a pyodbc connection."""
|
||||||
|
|
||||||
|
def __init__(self, conn_str: str, queue_table: str = "dbo.SyncQueue"):
|
||||||
|
self.conn_str = conn_str
|
||||||
|
self.queue_table = queue_table
|
||||||
|
# autocommit=True: usp_SyncApply manages its own transaction internally.
|
||||||
|
# An outer pyodbc transaction would conflict on ROLLBACK (SQL error 266).
|
||||||
|
self._conn = pyodbc.connect(conn_str, autocommit=True)
|
||||||
|
|
||||||
|
def insert_queue_row(self, row: QueueRow) -> None:
|
||||||
|
"""Idempotently enqueue ``row`` (dedup on SourceFile/Table/LogID).
|
||||||
|
|
||||||
|
``IF NOT EXISTS ... INSERT`` is a single statement, atomic under
|
||||||
|
autocommit. 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.execute(
|
||||||
|
"IF NOT EXISTS (SELECT 1 FROM dbo.SyncQueue "
|
||||||
|
"WHERE SourceFile=? AND SourceTable=? AND SourceLogID=?) "
|
||||||
|
"INSERT dbo.SyncQueue(SourceFile,SourceTable,SourceLogID,"
|
||||||
|
"TargetSchema,TargetTable,RecordID,OperateType,RowData,Status) "
|
||||||
|
"VALUES (?,?,?,?,?,?,?,?, 'pending')",
|
||||||
|
row.source_file,
|
||||||
|
row.source_table,
|
||||||
|
row.source_log_id,
|
||||||
|
row.source_file,
|
||||||
|
row.source_table,
|
||||||
|
row.source_log_id,
|
||||||
|
row.target_schema,
|
||||||
|
row.target_table,
|
||||||
|
row.record_id,
|
||||||
|
row.operate_type,
|
||||||
|
row.row_data,
|
||||||
|
)
|
||||||
|
# autocommit: statement already committed.
|
||||||
|
|
||||||
|
def call_apply(self, max_retries: int) -> None:
|
||||||
|
"""Drain the pending queue via the stored procedure.
|
||||||
|
|
||||||
|
``usp_SyncApply`` flips rows to ``applied`` (or ``error`` after retries).
|
||||||
|
"""
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
cur.execute("EXEC dbo.usp_SyncApply ?", max_retries)
|
||||||
|
|
||||||
|
def applied_log_ids(self, source_file: str) -> list[int]:
|
||||||
|
"""Return applied SourceLogIDs for ``source_file`` in ascending order."""
|
||||||
|
cur = self._conn.cursor()
|
||||||
|
cur.execute(
|
||||||
|
"SELECT SourceLogID FROM dbo.SyncQueue "
|
||||||
|
"WHERE SourceFile=? AND Status='applied' ORDER BY SourceLogID",
|
||||||
|
source_file,
|
||||||
|
)
|
||||||
|
return [r[0] for r in cur.fetchall()]
|
||||||
|
|
||||||
|
def close(self) -> None:
|
||||||
|
"""Close the underlying pyodbc connection."""
|
||||||
|
self._conn.close()
|
||||||
55
tests/test_sql_writer.py
Normal file
55
tests/test_sql_writer.py
Normal file
@@ -0,0 +1,55 @@
|
|||||||
|
"""Integration test for sync.sql_writer.SqlWriter.
|
||||||
|
|
||||||
|
Validates the dedup INSERT (same SourceLogID is only ever enqueued once), the
|
||||||
|
``applied_log_ids`` query (rows flipped to ``applied`` come back in order), and
|
||||||
|
that ``call_apply`` invokes the proc without raising.
|
||||||
|
|
||||||
|
conn_str is read from the gitignored ``config.yaml`` via ``load_config``; no
|
||||||
|
credentials are hardcoded here. The test self-cleans using a throwaway
|
||||||
|
``SourceFile='sqlw_test.accdb'`` marker so no residue is left on SyncQueue.
|
||||||
|
"""
|
||||||
|
import os
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from sync.sql_writer import SqlWriter, QueueRow
|
||||||
|
from sync.config import load_config
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.integration
|
||||||
|
def test_insert_dedup_and_applied_ids():
|
||||||
|
if not os.environ.get("RUN_INTEGRATION"):
|
||||||
|
pytest.skip("integration")
|
||||||
|
cfg = load_config("config.yaml")
|
||||||
|
w = SqlWriter(cfg.sql_server.conn_str, "dbo.SyncQueue")
|
||||||
|
try:
|
||||||
|
cur = w._conn.cursor()
|
||||||
|
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'")
|
||||||
|
qr = QueueRow(
|
||||||
|
source_file="sqlw_test.accdb",
|
||||||
|
source_table="T",
|
||||||
|
record_id="7",
|
||||||
|
target_schema="sync_test",
|
||||||
|
target_table="T_YEAR2026",
|
||||||
|
source_log_id=100,
|
||||||
|
operate_type="Insert",
|
||||||
|
row_data='{"ID":7}',
|
||||||
|
)
|
||||||
|
w.insert_queue_row(qr)
|
||||||
|
w.insert_queue_row(qr) # duplicate must be deduped (ignored)
|
||||||
|
cur.execute(
|
||||||
|
"SELECT COUNT(*) FROM dbo.SyncQueue "
|
||||||
|
"WHERE SourceFile='sqlw_test.accdb' AND SourceLogID=100"
|
||||||
|
)
|
||||||
|
assert cur.fetchone()[0] == 1
|
||||||
|
|
||||||
|
# call_apply must run without error (proc logic validated elsewhere).
|
||||||
|
w.call_apply(max_retries=5)
|
||||||
|
|
||||||
|
cur.execute(
|
||||||
|
"UPDATE dbo.SyncQueue SET Status='applied' "
|
||||||
|
"WHERE SourceFile='sqlw_test.accdb'"
|
||||||
|
)
|
||||||
|
assert w.applied_log_ids("sqlw_test.accdb") == [100]
|
||||||
|
finally:
|
||||||
|
cur.execute("DELETE dbo.SyncQueue WHERE SourceFile='sqlw_test.accdb'")
|
||||||
|
w.close()
|
||||||
Reference in New Issue
Block a user