using System.IO.Compression; using System.Text.Json; using System.Text.Json.Nodes; using FrameworkBLL.Connection; using FrameworkBLL.Engine; using FrameworkDAL.CustomCode.DataSync; using FrameworkDAL.CustomCode.Server; using FrameworkDAL.DTO.DataSync; using GB5Shared.DTO.Framework.Login; using GB5Shared.GenerateAutoNumber; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.DataSync { public class DataSyncImportBLL : IDataSyncImportBLL { private readonly IServerDAL _serverDal; private readonly IExternalDbConnectionFactory _connFactory; private readonly ISyncEngine _engine; private readonly IDataSyncRunLogDAL _runLogDal; private readonly AutoNumber _autoNumber; private readonly ILogger _logger; public DataSyncImportBLL( IServerDAL serverDal, IExternalDbConnectionFactory connFactory, ISyncEngine engine, IDataSyncRunLogDAL runLogDal, AutoNumber autoNumber, ILogger logger) { _serverDal = serverDal; _connFactory = connFactory; _engine = engine; _runLogDal = runLogDal; _autoNumber = autoNumber; _logger = logger; } // ── Validate ────────────────────────────────────────────────────────── public async Task ValidateImportFileAsync( Stream zipStream, LoginDTO login, CancellationToken ct) { GB5Trace.Step("validate-import-file", new { }); using var zip = new ZipArchive(zipStream, ZipArchiveMode.Read, leaveOpen: true); var manifest = ReadManifest(zip); var preview = new SyncImportPreviewDTO { SyncJobId = manifest["SyncJobId"]?.GetValue() ?? 0, SyncJobDescription = manifest["SyncJobDescription"]?.ToString() ?? string.Empty, SourceDBInstance = manifest["SourceDBInstance"]?.ToString() ?? string.Empty, ExportedAt = manifest["ExportedAt"]?.ToString() ?? string.Empty, ExportMode = manifest["ExportMode"]?.ToString() ?? string.Empty, ConflictPolicy = manifest["ConflictPolicy"]?.GetValue() ?? 0, }; var tables = manifest["Tables"]?.AsArray() ?? new JsonArray(); foreach (var t in tables) { var tableName = t?["DBObjectName"]?.ToString() ?? string.Empty; var rowCount = t?["RowCount"]?.GetValue() ?? 0; var pk = t?["PrimaryKeyColumn"]?.ToString() ?? string.Empty; var modDate = t?["ModifiedDateColumn"]?.ToString() ?? string.Empty; var tableEntry = zip.GetEntry($"{tableName}.json"); var schemaStatus = tableEntry is null ? "File missing in ZIP" : "OK"; preview.Tables.Add(new SyncImportTablePreviewDTO { TableName = tableName, RowCount = rowCount, SchemaStatus = schemaStatus, PrimaryKeyColumn = pk, ModifiedDateColumn = modDate, }); } return preview; } // ── Import ──────────────────────────────────────────────────────────── public async Task ImportFromFileAsync( Stream zipStream, int targetDbInstanceId, byte conflictPolicy, LoginDTO login, CancellationToken ct) { GB5Trace.Step("start-import-from-file", new { targetDbInstanceId, conflictPolicy }); var autoNum = await _autoNumber.GetAutoNumber(1, AUTONUMBERCONSTANT.DSYNCRUNLOG, login).ConfigureAwait(false); var runLog = new SyncRunLogDTO { RunId = autoNum.StartNumber, SyncJobId = 0, // file import — no source job on this system TriggeredBy = 2, // 2 = FileImport TriggeredByUserId = login.UserId }; await _runLogDal.SaveSyncRunLog(runLog, login, ct).ConfigureAwait(false); var runId = runLog.RunId; var tgtInstance = await _serverDal.GetDBInstanceById(targetDbInstanceId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Target DBInstance {targetDbInstanceId} not found."); var tgtDbType = (DatabaseType)tgtInstance.DBTypeId; var tgtSchema = tgtDbType == DatabaseType.PostgreSQL ? "public" : "dbo"; var policy = (ConflictPolicy)conflictPolicy; int totalRows = 0, insertedRows = 0; byte finalStatus = 2; string? errorDetail = null; try { using var zip = new ZipArchive(zipStream, ZipArchiveMode.Read, leaveOpen: true); var manifest = ReadManifest(zip); var tables = manifest["Tables"]?.AsArray() ?? new JsonArray(); // Manifest ConflictPolicy is advisory — caller-supplied policy wins await using var tgtConn = await _connFactory.OpenAsync(tgtInstance, ct).ConfigureAwait(false); foreach (var t in tables) { var tableName = t?["DBObjectName"]?.ToString() ?? string.Empty; var pk = t?["PrimaryKeyColumn"]?.ToString() ?? string.Empty; var modDate = t?["ModifiedDateColumn"]?.ToString(); if (string.IsNullOrWhiteSpace(tableName) || string.IsNullOrWhiteSpace(pk)) continue; var tableEntry = zip.GetEntry($"{tableName}.json"); if (tableEntry is null) { _logger.LogWarning("Import: {Table}.json not found in ZIP — skipping", tableName); continue; } List> rows; try { await using var es = tableEntry.Open(); rows = await JsonSerializer.DeserializeAsync>>(es, cancellationToken: ct).ConfigureAwait(false) is { } raw ? raw.Select(r => (IDictionary)r.ToDictionary(kv => kv.Key, kv => (object?)UnboxJsonElement(kv.Value))).ToList() : new List>(); } catch (Exception ex) { _logger.LogError(ex, "Import: failed deserialising {Table}", tableName); GB5Trace.MarkFailed("import-deserialise-failed", ex); finalStatus = 4; errorDetail = $"Table {tableName}: {ex.Message}"; continue; } const int chunkSize = 500; for (int i = 0; i < rows.Count; i += chunkSize) { var chunk = rows.Skip(i).Take(chunkSize).ToList(); var written = await _engine.ApplyChunkToTargetAsync( tgtConn, tgtSchema, tableName, chunk, pk, policy, modDate, DatabaseType.SqlServer, // source type is unknown from file; assume SQL Server for type normalizer tgtDbType, ct).ConfigureAwait(false); insertedRows += written; totalRows += chunk.Count; GB5Trace.Step("import-chunk", new { runId, tableName, chunkIndex = i / chunkSize, rowCount = chunk.Count }); } _logger.LogInformation("Import: {Table} — {Rows} rows written (Run {RunId})", tableName, rows.Count, runId); } } catch (Exception ex) { finalStatus = 3; // Failed errorDetail = ex.Message; _logger.LogError(ex, "ImportFromFileAsync failed for Run {RunId}", runId); GB5Trace.MarkFailed("import-from-file-failed", ex); } finally { var completed = new SyncRunLogDTO { RunId = runId, Status = finalStatus, TotalRows = totalRows, InsertedRows = insertedRows, UpdatedRows = 0, DeletedRows = 0, SkippedRows = 0, ErrorDetail = errorDetail }; try { // Use CancellationToken.None — the original ct may already be cancelled, // but the run log must always be written regardless. await _runLogDal.UpdateSyncRunLogComplete(completed, login, CancellationToken.None).ConfigureAwait(false); } catch (Exception finEx) { _logger.LogError(finEx, "Failed to update run log for Run {RunId} — status {Status}", runId, finalStatus); } GB5Trace.Step("import-complete", new { runId, finalStatus, totalRows }); } return runId; } // ── Helpers ─────────────────────────────────────────────────────────── private static JsonObject ReadManifest(ZipArchive zip) { var entry = zip.GetEntry("manifest.json") ?? throw new InvalidOperationException("manifest.json not found in ZIP. This does not appear to be a DataSync export file."); using var stream = entry.Open(); var node = JsonNode.Parse(stream) ?? throw new InvalidOperationException("manifest.json is empty or invalid JSON."); return node.AsObject(); } private static object? UnboxJsonElement(JsonElement el) => el.ValueKind switch { JsonValueKind.String => el.GetString(), JsonValueKind.Number => el.TryGetInt64(out var l) ? l : el.GetDouble(), JsonValueKind.True => true, JsonValueKind.False => false, JsonValueKind.Null => null, _ => el.ToString() }; } }