using GB5Shared.DTO.Framework.Login; using GB5Shared.QueryExecutor; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using IceImportBLL.IceImportRunNotifier; using IceImportDAL.CustomCode.IceImportRun; using IceImportDAL.DTO.IceImportRun; using Microsoft.Extensions.Logging; namespace IceImportBLL.IceImportRun; public class IceImportRunBLL : IIceImportRunBLL { // CK_TICEIMPORTRUN_STATUS (Phase 0 migration) allows only these terminal values — no // PartialFailure. IceImportRunExecutionService maps a run with any failed/skipped rows to // 'Failed' (with RowTotalCount/SuccessCount/FailedCount plus a summary ErrorMessage carrying the // partial-failure detail) rather than inventing a status value the CHECK constraint would reject. private static readonly HashSet TerminalStatuses = new(StringComparer.OrdinalIgnoreCase) { "Success", "Failed", "Cancelled" }; private readonly IIceImportRunDAL _dal; private readonly IIceImportRunNotifier _notifier; private readonly ILogger _logger; public IceImportRunBLL(IIceImportRunDAL dal, IIceImportRunNotifier notifier, ILogger logger) { _dal = dal; _notifier = notifier; _logger = logger; } public async Task EnqueueRunAsync( int iceMapId, string triggeredBy, string? sourceFileStorageKey, int? sourceProfileId, LoginDTO login, CancellationToken ct, long? retryOfRunId = null) { try { GB5Trace.Step("validate-enqueue-run", new { iceMapId, triggeredBy, retryOfRunId }); var runId = await _dal .InsertRunAsync(iceMapId, triggeredBy, sourceProfileId, retryOfRunId, login, ct) .ConfigureAwait(false); GB5Trace.Step("enqueue-run", new { iceMapId, runId, triggeredBy }); if (!string.IsNullOrWhiteSpace(sourceFileStorageKey)) { var artifactId = await _dal .InsertArtifactAsync(runId, "SourceFile", sourceFileStorageKey, login, ct) .ConfigureAwait(false); await _dal.UpdateSourceFileArtifactAsync(runId, artifactId, login, ct).ConfigureAwait(false); } return runId; } catch (Exception ex) { GB5Trace.MarkFailed("enqueue-import-run-failed", ex); _logger.LogError(ex, "EnqueueRunAsync failed for IceMapId {IceMapId} TriggeredBy {TriggeredBy}", iceMapId, triggeredBy); throw; } } public async Task UpdateStageAsync(long runId, string currentStage, LoginDTO login, CancellationToken ct) { GB5Trace.Step("update-stage", new { runId, currentStage }); var rows = await _dal.UpdateStageAsync(runId, currentStage, login, ct).ConfigureAwait(false); if (rows == 0) _logger.LogWarning( "UpdateStageAsync: RunId {RunId} was not in Queued/Running state (already cancelled/completed?) — stage {Stage} not applied.", runId, currentStage); } public async Task AppendLogAsync(long runId, string logLevel, string message, LoginDTO login, CancellationToken ct) { _ = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {runId} not found."); await _dal.InsertLogAsync(runId, logLevel, message, login, ct).ConfigureAwait(false); await _notifier.NotifyLogAppendedAsync(runId, logLevel, message, ct).ConfigureAwait(false); } public async Task CompleteRunAsync( long runId, string finalStatus, int rowTotalCount, int successCount, int failedCount, string? errorMessage, LoginDTO login, CancellationToken ct) { try { if (!TerminalStatuses.Contains(finalStatus)) throw new InvalidOperationException( $"finalStatus must be one of: {string.Join(", ", TerminalStatuses)}"); GB5Trace.Step("complete-import-run", new { runId, finalStatus, rowTotalCount, successCount, failedCount }); var rows = await _dal .CompleteRunAsync(runId, finalStatus, rowTotalCount, successCount, failedCount, errorMessage, login, ct) .ConfigureAwait(false); if (rows == 0) _logger.LogWarning( "CompleteRunAsync: RunId {RunId} was not in Queued/Running state (already completed/cancelled concurrently) — status {Status} not applied.", runId, finalStatus); await _notifier .NotifyCompletedAsync(runId, finalStatus, rowTotalCount, successCount, failedCount, errorMessage, ct) .ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("complete-import-run-failed", ex); _logger.LogError(ex, "CompleteRunAsync failed for RunId {RunId}", runId); throw; } } public async Task SaveArtifactAsync(long runId, string artifactType, string storageKey, LoginDTO login, CancellationToken ct) { _ = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {runId} not found."); GB5Trace.Step("save-import-run-artifact", new { runId, artifactType }); return await _dal.InsertArtifactAsync(runId, artifactType, storageKey, login, ct).ConfigureAwait(false); } public async Task SaveRowResultsAsync(long runId, IEnumerable results, LoginDTO login, CancellationToken ct) { _ = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {runId} not found."); var list = results.ToList(); GB5Trace.Step("save-import-row-results", new { runId, Count = list.Count }); await _dal.InsertRowResultsAsync(runId, list, login, ct).ConfigureAwait(false); } public async Task> GetHistoryAsync( int iceMapId, int page, int pageSize, LoginDTO login, CancellationToken ct) { var (offset, size) = NormalizePaging(page, pageSize); return await _dal.GetHistoryAsync(iceMapId, offset, size, login, ct).ConfigureAwait(false); } public async Task GetDetailAsync(long runId, LoginDTO login, CancellationToken ct) { var run = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false); if (run is null) return null; var logs = await _dal.GetLogsAsync(runId, login, ct).ConfigureAwait(false); var artifacts = await _dal.GetArtifactsAsync(runId, login, ct).ConfigureAwait(false); return new IceImportRunDetailDTO { Run = run, Logs = logs, Artifacts = artifacts }; } public async Task> GetRowErrorsAsync( long runId, string? status, int page, int pageSize, LoginDTO login, CancellationToken ct) { _ = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {runId} not found."); var (offset, size) = NormalizePaging(page, pageSize); return await _dal.GetRowErrorsAsync(runId, status, offset, size, login, ct).ConfigureAwait(false); } public async Task CancelRunAsync(long runId, LoginDTO login, CancellationToken ct) { try { GB5Trace.Step("cancel-import-run", new { runId }); var rows = await _dal.CancelRunAsync(runId, login, ct).ConfigureAwait(false); if (rows == 0) throw new InvalidOperationException( $"IceImport run {runId} was not found or is not in a cancellable (Queued/Running) state."); await _notifier .NotifyCompletedAsync(runId, "Cancelled", 0, 0, 0, "Run was cancelled by user request.", ct) .ConfigureAwait(false); return SuccessResponse.UpdateSuccess; } catch (Exception ex) { GB5Trace.MarkFailed("cancel-import-run-failed", ex); _logger.LogError(ex, "CancelRunAsync failed for RunId {RunId}", runId); throw; } } public async Task GetByIdAsync(long runId, LoginDTO login, CancellationToken ct) => await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false); public async Task GetArtifactByTypeAsync( long runId, string artifactType, LoginDTO login, CancellationToken ct) { _ = await _dal.GetByIdAsync(runId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {runId} not found."); return await _dal.GetLatestArtifactByTypeAsync(runId, artifactType, login, ct).ConfigureAwait(false); } public async Task PrepareRetryAsync(long originalRunId, LoginDTO login, CancellationToken ct) { var original = await _dal.GetByIdAsync(originalRunId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException($"IceImport run {originalRunId} not found."); if (!string.Equals(original.Status, "Failed", StringComparison.OrdinalIgnoreCase) && !string.Equals(original.Status, "Cancelled", StringComparison.OrdinalIgnoreCase)) throw new InvalidOperationException( $"IceImport run {originalRunId} is in status '{original.Status}' — only Failed/Cancelled runs can be retried."); var sourceArtifact = await _dal .GetLatestArtifactByTypeAsync(originalRunId, "SourceFile", login, ct) .ConfigureAwait(false) ?? throw new InvalidOperationException( $"IceImport run {originalRunId} has no recorded SourceFile artifact — cannot retry."); GB5Trace.Step("prepare-retry-import-run", new { originalRunId, original.IceMapId }); return await EnqueueRunAsync( original.IceMapId, $"Retry:{login.UserId}", sourceArtifact.StorageKey, original.SourceProfileId, login, ct, retryOfRunId: originalRunId).ConfigureAwait(false); } private static (int Offset, int PageSize) NormalizePaging(int page, int pageSize) { var size = pageSize <= 0 ? 20 : Math.Min(pageSize, 200); var safePage = page <= 0 ? 1 : page; return ((safePage - 1) * size, size); } }