diff --git a/src/sync/sql_writer.py b/src/sync/sql_writer.py new file mode 100644 index 0000000..9212d6c --- /dev/null +++ b/src/sync/sql_writer.py @@ -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() diff --git a/tests/test_sql_writer.py b/tests/test_sql_writer.py new file mode 100644 index 0000000..6fc4101 --- /dev/null +++ b/tests/test_sql_writer.py @@ -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()