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
}
}