using System.ComponentModel.DataAnnotations; using System.Text.Json; using GB5Shared.ActionProcessor; using GB5Shared.DaprCache; using GB5Shared.DTO.Framework.AutoNumber; using GB5Shared.DTO.Framework.Criteria; using GB5Shared.DTO.Framework.Login; using GB5Shared.EntityHandler; using GB5Shared.GenerateAutoNumber; using GB5Shared.Resource.Response; using GB5Shared.Telemetry; using GB5Shared.Validation; using MarketingDAL.CustomCode.Communication; using MarketingDAL.DTO.Communication; using Microsoft.Extensions.Logging; using Newtonsoft.Json; using static GB5Shared.GB5Constant.Constant; namespace MarketingBLL.Communication; public class CommunicationBLL : ICommunicationBLL { private readonly ICommunicationDAL _CommunicationDAL; private readonly AutoNumber _AutoNumber; private readonly IValidation _Validation; private readonly KeyInvalidate _KeyInvalidate; private readonly BaseEntityAppService _BaseEntityAppService; private readonly IEventActionRunDAL _EventActionRunDAL; private readonly IActionOutboxDAL _ActionOutboxDAL; private readonly ILogger _Logger; public CommunicationBLL( ICommunicationDAL communicationDAL, AutoNumber autoNumber, IValidation validation, KeyInvalidate keyInvalidate, BaseEntityAppService baseEntityAppService, IEventActionRunDAL eventActionRunDAL, IActionOutboxDAL actionOutboxDAL, ILogger logger) { _CommunicationDAL = communicationDAL; _AutoNumber = autoNumber; _Validation = validation; _KeyInvalidate = keyInvalidate; _BaseEntityAppService = baseEntityAppService; _EventActionRunDAL = eventActionRunDAL; _ActionOutboxDAL = actionOutboxDAL; _Logger = logger; } public async Task GetCommunication(int communicationId, LoginDTO loginDTO, CancellationToken ct) { try { return await _CommunicationDAL.GetCommunication(communicationId, loginDTO, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("get-communication-failed", ex); _Logger.LogError(ex, "GetCommunication failed for CommunicationId {CommunicationId}", communicationId); throw; } } public async Task SaveCommunication(CommunicationDTO communicationDTO, LoginDTO loginDTO, CancellationToken ct) { if (communicationDTO == null) throw new ArgumentNullException(nameof(communicationDTO)); AutoNumberDTO? autoNumberDTO = null; try { GB5Trace.Step("validate-communication", new { communicationDTO.CommunicationId }); await _Validation.NotEmpty(communicationDTO.CommunicationCode, nameof(communicationDTO.CommunicationCode)); await _Validation.NotEmpty(communicationDTO.CommunicationName, nameof(communicationDTO.CommunicationName)); await _Validation.NotEmpty(communicationDTO.BodyText, nameof(communicationDTO.BodyText)); await _Validation.Range((byte)1, (byte)4, communicationDTO.Channel, nameof(communicationDTO.Channel)); await _Validation.Range((byte)1, (byte)3, communicationDTO.CommunicationStatus, nameof(communicationDTO.CommunicationStatus)); if (communicationDTO.CampaignId == 0) communicationDTO.CampaignId = -1; if (communicationDTO.SegmentId == 0) communicationDTO.SegmentId = -1; if (communicationDTO.PersonaId == 0) communicationDTO.PersonaId = -1; communicationDTO.ModifiedById = loginDTO.UserId; communicationDTO.ModifiedOn = DateTime.UtcNow; bool isNew = communicationDTO.CommunicationId == 0; if (isNew) { autoNumberDTO = await _AutoNumber.GetNumberAsync(1, "COMMUNICATION", loginDTO).ConfigureAwait(false); communicationDTO.CommunicationId = autoNumberDTO.StartNumber; communicationDTO.CreatedById = loginDTO.UserId; communicationDTO.CreatedOn = DateTime.UtcNow; } GB5Trace.Step("save-communication", new { communicationDTO.CommunicationId, isNew }); await _BaseEntityAppService.ExecuteSaveAsync( EntityConstant.OBJECTMARKETINGCOMMUNICATION, isNew ? EventTypeConstant.SAVEMARKETINGCOMMUNICATIONEVENTTYPEID : EventTypeConstant.UPDATEMARKETINGCOMMUNICATIONEVENTTYPEID, communicationDTO, loginDTO, async tx => { _ = isNew ? await _CommunicationDAL.SaveCommunication(communicationDTO, loginDTO, tx, ct).ConfigureAwait(false) : await _CommunicationDAL.UpdateCommunication(communicationDTO, loginDTO, tx, ct).ConfigureAwait(false); return communicationDTO.CommunicationId; }, null, -1, -1, null, isNew).ConfigureAwait(false); GB5Trace.Step("event-publish", new { EventTypeId = isNew ? EventTypeConstant.SAVEMARKETINGCOMMUNICATIONEVENTTYPEID : EventTypeConstant.UPDATEMARKETINGCOMMUNICATIONEVENTTYPEID }); var cacheKey = new CacheKeyGeneration().KeyGeneration( communicationDTO.CommunicationId, EntityConstant.OBJECTMARKETINGCOMMUNICATION, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.AllInvalidateCache(cacheKey).ConfigureAwait(false); _Logger.LogInformation("Communication {CommunicationId} {Action} by user {UserId}", communicationDTO.CommunicationId, isNew ? "saved" : "updated", loginDTO.UserId); return isNew ? $"{SuccessResponse.SaveSuccessMessage} {communicationDTO.CommunicationId}" : $"{SuccessResponse.UpdateSuccessMessage} {communicationDTO.CommunicationId}"; } catch (ValidationException vex) { GB5Trace.MarkFailed("save-communication-failed", vex); if (autoNumberDTO != null && communicationDTO.CommunicationId > 0) await _AutoNumber.RollbackAutoNumber("COMMUNICATION", communicationDTO.CommunicationId, loginDTO).ConfigureAwait(false); throw new Exception(vex.Message); } catch (Exception ex) { GB5Trace.MarkFailed("save-communication-failed", ex); _Logger.LogError(ex, "SaveCommunication failed for CommunicationId {CommunicationId}", communicationDTO.CommunicationId); if (autoNumberDTO != null && communicationDTO.CommunicationId > 0) await _AutoNumber.RollbackAutoNumber("COMMUNICATION", communicationDTO.CommunicationId, loginDTO).ConfigureAwait(false); throw; } } public async Task GetSelectListCommunication(CriteriaDTO criteriaDTO, LoginDTO loginDTO, CancellationToken ct) { try { return await _CommunicationDAL.GetSelectListCommunication(criteriaDTO, loginDTO, ct).ConfigureAwait(false); } catch (Exception ex) { GB5Trace.MarkFailed("get-selectlist-communication-failed", ex); _Logger.LogError(ex, "GetSelectListCommunication failed"); throw; } } public async Task DeleteCommunication(int communicationId, LoginDTO loginDTO, CancellationToken ct) { try { GB5Trace.Step("delete-communication", new { communicationId }); var result = await _CommunicationDAL.DeleteCommunication(communicationId, loginDTO, ct).ConfigureAwait(false); var cacheKey = new CacheKeyGeneration().KeyGeneration( communicationId, EntityConstant.OBJECTMARKETINGCOMMUNICATION, CacheKeyLevel.CLIENT_LEVEL, loginDTO); await _KeyInvalidate.AllInvalidateCache(cacheKey).ConfigureAwait(false); _Logger.LogInformation("Communication {CommunicationId} deleted by user {UserId}", communicationId, loginDTO.UserId); return result; } catch (Exception ex) { GB5Trace.MarkFailed("delete-communication-failed", ex); _Logger.LogError(ex, "DeleteCommunication failed for CommunicationId {CommunicationId}", communicationId); throw; } } /// 1=Email, 2=SMS, 3=WhatsApp, 4=Social — mirrors CommunicationDTO.Channel. private static class ChannelValue { public const byte Email = 1; public const byte Sms = 2; public const byte WhatsApp = 3; public const byte Social = 4; } /// ActionType values IActionHandler implementations are registered for /// (FrameworkBLL/ActionProcessor/Handlers) — 0=Email, 1=SMS. No WhatsApp/Social handler /// exists anywhere in this codebase yet. private static class ActionTypeValue { public const int Email = 0; public const int Sms = 1; } /// /// Dispatches a Draft/Scheduled Communication to every Lead in its target Segment, via the /// same real, working Action/Outbox mechanism ECP's CorrespondenceBLL.QueueDeliveryAsync /// uses (GB5Shared.ActionProcessor — a thin, GB5Shared-level contract, so this carries none /// of the cross-module ProjectReference/assembly-load risk that ruled out /// FrameworkBLL.MessageHubGenerator in this method's previous investigation). /// /// One TEVENTACTIONRUN + TACTIONOUTBOX row per resolved recipient, not one row for the whole /// audience — required for SMS (Twilio's MessageResource.CreateAsync takes exactly one `to` /// per call; SmsActionHandler.HandleAsync would treat a combined list as one malformed /// number) and kept for Email too for consistency and per-recipient failure tracking. /// /// /// ActionId=-1 (not a real MACTION row) — matches ECP's own precedent exactly /// (CorrespondenceBLL.QueueDeliveryAsync's comment): this is an explicit, user-triggered send, /// not a rule EventSubBLL evaluates against a configured MACTION row, so -1 follows the /// "-1 = not applicable" convention used throughout GB5. /// /// /// WhatsApp/Social remain genuinely unimplemented — no IActionHandler exists for those /// ActionTypes anywhere in this codebase (only Email/SMS/Webhook/Correspondence do) — this /// throws NotImplementedException for those two channels rather than silently no-op'ing, /// matching GB5Shared/RecordingIntelligence/NotConfiguredAICapabilityService's pattern. /// /// public async Task SendCommunicationAsync(int communicationId, LoginDTO loginDTO, CancellationToken ct) { GB5Trace.Step("validate-send-communication", new { communicationId }); var json = await _CommunicationDAL.GetCommunication(communicationId, loginDTO, ct).ConfigureAwait(false); var communicationDTO = JsonConvert.DeserializeObject(json) ?? throw new InvalidOperationException($"Communication {communicationId} not found."); if (communicationDTO.CommunicationStatus == 3) throw new InvalidOperationException($"Communication {communicationId} has already been sent."); int actionType = communicationDTO.Channel switch { ChannelValue.Email => ActionTypeValue.Email, ChannelValue.Sms => ActionTypeValue.Sms, ChannelValue.WhatsApp => throw new NotImplementedException( "WhatsApp dispatch is not implemented — no IActionHandler is registered for a " + "WhatsApp ActionType anywhere in this codebase (only Email/SMS/Webhook/" + "Correspondence exist, per FrameworkBLL/ActionProcessor/Handlers)."), ChannelValue.Social => throw new NotImplementedException( "Social dispatch is not implemented — no IActionHandler is registered for a " + "Social ActionType anywhere in this codebase."), _ => throw new InvalidOperationException($"Unknown Channel {communicationDTO.Channel}.") }; if (communicationDTO.SegmentId <= 0) throw new InvalidOperationException( "Communication must target a real Segment (SegmentId) to resolve a send audience."); try { GB5Trace.Step("resolve-communication-audience", new { communicationDTO.SegmentId }); var contacts = await _CommunicationDAL .GetCommunicationAudienceContacts(communicationDTO.SegmentId, loginDTO, ct) .ConfigureAwait(false); var recipients = actionType == ActionTypeValue.Email ? contacts.Select(c => c.Mail).Where(m => !string.IsNullOrWhiteSpace(m)).Distinct().ToList() : contacts.Select(c => c.Mobile).Where(m => !string.IsNullOrWhiteSpace(m)).Distinct().ToList(); if (recipients.Count == 0) { _Logger.LogWarning( "SendCommunicationAsync: no recipients with a usable {Channel} contact found for Communication {CommunicationId} Segment {SegmentId}", actionType == ActionTypeValue.Email ? "email" : "mobile", communicationId, communicationDTO.SegmentId); return $"No recipients found for Segment {communicationDTO.SegmentId} — nothing sent."; } GB5Trace.Step("enqueue-communication-dispatch", new { communicationId, RecipientCount = recipients.Count }); foreach (var recipient in recipients) { var correlationKey = $"communication-{communicationId}-{Guid.NewGuid():N}"; var actionRunId = await _EventActionRunDAL.InsertAsync(new EventActionRunDTO { ActionId = -1, JobExecutionId = -1, EventTypeId = EventTypeConstant.SAVEMARKETINGCOMMUNICATIONEVENTTYPEID, Payload = json, CorrelationKey = correlationKey, }, loginDTO, ct).ConfigureAwait(false); var actionEventDto = new ActionEventDto { ActionRunId = actionRunId, ActionId = -1, ActionType = actionType, TenantId = loginDTO.ClientId, DatabaseName = loginDTO.DatabaseName, SendTo = recipient, ToDeliveryType = 3, // direct value, already resolved above TemplateId = -1, // Subject/Body travel in Payload — no MMAILTEMPLATE lookup CorrelationKey = correlationKey, Payload = actionType == ActionTypeValue.Email ? System.Text.Json.JsonSerializer.SerializeToElement(new { Subject = communicationDTO.Subject ?? string.Empty, Body = communicationDTO.BodyText, IsHtml = false, }) : System.Text.Json.JsonSerializer.SerializeToElement(communicationDTO.BodyText), }; var partition = Math.Abs(loginDTO.ClientId % 5); await _ActionOutboxDAL.InsertAsync(new ActionOutboxDTO { ActionRunId = actionRunId, DestinationTopic = $"action-exec-p{partition}", Payload = System.Text.Json.JsonSerializer.Serialize(actionEventDto), CorrelationKey = correlationKey, }, loginDTO, ct).ConfigureAwait(false); } communicationDTO.CommunicationStatus = 3; // Sent communicationDTO.SentOn = DateTime.UtcNow; await SaveCommunication(communicationDTO, loginDTO, ct).ConfigureAwait(false); _Logger.LogInformation( "Communication {CommunicationId} dispatched to {RecipientCount} recipients by user {UserId}", communicationId, recipients.Count, loginDTO.UserId); return $"Communication queued for delivery to {recipients.Count} recipient(s)."; } catch (Exception ex) { GB5Trace.MarkFailed("send-communication-failed", ex); _Logger.LogError(ex, "SendCommunicationAsync failed for CommunicationId {CommunicationId}", communicationId); throw; } } }