refine message hub

This commit is contained in:
Jicheng Lu 2025-07-29 13:55:53 -05:00
parent 8391fa9293
commit 7d3579c1ec
14 changed files with 21 additions and 21 deletions

View file

@ -120,7 +120,7 @@ public class RoleDialogModel : ITrackableMessage
[JsonIgnore(Condition = JsonIgnoreCondition.Always)]
public bool IsStreaming { get; set; }
private RoleDialogModel()
public RoleDialogModel()
{
}

View file

@ -1,7 +1,7 @@
namespace BotSharp.Abstraction.MessageHub.Models;
public class HubObserveData : ObserveDataBase
public class HubObserveData<TData> : ObserveDataBase where TData : class, new()
{
public string EventName { get; set; } = null!;
public RoleDialogModel Data { get; set; } = null!;
public TData Data { get; set; } = null!;
}

View file

@ -44,7 +44,7 @@ public class ConversationPlugin : IBotSharpPlugin, IBotSharpAppPlugin
return settingService.Bind<GoogleApiSettings>("GoogleApi");
});
services.AddSingleton<MessageHub<HubObserveData>>();
services.AddSingleton<MessageHub<HubObserveData<RoleDialogModel>>>();
services.AddScoped<IConversationStorage, ConversationStorage>();
services.AddScoped<IConversationService, ConversationService>();
@ -72,8 +72,8 @@ public class ConversationPlugin : IBotSharpPlugin, IBotSharpAppPlugin
public void Configure(IApplicationBuilder app)
{
var services = app.ApplicationServices;
var queue = services.GetRequiredService<MessageHub<HubObserveData>>();
var logger = services.GetRequiredService<ILogger<MessageHub<HubObserveData>>>();
var queue = services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var logger = services.GetRequiredService<ILogger<MessageHub<HubObserveData<RoleDialogModel>>>>();
queue.Events.Subscribe(new ConversationObserver(logger));
}
}

View file

@ -18,7 +18,7 @@ public class GetWeatherFn : IFunctionCallback
public async Task<bool> Execute(RoleDialogModel message)
{
var conv = _services.GetRequiredService<IConversationService>();
var messageHub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var messageHub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
await Task.Delay(1000);

View file

@ -2,7 +2,7 @@ using System.Reactive.Subjects;
namespace BotSharp.Core.MessageHub;
public class MessageHub<T> where T : class
public class MessageHub<T> where T : class, new()
{
private readonly ILogger<MessageHub<T>> _logger;
private readonly ISubject<T> _observable = Subject.Synchronize(new Subject<T>());

View file

@ -1,6 +1,6 @@
namespace BotSharp.Core.MessageHub.Observers;
public class ConversationObserver : IObserver<HubObserveData>
public class ConversationObserver : IObserver<HubObserveData<RoleDialogModel>>
{
private readonly ILogger _logger;
private IServiceProvider _services;
@ -20,11 +20,11 @@ public class ConversationObserver : IObserver<HubObserveData>
_logger.LogError(error, $"{nameof(ConversationObserver)} receives error notification: {error.Message}");
}
public void OnNext(HubObserveData value)
public void OnNext(HubObserveData<RoleDialogModel> value)
{
_services = value.ServiceProvider;
var progress = _services.GetRequiredService<IConversationProgressService>();
if (value.EventName == ChatEvent.OnIndicationReceived
&& progress.OnFunctionExecuting != null)
{

View file

@ -25,7 +25,7 @@ public partial class RoutingService
clonedMessage.FunctionName = name;
clonedMessage.Indication = await funcExecutor.GetIndicatorAsync(message);
var messageHub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var messageHub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
messageHub.Push(new()
{
EventName = ChatEvent.OnIndicationReceived,

View file

@ -212,7 +212,7 @@ public class ChatCompletionProvider : IChatCompletion
var chatClient = client.GetChatClient(_model);
var (prompt, messages, options) = PrepareOptions(agent, conversations);
var hub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var hub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var messageId = conversations.LastOrDefault()?.MessageId ?? string.Empty;
var contentHooks = _services.GetHooks<IContentGeneratingHook>(agent.Id);

View file

@ -36,8 +36,8 @@ public class ChatHubPlugin : IBotSharpPlugin, IBotSharpAppPlugin
public void Configure(IApplicationBuilder app)
{
var services = app.ApplicationServices;
var queue = services.GetRequiredService<MessageHub<HubObserveData>>();
var logger = services.GetRequiredService<ILogger<MessageHub<HubObserveData>>>();
var queue = services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var logger = services.GetRequiredService<ILogger<MessageHub<HubObserveData<RoleDialogModel>>>>();
queue.Events.Subscribe(new ChatHubObserver(logger));
}
}

View file

@ -6,7 +6,7 @@ using System.Runtime.CompilerServices;
namespace BotSharp.Plugin.ChatHub.Observers;
public class ChatHubObserver : IObserver<HubObserveData>
public class ChatHubObserver : IObserver<HubObserveData<RoleDialogModel>>
{
private readonly ILogger _logger;
private IServiceProvider _services;
@ -26,7 +26,7 @@ public class ChatHubObserver : IObserver<HubObserveData>
_logger.LogError(error, $"{nameof(ChatHubObserver)} receives error notification: {error.Message}");
}
public void OnNext(HubObserveData value)
public void OnNext(HubObserveData<RoleDialogModel> value)
{
_services = value.ServiceProvider;

View file

@ -179,7 +179,7 @@ public class ChatCompletionProvider : IChatCompletion
var chatClient = client.GetChatClient(_model);
var (prompt, messages, options) = PrepareOptions(agent, conversations);
var hub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var hub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var messageId = conversations.LastOrDefault()?.MessageId ?? string.Empty;
var contentHooks = _services.GetHooks<IContentGeneratingHook>(agent.Id);

View file

@ -183,7 +183,7 @@ public class ChatCompletionProvider : IChatCompletion
_logger.LogInformation(agent.Instruction);
}
var hub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var hub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var messageId = conversations.LastOrDefault()?.MessageId ?? string.Empty;
hub.Push(new()

View file

@ -188,7 +188,7 @@ public class ChatCompletionProvider : IChatCompletion
var chatClient = client.GetChatClient(_model);
var (prompt, messages, options) = PrepareOptions(agent, conversations);
var hub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var hub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
var messageId = conversations.LastOrDefault()?.MessageId ?? string.Empty;
var contentHooks = _services.GetHooks<IContentGeneratingHook>(agent.Id);

View file

@ -152,7 +152,7 @@ public class ChatCompletionProvider : IChatCompletion
var client = new SparkDeskClient(appId: _settings.AppId, apiKey: _settings.ApiKey, apiSecret: _settings.ApiSecret);
var (prompt, messages, funcall) = PrepareOptions(agent, conversations);
var messageId = conversations.LastOrDefault()?.MessageId ?? string.Empty;
var hub = _services.GetRequiredService<MessageHub<HubObserveData>>();
var hub = _services.GetRequiredService<MessageHub<HubObserveData<RoleDialogModel>>>();
hub.Push(new()
{