using System.Data.Common; using System.Runtime.CompilerServices; using Dapper; using Microsoft.Extensions.Logging; using Polly; using Polly.Retry; namespace FrameworkBLL.Engine { public class SyncEngine : ISyncEngine { private readonly ILogger _logger; // Polly retry: 3 attempts, exponential 500ms → 1s → 2s on transient DbException. // Watermarks are NOT updated for failed tables — next run re-processes the same window. private static readonly AsyncRetryPolicy _retryPolicy = Policy .Handle() .WaitAndRetryAsync(3, attempt => TimeSpan.FromMilliseconds(500 * Math.Pow(2, attempt - 1))); public SyncEngine(ILogger logger) { _logger = logger; } // ── Streaming ───────────────────────────────────────────────────────── public async IAsyncEnumerable>> StreamPagedAsync( DbConnection conn, string selectSql, int chunkSize, IDictionary baseParams, [EnumeratorCancellation] CancellationToken ct) { int offset = 0; while (true) { ct.ThrowIfCancellationRequested(); var dynParams = new DynamicParameters(baseParams); dynParams.Add("ChunkOffset", offset); dynParams.Add("ChunkSize", chunkSize); var rows = (await conn.QueryAsync( new CommandDefinition(selectSql, dynParams, cancellationToken: ct)) .ConfigureAwait(false)) .Select(r => { var dict = (System.Collections.Generic.IDictionary)r; return (IDictionary)dict .ToDictionary(kv => kv.Key, kv => (object?)kv.Value); }) .ToList(); if (rows.Count == 0) yield break; yield return rows; if (rows.Count < chunkSize) yield break; offset += chunkSize; } } // ── Apply chunk with retry ──────────────────────────────────────────── public async Task ApplyChunkToTargetAsync( DbConnection targetConn, string schemaName, string tableName, IEnumerable> chunk, string pkColumn, ConflictPolicy policy, string? modifiedDateColumn, DatabaseType srcDbType, DatabaseType tgtDbType, CancellationToken ct) { var chunkList = chunk.ToList(); return await _retryPolicy.ExecuteAsync(async () => await ConflictResolver.ApplyAsync( targetConn, schemaName, tableName, chunkList, pkColumn, policy, modifiedDateColumn, srcDbType, tgtDbType, ct) .ConfigureAwait(false)); } // ── Delete detection ───────────────────────────────────────────────── public async Task ApplyDeletesAsync( DbConnection sourceConn, DbConnection targetConn, string schemaName, string tableName, string pkColumn, DateTime sinceUtc, DatabaseType sourceDbType, DatabaseType tgtDbType, CancellationToken ct) { if (sourceDbType != DatabaseType.SqlServer) { // Non-SQL-Server: propagate ISDELETED=1 rows via the UPDATE flow. // Full hard-delete via PostgreSQL logical replication is Phase 4. return 0; } // SQL Server: read CDC table for hard-delete events using the configured PK column var cdcSql = SyncQueryBuilder.BuildCdcSelectSql(tableName, pkColumn, sinceUtc); IEnumerable deletedPks; try { deletedPks = await sourceConn.QueryAsync( new CommandDefinition(cdcSql, new { SinceUtc = sinceUtc }, cancellationToken: ct)) .ConfigureAwait(false); } catch (Exception ex) { // CDC table may not exist for every table — log and skip _logger.LogWarning(ex, "CDC query failed for {Schema}.{Table} — skipping hard-delete propagation", schemaName, tableName); return 0; } var pkList = deletedPks.ToList(); if (pkList.Count == 0) return 0; // Use target RDBMS quoting for the DELETE statement var deleteSql = $"DELETE FROM {SyncQueryBuilder.QuoteTable(schemaName, tableName, tgtDbType)} WHERE {SyncQueryBuilder.Quote(pkColumn, tgtDbType)} IN @PkValues"; return await targetConn.ExecuteAsync( new CommandDefinition(deleteSql, new { PkValues = pkList }, cancellationToken: ct)) .ConfigureAwait(false); } } }