using System; using System.Collections.Generic; using System.Data; using System.IO; using System.Linq; using System.Text; using System.Threading; using System.Threading.Tasks; using AdminDAL.CustomCode.BulkIceImport; using AdminDAL.DTO.BulkIceImport; using GB5Shared.DTO.Framework.Login; using GB5Shared.FileImport; using GB5Shared.QueryExecutor; using GB5Shared.Resource.Response; using Microsoft.Extensions.Configuration; using Newtonsoft.Json; namespace AdminBLL.BulkIceImport { public class BulkIceImportBLL : IBulkIceImportBLL { private readonly IBulkIceImportDAL _dal; private readonly IFileToDataTableService _fileReader; private readonly IQueryExecutor _qe; private readonly IEnumerable _postActions; private readonly IConfiguration _config; public BulkIceImportBLL( IBulkIceImportDAL dal, IFileToDataTableService fileReader, IQueryExecutor qe, IEnumerable postActions, IConfiguration config) { _dal = dal; _fileReader = fileReader; _qe = qe; _postActions = postActions; _config = config; } // ── Public entry point ─────────────────────────────────────────────── public async Task BulkImportAsync( BulkIceImportRequestDTO request, LoginDTO login, CancellationToken ct = default) { ValidateRequest(request); var configList = (await _dal .GetBulkIceConfigAsync(request.BulkIceId, login, ct) .ConfigureAwait(false)).ToList(); ValidateConfig(configList, request); string guid = Guid.NewGuid().ToString("N"); // no dashes var stagingCreated = new bool[5]; var stagingNames = new string[configList.Count]; bool outputCreated = false; try { // ── Phase 1: decode Base64 → staging tables (outside transaction) ── for (int i = 0; i < configList.Count; i++) { var file = i < request.Files.Length ? request.Files[i] : null; if (file is null || string.IsNullOrWhiteSpace(file.Base64Content)) continue; byte[] bytes = Convert.FromBase64String(file.Base64Content); DataTable dt = await _fileReader .ReadFromBytesAsync(bytes, request.ImportFileType, ct) .ConfigureAwait(false); stagingNames[i] = $"BULK_FILE_IMPORT_{i + 1}_{guid}"; await _dal.DumpFileToStagingAsync(stagingNames[i], dt, login, ct).ConfigureAwait(false); stagingCreated[i] = true; } // ── Phase 2: transaction — copy, execute, validate, commit ──── var trans = await _qe.BeginTransactionAsync(login).ConfigureAwait(false); try { string procedureCall = configList[0].CustomProcedure; for (int i = 0; i < configList.Count; i++) { if (!stagingCreated[i]) continue; await _dal.CopyToTempTableAsync( configList[i].TempTableName, stagingNames[i], login, ct) .ConfigureAwait(false); procedureCall = procedureCall.Replace( $":table{i + 1}", stagingNames[i], StringComparison.OrdinalIgnoreCase); } procedureCall = procedureCall.Replace( ":outputtablename", $"OUTPUT_{guid}", StringComparison.OrdinalIgnoreCase); var validation = await _dal .RunValidationAsync(request.BulkIceId, login.UserId, login, trans, ct) .ConfigureAwait(false); if (validation.HasErrors) { await _qe.RollbackAsync(trans).ConfigureAwait(false); return SerializeValidationFailure(validation); } await _dal.ExecuteImportProcedureAsync(procedureCall, login, trans, ct) .ConfigureAwait(false); if (configList[0].ToType != ToTypeEnum.DirectMethod) outputCreated = true; await _qe.CommitAsync(trans).ConfigureAwait(false); return await ExecutePostActionAsync(configList[0], guid, login, ct) .ConfigureAwait(false); } catch { await _qe.RollbackAsync(trans).ConfigureAwait(false); throw; } } finally { await _dal.DropStagingTablesAsync(guid, stagingCreated, login, ct).ConfigureAwait(false); if (outputCreated) await _dal.DropOutputTableAsync(guid, login, ct).ConfigureAwait(false); } } // ── Post-import action dispatch ────────────────────────────────────── private async Task ExecutePostActionAsync( BulkIceImportDTO config, string guid, LoginDTO login, CancellationToken ct) { switch (config.ToType) { case ToTypeEnum.DirectMethod: return SuccessResponse.SaveSuccessMessage; case ToTypeEnum.ReturnFile: { DataTable dt = await _dal.GetOutputTableDataAsync(guid, login, ct).ConfigureAwait(false); string csv = DataTableToCsv(dt); string logPath = _config["LogPath"] ?? Path.Combine(Path.GetTempPath(), "BulkIce"); string dir = Path.Combine(logPath, "CAIFileGenerate", guid); Directory.CreateDirectory(dir); string fileName = config.BulkIceCode + ".csv"; await File.WriteAllTextAsync(Path.Combine(dir, fileName), csv, ct).ConfigureAwait(false); string destDir = Path.Combine(logPath, "AttachmentFolder"); Directory.CreateDirectory(destDir); File.Copy(Path.Combine(dir, fileName), Path.Combine(destDir, guid + fileName), overwrite: true); return guid + fileName; } default: { string actionKey = string.IsNullOrEmpty(config.PostActionKey) ? config.ToType.ToString().ToUpperInvariant() : config.PostActionKey.ToUpperInvariant(); var action = _postActions.FirstOrDefault(a => string.Equals(a.ActionKey, actionKey, StringComparison.OrdinalIgnoreCase)); if (action is null) throw new InvalidOperationException( $"No IPostImportAction registered for key '{actionKey}'. " + "Register the implementation in AdminSL/Program.cs."); DataTable output = await _dal.GetOutputTableDataAsync(guid, login, ct).ConfigureAwait(false); string outCsv = DataTableToCsv(output); string outPath = _config["LogPath"] ?? Path.Combine(Path.GetTempPath(), "BulkIce"); string outDir = Path.Combine(outPath, "AttachmentFolder"); Directory.CreateDirectory(outDir); string outFileName = config.BulkIceCode + ".csv"; string outFullPath = Path.Combine(outDir, guid + outFileName); await File.WriteAllTextAsync(outFullPath, outCsv, ct).ConfigureAwait(false); await action.ExecuteAsync(guid + outFileName, config.SkipOverWrite, config.ApprovalDuringImport, login, ct).ConfigureAwait(false); return SuccessResponse.SaveSuccessMessage; } } } // ── Validation ─────────────────────────────────────────────────────── private static void ValidateRequest(BulkIceImportRequestDTO request) { if (request is null) throw new ArgumentNullException(nameof(request)); //if (request.BulkIceId <= 0) throw new ArgumentException("BulkIceId must be greater than 0."); if (request.Files is null || request.Files.Length == 0) throw new ArgumentException("At least one file must be provided."); if (request.Files.Length > 5) throw new ArgumentException("A maximum of 5 files is supported."); } private static void ValidateConfig(List configList, BulkIceImportRequestDTO request) { if (configList.Count == 0) throw new InvalidOperationException("No BulkIce configuration found for the given BulkIceId."); int expected = configList[0].NumberOfFiles; if (expected == 0) throw new InvalidOperationException("NumberOfFiles cannot be 0 in BulkIce configuration."); if (configList.Count != expected) throw new InvalidOperationException("Mismatch between NumberOfFiles and BulkIceDetail row count."); if (expected > 5) throw new InvalidOperationException("NumberOfFiles cannot exceed 5."); for (int i = 0; i < expected; i++) { if (i >= request.Files.Length || string.IsNullOrWhiteSpace(request.Files[i]?.Base64Content)) throw new ArgumentException($"File {i + 1} is required but was not provided."); } if (string.IsNullOrWhiteSpace(configList[0].CustomProcedure)) throw new InvalidOperationException("No stored procedure is defined for this BulkIce import."); } // ── Serialisation ──────────────────────────────────────────────────── // Backward-compatible validation failure envelope expected by legacy UI parser. private static string SerializeValidationFailure(BulkIceValidationResultDTO result) { string body = result.RowErrors.Count > 0 ? JsonConvert.SerializeObject(result.RowErrors) : JsonConvert.SerializeObject(result.Errors); return JsonConvert.SerializeObject(new { Id = -1, ValidationBody = body + "@#$" }); } private static string DataTableToCsv(DataTable table) { var sb = new StringBuilder(); sb.AppendLine(string.Join(",", table.Columns.Cast().Select(c => $"\"{c.ColumnName}\""))); foreach (DataRow row in table.Rows) sb.AppendLine(string.Join(",", table.Columns.Cast() .Select(c => $"\"{row[c]?.ToString()?.Replace("\"", "\"\"")}\""))); return sb.ToString(); } } }