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