using System; using System.Collections.Generic; using System.Data.Common; using System.Diagnostics; using System.Text; using System.Text.Json; using System.Threading.Tasks; using Dapr.Client; using GB5Shared.DTO.Framework.Login; using GB5Shared.DTO.PubSub; using GB5Shared.GB5Constant; using GB5Shared.Query.OutBox; using GB5Shared.QueryExecutor; using GB5Shared.Telemetry; using GB5Shared.Telemetry.Dapr; using Microsoft.Extensions.Logging; using OpenTelemetry.Trace; using static GB5Shared.GB5Constant.Constant; namespace GB5Shared.PubSub.OutBox { public sealed class OutBox : IOutBox { private readonly DaprClient _dapr; private readonly IQueryExecutor _queryExecutor; private readonly ILogger _logger; public OutBox(DaprClient dapr, IQueryExecutor queryExecutor, ILogger logger) { _dapr = dapr ?? throw new ArgumentNullException(nameof(dapr)); _queryExecutor = queryExecutor ?? throw new ArgumentNullException(nameof(queryExecutor)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } // ------------------------------------------------------------------ // WRITE TO OUTBOX (inside business transaction) // ------------------------------------------------------------------ public async Task PublishEventAsync(OutboxDTO dto, LoginDTO login, DbTransaction Trans) { if (dto == null) throw new ArgumentNullException(nameof(dto)); if (login == null) throw new ArgumentNullException(nameof(login)); dto.CorrelationKey ??= Guid.NewGuid().ToString(); await _queryExecutor.ExecuteAsync(login, OutBoxQB.InsertOutBox, new { dto.ObjectTypeId, dto.ObjectId, dto.EventTypeId, dto.Version, dto.CorrelationKey, userid = login.UserId, tenantid = login.ClientId, Payload = JsonSerializer.Serialize(dto.Payload), Status = OUTBOXSTATUS.PENDING, dto.HeaderRowGuid, ContextJson = dto.ContextJson }); } // ------------------------------------------------------------------ // BACKGROUND PUBLISHER // ------------------------------------------------------------------ public async Task PublishPendingEventsAsync( string pubsubName, LoginDTO login, int batchSize = 100) { if (login == null) throw new ArgumentNullException(nameof(login)); IEnumerable events; try { events = await _queryExecutor.QueryAsync( login.ConnectionDatabaseName, OutBoxQB.fetchSql, new { BatchSize = batchSize, Pending = OUTBOXSTATUS.PENDING }); } catch (Exception ex) { _logger.LogError(ex, "[outbox] FETCH FAILED | TENANTID={TenantId} DB={Db} | Error={Error}", login.ClientId, login.ConnectionDatabaseName, ex.Message); throw; } int published = 0, failed = 0; foreach (var evt in events) { // Extract identifiers up-front; any cast failure is caught below. int outboxId = -1; int eventTypeId = -1; int objectTypeId= -1; int objectId = -1; int tenantId = -1; int userId = -1; string topic = "UNKNOWN"; Activity? publishActivity = null; try { outboxId = Convert.ToInt32(evt.OUTBOXID); eventTypeId = Convert.ToInt32(evt.EVENTTYPEID); objectTypeId = Convert.ToInt32(evt.OBJECTTYPEID); objectId = Convert.ToInt32(evt.OBJECTID); tenantId = Convert.ToInt32(evt.TENANTID); userId = Convert.ToInt32(evt.USERID); topic = $"EVENTTYPEID:{eventTypeId}"; var correlationKey = evt.CORRELATIONKEY == null ? null : Convert.ToString(evt.CORRELATIONKEY); _logger.LogInformation( "[outbox] Processing | OUTBOXID={OutboxId} | Topic={Topic} | OBJECTTYPEID={ObjectTypeId} | OBJECTID={ObjectId} | TENANTID={TenantId}", outboxId, topic, objectTypeId, objectId, tenantId); // ── 1. Publish main domain event ────────────────────────────── publishActivity = DaprEventTracing.StartPublishActivity(topic, pubsubName); publishActivity?.SetTag("gb5.event.type_id", eventTypeId); publishActivity?.SetTag("gb5.tenant.id", tenantId); publishActivity?.SetTag("gb5.object.type_id", objectTypeId); publishActivity?.SetTag("gb5.object.id", objectId); publishActivity?.SetTag("gb5.outbox.id", outboxId); publishActivity?.SetTag("gb5.user.id", userId); publishActivity?.SetTag("gb5.correlation.key", correlationKey ?? string.Empty); // login.ConnectionDatabaseName is the MSERVER connection name that was used // to read this TOUTBOX row. Carry it in the message so the framework // subscriber connects to the exact same tenant DB — never re-derives it from // MSERVERCONFIG, which can have multiple STATUS=1 records for the same ClientId. publishActivity?.SetTag("gb5.connection.name", login.ConnectionDatabaseName ?? string.Empty); await _dapr.PublishEventAsync( pubsubName, topic, new OutboxEventMessage { EventTypeId = eventTypeId, ObjectTypeId = objectTypeId, ObjectId = objectId, TenantId = tenantId, UserId = userId, Payload = evt.PAYLOAD == null ? null : Convert.ToString(evt.PAYLOAD), CorrelationKey = correlationKey, HeaderRowGuid = evt.HEADERROWGUID == null ? (Guid?)null : (Guid)evt.HEADERROWGUID, ConnectionName = login.ConnectionDatabaseName ?? string.Empty, ContextJson = string.IsNullOrWhiteSpace(Convert.ToString(evt.CONTEXTJSON)) ? "{}" : Convert.ToString(evt.CONTEXTJSON)!, // Carry this producer span's W3C context in the message itself so // EventActionSubscribeController can parent its consumer span correctly // regardless of Dapr broker/CloudEvent metadata propagation. TraceParent = publishActivity?.Id, TraceState = publishActivity?.TraceStateString }); publishActivity?.SetTag("gb5.outbox.publish_status", "PUBLISHED"); publishActivity?.SetStatus(ActivityStatusCode.Ok); _logger.LogInformation( "[outbox] Dapr publish OK | OUTBOXID={OutboxId} | Topic={Topic}", outboxId, topic); // ── 2. Publish attachment-resolve (fire-and-forget, non-fatal) ─ if (evt.HEADERROWGUID != null) { var headerRowGuid = (Guid)evt.HEADERROWGUID; using var attachResolveActivity = DaprEventTracing.StartPublishActivity("ATTACHMENT-RESOLVE", pubsubName); try { attachResolveActivity?.SetTag("gb5.outbox.id", outboxId); attachResolveActivity?.SetTag("gb5.object.header_type_id", objectTypeId); attachResolveActivity?.SetTag("gb5.object.id", objectId); attachResolveActivity?.SetTag("gb5.tenant.id", tenantId); attachResolveActivity?.SetTag("gb5.header_row_guid", headerRowGuid.ToString()); await _dapr.PublishEventAsync( pubsubName, "ATTACHMENT-RESOLVE", new AttachmentResolveMessage { TenantId = tenantId, ObjectHeaderTypeId = objectTypeId, HeaderRowGuid = headerRowGuid, ObjectId = objectId, // Carry this producer span's W3C context so the subscriber // (AttachmentResolveSubBLL) can parent its consumer span onto // this same trace instead of starting an orphaned one. TraceParent = attachResolveActivity?.Id, TraceState = attachResolveActivity?.TraceStateString }); attachResolveActivity?.SetStatus(ActivityStatusCode.Ok); _logger.LogInformation( "[outbox] Attachment-resolve publish OK | OUTBOXID={OutboxId} | HeaderRowGuid={HeaderRowGuid} | ObjectId={ObjectId}", outboxId, headerRowGuid, objectId); } catch (Exception attEx) { attachResolveActivity?.RecordException(attEx); attachResolveActivity?.SetStatus(ActivityStatusCode.Error, attEx.Message); // Non-fatal: main event already published. Log and continue. _logger.LogWarning(attEx, "[outbox] WARNING attachment-resolve publish FAILED (non-fatal) | OUTBOXID={OutboxId} | HeaderRowGuid={HeaderRowGuid} | ObjectId={ObjectId} | Error={Error}", outboxId, headerRowGuid, objectId, attEx.Message); } } // ── 3. Mark PUBLISHED (STATUS = 1) ──────────────────────────── await _queryExecutor.ExecuteAsync(login, OutBoxQB.MarkPublished, new { Published = OUTBOXSTATUS.PUBLISHED, OutboxId = outboxId }); _logger.LogInformation( "[outbox] STATUS=1 PUBLISHED | OUTBOXID={OutboxId} | EVENTTYPEID={EventTypeId} | OBJECTID={ObjectId}", outboxId, eventTypeId, objectId); published++; } catch (Exception ex) when (IsDaprSidecarUnavailable(ex)) { // Dapr sidecar is not reachable — transient infrastructure failure. // Row stays PENDING so it will be picked up automatically when Dapr comes back. // Abort the rest of this batch: no point retrying other rows right now. publishActivity?.SetTag("gb5.outbox.publish_status", "SIDECAR_UNAVAILABLE"); publishActivity?.SetStatus(ActivityStatusCode.Error, "Dapr sidecar unavailable"); _logger.LogWarning( "[outbox] DAPR SIDECAR UNAVAILABLE — batch aborted, rows left as PENDING for auto-retry | OUTBOXID={OutboxId} | Hint=Start the app with 'dapr run' so the sidecar is present | Error={Error}", outboxId, BuildFullExceptionChain(ex)); break; } catch (Exception ex) { failed++; publishActivity?.SetTag("gb5.outbox.publish_status", "FAILED"); publishActivity?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "[outbox] ERROR — OUTBOXID={OutboxId} | EVENTTYPEID={EventTypeId} | OBJECTID={ObjectId} | TENANTID={TenantId} | Error={Error}", outboxId, eventTypeId, objectId, tenantId, BuildFullExceptionChain(ex)); // Mark FAILED (STATUS = 2) — wrap so a DB failure here is also visible try { await _queryExecutor.ExecuteAsync(login, OutBoxQB.MarkFailed, new { Failed = OUTBOXSTATUS.FAILED, OutboxId = outboxId }); _logger.LogError( "[outbox] STATUS=2 FAILED set | OUTBOXID={OutboxId}", outboxId); } catch (Exception markEx) { _logger.LogCritical(markEx, "[outbox] CRITICAL could not write STATUS=2 | OUTBOXID={OutboxId} | Error={Error}", outboxId, markEx.Message); } } finally { publishActivity?.Dispose(); } } if (published > 0 || failed > 0) _logger.LogInformation( "[outbox] Batch done | TENANTID={TenantId} | Published={Published} | Failed={Failed}", login.ClientId, published, failed); } // ------------------------------------------------------------------ // HELPERS // ------------------------------------------------------------------ /// /// Returns true when the Dapr sidecar is simply not reachable (SocketException /// anywhere in the chain). This is a transient infrastructure failure — the row /// should stay PENDING and be retried once the sidecar is running, not be /// permanently marked FAILED. /// private static bool IsDaprSidecarUnavailable(Exception ex) { var cur = ex; while (cur != null) { if (cur is System.Net.Sockets.SocketException) return true; cur = cur.InnerException; } return false; } /// /// Walks the full InnerException chain and returns every message on one line. /// Dapr wraps the real gRPC error inside InnerException — without this we only /// see "Publish operation failed: the Dapr endpoint indicated a failure." /// With this we see the actual gRPC status code and message that caused it. /// private static string BuildFullExceptionChain(Exception ex) { var sb = new StringBuilder(); var cur = ex; int depth = 0; while (cur != null) { if (depth > 0) sb.Append(" --> "); sb.Append('[').Append(cur.GetType().Name).Append("] ").Append(cur.Message); cur = cur.InnerException; depth++; } return sb.ToString(); } // ------------------------------------------------------------------ // QUERY BY CORRELATION KEY (debug / tracing) // ------------------------------------------------------------------ public async Task> GetEventsByCorrelationAsync( string correlationKey, LoginDTO login) { if (string.IsNullOrWhiteSpace(correlationKey)) throw new ArgumentException("CorrelationKey required"); try { return await _queryExecutor.QueryAsync( login, OutBoxQB.GetEventsByCorrelationKey, new { CorrelationKey = correlationKey }); } catch (Exception ex) { throw new Exception($"GetEventsByCorrelationAsync failed: {ex.Message}", ex); } } } }