using GB5Shared.IceMap; using GB5Shared.DTO.Ice; using GB5Shared.DTO.Framework.Login; using GB5Shared.Telemetry; using IceImportBLL.IceImportRun; using IceImportBLL.IceImportSecretResolver; using IceImportBLL.IceImportUploadStagingService; using IceImportDAL.CustomCode.IceMapFtpSource; using IceImportDAL.CustomCode.RemoteFileSources; using IceImportDAL.DTO.RemoteFileSources; using IceImportDAL.DTO.ScheduledImportService; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Newtonsoft.Json; namespace IceImportBLL.ScheduledImportService; public class ScheduledImportService : IScheduledImportService { private readonly IIceMapBLL _iceMapBLL; private readonly IIceMapFtpSourceDAL _ftpSourceDal; private readonly IIceImportSecretResolver _secretResolver; private readonly IEnumerable _sources; private readonly IIceImportUploadStagingService _staging; private readonly IIceImportRunBLL _runBll; private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public ScheduledImportService( IIceMapBLL iceMapBLL, IIceMapFtpSourceDAL ftpSourceDal, IIceImportSecretResolver secretResolver, IEnumerable sources, IIceImportUploadStagingService staging, IIceImportRunBLL runBll, IServiceScopeFactory scopeFactory, ILogger logger) { _iceMapBLL = iceMapBLL; _ftpSourceDal = ftpSourceDal; _secretResolver = secretResolver; _sources = sources; _staging = staging; _runBll = runBll; _scopeFactory = scopeFactory; _logger = logger; } public async Task RunScheduledImportAsync( int iceMapId, string triggeredBy, LoginDTO login, CancellationToken ct) { var result = new ScheduledImportResultDTO(); try { GB5Trace.Step("validate-scheduled-import", new { iceMapId, triggeredBy }); await EnsureIceMapExistsAsync(iceMapId, login).ConfigureAwait(false); var rawConfig = await _ftpSourceDal.GetByIceMapIdAsync(iceMapId, login, ct).ConfigureAwait(false) ?? throw new InvalidOperationException( $"IceMap {iceMapId} has no MICEMAPFTPSOURCE configuration — call SaveIceMapFtpSource first."); var source = _sources.FirstOrDefault(s => s.ConnectionType == rawConfig.ConnectionType) ?? throw new NotSupportedException( $"No IRemoteFileSource is registered for ConnectionType {rawConfig.ConnectionType}."); var resolvedConfig = await _secretResolver.ResolveAsync(rawConfig, ct).ConfigureAwait(false); GB5Trace.Step("list-remote-files", new { iceMapId, rawConfig.ConnectionType }); var files = await source.ListFilesAsync(resolvedConfig, ct).ConfigureAwait(false); result.FilesFound = files.Count; foreach (var file in files) { ct.ThrowIfCancellationRequested(); await ProcessOneFileAsync(iceMapId, triggeredBy, source, resolvedConfig, file, login, ct, result) .ConfigureAwait(false); } return result; } catch (Exception ex) { GB5Trace.MarkFailed("run-scheduled-import-failed", ex); _logger.LogError(ex, "IceImport: RunScheduledImportAsync failed for IceMapId {IceMapId}", iceMapId); throw; } } private async Task EnsureIceMapExistsAsync(int iceMapId, LoginDTO login) { // Same ownership/tenant discipline IceImportRunExecutionService.LoadMapDefinitionAsync // already uses — GetIceMap resolves (or fails to resolve) against the caller's own tenant // connection, so a successful deserialize here is proof IceMapId belongs to this tenant. var json = await _iceMapBLL.GetIceMap(iceMapId, login).ConfigureAwait(false); _ = JsonConvert.DeserializeObject(json) ?? throw new InvalidOperationException($"IceMap {iceMapId} was not found."); } private async Task ProcessOneFileAsync( int iceMapId, string triggeredBy, IRemoteFileSource source, RemoteFileSourceConfigDTO resolvedConfig, RemoteFileInfoDTO file, LoginDTO login, CancellationToken ct, ScheduledImportResultDTO result) { try { var alreadyProcessed = await _ftpSourceDal .IsFileAlreadyProcessedAsync(iceMapId, file.FileName, file.LastModifiedUtc, login, ct) .ConfigureAwait(false); if (alreadyProcessed) { result.FilesSkippedAlreadyProcessed++; return; } GB5Trace.Step("download-remote-file", new { iceMapId, file.FileName }); string storageKey; var downloaded = await source.DownloadAsync(resolvedConfig, file, ct).ConfigureAwait(false); await using (downloaded) { // Stages the downloaded bytes through the exact same // IIceImportUploadStagingService a manual UploadImportFile call would use, so the // rest of the pipeline (IIceImportRunExecutionService) can't tell the difference // between a browser upload and a remote-source pull — zero duplicated import logic. var uploadToken = await _staging .StageAsync(downloaded, file.FileName, iceMapId, login, ct) .ConfigureAwait(false); // Immediately unprotects the token just issued to recover its StorageKey — same // "stage then immediately open once for the StorageKey" dance CommitImport performs // on the client's pre-existing upload token. The opened stream is discarded unread; // only StorageKey is needed here. var staged = await _staging.OpenStagedAsync(uploadToken, iceMapId, login, ct).ConfigureAwait(false); await using var openedStream = staged.Stream; storageKey = staged.StorageKey; } GB5Trace.Step("enqueue-scheduled-run", new { iceMapId, file.FileName, triggeredBy }); var runId = await _runBll .EnqueueRunAsync(iceMapId, triggeredBy, storageKey, null, login, ct) .ConfigureAwait(false); await _ftpSourceDal .MarkFileProcessedAsync(iceMapId, file.FileName, file.LastModifiedUtc, runId, login, ct) .ConfigureAwait(false); KickOffBackgroundExecution(runId, iceMapId, storageKey, login); result.RunIds.Add(runId); result.FilesProcessed++; try { await source.PostProcessAsync(resolvedConfig, file, ct).ConfigureAwait(false); } catch (Exception postProcessEx) { // Best-effort per IRemoteFileSource's contract — a post-process failure (e.g. // archive-folder permissions) must never be allowed to undo an import that was // already enqueued and is now running. GB5Trace.RecordError(postProcessEx, "ice-import-ftp-postprocess-failed"); _logger.LogWarning(postProcessEx, "IceImport: PostProcessAsync failed for IceMapId {IceMapId} file {FileName} (RunId {RunId} was still enqueued).", iceMapId, file.FileName, runId); } } catch (Exception fileEx) { // One bad file's error must not abort the rest of the batch (mirrors // IceImportRunExecutionService's per-stage isolation philosophy, applied per-file here). GB5Trace.RecordError(fileEx, "ice-import-scheduled-file-failed"); _logger.LogError(fileEx, "IceImport: RunScheduledImport failed for IceMapId {IceMapId} file {FileName}", iceMapId, file.FileName); result.Errors.Add($"{file.FileName}: {fileEx.Message}"); } } // Mirrors IceImportSL.Endpoints.Import.CommitImport's detached-Task-on-a-fresh-scope pattern exactly — // see that file's doc comment for why a fresh IServiceScopeFactory scope (not this method's own // scope) is required: whatever scope RunScheduledImportAsync itself is running in will be // disposed the moment the HTTP request that triggered it returns, long before the pipeline // finishes. private void KickOffBackgroundExecution(long runId, int iceMapId, string storageKey, LoginDTO login) { var scopeFactory = _scopeFactory; var logger = _logger; _ = Task.Run(async () => { await using var scope = scopeFactory.CreateAsyncScope(); var executionService = scope.ServiceProvider.GetRequiredService(); try { await executionService .ExecuteTrackedRunAsync(runId, iceMapId, uploadToken: null, sourceStorageKey: storageKey, 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} (RunScheduledImport)", runId); } }, CancellationToken.None); } }