using System.IO.Compression; using System.Text; using System.Text.Json; using System.Text.Json.Serialization; using FrameworkBLL.Connection; using FrameworkBLL.Engine; using FrameworkDAL.CustomCode.DataSync; using FrameworkDAL.CustomCode.Dataset; using FrameworkDAL.CustomCode.Server; using FrameworkDAL.DTO.DataSync; using FrameworkDAL.DTO.Dataset; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using Microsoft.Extensions.Logging; namespace FrameworkBLL.DataSync { public class DataSyncExportBLL : IDataSyncExportBLL { private static readonly JsonSerializerOptions _jsonOptions = new() { DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull }; private static readonly JsonSerializerOptions _prettyJson = new() { WriteIndented = true }; private static readonly byte[] _comma = Encoding.UTF8.GetBytes(","); private static readonly byte[] _arrayOpen = Encoding.UTF8.GetBytes("["); private static readonly byte[] _arrayClose = Encoding.UTF8.GetBytes("]"); private readonly IDataSyncJobDAL _syncJobDal; private readonly IDatasetDAL _datasetDal; private readonly IServerDAL _serverDal; private readonly IExternalDbConnectionFactory _connFactory; private readonly ISyncEngine _engine; private readonly ILogger _logger; public DataSyncExportBLL( IDataSyncJobDAL syncJobDal, IDatasetDAL datasetDal, IServerDAL serverDal, IExternalDbConnectionFactory connFactory, ISyncEngine engine, ILogger logger) { _syncJobDal = syncJobDal; _datasetDal = datasetDal; _serverDal = serverDal; _connFactory = connFactory; _engine = engine; _logger = logger; } public async Task ExportToZipAsync(SyncExportRequestDTO req, LoginDTO login, CancellationToken ct) { GB5Trace.Step("validate-export", new { req.SyncJobId, req.ExportMode }); var jobJson = await _syncJobDal.GetSyncJob(req.SyncJobId, login, ct).ConfigureAwait(false); var job = Newtonsoft.Json.JsonConvert.DeserializeObject(jobJson) ?? throw new InvalidOperationException($"SyncJob {req.SyncJobId} not found."); var srcInstance = await _serverDal.GetDBInstanceById(job.SourceDBInstanceId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Source DBInstance {job.SourceDBInstanceId} not found."); var datasetJson = await _datasetDal.GetDataset(job.DatasetId, login).ConfigureAwait(false); var dataset = Newtonsoft.Json.JsonConvert.DeserializeObject(datasetJson) ?? throw new InvalidOperationException($"Dataset {job.DatasetId} not found."); var srcDbType = (DatabaseType)srcInstance.DBTypeId; var srcSchema = DefaultSchema(srcDbType); DateTime? fromDate = req.ExportMode switch { "Incremental" => job.LastInsertedTill ?? job.LastUpdatedTill, "Custom" => req.FromDate, _ => null // Full — no watermark }; GB5Trace.Step("start-export", new { req.SyncJobId, req.ExportMode, fromDate }); var ms = new MemoryStream(); using (var zip = new ZipArchive(ms, ZipArchiveMode.Create, leaveOpen: true)) { var tableMeta = new List(); await using var srcConn = await _connFactory.OpenAsync(srcInstance, ct).ConfigureAwait(false); foreach (var detail in (dataset.DatasetDetailArray ?? new List()) .OrderBy(d => d.DatasetDetailSlNo)) { if (string.IsNullOrWhiteSpace(detail.DBObjectName) || string.IsNullOrWhiteSpace(detail.PrimaryKeyColumn)) continue; var chunkSize = detail.ChunkSize > 0 ? detail.ChunkSize : 1000; var selectSql = SyncQueryBuilder.BuildSelectSql( schemaName: srcSchema, tableName: detail.DBObjectName, pkColumn: detail.PrimaryKeyColumn, createdDateCol: detail.CreatedDateColumn, modifiedDateCol: detail.ModifiedDateColumn, lastInsertedTill: fromDate, lastUpdatedTill: fromDate, extraWhere: job.QueryCondition, dbType: srcDbType); var baseParams = new Dictionary(); if (fromDate.HasValue && !string.IsNullOrWhiteSpace(detail.CreatedDateColumn)) baseParams["LastInsertedTill"] = fromDate.Value; if (fromDate.HasValue && !string.IsNullOrWhiteSpace(detail.ModifiedDateColumn)) baseParams["LastUpdatedTill"] = fromDate.Value; // Stream rows directly into the ZIP entry — never materialise the full table. // Each chunk is serialised as individual JSON objects inside a JSON array, // keeping memory bounded to one chunk at a time regardless of table size. var entry = zip.CreateEntry($"{detail.DBObjectName}.json", CompressionLevel.Optimal); await using var entryStream = entry.Open(); int rowCount = 0; try { await entryStream.WriteAsync(_arrayOpen, ct).ConfigureAwait(false); bool firstRow = true; await foreach (var chunk in _engine.StreamPagedAsync(srcConn, selectSql, chunkSize, baseParams, ct)) { foreach (var row in chunk) { if (!firstRow) await entryStream.WriteAsync(_comma, ct).ConfigureAwait(false); await JsonSerializer.SerializeAsync(entryStream, row, _jsonOptions, ct).ConfigureAwait(false); firstRow = false; rowCount++; } } await entryStream.WriteAsync(_arrayClose, ct).ConfigureAwait(false); } catch (Exception ex) { _logger.LogError(ex, "Export: failed reading table {Table} for SyncJob {SyncJobId}", detail.DBObjectName, req.SyncJobId); GB5Trace.MarkFailed("export-table-failed", ex); throw; } tableMeta.Add(new { detail.DBObjectName, RowCount = rowCount, detail.PrimaryKeyColumn, detail.ModifiedDateColumn }); GB5Trace.Step("export-table-done", new { tableName = detail.DBObjectName, rowCount }); } // Write manifest.json var manifest = new { ExportVersion = "1.0", ExportedAt = DateTime.UtcNow.ToString("O"), req.SyncJobId, SyncJobDescription = job.Description, SourceDBInstance = srcInstance.DBInstanceName, req.ExportMode, IncrementalFrom = fromDate?.ToString("O"), ConflictPolicy = job.ConflictPolicy, Tables = tableMeta }; var manifestEntry = zip.CreateEntry("manifest.json", CompressionLevel.NoCompression); await using var manifestStream = manifestEntry.Open(); var manifestBytes = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(manifest, _prettyJson)); await manifestStream.WriteAsync(manifestBytes, ct).ConfigureAwait(false); } ms.Position = 0; GB5Trace.Step("export-complete", new { req.SyncJobId, zipBytes = ms.Length }); return ms; } private static string DefaultSchema(DatabaseType dbType) => dbType switch { DatabaseType.PostgreSQL => "public", _ => "dbo" }; } }