using GB5Shared.ActionProcessor; using FrameworkDAL.CustomCode.ActionProcessor; using FrameworkDAL.DTO.EIPConversation; using GB5Shared.DTO.Framework.Login; using Microsoft.Extensions.Logging; using System; using System.Diagnostics; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.ActionHandler { /// /// Routes an in-app notification request through the ActionProcessor outbox /// (ActionType=3). The background worker picks it up and delivers it via /// the configured notification provider. /// public class NotificationActionHandler : IEIPActionHandler { public string ActionCode => "SEND_NOTIFICATION"; private const int Partitions = 5; private readonly IEventActionRunDAL _runDal; private readonly IActionOutboxDAL _outboxDal; private readonly ILogger _logger; public NotificationActionHandler( IEventActionRunDAL runDal, IActionOutboxDAL outboxDal, ILogger logger) { _runDal = runDal ?? throw new ArgumentNullException(nameof(runDal)); _outboxDal = outboxDal ?? throw new ArgumentNullException(nameof(outboxDal)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } public bool CanHandle(string actionCode) => string.Equals(actionCode, ActionCode, StringComparison.OrdinalIgnoreCase); public async Task HandleAsync( EIPActionContextDTO context, LoginDTO loginDTO, CancellationToken ct) { if (context.TenantId is null) { _logger.LogError("NotificationActionHandler: TenantId is missing."); return new EIPActionResultDTO { IsSuccess = false, Message = "Invalid TenantId." }; } var tenantId = context.TenantId.Value; try { // Insert run first — IDENTITY generates ActionRunId int actionRunId; try { actionRunId = await _runDal.InsertAsync(new EventActionRunDTO { ActionId = -1, JobExecutionId = -1, EventTypeId = -1, TenantId = tenantId, Payload = JsonSerializer.Serialize(context.Payload) }, loginDTO, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.run.id", actionRunId); Activity.Current?.SetTag("gb5.tenant.id", tenantId); Activity.Current?.SetTag("gb5.conversation.id", context.ConversationId); } catch (Exception insertEx) { Activity.Current?.SetStatus(ActivityStatusCode.Error, insertEx.Message); _logger.LogError(insertEx, "NotificationActionHandler: TEVENTACTIONRUN INSERT FAILED | TenantId={TenantId} ConversationId={ConvId} UserIdentifier={User}", tenantId, context.ConversationId, context.UserIdentifier); return new EIPActionResultDTO { IsSuccess = false, Message = "Failed to create action run.", ErrorMessage = insertEx.ToString() }; } var correlationKey = $"EIP-NOTIFY-{actionRunId}-{Guid.NewGuid():N}"; var eventDto = new ActionEventDto { ActionRunId = actionRunId, ActionId = -1, ActionType = 5, // Notification (EIP) — Correspondence uses 3 TenantId = tenantId, DatabaseName = loginDTO.DatabaseName, CorrelationKey = correlationKey, SendTo = context.Payload.TryGetValue("RecipientId", out var r) ? r?.ToString() : context.UserIdentifier, Payload = JsonSerializer.SerializeToElement(context.Payload) }; var partition = tenantId % Partitions; var destinationTopic = $"action-exec-p{partition}"; try { await _outboxDal.InsertAsync(new ActionOutboxDTO { ActionRunId = actionRunId, DestinationTopic = destinationTopic, Payload = JsonSerializer.Serialize(eventDto), CorrelationKey = correlationKey, TenantId = tenantId }, loginDTO, ct).ConfigureAwait(false); Activity.Current?.SetTag("gb5.outbox.topic", destinationTopic); } catch (Exception outboxEx) { Activity.Current?.SetStatus(ActivityStatusCode.Error, outboxEx.Message); _logger.LogError(outboxEx, "NotificationActionHandler: TACTIONOUTBOX INSERT FAILED | ActionRunId={RunId} TenantId={TenantId} ConversationId={ConvId} Topic={Topic} CorrelationKey={CorrelationKey}", actionRunId, tenantId, context.ConversationId, destinationTopic, correlationKey); return new EIPActionResultDTO { IsSuccess = false, Message = "Failed to queue notification.", ErrorMessage = outboxEx.ToString() }; } _logger.LogInformation( "NotificationActionHandler: queued | ActionRunId={RunId} TenantId={TenantId} Partition={P} ConversationId={ConvId} CorrelationKey={CorrelationKey}", actionRunId, tenantId, partition, context.ConversationId, correlationKey); return new EIPActionResultDTO { IsSuccess = true, Message = "Notification queued." }; } catch (OperationCanceledException) { Activity.Current?.SetStatus(ActivityStatusCode.Error, "Cancelled"); _logger.LogWarning("NotificationActionHandler: cancelled | ConversationId={Id} TenantId={TenantId}", context.ConversationId, context.TenantId); return new EIPActionResultDTO { IsSuccess = false, Message = "Operation cancelled." }; } catch (Exception ex) { Activity.Current?.SetStatus(ActivityStatusCode.Error, ex.Message); _logger.LogError(ex, "NotificationActionHandler: unexpected failure | ConversationId={Id} TenantId={TenantId}", context.ConversationId, context.TenantId); return new EIPActionResultDTO { IsSuccess = false, Message = ex.Message, ErrorMessage = ex.ToString() }; } } } }