using FrameworkDAL.CustomCode.EventSub; using FrameworkDAL.DTO.EventSub; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkBLL.EventSub { public interface IDaprSubscriptionBLL { /// /// Returns one subscription entry per active EventTypeId in MEVENTTYPE, /// all routed to /eventsub/handle. Called by Dapr at startup via /// GET /dapr/subscribe. /// Task> GetSubscriptionsAsync(CancellationToken ct = default); } public sealed class DaprSubscriptionBLL : IDaprSubscriptionBLL { private readonly IDaprSubscriptionDAL _dal; private readonly ILogger _logger; public DaprSubscriptionBLL(IDaprSubscriptionDAL dal, ILogger logger) { _dal = dal; _logger = logger; } public async Task> GetSubscriptionsAsync( CancellationToken ct = default) { const int MaxAttempts = 3; const int RetryDelayMs = 2_000; Exception? lastEx = null; for (int attempt = 1; attempt <= MaxAttempts; attempt++) { try { var eventTypeIds = await _dal.GetActiveEventTypeIdsAsync(ct).ConfigureAwait(false); // One subscription per EventTypeId — routed to /eventsub/handle. // EventActionSubscribeController fans out to both EventSubBLL (actions) // and EventLogSubBLL (audit) internally. Dapr does not support two // routes for the same pubsubname+topic; duplicate entries are deduped // and only the last one is used. var subscriptions = new List(eventTypeIds.Count()); foreach (var id in eventTypeIds) { subscriptions.Add(new DaprSubscriptionDTO { PubsubName = PUBLISHTYPE.PUBSUB, Topic = $"EVENTTYPEID:{id}", Route = "/eventsub/handle" }); } // Fixed subscription — routes ATTACHMENT-RESOLVE messages to the attachment // resolve controller regardless of EventTypeId configuration. subscriptions.Add(new DaprSubscriptionDTO { PubsubName = PUBLISHTYPE.PUBSUB, Topic = "ATTACHMENT-RESOLVE", Route = "/attachment/resolve" }); // Fixed subscription — SqlWorkbench publishes this literal topic name (not // an EVENTTYPEID:{id} topic, and not driven by any MEVENTTYPE row), so it // needs its own static entry here, same as ATTACHMENT-RESOLVE above. subscriptions.Add(new DaprSubscriptionDTO { PubsubName = PUBLISHTYPE.PUBSUB, Topic = "sqlworkbench.changerequest.applied", Route = "/versionsync/handle" }); _logger.LogInformation( "DaprSubscription: {EventTypeCount} EventTypeId(s) → {SubCount} subscription(s) registered (includes ATTACHMENT-RESOLVE, sqlworkbench.changerequest.applied)", eventTypeIds.Count(), subscriptions.Count); return subscriptions; } catch (OperationCanceledException) when (ct.IsCancellationRequested) { throw; } catch (Exception ex) { lastEx = ex; _logger.LogWarning(ex, "DaprSubscription: attempt {Attempt}/{Max} failed — {Error}", attempt, MaxAttempts, ex.Message); if (attempt < MaxAttempts) await Task.Delay(RetryDelayMs, ct).ConfigureAwait(false); } } // All retries exhausted — throw so the controller returns 503. // Dapr then retries GET /dapr/subscribe instead of locking in zero subscriptions. throw new InvalidOperationException( $"DaprSubscription: failed to load EventTypeIds after {MaxAttempts} attempts", lastEx); } } }