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;
}
}
}
}