using System.Data.Common; using System.Text; using Dapper; namespace FrameworkBLL.Engine { /// /// Applies one chunk of source rows to a target table using the specified conflict policy. /// All SQL generated here is runtime-dynamic (column names from INFORMATION_SCHEMA) — this is /// the documented exception to the QB-only SQL rule. /// internal static class ConflictResolver { // Rows per multi-row INSERT statement — keeps parameter count low while reducing round-trips. private const int BatchSize = 25; internal static async Task ApplyAsync( DbConnection target, string schemaName, string tableName, IEnumerable> chunk, string pkColumn, ConflictPolicy policy, string? modifiedDateColumn, DatabaseType srcDbType, DatabaseType tgtDbType, CancellationToken ct) { // Normalize source values to target RDBMS types before any write var rows = NormalizeRows(chunk.ToList(), srcDbType, tgtDbType); if (rows.Count == 0) return 0; return policy switch { ConflictPolicy.Overwrite => await ApplyOverwriteAsync(target, schemaName, tableName, rows, pkColumn, tgtDbType, ct).ConfigureAwait(false), ConflictPolicy.Skip => await ApplySkipAsync(target, schemaName, tableName, rows, pkColumn, tgtDbType, ct).ConfigureAwait(false), ConflictPolicy.LastWriteWins => await ApplyLastWriteWinsAsync(target, schemaName, tableName, rows, pkColumn, modifiedDateColumn, tgtDbType, ct).ConfigureAwait(false), _ => throw new ArgumentOutOfRangeException(nameof(policy), policy, null) }; } // ── Overwrite: DELETE then INSERT — wrapped in a transaction ───────── private static async Task ApplyOverwriteAsync( DbConnection target, string schema, string table, List> rows, string pk, DatabaseType dbType, CancellationToken ct) { var pkValues = rows.Select(r => r[pk]).ToList(); var columns = rows[0].Keys.ToList(); var deleteSql = $"DELETE FROM {SyncQueryBuilder.QuoteTable(schema, table, dbType)} WHERE {SyncQueryBuilder.Quote(pk, dbType)} IN @PkValues"; await using var tx = await target.BeginTransactionAsync(ct).ConfigureAwait(false); try { await target.ExecuteAsync( new CommandDefinition(deleteSql, new { PkValues = pkValues }, transaction: tx, cancellationToken: ct)) .ConfigureAwait(false); var inserted = await BatchInsertAsync(target, schema, table, columns, rows, dbType, tx, ct).ConfigureAwait(false); await tx.CommitAsync(ct).ConfigureAwait(false); return inserted; } catch { await tx.RollbackAsync(ct).ConfigureAwait(false); throw; } } // ── Skip: INSERT only if PK absent ─────────────────────────────────── private static async Task ApplySkipAsync( DbConnection target, string schema, string table, List> rows, string pk, DatabaseType dbType, CancellationToken ct) { var pkValues = rows.Select(r => r[pk]).ToList(); var columns = rows[0].Keys.ToList(); var existsSql = $"SELECT {SyncQueryBuilder.Quote(pk, dbType)} FROM {SyncQueryBuilder.QuoteTable(schema, table, dbType)} WHERE {SyncQueryBuilder.Quote(pk, dbType)} IN @PkValues"; // NOTE: deliberately not QueryAsync — Dapper special-cases a requested type of // `object` and returns a DapperRow wrapper for the whole row rather than the raw scalar, // so a naive `.ToString()` on it never matches the chunk's own PK values (found via a // real SQLite-backed test — see FrameworkTests/ConflictResolverTests.cs). Read the row // as a dictionary and pull the PK column out explicitly instead, mirroring the pattern // ApplyLastWriteWinsAsync already uses below for the same reason. var existingRaw = await target.QueryAsync( new CommandDefinition(existsSql, new { PkValues = pkValues }, cancellationToken: ct)) .ConfigureAwait(false); var existing = existingRaw .Select(r => { var d = (IDictionary)r; return d.TryGetValue(pk, out var v) ? v?.ToString() : null; }) .Where(s => s is not null) .ToHashSet(); var newRows = rows.Where(r => !existing.Contains(r[pk]?.ToString())).ToList(); if (newRows.Count == 0) return 0; return await BatchInsertAsync(target, schema, table, columns, newRows, dbType, null, ct).ConfigureAwait(false); } // ── LastWriteWins: write only if source is newer than target ───────── private static async Task ApplyLastWriteWinsAsync( DbConnection target, string schema, string table, List> rows, string pk, string? modifiedCol, DatabaseType dbType, CancellationToken ct) { if (string.IsNullOrWhiteSpace(modifiedCol)) return await ApplyOverwriteAsync(target, schema, table, rows, pk, dbType, ct).ConfigureAwait(false); var pkValues = rows.Select(r => r[pk]).ToList(); var columns = rows[0].Keys.ToList(); var fetchSql = $"SELECT {SyncQueryBuilder.Quote(pk, dbType)}, {SyncQueryBuilder.Quote(modifiedCol, dbType)} FROM {SyncQueryBuilder.QuoteTable(schema, table, dbType)} WHERE {SyncQueryBuilder.Quote(pk, dbType)} IN @PkValues"; var targetRows = await target.QueryAsync( new CommandDefinition(fetchSql, new { PkValues = pkValues }, cancellationToken: ct)) .ConfigureAwait(false); var targetMods = targetRows.ToDictionary( r => { var d = (IDictionary)r; return d.TryGetValue(pk, out var v) ? v?.ToString() ?? string.Empty : string.Empty; }, r => { var d = (IDictionary)r; return d.TryGetValue(modifiedCol!, out var v) ? ToDateTimeOrNull(v) : null; }); var toInsert = new List>(); var toUpdate = new List>(); foreach (var row in rows) { var key = row[pk]?.ToString() ?? string.Empty; if (!targetMods.TryGetValue(key, out var targetTs)) toInsert.Add(row); else if (ToDateTimeOrNull(row[modifiedCol]) is DateTime srcTs && (!targetTs.HasValue || srcTs > targetTs.Value)) toUpdate.Add(row); // else: source is same age or older — skip } int written = 0; if (toInsert.Count > 0) written += await BatchInsertAsync(target, schema, table, columns, toInsert, dbType, null, ct).ConfigureAwait(false); foreach (var row in toUpdate) written += await UpsertRowAsync(target, schema, table, columns, pk, row, dbType, ct).ConfigureAwait(false); return written; } // ── Helpers ────────────────────────────────────────────────────────── // Batches rows into multi-row INSERT statements (BatchSize rows per statement). // Uniquely names parameters @colname_rowindex to avoid collisions within a batch. private static async Task BatchInsertAsync( DbConnection conn, string schema, string table, List columns, List> rows, DatabaseType dbType, DbTransaction? tx, CancellationToken ct) { var tableFqn = SyncQueryBuilder.QuoteTable(schema, table, dbType); var colList = string.Join(", ", columns.Select(c => SyncQueryBuilder.Quote(c, dbType))); int total = 0; for (int batchStart = 0; batchStart < rows.Count; batchStart += BatchSize) { var batch = rows.Skip(batchStart).Take(BatchSize).ToList(); var sb = new StringBuilder(); sb.Append($"INSERT INTO {tableFqn} ({colList}) VALUES "); var valueRows = new List(batch.Count); var dynParams = new DynamicParameters(); for (int j = 0; j < batch.Count; j++) { var paramRefs = columns.Select(c => $"@{c}_{j}"); valueRows.Add($"({string.Join(", ", paramRefs)})"); foreach (var col in columns) dynParams.Add($"{col}_{j}", batch[j].TryGetValue(col, out var v) ? v : null); } sb.Append(string.Join(", ", valueRows)); total += await conn.ExecuteAsync( new CommandDefinition(sb.ToString(), dynParams, transaction: tx, cancellationToken: ct)) .ConfigureAwait(false); } return total; } private static async Task UpsertRowAsync( DbConnection conn, string schema, string table, List columns, string pk, IDictionary row, DatabaseType dbType, CancellationToken ct) { var setClauses = string.Join(", ", columns .Where(c => c != pk) .Select(c => $"{SyncQueryBuilder.Quote(c, dbType)} = @{c}")); var updateSql = $"UPDATE {SyncQueryBuilder.QuoteTable(schema, table, dbType)} SET {setClauses} WHERE {SyncQueryBuilder.Quote(pk, dbType)} = @{pk}"; return await conn.ExecuteAsync( new CommandDefinition(updateSql, new DynamicParameters(row), cancellationToken: ct)) .ConfigureAwait(false); } // Coerces a modified-date column value into DateTime regardless of exactly which CLR // type the ADO provider boxed it as. A hard "(DateTime?)v" cast is not safe here — Npgsql // returns DateTimeOffset for `timestamptz` columns (not DateTime), and some providers/ // configurations return the raw textual representation — either would throw // InvalidCastException (or, for the "is DateTime" pattern match on the source side, // silently and permanently skip every row) with a plain cast instead of a real comparison. private static DateTime? ToDateTimeOrNull(object? v) => v switch { null or DBNull => null, DateTime dt => dt, DateTimeOffset dto => dto.UtcDateTime, string s when DateTime.TryParse( s, System.Globalization.CultureInfo.InvariantCulture, System.Globalization.DateTimeStyles.RoundtripKind, out var parsed) => parsed, _ => null }; // Normalizes each column value from source RDBMS types to target RDBMS types. // No-op when source and target are the same engine. private static List> NormalizeRows( List> rows, DatabaseType srcDbType, DatabaseType tgtDbType) { if (srcDbType == tgtDbType) return rows; return rows .Select(row => (IDictionary)row.ToDictionary( kv => kv.Key, kv => TypeNormalizer.NormalizeCLR(kv.Value, srcDbType, tgtDbType))) .ToList(); } } }