using FastEndpoints; using GB5Shared.IceMap; using GB5Shared.DTO.Ice; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ResponseStandard; using GB5Shared.FastEndPoint; using IceImportBLL.IceImportRun; using IceImportBLL.IceImportUploadStagingService; using IceImportSL.Parameters.Import; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using static GB5Shared.GB5Constant.Constant; namespace IceImportSL.Endpoints.Import; // Phase 2: enqueue-and-return. Inserts ICEIMPORT.TICEIMPORTRUN{Status=Queued} via // IIceImportRunBLL.EnqueueRunAsync, returns { RunId } immediately, then runs the actual // Parse -> MapAndTransform -> Validate -> Commit pipeline as a detached background Task on its own // IServiceScopeFactory-created DI scope — mirroring FrameworkBLL.DataSync.DataSyncRunBLL's // "new scope, not the request scope" pattern (the request scope — and this endpoint instance — is // disposed the instant ExecuteAsync returns, so the background task must resolve everything fresh). // // Phase 1's synchronous IIceImportPipeline.ExecuteAsync call is gone from here — run-tracked, // stage-by-stage execution now lives in IceImportRunExecutionService (BLL), which this background // Task resolves and invokes. public class CommitImport : BaseEndpoint> { private readonly IIceImportUploadStagingService _staging; private readonly IIceImportRunBLL _runBll; private readonly IIceMapBLL _iceMapBLL; private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public CommitImport( IIceImportUploadStagingService staging, IIceImportRunBLL runBll, IIceMapBLL iceMapBLL, IServiceScopeFactory scopeFactory, ILogger logger) { _staging = staging; _runBll = runBll; _iceMapBLL = iceMapBLL; _scopeFactory = scopeFactory; _logger = logger; } public override void Configure() { Post("/IceImport/CommitImport"); AllowAnonymous(); } // Mutation endpoint — never cached. protected override string? GetCacheKey(CommitImportParameters req, LoginDTO login) => null; protected override async Task> ExecuteAsync( CommitImportParameters req, LoginDTO login, CancellationToken ct) { try { if (string.IsNullOrWhiteSpace(req.UploadToken)) throw new ArgumentException("UploadToken is required."); // Validate the token up front (tenant/IceMapId/expiry) and capture its StorageKey so // EnqueueRunAsync can record a "SourceFile" artifact immediately — needed later by // RetryImportRun even if the run never gets past Queued. The stream itself is opened // and immediately discarded here; the background task re-opens its own via // IIceImportUploadStagingService.OpenStagedAsync inside a fresh scope rather than // inheriting this request-scoped stream (which would be disposed the moment this method // returns, long before the background pipeline finishes reading it). string storageKey; string fileName; { var staged = await _staging .OpenStagedAsync(req.UploadToken, req.IceMapId, login, ct) .ConfigureAwait(false); await using var _ = staged.Stream; storageKey = staged.StorageKey; fileName = staged.FileName; } var json = await _iceMapBLL.GetIceMap(req.IceMapId, login).ConfigureAwait(false); var mapDefinition = JsonConvert.DeserializeObject(json) ?? throw new InvalidOperationException($"IceMap {req.IceMapId} was not found."); // MICEMAP.ICESOURCETYPE: 0-Excel, 1-Json, 2-Xml, 3-Dbf, 4-Csv, 5-Mdb, 6-Sql. Only Excel // is handled by /IceImport/CommitImport today — ExecuteTrackedRunAsync below must never // be reached for anything else, or for an Excel map whose uploaded file isn't actually // .xls/.xlsx. if (mapDefinition.IceMapIceSourceType != 0) { throw new NotSupportedException( $"Given Type ({mapDefinition.IceMapIceSourceType}) is not supported by /IceImport/CommitImport — only Excel (ICESOURCETYPE=0) is currently supported."); } var ext = Path.GetExtension(fileName); if (!string.Equals(ext, ".xls", StringComparison.OrdinalIgnoreCase) && !string.Equals(ext, ".xlsx", StringComparison.OrdinalIgnoreCase)) { throw new ArgumentException( $"This IceMap's ICESOURCETYPE is Excel (0), but the uploaded file '{fileName}' is not a .xls/.xlsx file."); } var runId = await _runBll .EnqueueRunAsync(req.IceMapId, $"Manual:{login.UserId}", storageKey, null, login, ct) .ConfigureAwait(false); var scopeFactory = _scopeFactory; var logger = _logger; var uploadToken = req.UploadToken; var iceMapId = req.IceMapId; _ = Task.Run(async () => { await using var scope = scopeFactory.CreateAsyncScope(); var executionService = scope.ServiceProvider.GetRequiredService(); try { await executionService .ExecuteTrackedRunAsync(runId, iceMapId, uploadToken, null, login, CancellationToken.None) .ConfigureAwait(false); } catch (Exception ex) { // ExecuteTrackedRunAsync already wraps its own body in try/catch and always // calls CompleteRunAsync — this is the last-resort net if even that fails. logger.LogError(ex, "IceImport background execution failed for RunId {RunId}", runId); } }, CancellationToken.None); //await using var scope = scopeFactory.CreateAsyncScope(); //var executionService = scope.ServiceProvider.GetRequiredService(); //try //{ // await executionService // .ExecuteTrackedRunAsync(runId, iceMapId, uploadToken, null, login, CancellationToken.None) // .ConfigureAwait(false); //} //catch (Exception ex) //{ // // ExecuteTrackedRunAsync already wraps its own body in try/catch and always // // calls CompleteRunAsync — this is the last-resort net if even that fails. // logger.LogError(ex, "IceImport background execution failed for RunId {RunId}", runId); //} return await GB5Shared.ResponseStandard.Response.CreateSuccessResponse( new { RunId = runId }, CacheKeyLevel.NOT_REQUIRED, login); } catch (Exception ex) { return await GB5Shared.ResponseStandard.Response.CreateExceptionError( ex, CacheKeyLevel.NOT_REQUIRED, login, ex.Message, 500); } } }