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

145 lines
5.2 KiB
C#
Raw Normal View History

2025-06-17 18:17:28 +00:00
using BotSharp.Abstraction.Conversations.Dtos;
2025-07-28 22:41:28 +00:00
using BotSharp.Abstraction.Conversations.Enums;
using BotSharp.Abstraction.MessageHub.Models;
2025-07-31 22:48:51 +00:00
using BotSharp.Abstraction.MessageHub.Observers;
2025-06-17 18:17:28 +00:00
using BotSharp.Abstraction.SideCar;
2025-07-29 15:32:55 +00:00
using System.Runtime.CompilerServices;
2025-06-17 18:17:28 +00:00
namespace BotSharp.Plugin.ChatHub.Observers;
2025-08-01 18:23:46 +00:00
public class ChatHubObserver : BotSharpObserverBase<HubObserveData<RoleDialogModel>>
2025-06-17 18:17:28 +00:00
{
2025-08-01 18:23:46 +00:00
private readonly IServiceProvider _services;
2025-08-06 16:27:26 +00:00
private readonly BotSharpOptions _options;
private readonly ILogger _logger;
2025-06-17 18:17:28 +00:00
2025-07-31 22:48:51 +00:00
public ChatHubObserver(
2025-08-01 18:23:46 +00:00
IServiceProvider services,
2025-08-06 16:27:26 +00:00
BotSharpOptions options,
2025-08-01 18:23:46 +00:00
ILogger<ChatHubObserver> logger) : base()
2025-06-17 18:17:28 +00:00
{
2025-08-01 18:23:46 +00:00
_services = services;
2025-08-06 16:27:26 +00:00
_options = options;
2025-06-17 18:17:28 +00:00
_logger = logger;
}
2025-08-01 18:23:46 +00:00
public override string Name => nameof(ChatHubObserver);
2025-07-31 22:48:51 +00:00
2025-08-01 18:23:46 +00:00
public override void OnCompleted()
2025-06-17 18:17:28 +00:00
{
2025-06-18 03:23:43 +00:00
_logger.LogWarning($"{nameof(ChatHubObserver)} receives complete notification.");
2025-06-17 18:17:28 +00:00
}
2025-08-01 18:23:46 +00:00
public override void OnError(Exception error)
2025-06-17 18:17:28 +00:00
{
_logger.LogError(error, $"{nameof(ChatHubObserver)} receives error notification: {error.Message}");
}
2025-08-01 18:23:46 +00:00
public override void OnNext(HubObserveData<RoleDialogModel> value)
2025-06-17 18:17:28 +00:00
{
2025-06-17 22:57:57 +00:00
var message = value.Data;
var model = new ChatResponseDto();
2025-07-28 22:41:28 +00:00
var action = new ConversationSenderActionModel();
var conv = _services.GetRequiredService<IConversationService>();
2025-06-17 22:57:57 +00:00
2025-07-28 22:41:28 +00:00
switch (value.EventName)
2025-06-17 22:57:57 +00:00
{
2025-07-28 22:41:28 +00:00
case ChatEvent.BeforeReceiveLlmStreamMessage:
2025-07-29 16:17:33 +00:00
if (!AllowSendingMessage()) return;
2025-07-28 22:41:28 +00:00
model = new ChatResponseDto()
2025-06-18 03:23:43 +00:00
{
2025-07-28 22:41:28 +00:00
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Text = string.Empty,
Sender = new()
{
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
action = new ConversationSenderActionModel
{
ConversationId = conv.ConversationId,
SenderAction = SenderActionEnum.TypingOn
};
2025-07-29 15:32:55 +00:00
SendEvent(ChatEvent.OnSenderActionGenerated, conv.ConversationId, action);
2025-07-28 22:41:28 +00:00
break;
case ChatEvent.OnReceiveLlmStreamMessage:
2025-07-29 16:17:33 +00:00
if (!AllowSendingMessage()) return;
2025-07-28 22:41:28 +00:00
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
}
};
break;
case ChatEvent.AfterReceiveLlmStreamMessage:
2025-07-29 16:17:33 +00:00
if (!AllowSendingMessage()) return;
2025-07-28 22:41:28 +00:00
model = new ChatResponseDto()
{
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Text = message.Content,
Sender = new()
{
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
break;
case ChatEvent.OnIndicationReceived:
model = new ChatResponseDto
2025-06-17 22:57:57 +00:00
{
2025-07-28 22:41:28 +00:00
ConversationId = conv.ConversationId,
MessageId = message.MessageId,
Indication = message.Indication,
Sender = new()
{
FirstName = "AI",
LastName = "Assistant",
Role = AgentRole.Assistant
}
};
2025-07-31 22:48:51 +00:00
2025-08-01 19:41:17 +00:00
#if DEBUG
2025-07-31 22:48:51 +00:00
_logger.LogCritical($"Receiving {value.EventName} ({value.Data.Indication}) in {nameof(ChatHubObserver)} - {conv.ConversationId}");
2025-08-01 19:41:17 +00:00
#endif
2025-07-28 22:41:28 +00:00
break;
2025-06-17 22:57:57 +00:00
}
2025-07-29 15:32:55 +00:00
SendEvent(value.EventName, model.ConversationId, model);
2025-06-17 18:17:28 +00:00
}
2025-07-24 21:00:37 +00:00
private bool AllowSendingMessage()
{
var sidecar = _services.GetService<IConversationSideCar>();
return sidecar == null || !sidecar.IsEnabled;
}
2025-07-29 15:32:55 +00:00
#region Private methods
private void SendEvent<T>(string @event, string conversationId, T data, [CallerMemberName] string callerName = "")
2025-06-17 18:17:28 +00:00
{
2025-07-29 15:32:55 +00:00
var user = _services.GetRequiredService<IUserIdentity>();
2025-08-06 16:27:26 +00:00
var json = JsonSerializer.Serialize(data, _options.JsonSerializerOptions);
EventEmitter.SendChatEvent(_services, _logger, @event, conversationId, user?.Id, json, nameof(ChatHubObserver), callerName)
2025-07-29 15:32:55 +00:00
.ConfigureAwait(false).GetAwaiter().GetResult();
2025-06-17 18:17:28 +00:00
}
2025-07-29 15:32:55 +00:00
#endregion
2025-06-17 18:17:28 +00:00
}