using System; using System.Collections.Generic; using System.Data.OleDb; using System.Diagnostics; using System.Globalization; using System.Linq; using Microsoft.Data.SqlClient; using Serilog; using SyncEngine.Access; using SyncEngine.Configuration; using SyncEngine.Logging; using SyncEngine.Metadata; using SyncEngine.Sync; namespace SyncEngine.SqlServer { public sealed class SqlDataSync { /// SQL Server allows 2100 parameters; leave headroom for @seen on touch updates. private const int MaxInClauseParameters = 2000; private readonly AccessConnectionFactory _accessFactory; private readonly SqlServerConnectionFactory _sqlFactory; private readonly SyncOptions _options; private readonly SyncMetadataRepository _metadata; private readonly SoftDeleteTracker _softDelete; private readonly ILogger _logger; public SqlDataSync( AccessConnectionFactory accessFactory, SqlServerConnectionFactory sqlFactory, SyncOptions options, SyncMetadataRepository metadata, SoftDeleteTracker softDelete, ILogger logger) { _accessFactory = accessFactory; _sqlFactory = sqlFactory; _options = options; _metadata = metadata; _softDelete = softDelete; _logger = logger; } public void SyncTables( SqlConnection sqlConn, IList tables, Guid runId, SyncStatistics stats, DateTime runStart) { using (var accessConn = _accessFactory.Create()) { accessConn.Open(); foreach (var table in tables) { SyncTable(sqlConn, accessConn, table, runId, stats, runStart); } } } private void SyncTable( SqlConnection sqlConn, OleDbConnection accessConn, AccessTableSchema table, Guid runId, SyncStatistics stats, DateTime runStart) { var pkColumns = DataSyncPrimaryKeyHelper.GetPrimaryKeyColumns(table); if (pkColumns.Count == 0) { _logger.Warning("Skipping {Table}: no primary key detected", table.TableName); stats.TablesSkipped++; return; } if (pkColumns.Count > 1) { _logger.Debug( "Table {Table} has composite primary key ({Columns})", table.TableName, string.Join(", ", pkColumns)); } var sw = Stopwatch.StartNew(); _logger.Information("Syncing table {Table}...", table.TableName); var qualified = _sqlFactory.Qualify(table.TableName); var dataColumns = table.Columns.Select(c => c.Name).ToList(); var inserted = 0; var updated = 0; var batchCount = 0; var orderBy = DataSyncPrimaryKeyHelper.BuildOrderBy(pkColumns); Dictionary lastRow = null; var where = string.Empty; while (true) { var batch = ReadBatch(accessConn, table, orderBy, where, _options.BatchSize); if (batch.Count == 0) { break; } batchCount++; if (!_options.DryRun) { BatchProcessor.ExecuteWithRetry( () => UpsertBatch(sqlConn, qualified, table, pkColumns, dataColumns, batch, runStart, ref inserted, ref updated), _options.MaxRetries, _logger, table.TableName); } lastRow = batch[batch.Count - 1]; where = DataSyncPrimaryKeyHelper.BuildAccessResumeWhere(pkColumns, lastRow); } if (!_options.DryRun) { var checkpoint = lastRow == null ? null : (object)DataSyncPrimaryKeyHelper.BuildPkKey(lastRow, pkColumns); _metadata.SaveCheckpoint(sqlConn, runId, table.TableName, checkpoint, inserted, updated, 0); ReseedIdentityIfNeeded(sqlConn, qualified, table); } var softDeleted = _softDelete.MarkMissingRows(sqlConn, qualified, table.PrimaryKeyColumn, runStart); stats.RowsInserted += inserted; stats.RowsUpdated += updated; stats.RowsSoftDeleted += softDeleted; stats.TablesProcessed++; _logger.Information( "Finished {Table} in {Elapsed:F1}s: +{Ins} ~{Upd} -{Del} ({Batches} batches)", table.TableName, sw.Elapsed.TotalSeconds, inserted, updated, softDeleted, batchCount); } private void ReseedIdentityIfNeeded(SqlConnection sqlConn, string qualifiedTable, AccessTableSchema table) { var identityCol = SqlIdentityHelper.GetIdentityColumn(table); if (identityCol == null) { return; } var maxValue = SqlIdentityHelper.GetMaxColumnValue(sqlConn, qualifiedTable, identityCol.Name); SqlIdentityHelper.ReseedIdentity(sqlConn, qualifiedTable, maxValue ?? 0); } private static List> ReadBatch( OleDbConnection conn, AccessTableSchema table, string orderBy, string whereClause, int batchSize) { var result = new List>(); var sql = string.Format( "SELECT TOP {0} * FROM [{1}]{2} ORDER BY {3}", batchSize, table.TableName, whereClause, orderBy); using (var cmd = new OleDbCommand(sql, conn)) using (var reader = cmd.ExecuteReader()) { while (reader.Read()) { var row = new Dictionary(StringComparer.OrdinalIgnoreCase); for (var i = 0; i < reader.FieldCount; i++) { row[reader.GetName(i)] = reader.IsDBNull(i) ? null : reader.GetValue(i); } result.Add(row); } } return result; } private void UpsertBatch( SqlConnection sqlConn, string qualifiedTable, AccessTableSchema table, IList pkColumns, IList dataColumns, List> batch, DateTime runStart, ref int inserted, ref int updated) { using (var tx = sqlConn.BeginTransaction()) { try { var identityCol = SqlIdentityHelper.GetIdentityColumn(table); var useIdentityInsert = identityCol != null; if (useIdentityInsert) { SqlIdentityHelper.SetIdentityInsert(sqlConn, tx, qualifiedTable, true); } var existing = LoadExistingRowHashes(sqlConn, tx, qualifiedTable, pkColumns, batch); foreach (var row in batch) { var values = dataColumns.Select(c => row.ContainsKey(c) ? row[c] : null).ToList(); var hash = ChangeDetector.ComputeRowHash(values); var pkKey = DataSyncPrimaryKeyHelper.BuildPkKey(row, pkColumns); if (!existing.TryGetValue(pkKey, out var existingHash)) { InsertRow(sqlConn, tx, qualifiedTable, dataColumns, row, hash, runStart); inserted++; } else if (!string.Equals(existingHash, hash, StringComparison.Ordinal)) { if (UpdateRowIfChanged(sqlConn, tx, qualifiedTable, pkColumns, dataColumns, row, hash, runStart)) { updated++; } } } // Every row in the batch was read from Access this run — refresh last_seen even when // hash comparison skipped the data UPDATE (e.g. CHAR(64) hash padding in SQL Server). MarkRowsSeen(sqlConn, tx, qualifiedTable, pkColumns, batch, runStart); if (useIdentityInsert) { SqlIdentityHelper.SetIdentityInsert(sqlConn, tx, qualifiedTable, false); } tx.Commit(); } catch { tx.Rollback(); throw; } } } private Dictionary LoadExistingRowHashes( SqlConnection conn, SqlTransaction tx, string table, IList pkColumns, List> batch) { var result = new Dictionary(StringComparer.OrdinalIgnoreCase); if (batch.Count == 0) { return result; } var firstCol = pkColumns[0]; var firstValues = batch .Select(r => r.ContainsKey(firstCol) ? r[firstCol] : null) .Distinct() .ToList(); var selectColumns = string.Join( ", ", pkColumns.Select(c => "[" + c + "]").Concat(new[] { "_sync_row_hash" })); for (var offset = 0; offset < firstValues.Count; offset += MaxInClauseParameters) { var chunk = firstValues.Skip(offset).Take(MaxInClauseParameters).ToList(); var paramNames = chunk.Select((_, i) => "@pk" + i).ToList(); var sql = string.Format( "SELECT {0} FROM {1} WHERE [{2}] IN ({3})", selectColumns, table, firstCol, string.Join(", ", paramNames)); using (var cmd = new SqlCommand(sql, conn, tx)) { for (var i = 0; i < chunk.Count; i++) { cmd.Parameters.AddWithValue("@pk" + i, chunk[i] ?? DBNull.Value); } cmd.CommandTimeout = _options.CommandTimeoutSeconds; using (var reader = cmd.ExecuteReader()) { while (reader.Read()) { var keyRow = new Dictionary(StringComparer.OrdinalIgnoreCase); for (var i = 0; i < pkColumns.Count; i++) { keyRow[pkColumns[i]] = reader.IsDBNull(i) ? null : reader.GetValue(i); } result[DataSyncPrimaryKeyHelper.BuildPkKey(keyRow, pkColumns)] = ReadStoredRowHash(reader, pkColumns.Count); } } } } return result; } private static string ReadStoredRowHash(Microsoft.Data.SqlClient.SqlDataReader reader, int hashColumnIndex) { if (reader.IsDBNull(hashColumnIndex)) { return null; } // _sync_row_hash may be CHAR(64) on older mirrors — trim trailing spaces before compare. return Convert.ToString(reader.GetValue(hashColumnIndex), CultureInfo.InvariantCulture).TrimEnd(); } private void MarkRowsSeen( SqlConnection conn, SqlTransaction tx, string table, IList pkColumns, IList> rows, DateTime runStart) { if (rows.Count == 0) { return; } var now = DateTime.UtcNow; var whereClause = DataSyncPrimaryKeyHelper.BuildSqlWhereClause(pkColumns, "@pk"); var sql = string.Format( "UPDATE {0} SET _sync_last_seen_at = @seen, _sync_is_deleted = 0, _sync_updated_at = @now WHERE {1}", table, whereClause); foreach (var row in rows) { using (var cmd = new SqlCommand(sql, conn, tx)) { cmd.Parameters.AddWithValue("@seen", runStart); cmd.Parameters.AddWithValue("@now", now); DataSyncPrimaryKeyHelper.BindPrimaryKeyParameters(cmd, row, pkColumns, "@pk"); cmd.CommandTimeout = _options.CommandTimeoutSeconds; cmd.ExecuteNonQuery(); } } } private void InsertRow( SqlConnection conn, SqlTransaction tx, string table, IList dataColumns, Dictionary row, string hash, DateTime runStart) { var cols = new List(dataColumns) { "_sync_row_hash", "_sync_is_deleted", "_sync_last_seen_at", "_sync_updated_at" }; var paramNames = cols.Select((c, i) => "@p" + i).ToList(); var sql = string.Format( "INSERT INTO {0} ({1}) VALUES ({2})", table, string.Join(", ", cols.Select(c => "[" + c + "]")), string.Join(", ", paramNames)); using (var cmd = new SqlCommand(sql, conn, tx)) { for (var i = 0; i < dataColumns.Count; i++) { cmd.Parameters.AddWithValue("@p" + i, row.ContainsKey(dataColumns[i]) ? (row[dataColumns[i]] ?? DBNull.Value) : DBNull.Value); } cmd.Parameters.AddWithValue("@p" + dataColumns.Count, hash); cmd.Parameters.AddWithValue("@p" + (dataColumns.Count + 1), false); cmd.Parameters.AddWithValue("@p" + (dataColumns.Count + 2), runStart); cmd.Parameters.AddWithValue("@p" + (dataColumns.Count + 3), DateTime.UtcNow); cmd.CommandTimeout = _options.CommandTimeoutSeconds; cmd.ExecuteNonQuery(); } } private bool UpdateRowIfChanged( SqlConnection conn, SqlTransaction tx, string table, IList pkColumns, IList dataColumns, Dictionary row, string hash, DateTime runStart) { var pkSet = new HashSet(pkColumns, StringComparer.OrdinalIgnoreCase); var updateColumns = dataColumns.Where(c => !pkSet.Contains(c)).ToList(); var sets = updateColumns.Select((c, i) => string.Format("[{0}] = @p{1}", c, i)).ToList(); sets.Add("_sync_row_hash = @hash"); sets.Add("_sync_is_deleted = 0"); sets.Add("_sync_last_seen_at = @seen"); sets.Add("_sync_updated_at = @now"); var whereClause = DataSyncPrimaryKeyHelper.BuildSqlWhereClause(pkColumns, "@pk"); var sql = string.Format( "UPDATE {0} SET {1} WHERE {2} AND (_sync_row_hash IS NULL OR _sync_row_hash <> @hash)", table, string.Join(", ", sets), whereClause); using (var cmd = new SqlCommand(sql, conn, tx)) { for (var i = 0; i < updateColumns.Count; i++) { cmd.Parameters.AddWithValue( "@p" + i, row.ContainsKey(updateColumns[i]) ? (row[updateColumns[i]] ?? DBNull.Value) : DBNull.Value); } cmd.Parameters.AddWithValue("@hash", hash); cmd.Parameters.AddWithValue("@seen", runStart); cmd.Parameters.AddWithValue("@now", DateTime.UtcNow); DataSyncPrimaryKeyHelper.BindPrimaryKeyParameters(cmd, row, pkColumns, "@pk"); cmd.CommandTimeout = _options.CommandTimeoutSeconds; return cmd.ExecuteNonQuery() > 0; } } } }