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