using System; using System.Collections.Generic; using System.Data; using System.Data.Common; using System.Linq; using System.Text; using System.Threading; using System.Threading.Tasks; using AdminDAL.CustomCode.BulkInsert; using AdminDAL.DTO.BulkIceImport; using AdminDAL.Query.BulkIceImport; using Dapper; using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; namespace AdminDAL.CustomCode.BulkIceImport { public class BulkIceImportDAL : IBulkIceImportDAL { private readonly IQueryExecutor _qe; private readonly IBulkInsertService _bulkInsert; public BulkIceImportDAL(IQueryExecutor qe, IBulkInsertService bulkInsert) { _qe = qe; _bulkInsert = bulkInsert; } // ── Config ────────────────────────────────────────────────────────── public async Task> GetBulkIceConfigAsync( int bulkIceId, LoginDTO login, CancellationToken ct = default) { return await _qe .QueryAsync(login, BulkIceImportQB.GET_BULKICE_CONFIG, new { BulkIceId = bulkIceId,TenantId=login.ClientId }, cancellationToken: ct) .ConfigureAwait(false); } // ── Staging dump (outside transaction) ────────────────────────────── public async Task DumpFileToStagingAsync( string stagingTableName, DataTable data, LoginDTO login, CancellationToken ct = default) { string ddl = BuildCreateTableDdl(stagingTableName, data); await _qe.ExecuteAsync(login, ddl, null!, cancellationToken: ct).ConfigureAwait(false); await _bulkInsert.InsertAsync(stagingTableName, data, login, ct).ConfigureAwait(false); } // ── Copy staging → persistent temp (inside transaction) ───────────── public async Task CopyToTempTableAsync( string tempTableName, string stagingTableName, LoginDTO login, CancellationToken ct = default) { // Best-effort reset: matches legacy behaviour — log warning, keep going on failure. try { string dropSql = $"IF OBJECT_ID(N'{tempTableName}') IS NOT NULL DROP TABLE {tempTableName}"; await _qe.ExecuteAsync(login, dropSql, null, null, ct).ConfigureAwait(false); string createSql = $"IF OBJECT_ID(N'{tempTableName}') IS NULL " + $"SELECT * INTO {tempTableName} FROM {stagingTableName} WHERE 1 = 0"; await _qe.ExecuteAsync(login, createSql, null, null, ct).ConfigureAwait(false); } catch { // Non-critical reset failure — proceed to insert; let that surface any real error. } string insertSql = $"INSERT INTO {tempTableName} SELECT * FROM {stagingTableName}"; await _qe.ExecuteAsync(login, insertSql, null, null, ct).ConfigureAwait(false); } // ── Execute stored procedure (inside transaction) ──────────────────── public async Task ExecuteImportProcedureAsync( string procedureCall, LoginDTO login, DbTransaction tx, CancellationToken ct = default) { await _qe.ExecuteAsync(login, procedureCall, null, tx, ct).ConfigureAwait(false); } // ── Output table read ──────────────────────────────────────────────── public async Task GetOutputTableDataAsync( string guid, LoginDTO login, CancellationToken ct = default) { string sql = BulkIceImportQB.GetOutputTableSql(guid); var rows = (await _qe.QueryAsync(login, sql, null!, cancellationToken: ct).ConfigureAwait(false)) .Cast>() .ToList(); var table = new DataTable(); if (rows.Count == 0) return table; foreach (var key in rows[0].Keys) table.Columns.Add(key, typeof(string)); foreach (var row in rows) { var dr = table.NewRow(); foreach (var kvp in row) dr[kvp.Key] = kvp.Value?.ToString() ?? string.Empty; table.Rows.Add(dr); } return table; } // ── Validation (inside transaction) ────────────────────────────────── public async Task RunValidationAsync( int bulkIceId, int userId, LoginDTO login, DbTransaction tx, CancellationToken ct = default) { var param = new { BulkIceId = bulkIceId, UserId = userId }; try { var errors = (await _qe .QueryAsync( login, BulkIceImportQB.VALIDATION_BATCH_NO_ROWNUM, param, cancellationToken: ct) .ConfigureAwait(false)).ToList(); return new BulkIceValidationResultDTO { Errors = errors }; } catch (Exception ex) when (ContainsColumnMismatch(ex)) { // Stored procedure returns a ROWNUMBER column — retry with the wider CREATE TABLE. var rowErrors = (await _qe .QueryAsync( login, BulkIceImportQB.VALIDATION_BATCH_WITH_ROWNUM, param, cancellationToken: ct) .ConfigureAwait(false)).ToList(); return new BulkIceValidationResultDTO { RowErrors = rowErrors }; } } // ── Cleanup (always, never throws) ─────────────────────────────────── public async Task DropStagingTablesAsync( string guid, bool[] tablesCreated, LoginDTO login, CancellationToken ct = default) { for (int i = 0; i < tablesCreated.Length; i++) { if (!tablesCreated[i]) continue; try { string tableName = $"BULK_FILE_IMPORT_{i + 1}_{guid}"; string sql = $"IF OBJECT_ID(N'{tableName}') IS NOT NULL DROP TABLE {tableName}"; await _qe.ExecuteAsync(login, sql, null!, cancellationToken: ct).ConfigureAwait(false); } catch { // Swallow — cleanup must never interrupt the finally block. } } } public async Task DropOutputTableAsync( string guid, LoginDTO login, CancellationToken ct = default) { try { string tableName = $"OUTPUT_{guid}"; string sql = $"IF OBJECT_ID(N'{tableName}') IS NOT NULL DROP TABLE {tableName}"; await _qe.ExecuteAsync(login, sql, null!, cancellationToken: ct).ConfigureAwait(false); } catch { // Swallow — cleanup must never throw. } } // ── Helpers ────────────────────────────────────────────────────────── private static string BuildCreateTableDdl(string tableName, DataTable data) { var sb = new StringBuilder(); sb.Append($"CREATE TABLE {tableName} ("); for (int i = 0; i < data.Columns.Count; i++) { if (i > 0) sb.Append(", "); sb.Append($"[{data.Columns[i].ColumnName}] NVARCHAR(MAX) NULL"); } sb.Append(")"); return sb.ToString(); } private static bool ContainsColumnMismatch(Exception ex) { const string msg = "Column name or number of supplied values does not match table definition"; return ex.Message.Contains(msg, StringComparison.OrdinalIgnoreCase) || (ex.InnerException?.Message.Contains(msg, StringComparison.OrdinalIgnoreCase) ?? false); } } }