using FrameworkBLL.EventSub; using GB5Shared.DTO.PubSub; using GB5Shared.Telemetry.Dapr; using Microsoft.AspNetCore.Mvc; using Microsoft.Extensions.Logging; using System; using System.Diagnostics; using System.Threading; using System.Threading.Tasks; namespace FrameworkSL.Controllers.EventSub { /// /// Dapr pub/sub subscriber — single entry point for every EVENTTYPEID:{id} topic. /// /// Fan-out: one Dapr subscription → two BLL calls in sequence: /// 1. EventSubBLL.ProcessEventAsync → writes TEVENTACTIONRUN + TACTIONOUTBOX /// 2. EventLogSubBLL.RecordEventAsync → writes audit row to TEVENTLOG /// /// Two separate Dapr subscriptions for the same pubsubname+topic are deduped by /// the sidecar (last one wins), so fan-out must happen here instead. /// /// Subscriptions registered dynamically by DaprSubscribeController (GET /dapr/subscribe). /// [ApiController] public class EventActionSubscribeController : ControllerBase { private readonly IEventSubBLL _bll; private readonly IEventLogSubBLL _logBll; private readonly ILogger _logger; public EventActionSubscribeController( IEventSubBLL bll, IEventLogSubBLL logBll, ILogger logger) { _bll = bll; _logBll = logBll; _logger = logger; } [HttpPost("/eventsub/handle")] public async Task HandleEntityEvent( [FromBody] OutboxEventMessage message, CancellationToken ct) { var topic = $"EVENTTYPEID:{message.EventTypeId}"; // Parent from the publisher's own span (embedded in the message body — see // OutboxEventMessage.TraceParent/TraceState) so this consumer span always links // back into the originating entity-save trace in Zipkin, regardless of whether // the Dapr broker/CloudEvent envelope forwards custom metadata through. using var subActivity = DaprEventTracing.StartSubscribeActivity( topic, "pubsub", message.TraceParent, message.TraceState); subActivity?.SetTag("gb5.event.type_id", message.EventTypeId); subActivity?.SetTag("gb5.tenant.id", message.TenantId); subActivity?.SetTag("gb5.object.id", message.ObjectId); subActivity?.SetTag("gb5.object.type_id", message.ObjectTypeId); subActivity?.SetTag("gb5.correlation.key", message.CorrelationKey ?? string.Empty); subActivity?.SetTag("gb5.connection.name", message.ConnectionName ?? "(missing)"); // "Who subscribed / what subscribed" — gb5.subscriber answers "who"; the // payload_length/payload pair answers "what data did it receive", both directly // on this span (gb5.actions.count/ids/types, set later in EventSubBLL, then answer // "what did it DO with that data" — who/what resolved as a result). subActivity?.SetTag("gb5.subscriber", "EventActionSubscribe"); subActivity?.SetTag("gb5.subscribe.payload_length", message.Payload?.Length ?? 0); subActivity?.SetTag("gb5.subscribe.payload", GB5Shared.Telemetry.GB5Trace.Preview(message.Payload)); _logger.LogInformation( "EventActionSubscribe: received | EventTypeId={EventTypeId} TenantId={TenantId} ObjectId={ObjectId}", message.EventTypeId, message.TenantId, message.ObjectId); // Capture the Framework service base URL from the incoming Dapr HTTP request. // Pattern: same as WIP capturing login.RequestUrl in BaseEndpoint. // "http://192.168.0.112:5000" — used by EventSubBLL to build email approval links. // Zero config: automatically correct for any deployment (dev, staging, prod). var frameworkBaseUrl = $"{Request.Scheme}://{Request.Host}"; try { // 1. Queue actions (TEVENTACTIONRUN + TACTIONOUTBOX) await _bll.ProcessEventAsync(message, frameworkBaseUrl, ct).ConfigureAwait(false); // 2. Write audit log (TEVENTLOG) — non-critical; log but don't NACK on failure try { await _logBll.RecordEventAsync(message, ct).ConfigureAwait(false); } catch (Exception exLog) { _logger.LogWarning(exLog, "EventActionSubscribe: event log write failed (non-critical) | EventTypeId={EventTypeId}", message.EventTypeId); } subActivity?.SetStatus(ActivityStatusCode.Ok); return Ok(); // 200 = ACK to Dapr } catch (Exception ex) { subActivity?.AddException(ex); subActivity?.SetTag("gb5.error.type", ex.GetType().Name); subActivity?.SetTag("gb5.error.message", ex.Message); subActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "EventActionSubscribe: failed | EventTypeId={EventTypeId} TenantId={TenantId}", message.EventTypeId, message.TenantId); return StatusCode(500); // 500 = NACK → Dapr retries } } } }