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;
}
}
}
}