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.EventLogPublish; using GB5Shared.GenerateAutoNumber; using GB5Shared.Telemetry; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.DataSync { // IDataSyncHubContext is a thin wrapper interface so FrameworkBLL can push SignalR // events without taking a compile-time dependency on FrameworkSL hub types. public interface IDataSyncHubContext { Task PushProgressAsync(int runId, string tableName, int chunkNumber, int rowsInChunk, int totalRowsSoFar, CancellationToken ct); Task PushRunCompletedAsync(int runId, int syncJobId, byte status, int totalRows, int insertedRows, int updatedRows, int deletedRows, int skippedRows, string? errorDetail, DateTime startedAt, CancellationToken ct); Task PushErrorAsync(int runId, string message, CancellationToken ct); } public class DataSyncRunBLL : IDataSyncRunBLL { private readonly IDataSyncJobDAL _syncJobDal; private readonly IDataSyncRunLogDAL _runLogDal; private readonly IDatasetDAL _datasetDal; private readonly IServerDAL _serverDal; private readonly IExternalDbConnectionFactory _connFactory; private readonly ISyncEngine _engine; private readonly AutoNumber _autoNumber; private readonly EventLogPublish _eventLog; private readonly IDataSyncHubContext? _hubContext; private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public DataSyncRunBLL( IDataSyncJobDAL syncJobDal, IDataSyncRunLogDAL runLogDal, IDatasetDAL datasetDal, IServerDAL serverDal, IExternalDbConnectionFactory connFactory, ISyncEngine engine, AutoNumber autoNumber, EventLogPublish eventLog, IServiceScopeFactory scopeFactory, ILogger logger, IDataSyncHubContext? hubContext = null) { _syncJobDal = syncJobDal; _runLogDal = runLogDal; _datasetDal = datasetDal; _serverDal = serverDal; _connFactory = connFactory; _engine = engine; _autoNumber = autoNumber; _eventLog = eventLog; _scopeFactory = scopeFactory; _logger = logger; _hubContext = hubContext; } // ── Read ────────────────────────────────────────────────────────────── public async Task GetSyncRunLog(int runId, LoginDTO login, CancellationToken ct) => await _runLogDal.GetSyncRunLog(runId, login, ct).ConfigureAwait(false); public async Task GetSyncRunLogList(int syncJobId, int offset, int pageSize, LoginDTO login, CancellationToken ct) => await _runLogDal.GetSyncRunLogList(syncJobId, offset, pageSize, login, ct).ConfigureAwait(false); // ── Main Orchestration ──────────────────────────────────────────────── public async Task ExecuteSyncJobAsync( int syncJobId, byte triggeredBy, int? triggeredByUserId, LoginDTO login, CancellationToken ct) { GB5Trace.Step("validate-syncjob", new { syncJobId }); // Concurrency guard — skip if job is already running var locked = await _syncJobDal.LockSyncJobForRun(syncJobId, login, ct).ConfigureAwait(false); if (locked == 0) { _logger.LogWarning("SyncJob {SyncJobId} is already running — skipping this trigger", syncJobId); throw new InvalidOperationException($"SyncJob {syncJobId} is already running."); } // Create run log entry (STATUS=1 Running); RunId comes from AutoNumber, not DB identity var autoNum = await _autoNumber.GetAutoNumber(1, AUTONUMBERCONSTANT.DSYNCRUNLOG, login); var runLog = new SyncRunLogDTO { RunId = autoNum.StartNumber, SyncJobId = syncJobId, TriggeredBy = triggeredBy, TriggeredByUserId = triggeredByUserId }; await _runLogDal.SaveSyncRunLog(runLog, login, ct).ConfigureAwait(false); var runId = runLog.RunId; GB5Trace.Step("start-run", new { syncJobId, runId }); // Fire execution in a new DI scope — the request scope will be disposed by the time // the background task runs, so we must not capture 'this' or any injected services. var scopeFactory = _scopeFactory; // singleton: safe to capture across scope boundaries var logger = _logger; // ILogger: backed by singleton infrastructure _ = Task.Run(async () => { await using var scope = scopeFactory.CreateAsyncScope(); var freshBll = scope.ServiceProvider.GetRequiredService(); try { await freshBll.RunBackgroundAsync(syncJobId, runId, login, CancellationToken.None) .ConfigureAwait(false); } catch (Exception ex) { logger.LogError(ex, "Background execution failed for SyncJob {SyncJobId} Run {RunId}", syncJobId, runId); } }, CancellationToken.None); return runId; } // Entry point called from background Task.Run on a fresh scope instance. public Task RunBackgroundAsync(int syncJobId, int runId, LoginDTO login, CancellationToken ct) => ExecuteRunAsync(syncJobId, runId, login, ct); // ── Private Execution ───────────────────────────────────────────────── private async Task ExecuteRunAsync(int syncJobId, int runId, LoginDTO login, CancellationToken ct) { int totalRows = 0, insertedRows = 0, updatedRows = 0, deletedRows = 0, skippedRows = 0; byte finalStatus = 2; // 2=Success string? errorDetail = null; var runAt = DateTime.UtcNow; try { // Load config var jobJson = await _syncJobDal.GetSyncJob(syncJobId, login, ct).ConfigureAwait(false); var job = JsonConvert.DeserializeObject(jobJson) ?? throw new InvalidOperationException($"SyncJob {syncJobId} not found."); var srcInstance = await _serverDal.GetDBInstanceById(job.SourceDBInstanceId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Source DBInstance {job.SourceDBInstanceId} not found."); var tgtInstance = await _serverDal.GetDBInstanceById(job.TargetDBInstanceId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"Target DBInstance {job.TargetDBInstanceId} not found."); var datasetJson = await _datasetDal.GetDataset(job.DatasetId, login).ConfigureAwait(false); var dataset = JsonConvert.DeserializeObject(datasetJson) ?? throw new InvalidOperationException($"Dataset {job.DatasetId} not found."); var srcDbType = (DatabaseType)srcInstance.DBTypeId; var tgtDbType = (DatabaseType)tgtInstance.DBTypeId; var policy = (ConflictPolicy)job.ConflictPolicy; // Default schema per RDBMS: SQL Server = dbo, PostgreSQL = public var srcSchema = DefaultSchema(srcDbType); var tgtSchema = DefaultSchema(tgtDbType); GB5Trace.Step("load-config", new { syncJobId, datasetId = job.DatasetId, tableCount = dataset.DatasetDetailArray?.Count ?? 0 }); await using var srcConn = await _connFactory.OpenAsync(srcInstance, ct).ConfigureAwait(false); await using var tgtConn = await _connFactory.OpenAsync(tgtInstance, ct).ConfigureAwait(false); var tableErrors = new List(); foreach (var detail in (dataset.DatasetDetailArray ?? new List()) .OrderBy(d => d.DatasetDetailSlNo)) { if (string.IsNullOrWhiteSpace(detail.DBObjectName) || string.IsNullOrWhiteSpace(detail.PrimaryKeyColumn)) { _logger.LogWarning("Skipping DatasetDetail {Id}: DBObjectName or PrimaryKeyColumn is empty", detail.DatasetDetailId); skippedRows++; 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: job.LastInsertedTill, lastUpdatedTill: job.LastUpdatedTill, extraWhere: job.QueryCondition, dbType: srcDbType); var baseParams = new Dictionary(); if (job.LastInsertedTill.HasValue && !string.IsNullOrWhiteSpace(detail.CreatedDateColumn)) baseParams["LastInsertedTill"] = job.LastInsertedTill.Value; if (job.LastUpdatedTill.HasValue && !string.IsNullOrWhiteSpace(detail.ModifiedDateColumn)) baseParams["LastUpdatedTill"] = job.LastUpdatedTill.Value; try { int chunkNumber = 0; await foreach (var chunk in _engine.StreamPagedAsync(srcConn, selectSql, chunkSize, baseParams, ct)) { chunkNumber++; var chunkList = chunk.ToList(); var written = await _engine.ApplyChunkToTargetAsync( tgtConn, tgtSchema, detail.DBObjectName!, chunkList, detail.PrimaryKeyColumn!, policy, detail.ModifiedDateColumn, srcDbType, tgtDbType, ct).ConfigureAwait(false); insertedRows += written; totalRows += chunkList.Count; GB5Trace.Step("process-chunk", new { runId, tableName = detail.DBObjectName, chunkNumber, rowCount = chunkList.Count }); if (_hubContext is not null) await _hubContext.PushProgressAsync(runId, detail.DBObjectName!, chunkNumber, chunkList.Count, totalRows, ct).ConfigureAwait(false); } // Delete detection if (job.LastDeleteTill.HasValue) { var deletes = await _engine.ApplyDeletesAsync( srcConn, tgtConn, tgtSchema, detail.DBObjectName!, detail.PrimaryKeyColumn!, job.LastDeleteTill.Value, srcDbType, tgtDbType, ct).ConfigureAwait(false); deletedRows += deletes; } } catch (Exception ex) { finalStatus = 4; // 4=PartialFailure var msg = $"Table {detail.DBObjectName}: {ex.Message}"; tableErrors.Add(msg); _logger.LogError(ex, "Sync failed for table {Table} in Run {RunId}", detail.DBObjectName, runId); GB5Trace.MarkFailed("process-table-failed", ex); } } if (tableErrors.Count > 0) errorDetail = string.Join(Environment.NewLine, tableErrors); // Update watermarks only on full success if (finalStatus == 2) { var now = DateTime.UtcNow; GB5Trace.Step("update-watermarks", new { syncJobId, now }); await _syncJobDal.UpdateWatermarks(syncJobId, now, now, now, login, ct).ConfigureAwait(false); } } catch (Exception ex) { finalStatus = 3; // 3=Failed errorDetail = ex.Message; GB5Trace.MarkFailed("syncjob-run-failed", ex); _logger.LogError(ex, "SyncJob {SyncJobId} Run {RunId} failed", syncJobId, runId); } finally { // Complete the run log var completedLog = new SyncRunLogDTO { RunId = runId, Status = finalStatus, TotalRows = totalRows, InsertedRows = insertedRows, UpdatedRows = updatedRows, DeletedRows = deletedRows, SkippedRows = skippedRows, ErrorDetail = errorDetail }; await _runLogDal.UpdateSyncRunLogComplete(completedLog, login, CancellationToken.None) .ConfigureAwait(false); await _syncJobDal.UpdateLastRun(syncJobId, runId, finalStatus, runAt, login, CancellationToken.None) .ConfigureAwait(false); GB5Trace.Step("event-publish", new { EventTypeConstant.DSYNCRUNCOMPLETEDEVENTTYPEID }); await _eventLog.PublishEventLogAsync( finalStatus == 2 ? $"DataSync Run Completed: SyncJobId={syncJobId} RunId={runId} Rows={totalRows}" : $"DataSync Run {(finalStatus == 3 ? "Failed" : "PartialFailure")}: SyncJobId={syncJobId} RunId={runId}", new { SyncJobId = syncJobId, RunId = runId, Status = finalStatus }, finalStatus == 2 ? EventTypeConstant.DSYNCRUNCOMPLETEDEVENTTYPEID : EventTypeConstant.DSYNCRUNFAILEDEVENTTYPEID, runId, login, "eventlog-topic" ).ConfigureAwait(false); if (_hubContext is not null) await _hubContext.PushRunCompletedAsync( runId, syncJobId, finalStatus, totalRows, insertedRows, updatedRows, deletedRows, skippedRows, errorDetail, runAt, CancellationToken.None).ConfigureAwait(false); } } private static string DefaultSchema(DatabaseType dbType) => dbType switch { DatabaseType.PostgreSQL => "public", DatabaseType.MySQL => string.Empty, _ => "dbo" }; } }