BotSharp/src/Plugins/BotSharp.Plugin.ChatHub/Observers/ChatHubObserver.cs

164 lines
5.9 KiB
C#
Raw Normal View History

2025-06-17 18:17:28 +00:00
using BotSharp.Abstraction.Conversations.Dtos;
using BotSharp.Abstraction.Observables.Models;
using BotSharp.Abstraction.SideCar;
using BotSharp.Plugin.ChatHub.Hooks;
using Microsoft.AspNetCore.SignalR;
namespace BotSharp.Plugin.ChatHub.Observers;
public class ChatHubObserver : IObserver<HubObserveData>
{
private readonly ILogger _logger;
private IServiceProvider _services;
2025-06-17 22:57:57 +00:00
private const string BEFORE_RECEIVE_LLM_STREAM_MESSAGE = "BeforeReceiveLlmStreamMessage";
private const string ON_RECEIVE_LLM_STREAM_MESSAGE = "OnReceiveLlmStreamMessage";
private const string AFTER_RECEIVE_LLM_STREAM_MESSAGE = "AfterReceiveLlmStreamMessage";
2025-06-17 18:17:28 +00:00
private const string GENERATE_SENDER_ACTION = "OnSenderActionGenerated";
public ChatHubObserver(ILogger logger)
{
_logger = logger;
}
public void OnCompleted()
{
2025-06-18 03:23:43 +00:00
_logger.LogWarning($"{nameof(ChatHubObserver)} receives complete notification.");
2025-06-17 18:17:28 +00:00
}
public void OnError(Exception error)
{
_logger.LogError(error, $"{nameof(ChatHubObserver)} receives error notification: {error.Message}");
}
public void OnNext(HubObserveData value)
{
_services = value.ServiceProvider;
2025-06-17 22:57:57 +00:00
2025-06-18 03:23:43 +00:00
if (!AllowSendingMessage()) return;
2025-06-17 22:57:57 +00:00
var message = value.Data;
var model = new ChatResponseDto();
2025-06-18 03:23:43 +00:00
if (value.EventName == BEFORE_RECEIVE_LLM_STREAM_MESSAGE)
2025-06-17 22:57:57 +00:00
{
var conv = _services.GetRequiredService<IConversationService>();
model = new ChatResponseDto()
{
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Text = string.Empty,
Sender = new()
{
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
var action = new ConversationSenderActionModel
{
ConversationId = conv.ConversationId,
2025-06-18 03:23:43 +00:00
SenderAction = SenderActionEnum.TypingOn
2025-06-17 22:57:57 +00:00
};
GenerateSenderAction(conv.ConversationId, action).ConfigureAwait(false).GetAwaiter().GetResult();
}
2025-06-24 13:34:38 +00:00
else if (value.EventName == AFTER_RECEIVE_LLM_STREAM_MESSAGE && message.IsStreaming)
2025-06-17 22:57:57 +00:00
{
2025-06-24 13:34:38 +00:00
var conv = _services.GetRequiredService<IConversationService>();
model = new ChatResponseDto()
2025-06-18 03:23:43 +00:00
{
2025-06-24 13:34:38 +00:00
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Text = message.Content,
Sender = new()
2025-06-18 03:23:43 +00:00
{
2025-06-24 13:34:38 +00:00
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
2025-06-18 03:23:43 +00:00
2025-06-24 13:34:38 +00:00
var action = new ConversationSenderActionModel
{
ConversationId = conv.ConversationId,
SenderAction = SenderActionEnum.TypingOff
};
2025-06-18 03:23:43 +00:00
2025-06-24 13:34:38 +00:00
GenerateSenderAction(conv.ConversationId, action).ConfigureAwait(false).GetAwaiter().GetResult();
2025-06-17 22:57:57 +00:00
}
else if (value.EventName == ON_RECEIVE_LLM_STREAM_MESSAGE)
{
var conv = _services.GetRequiredService<IConversationService>();
model = new ChatResponseDto()
{
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Text = !string.IsNullOrEmpty(message.SecondaryContent) ? message.SecondaryContent : message.Content,
Function = message.FunctionName,
RichContent = message.SecondaryRichContent ?? message.RichContent,
Data = message.Data,
Sender = new()
{
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
}
OnReceiveAssistantMessage(value.EventName, model.ConversationId, model).ConfigureAwait(false).GetAwaiter().GetResult();
2025-06-17 18:17:28 +00:00
}
2025-06-17 22:57:57 +00:00
private async Task OnReceiveAssistantMessage(string @event, string conversationId, ChatResponseDto model)
2025-06-17 18:17:28 +00:00
{
try
{
var settings = _services.GetRequiredService<ChatHubSettings>();
var chatHub = _services.GetRequiredService<IHubContext<SignalRHub>>();
if (settings.EventDispatchBy == EventDispatchType.Group)
{
2025-06-17 22:57:57 +00:00
await chatHub.Clients.Group(conversationId).SendAsync(@event, model);
2025-06-17 18:17:28 +00:00
}
else
{
2025-06-17 22:57:57 +00:00
var user = _services.GetRequiredService<IUserIdentity>();
await chatHub.Clients.User(user.Id).SendAsync(@event, model);
2025-06-17 18:17:28 +00:00
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, $"Failed to receive assistant message in {nameof(ChatHubConversationHook)} (conversation id: {conversationId})");
}
}
private bool AllowSendingMessage()
{
var sidecar = _services.GetService<IConversationSideCar>();
2025-07-24 20:15:53 +00:00
return sidecar == null || !sidecar.IsEnabled;
2025-06-17 18:17:28 +00:00
}
private async Task GenerateSenderAction(string conversationId, ConversationSenderActionModel action)
{
try
{
var settings = _services.GetRequiredService<ChatHubSettings>();
var chatHub = _services.GetRequiredService<IHubContext<SignalRHub>>();
if (settings.EventDispatchBy == EventDispatchType.Group)
{
await chatHub.Clients.Group(conversationId).SendAsync(GENERATE_SENDER_ACTION, action);
}
else
{
2025-06-17 22:57:57 +00:00
var user = _services.GetRequiredService<IUserIdentity>();
await chatHub.Clients.User(user.Id).SendAsync(GENERATE_SENDER_ACTION, action);
2025-06-17 18:17:28 +00:00
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, $"Failed to generate sender action in {nameof(ChatHubConversationHook)} (conversation id: {conversationId})");
}
}
}