using FastEndpoints; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.Framework.ResponseStandard; using GB5Shared.FastEndPoint; using IceImportBLL.ScheduledImportService; using IceImportSL.Parameters.FtpSource; using static GB5Shared.GB5Constant.Constant; namespace IceImportSL.Endpoints.FtpSource; // Phase 4 — FTP/SFTP ingestion, manually callable (recurring trigger deferred pending JobEngine; // see the plan's "Key architecture decisions" #2). This endpoint IS the whole "run this import // from the remote source now" operation — everything (list -> filter -> download -> stage -> // enqueue -> background pipeline kickoff -> mark processed -> post-process) happens inside // IScheduledImportService.RunScheduledImportAsync, which owns its own IServiceScopeFactory for the // per-file background kickoff (see that interface's doc comment for why — orchestration option // (a), not CommitImport's option (b)). This endpoint stays a thin pass-through by design, same // spirit as every other BaseEndpoint in this module. // // Deliberately scheduler-agnostic per the plan: no dependency on SysJob/JobEngine internals, no // special "internal-only" caller marker beyond the module's usual AllowAnonymous() — anything that // can call an HTTP endpoint can trigger this (a human via Swagger today; JobEngine, once // provisioned, with zero rework later). public class RunScheduledImport : BaseEndpoint> { private readonly IScheduledImportService _service; public RunScheduledImport(IScheduledImportService service) => _service = service; public override void Configure() { Post("/IceImport/RunScheduledImport"); AllowAnonymous(); } // Action endpoint, not a plain read — never cached, same as CommitImport/RetryImportRun. protected override string? GetCacheKey(RunScheduledImportParameters req, LoginDTO login) => null; protected override async Task> ExecuteAsync( RunScheduledImportParameters req, LoginDTO login, CancellationToken ct) { try { var triggeredBy = string.IsNullOrWhiteSpace(req.TriggeredBy) ? "Manual" : req.TriggeredBy!; var result = await _service .RunScheduledImportAsync(req.IceMapId, triggeredBy, login, ct) .ConfigureAwait(false); var message = result.FilesFound == 0 ? "No new files found." : $"{result.FilesProcessed} file(s) processed, {result.FilesSkippedAlreadyProcessed} already up to date, {result.Errors.Count} error(s)."; return await GB5Shared.ResponseStandard.Response.CreateSuccessResponse(new { result.RunIds, result.FilesFound, result.FilesProcessed, result.FilesSkippedAlreadyProcessed, result.Errors, Message = message }, CacheKeyLevel.NOT_REQUIRED, login); } catch (Exception ex) { return await GB5Shared.ResponseStandard.Response.CreateExceptionError( ex, CacheKeyLevel.NOT_REQUIRED, login, ex.Message, 500); } } }