using FrameworkBLL.EIPConversation.EIPHandlers.ChannelHandler; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.ActionEngine; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.CapabilityEngine; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.ChannelNormalizer; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.FlowEngine; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.ResponseEngine; using FrameworkBLL.EIPConversation.EIPHandlers.EIPEngine.RoutingEngine; using FrameworkDAL.CustomCode.EIPConversation.EIPConversationLog; using FrameworkDAL.CustomCode.EIPConversation.EIPInteractionSession; using FrameworkDAL.CustomCode.EIPConversation.EIPSession; using FrameworkDAL.DTO.EIPConversation; using GB5Shared.DTO.Framework.Login; using Microsoft.Extensions.Logging; using System; using System.Collections.Generic; using System.Diagnostics; using System.Text.Json; using System.Threading; using System.Threading.Tasks; namespace FrameworkBLL.EIPConversation.EIPHandlers.EIPConversationEngine { public class EIPConversationEngine : IEIPConversationEngine { private readonly IEIPChannelNormalizer _channelNormalizer; private readonly IEIPRoutingEngine _routingEngine; private readonly IEIPFlowEngine _flowEngine; private readonly IEIPCapabilityEngine _capabilityEngine; private readonly IEIPActionEngine _actionEngine; private readonly IEIPResponseEngine _responseEngine; private readonly IEIPChannelHandler _channelHandler; private readonly IEIPSessionDAL _sessionDAL; private readonly IEIPInteractionSessionDAL _interactionSessionDAL; private readonly IEIPConversationLogDAL _conversationLogDAL; private readonly ILogger _logger; public EIPConversationEngine( IEIPChannelNormalizer channelNormalizer, IEIPRoutingEngine routingEngine, IEIPFlowEngine flowEngine, IEIPCapabilityEngine capabilityEngine, IEIPActionEngine actionEngine, IEIPResponseEngine responseEngine, IEIPChannelHandler channelHandler, IEIPSessionDAL sessionDAL, IEIPInteractionSessionDAL interactionSessionDAL, IEIPConversationLogDAL conversationLogDAL, ILogger logger) { _channelNormalizer = channelNormalizer ?? throw new ArgumentNullException(nameof(channelNormalizer)); _routingEngine = routingEngine ?? throw new ArgumentNullException(nameof(routingEngine)); _flowEngine = flowEngine ?? throw new ArgumentNullException(nameof(flowEngine)); _capabilityEngine = capabilityEngine ?? throw new ArgumentNullException(nameof(capabilityEngine)); _actionEngine = actionEngine ?? throw new ArgumentNullException(nameof(actionEngine)); _responseEngine = responseEngine ?? throw new ArgumentNullException(nameof(responseEngine)); _channelHandler = channelHandler ?? throw new ArgumentNullException(nameof(channelHandler)); _sessionDAL = sessionDAL ?? throw new ArgumentNullException(nameof(sessionDAL)); _interactionSessionDAL = interactionSessionDAL ?? throw new ArgumentNullException(nameof(interactionSessionDAL)); _conversationLogDAL = conversationLogDAL ?? throw new ArgumentNullException(nameof(conversationLogDAL)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); } public async Task ProcessIncomingMessageAsync( EIPConversationDTO payload, LoginDTO loginDTO, CancellationToken cancellationToken) { var stopwatch = Stopwatch.StartNew(); if (payload == null) throw new ArgumentNullException(nameof(payload)); if (loginDTO == null) throw new ArgumentNullException(nameof(loginDTO)); _logger.LogInformation("ENGINE PROCESS START"); try { // ---------------- TENANT VALIDATION ---------------- EnsureTenant(loginDTO, payload); if (loginDTO.UserId == 0) throw new UnauthorizedAccessException("UserId invalid."); payload.TenantId = loginDTO.ClientId; _logger.LogInformation( "ENGINE START | Tenant={Tenant} | User={User} | Channel={Channel} | ExternalUser={ExternalUser}", loginDTO.ClientId, loginDTO.UserId, payload.ChannelType, payload.UserIdentifier); // ---------------- PHASE 1: NORMALIZATION ---------------- var context = await _channelNormalizer.NormalizeAsync(payload, cancellationToken) ?? throw new InvalidOperationException("ChannelNormalizer returned null."); context.TenantId = loginDTO.ClientId; // EIPOtpService (and anything else downstream reading context.LoginDTO // rather than taking a LoginDTO parameter directly) needs the real, // fully-populated LoginDTO here -- DBConnectionStringCached resolves the // physical DB connection from LoginDTO.DatabaseName, so a null/default // context.LoginDTO throws a NullReferenceException the moment any // OTP-gated Capability step runs. This was never assigned anywhere in // the codebase before this fix. context.LoginDTO = loginDTO; if (context.ChannelType != EIPChannelType.WhatsApp) context.UserIdentifier ??= loginDTO.UserId.ToString(); context.PhaseTrace ??= new List(); context.Metadata ??= new Dictionary(); // Extract inbound location into context variables for LOCATION step access if (payload.Location != null) { context.Variables["_Location.Lat"] = payload.Location.Latitude.ToString(System.Globalization.CultureInfo.InvariantCulture); context.Variables["_Location.Lng"] = payload.Location.Longitude.ToString(System.Globalization.CultureInfo.InvariantCulture); context.Variables["_Location.Address"] = payload.Location.Address ?? string.Empty; } AddPhase(context, "PHASE1: NORMALIZATION", "SUCCESS", "Normalization completed"); _logger.LogInformation( "PHASE1 COMPLETE | Tenant={Tenant} | User={User} | ConversationId={ConversationId}", context.TenantId, context.UserIdentifier, context.ConversationId); // ---------------- SESSION RESTORE ---------------- // Load the active session to resume multi-turn conversations. var session = await _sessionDAL.GetActiveSessionAsync( context.UserIdentifier, loginDTO.ClientId, (int)context.ChannelType, loginDTO, cancellationToken).ConfigureAwait(false); if (session != null) { context.CurrentStepCode = session.CurrentStepCode; context.ConversationId = GuidToInt(session.SessionId); // ConversationId guard: if the client sends a non-zero ID that doesn't match // the active session, reject immediately. ConversationId == 0 means the // client hasn't received an ID yet (first message or fresh start) — allow it. if (payload.ConversationId != 0 && payload.ConversationId != context.ConversationId) { _logger.LogWarning( "SESSION MISMATCH | User={User} | Expected={Expected} | Received={Received}", context.UserIdentifier, context.ConversationId, payload.ConversationId); return new EIPExecutionResultDTO { TenantId = loginDTO.ClientId, ChannelType = context.ChannelType, UserIdentifier = context.UserIdentifier, FlowCode = string.Empty, ConversationId = context.ConversationId, ResponseMessage = "Session mismatch. The conversation ID you sent does not match your active session. " + "Please use the ConversationId returned by the previous response, or send 0 to start a new conversation.", IsCompleted = true, IsSuccess = false, PhaseTrace = context.PhaseTrace, Metadata = context.Metadata }; } if (!string.IsNullOrWhiteSpace(session.FlowCode)) context.FlowCode = session.FlowCode; if (!string.IsNullOrWhiteSpace(session.ContextJson)) { var restored = JsonSerializer.Deserialize>(session.ContextJson); if (restored != null) foreach (var kv in restored) context.Variables[kv.Key] = kv.Value; } _logger.LogInformation( "SESSION RESTORED | User={User} | Step={Step} | Flow={Flow} | ConversationId={Id}", context.UserIdentifier, session.CurrentStepCode, session.FlowCode, context.ConversationId); } // ---------------- PHASE 2: ROUTING ---------------- // Skip routing when an active session already has a FlowCode. // Mid-conversation messages (e.g. "15-05-2026") are inputs for the // current flow step — they must NOT be matched against routing rules. if (session == null || string.IsNullOrWhiteSpace(context.FlowCode)) { var routing = await _routingEngine.ResolveRoutingAsync(context, loginDTO, cancellationToken); context.FlowCode = routing?.IsMatched == true ? routing.TargetFlowCode : context.FlowCode; _logger.LogInformation("ROUTING COMPLETE | Flow={Flow}", context.FlowCode); } else { _logger.LogInformation("ROUTING SKIPPED | Session active | Flow={Flow}", context.FlowCode); } AddPhase(context, "PHASE2: ROUTING", "SUCCESS", $"FlowCode resolved: {context.FlowCode ?? "NONE"}"); // No flow resolved — DB has no catch-all routing rule configured. // Return a neutral response without any hardcoded flow-specific text. if (string.IsNullOrWhiteSpace(context.FlowCode)) { _logger.LogWarning( "ENGINE | No flow resolved and no catch-all rule in DB | Tenant={Tenant} | User={User} | Message='{Message}'", loginDTO.ClientId, context.UserIdentifier, context.NormalizedMessage); return new EIPExecutionResultDTO { TenantId = loginDTO.ClientId, ChannelType = context.ChannelType, UserIdentifier = context.UserIdentifier, FlowCode = string.Empty, ConversationId = 0, ResponseMessage = "I'm unable to process your request at this time. Please try again.", IsCompleted = true, IsSuccess = true, PhaseTrace = context.PhaseTrace, Metadata = context.Metadata }; } // ---------------- OBSERVABILITY: INTERACTION SESSION + CONVERSATION LOG ---------------- // Never let an audit-logging failure break the actual conversation — log and // continue with InteractionSessionId/ConversationLogId left at 0, which the // flow engine's step logger treats as "logging unavailable for this turn." try { var channelTypeByte = (byte)context.ChannelType; context.InteractionSessionId = await _interactionSessionDAL.GetOrCreateAsync( context.UserIdentifier, loginDTO.ClientId, channelTypeByte, context.FlowCode!, context.CurrentStepCode, loginDTO, cancellationToken) .ConfigureAwait(false); context.ConversationLogId = await _conversationLogDAL.GetOrCreateConversationAsync( context.InteractionSessionId, loginDTO.ClientId, channelTypeByte, $"ISESSION-{context.InteractionSessionId}", loginDTO, cancellationToken) .ConfigureAwait(false); await _conversationLogDAL.InsertDetailAsync( context.ConversationLogId, 1 /*Inbound*/, context.OriginalMessage ?? context.Message ?? string.Empty, loginDTO, cancellationToken).ConfigureAwait(false); } catch (Exception logEx) { _logger.LogError(logEx, "Observability logging failed (interaction session/conversation) — continuing without it | User={User}", context.UserIdentifier); } // ---------------- PHASE 3: FLOW ---------------- var flowResult = await _flowEngine.ExecuteFlowAsync( context, context.FlowCode, loginDTO, cancellationToken) ?? throw new InvalidOperationException("Flow execution returned null."); context.CurrentStepCode = flowResult.CurrentStepCode ?? context.CurrentStepCode; AddPhase(context, "PHASE3: FLOW", "SUCCESS", $"CurrentStep={context.CurrentStepCode}"); _logger.LogInformation( "FLOW COMPLETE | CurrentStep={Step} | NextStep={Next} | Completed={Completed}", context.CurrentStepCode, flowResult.NextStepCode, flowResult.IsCompleted); // ---------------- SESSION SAVE / EXPIRE ---------------- if (flowResult.IsCompleted) { await _sessionDAL.ExpireSessionAsync( context.UserIdentifier, loginDTO.ClientId, (int)context.ChannelType, loginDTO, cancellationToken).ConfigureAwait(false); _logger.LogInformation("SESSION EXPIRED | User={User} | Flow={Flow}", context.UserIdentifier, context.FlowCode); } else { var contextJson = JsonSerializer.Serialize(context.Variables); var sessionGuid = await _sessionDAL.UpsertSessionAsync(new EIPUserSessionDTO { UserIdentifier = context.UserIdentifier, TenantId = loginDTO.ClientId, ChannelType = (int)context.ChannelType, FlowCode = context.FlowCode ?? string.Empty, CurrentStepCode = flowResult.NextStepCode ?? context.CurrentStepCode ?? string.Empty, ContextJson = contextJson }, loginDTO, cancellationToken).ConfigureAwait(false); context.ConversationId = GuidToInt(sessionGuid); _logger.LogInformation( "SESSION SAVED | User={User} | NextStep={Step} | ConversationId={Id}", context.UserIdentifier, flowResult.NextStepCode, context.ConversationId); } if (context.InteractionSessionId != 0) { try { await _interactionSessionDAL.MarkStatusAsync( context.InteractionSessionId, sessionStatus: flowResult.IsCompleted ? (byte)5 /*Completed*/ : (byte)2 /*AwaitingInput*/, currentStep: flowResult.IsCompleted ? flowResult.CurrentStepCode : flowResult.NextStepCode, completed: flowResult.IsCompleted, loginDTO, cancellationToken).ConfigureAwait(false); } catch (Exception logEx) { _logger.LogError(logEx, "Failed to update interaction session status | InteractionSessionId={Id}", context.InteractionSessionId); } } // ---------------- PHASE 4: CAPABILITY ---------------- if (!string.IsNullOrWhiteSpace(flowResult.CapabilityCode)) { await _capabilityEngine.ExecuteCapabilityAsync( flowResult.CapabilityCode, context, loginDTO, cancellationToken); AddPhase(context, "PHASE4: CAPABILITY", "SUCCESS", $"Capability executed: {flowResult.CapabilityCode}"); } // ---------------- PHASE 5: ACTION ---------------- if (!string.IsNullOrWhiteSpace(flowResult.ActionCode)) { await _actionEngine.ExecuteAsync( new EIPActionContextDTO { ActionCode = flowResult.ActionCode, ConversationId = context.ConversationId, UserIdentifier = context.UserIdentifier, TenantId = context.TenantId }, loginDTO, cancellationToken); AddPhase(context, "PHASE5: ACTION", "SUCCESS", $"Action executed: {flowResult.ActionCode}"); } // ---------------- PHASE 6: RESPONSE ---------------- if (!string.IsNullOrWhiteSpace(flowResult.ResponseMessage) || flowResult.StructuredData != null || flowResult.QrCodeData != null || flowResult.RequestLocation || flowResult.Buttons?.Count > 0) { var responseContext = new EIPResponseContext { Recipient = context.UserIdentifier, Channel = GetChannelString(context.ChannelType), Message = flowResult.ResponseMessage ?? string.Empty, MessageKey = flowResult.MessageKey, ChannelType = context.ChannelType, FlowCode = context.FlowCode, PhaseTrace = context.PhaseTrace, Metadata = context.Metadata, TenantId = loginDTO.ClientId, StructuredData = flowResult.StructuredData, QrCodeData = flowResult.QrCodeData, RequestLocation = flowResult.RequestLocation, InputHint = flowResult.InputHint, RatingMax = flowResult.RatingMax, VoiceHint = flowResult.VoiceHint, Buttons = flowResult.Buttons?.Count > 0 ? flowResult.Buttons : null }; AddPhase(responseContext, "PHASE6: RESPONSE", "START", $"Channel={responseContext.Channel} | Recipient={responseContext.Recipient}"); await _responseEngine.GenerateAsync( responseContext, loginDTO, cancellationToken); // EIPResponseEngine may have overwritten responseContext.Message with a // resolved MEIPMESSAGETEMPLATE row (MessageKey lookup) — reflect that back // into flowResult so the API response and the outbound audit log record // what was actually sent, not the pre-lookup flow-JSON text. flowResult.ResponseMessage = responseContext.Message; AddPhase(responseContext, "PHASE6: RESPONSE", "SUCCESS", "Response sent successfully"); if (context.ConversationLogId != 0 && !string.IsNullOrWhiteSpace(flowResult.ResponseMessage)) { try { await _conversationLogDAL.InsertDetailAsync( context.ConversationLogId, 2 /*Outbound*/, flowResult.ResponseMessage, loginDTO, cancellationToken).ConfigureAwait(false); } catch (Exception logEx) { _logger.LogError(logEx, "Failed to log outbound conversation detail | ConversationLogId={Id}", context.ConversationLogId); } } } stopwatch.Stop(); _logger.LogInformation( "ENGINE COMPLETE | Tenant={Tenant} | DurationMs={Duration} | ConversationId={ConversationId}", loginDTO.ClientId, stopwatch.ElapsedMilliseconds, context.ConversationId); return new EIPExecutionResultDTO { TenantId = loginDTO.ClientId, ChannelType = context.ChannelType, UserIdentifier = context.UserIdentifier, FlowCode = context.FlowCode, ConversationId = context.ConversationId, CurrentStepCode = flowResult.CurrentStepCode, NextStepCode = flowResult.NextStepCode, ResponseMessage = flowResult.ResponseMessage, IsCompleted = flowResult.IsCompleted, IsSuccess = true, StructuredData = flowResult.StructuredData, QrCodeData = flowResult.QrCodeData, InputHint = flowResult.InputHint, RatingMax = flowResult.RatingMax, RequestLocation = flowResult.RequestLocation, VoiceHint = flowResult.VoiceHint, Buttons = flowResult.Buttons?.Count > 0 ? flowResult.Buttons : null, PhaseTrace = context.PhaseTrace, Metadata = context.Metadata }; } catch (Exception ex) { stopwatch.Stop(); _logger.LogError(ex, "ENGINE ERROR | Tenant={Tenant} | DurationMs={Duration}", loginDTO.ClientId, stopwatch.ElapsedMilliseconds); throw; } } // ---------------- HELPER: MAP ENUM TO STRING ---------------- private string GetChannelString(EIPChannelType channelType) => channelType switch { EIPChannelType.PostMan => "POSTMAN", EIPChannelType.WhatsApp => "WHATSAPP", EIPChannelType.Teams => "TEAMS", EIPChannelType.Telegram => "TELEGRAM", EIPChannelType.Slack => "SLACK", EIPChannelType.SMS => "SMS", EIPChannelType.AppChat => "APPCHAT", EIPChannelType.Voice => "VOICE", _ => "UNKNOWN" }; // ---------------- HELPER: ADD PHASE ---------------- private void AddPhase(dynamic context, string phaseName, string status, string description) { if (context.PhaseTrace == null) context.PhaseTrace = new List(); context.PhaseTrace.Add(new EIPExecutionPhaseDTO { PhaseName = phaseName, Status = status, Description = description, Timestamp = DateTime.UtcNow }); } // ---------------- SESSION GUID → INT ---------------- // Derives a stable positive int from the first 4 bytes of the session GUID. // Used only for display/tracking in the response — session lookup uses UserIdentifier. private static int GuidToInt(Guid g) { var raw = BitConverter.ToInt32(g.ToByteArray(), 0); return raw == int.MinValue ? int.MaxValue : Math.Abs(raw); } // ---------------- TENANT HELPER ---------------- private void EnsureTenant(LoginDTO loginDTO, EIPConversationDTO payload) { if (loginDTO.ClientId != 0) return; if (payload.TenantId == 0) throw new UnauthorizedAccessException("TenantId missing."); loginDTO.ClientId = payload.TenantId; _logger.LogWarning("ClientId defaulted from payload: {ClientId}", payload.TenantId); } // ---------------- MANUAL EXECUTION ---------------- public async Task ProcessManualRequestAsync( int conversationId, LoginDTO loginDTO, CancellationToken cancellationToken) { if (conversationId <= 0) throw new ArgumentException("Invalid conversationId", nameof(conversationId)); if (loginDTO == null) throw new ArgumentNullException(nameof(loginDTO)); _logger.LogInformation("MANUAL EXECUTION START | ConversationId={Id}", conversationId); return await Task.FromResult(new EIPExecutionResultDTO { ConversationId = conversationId, CurrentStepCode = "MANUAL_START", ResponseMessage = "Manual execution triggered successfully.", IsCompleted = false, IsSuccess = true, PhaseTrace = new List() }); } } }