using Dapr.Client; using FrameworkDAL.DTO.SchedulerTaskGenerator; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.PubSub; using GB5Shared.QueryExecutor; using GB5Shared.QueueReader; using GB5Shared.Telemetry; using GB5Shared.Telemetry.Dapr; using Microsoft.Extensions.Logging; using System; using System.Diagnostics; using System.Text.Json; using System.Threading; using System.Threading.Tasks; using static GB5Shared.GB5Constant.Constant; namespace FrameworkSL.Controllers.SchedulerTaskGenerator { /// /// TJOBQUEUE handler for MessageType="SCHEDULER_ACTION_EVENT" — publishes a completed /// scheduler job's action-dispatch event to Dapr pubsub, feeding the SAME subscriber /// already used for every other MACTION-driven action (mail/SMS/webhook/etc): /// /// SchedulerExecutionDeliveryService (on success) -> enqueue TJOBQUEUE row -> /// this handler -> Dapr publish, topic "EVENTTYPEID:{SchedulerJobEventTypeId}" -> /// EventActionSubscribeController (/eventsub/handle) -> EventSubBLL.ProcessEventAsync /// -> resolves MACTION via MEVENTTYPEACTION/MEVENTTYPEACTIONDETAIL (already wired by /// AutoSchedulerDAL when the job's action was saved) -> TEVENTACTIONRUN + TACTIONOUTBOX /// -> ActionOutboxDispatcherQuartzJob -> RabbitMQ action-exec-p{n} -> ActionProcessorWorker /// -> actual send. /// /// No new subscriber, no new MACTION lookup — this handler's only job is the Dapr /// publish. Retry/backoff/DLQ on publish failure is the existing generic /// BaseJobQueueProcessor infrastructure (GB5Shared/QueueReader) — the STATUS column drives /// that. ISPUBLISHED/PUBLISHEDON below are separate, explicit "did the Dapr publish itself /// succeed" bookkeeping (mirroring TOUTBOX's own ISPUBLISHED-style convention), stamped by /// this handler directly rather than the generic base processor. /// public sealed class SchedulerActionEventPublishHandler : IJobQueueHandler { private const string MarkPublishedSql = @" UPDATE TJOBQUEUE SET ISPUBLISHED = 1, PUBLISHEDON = GETUTCDATE() WHERE QUEUEID = @QueueId"; public string MessageType => SchedulerActionEventConstants.MessageType; private readonly DaprClient _dapr; private readonly IQueryExecutor _queryExecutor; private readonly ILogger _logger; public SchedulerActionEventPublishHandler( DaprClient dapr, IQueryExecutor queryExecutor, ILogger logger) { _dapr = dapr; _queryExecutor = queryExecutor; _logger = logger; } public async Task HandleAsync(string payload, LoginDTO login, long queueId, CancellationToken ct) { var message = JsonSerializer.Deserialize(payload) ?? throw new InvalidOperationException( "SchedulerActionEventPublishHandler: payload deserialized to null"); var topic = $"EVENTTYPEID:{message.EventTypeId}"; // Reconstruct the "deliver" span (a child of the original "JobScheduler" root) // that SchedulerExecutionDeliveryService persisted into message.TraceParent, so // the publish span below — and everything downstream of it — stays in the SAME // trace even though this runs in a separate poll tick / async continuation. var publishParentContext = message.TraceParent is not null && ActivityContext.TryParse(message.TraceParent, message.TraceState, isRemote: true, out var publishCtx) ? publishCtx : default; using var handlerActivity = publishParentContext != default ? GB5ActivitySources.Scheduler.StartActivity( "publish-action-event", ActivityKind.Internal, publishParentContext) : GB5ActivitySources.Scheduler.StartActivity("publish-action-event", ActivityKind.Internal); GB5Trace.Step("scheduler-action-event-publish", new { message.EventTypeId, message.ObjectId, message.TenantId }); try { // Opens its own child span (auto-nests under handlerActivity, now // Activity.Current) and, critically, re-points message.TraceParent/TraceState // at ITSELF — matching OutBox.PublishPendingEventsAsync's exact pattern — so // EventActionSubscribeController's existing StartSubscribeActivity(topic, // "pubsub", message.TraceParent, message.TraceState) parents from the publish // span specifically, not the broader handler span. using var publishActivity = DaprEventTracing.StartPublishActivity(topic, PUBLISHTYPE.PUBSUB); message.TraceParent = publishActivity?.Id ?? message.TraceParent; message.TraceState = publishActivity?.TraceStateString ?? message.TraceState; // "Who published / what published" — visible directly on this span so a // reviewer never has to leave Zipkin to answer either question. publishActivity?.SetTag("gb5.publisher", nameof(SchedulerActionEventPublishHandler)); publishActivity?.SetTag("gb5.publish.event_type_id", message.EventTypeId); publishActivity?.SetTag("gb5.publish.object_id", message.ObjectId); publishActivity?.SetTag("gb5.publish.tenant_id", message.TenantId); publishActivity?.SetTag("gb5.publish.correlation_key", message.CorrelationKey ?? string.Empty); publishActivity?.SetTag("gb5.publish.payload_length", message.Payload?.Length ?? 0); publishActivity?.SetTag("gb5.publish.payload", GB5Trace.Preview(message.Payload)); await _dapr.PublishEventAsync(PUBLISHTYPE.PUBSUB, topic, message, ct) .ConfigureAwait(false); await _queryExecutor.ExecuteAsync(login, MarkPublishedSql, new { QueueId = queueId }, cancellationToken: ct) .ConfigureAwait(false); _logger.LogInformation( "SchedulerActionEventPublishHandler: published | Topic={Topic} JobId={JobId} TenantId={TenantId} QueueId={QueueId}", topic, message.ObjectId, message.TenantId, queueId); } catch (Exception ex) { GB5Trace.MarkFailed("scheduler-action-event-publish-failed", ex); _logger.LogError(ex, "SchedulerActionEventPublishHandler: publish failed | Topic={Topic} JobId={JobId} TenantId={TenantId}", topic, message.ObjectId, message.TenantId); // Rethrow so BaseJobQueueProcessor applies its own retry/backoff/DLQ to the // TJOBQUEUE row (30s->2m->8m->30m->2h, then DLQ) — same as every other // message type on this queue. ISPUBLISHED stays 0/PUBLISHEDON stays NULL until // a retry actually succeeds. throw; } } } }