using System; using System.Text.Json; using System.Threading.Tasks; using FrameworkBLL.GOP; using FrameworkBLL.Promotion; using GB5Shared.DTO.GOP; using Microsoft.Extensions.Logging; using Quartz; namespace FrameworkSL.Controllers.Promotion { // ============================================================ // PromotionExchangeIngestQuartzJob — the destination-side "check my pending // packages" poller for Tier 2 (Exchange) delivery. On each fire: lists what's // waiting for this deployment's own identity (resolved server-side by the // Exchange from the configured API key, never a parameter here), downloads // and opens each one, then submits one GOP execution per bundle — the exact // same sequence ImportPromotionPackageService already runs for a manually // uploaded Tier 1 file, just sourced from the Exchange instead of an // IFormFile. Ingest -> Approval -> Import itself is unchanged either way. // // Never throws back to Quartz (outer try/catch is a final safety net, // matching OutboxPublishQuartzJob) and never lets one bad package abort the // rest of the batch (per-item try/catch, matching // ImportPromotionPackageService's own per-item isolation). // ============================================================ public class PromotionExchangeIngestQuartzJob : IJob { private readonly IPromotionExchangeClient _ExchangeClient; private readonly IPromotionBLL _PromotionBLL; private readonly IGopQueueBLL _GopQueueBLL; private readonly IPromotionSystemContext _SystemContext; private readonly ILogger _Logger; public PromotionExchangeIngestQuartzJob( IPromotionExchangeClient exchangeClient, IPromotionBLL promotionBLL, IGopQueueBLL gopQueueBLL, IPromotionSystemContext systemContext, ILogger logger) { _ExchangeClient = exchangeClient; _PromotionBLL = promotionBLL; _GopQueueBLL = gopQueueBLL; _SystemContext = systemContext; _Logger = logger; } public async Task Execute(IJobExecutionContext context) { try { var ct = context.CancellationToken; var login = _SystemContext.GetSystemLogin(); var pending = await _ExchangeClient.ListPendingAsync(ct).ConfigureAwait(false); if (pending.Count == 0) return; foreach (var item in pending) { try { var packageBytes = await _ExchangeClient.DownloadAsync(item.PackageId, ct).ConfigureAwait(false); if (packageBytes is null) { _Logger.LogWarning( "PromotionExchangeIngest: could not download PackageId {PackageId} from {Origin}", item.PackageId, item.OriginEnvironmentCode); continue; } var manifest = await _PromotionBLL.OpenPackageAsync(packageBytes, ct).ConfigureAwait(false); foreach (var bundle in manifest.Items) { try { await _GopQueueBLL.SubmitExecution( new GopSubmitRequestDTO { SourceCode = "PromotionIngest", PayloadJson = JsonSerializer.Serialize(bundle), Priority = 5 }, login).ConfigureAwait(false); } catch (Exception itemEx) { _Logger.LogError(itemEx, "PromotionExchangeIngest: SubmitExecution failed for {EntityTypeCode}/{OriginEntityId} in ExchangePackageId {PackageId}", bundle.EntityTypeCode, bundle.OriginEntityId, item.PackageId); } } } catch (Exception packageEx) { _Logger.LogError(packageEx, "PromotionExchangeIngest: failed processing ExchangePackageId {PackageId} from {Origin}", item.PackageId, item.OriginEnvironmentCode); } } } catch (Exception ex) { _Logger.LogError(ex, "PromotionExchangeIngestQuartzJob: unhandled failure"); } } } }