using FrameworkBLL.EIPConcept; using FrameworkBLL.MessageHub.MessageHubPlatforms; using FrameworkDAL.CustomCode.EIPConversation; using FrameworkDAL.DTO.EIPConversation; using FrameworkDAL.DTO.MessageHub.MessageHubGenerator; using GB5Shared.DTO.Framework.Login; using Microsoft.Extensions.Caching.Memory; using Microsoft.Extensions.Logging; using System.Diagnostics; namespace FrameworkBLL.EIPConversation { /// /// MessageHubConversationBLL /// Orchestrator Layer /// - Session handling /// - Flow execution delegation /// - Platform dispatch delegation /// - Conversation persistence /// public class EIPConceptBLL : IEIPConceptBLL { private readonly IConceptDAL _dal; private readonly MessageHubConversationDTO _config; private readonly ILogger _logger; private readonly IEIPConceptEngine _EIPConceptEngine; private readonly ConversationPlatformType _conversationPlatform; private readonly IMemoryCache _sessionCache; // Sliding TTL — evicts idle sessions after 30 minutes private static readonly MemoryCacheEntryOptions _sessionCacheOptions = new MemoryCacheEntryOptions().SetSlidingExpiration(TimeSpan.FromMinutes(30)); #region Constructor public EIPConceptBLL( IConceptDAL dal, MessageHubConversationDTO config, IEIPConceptEngine EIPConceptEngine, ConversationPlatformType conversationPlatform, IMemoryCache sessionCache, ILogger logger) { _dal = dal ?? throw new ArgumentNullException(nameof(dal)); _config = config ?? throw new ArgumentNullException(nameof(config)); _EIPConceptEngine = EIPConceptEngine ?? throw new ArgumentNullException(nameof(EIPConceptEngine)); _conversationPlatform = conversationPlatform ?? throw new ArgumentNullException(nameof(conversationPlatform)); _sessionCache = sessionCache ?? throw new ArgumentNullException(nameof(sessionCache)); _logger = logger ?? throw new ArgumentNullException(nameof(logger)); _logger.LogInformation("MessageHubConversationBLL initialized successfully."); } #endregion #region Public Entry Point public async Task HandleIncomingMessageAsync( string sessionToken, string userMessage, LoginDTO loginContext, string platformName = "POSTMAN", CancellationToken ct = default) { var stopwatch = Stopwatch.StartNew(); sessionToken ??= string.Empty; userMessage ??= string.Empty; _logger.LogInformation("Received incoming message | Session={Session} | Platform={Platform} | Message={Message}", sessionToken, platformName, userMessage); // Step 1: Get or create conversation session var conversation = await GetOrCreateConversationSessionAsync(loginContext, sessionToken, ct); _logger.LogInformation("Conversation session loaded | Session={Session} | ConversationId={ConversationId} | CurrentStep={StepId}", sessionToken, conversation.MessageHubConversationId, conversation.CurrentStepId); // Ensure conversation structure is initialized conversation.Context ??= new Dictionary(); conversation.Triggers ??= _config.Triggers ?? new(); conversation.Steps ??= _config.Steps ?? new(); // Step 2: Persist incoming user message if (!string.IsNullOrWhiteSpace(userMessage)) { _logger.LogDebug("Persisting user message to DAL."); await _dal.AddConversationMessageAsync( loginContext, conversation.MessageHubConversationId, 0, // 0 = user userMessage, null, ct); } // Step 3: Process conversation step string replyMessage; try { _logger.LogInformation("Processing conversation step | Session={Session} | StepId={StepId}", sessionToken, conversation.CurrentStepId); var stepStopwatch = Stopwatch.StartNew(); replyMessage = await _EIPConceptEngine.ProcessConversationStepAsync(userMessage, conversation, loginContext, ct); stepStopwatch.Stop(); _logger.LogInformation("Processed conversation step | StepId={StepId} | Duration={Duration}ms", conversation.CurrentStepId, stepStopwatch.ElapsedMilliseconds); } catch (Exception ex) { _logger.LogError(ex, "Error processing conversation step | Session={Session} | Message={Message}", sessionToken, userMessage); replyMessage = "Sorry, something went wrong. Please try again."; } // Step 4: Persist bot reply and dispatch to platform if (!string.IsNullOrWhiteSpace(replyMessage)) { _logger.LogDebug("Persisting bot reply and dispatching to platform."); await _dal.AddConversationMessageAsync( loginContext, conversation.MessageHubConversationId, 1, // 1 = bot replyMessage, null, ct); try { await _conversationPlatform.SendMessageToPlatformAsync(replyMessage, loginContext, platformName, ct); _logger.LogInformation("Reply sent to platform | Platform={Platform} | Session={Session}", platformName, sessionToken); } catch (Exception ex) { _logger.LogError(ex, "Failed to send reply to platform | Platform={Platform} | Session={Session}", platformName, sessionToken); } } else { _logger.LogWarning("No reply generated for incoming message | Session={Session}", sessionToken); } // Step 5: Update conversation last activity try { await _dal.UpdateConversationLastActivityAsync(loginContext, sessionToken, ct); _logger.LogDebug("Updated conversation last activity | Session={Session}", sessionToken); } catch (Exception ex) { _logger.LogError(ex, "Failed to update conversation last activity | Session={Session}", sessionToken); } // Step 6: Update session cache try { _sessionCache.Set(BuildSessionKey(loginContext.ClientId, sessionToken), conversation, _sessionCacheOptions); _logger.LogDebug("Session cache updated | Session={Session}", sessionToken); } catch (Exception ex) { _logger.LogError(ex, "Failed to update session cache | Session={Session}", sessionToken); } stopwatch.Stop(); _logger.LogInformation("HandleIncomingMessageAsync completed | Session={Session} | Duration={Duration}ms", sessionToken, stopwatch.ElapsedMilliseconds); return replyMessage; } #endregion #region Session Management private async Task GetOrCreateConversationSessionAsync( LoginDTO loginContext, string sessionToken, CancellationToken ct) { var key = BuildSessionKey(loginContext.ClientId, sessionToken); _logger.LogDebug("Building session key | Key={Key}", key); // Step 1: Check in-memory cache if (_sessionCache.TryGetValue(key, out MessageHubConversationDTO? cachedConversation) && cachedConversation is not null) { _logger.LogDebug("Conversation found in cache | Session={Session} | ConversationId={ConversationId}", sessionToken, cachedConversation.MessageHubConversationId); return cachedConversation; } // Step 2: Check persisted conversation in DAL var existingConversation = await _dal.GetConversationSessionAsync(loginContext, sessionToken, ct); if (existingConversation != null) { existingConversation.Context ??= new(); existingConversation.Steps ??= _config.Steps ?? new(); existingConversation.Triggers ??= _config.Triggers ?? new(); _sessionCache.Set(key, existingConversation, _sessionCacheOptions); _logger.LogInformation("Loaded existing conversation from DAL | Session={Session} | ConversationId={ConversationId}", sessionToken, existingConversation.MessageHubConversationId); return existingConversation; } // Step 3: Create new conversation var newConversation = new MessageHubConversationDTO { ClientId = loginContext.ClientId, SessionId = sessionToken, Status = 1, Context = new(), Steps = _config.Steps ?? new(), Triggers = _config.Triggers ?? new(), CurrentStepId = 0, IsActive = true, IsCompleted = false }; newConversation.MessageHubConversationId = await _dal.CreateConversationSessionAsync(loginContext, newConversation, ct); _sessionCache.Set(key, newConversation, _sessionCacheOptions); _logger.LogInformation("Created new conversation | Session={Session} | ConversationId={ConversationId}", sessionToken, newConversation.MessageHubConversationId); return newConversation; } private static string BuildSessionKey(int clientId, string sessionToken) { if (string.IsNullOrWhiteSpace(sessionToken)) sessionToken = "ANONYMOUS"; var safeToken = new string(sessionToken.Where(c => char.IsLetterOrDigit(c) || c == '_').ToArray()) .ToUpperInvariant(); return $"CONV_{clientId}_{safeToken}"; } #endregion } }