From fa88e9dcafe25abe34fa51f6c51c6246d9d755b6 Mon Sep 17 00:00:00 2001 From: "nick.yi" Date: Thu, 24 Sep 2026 09:48:14 +0800 Subject: [PATCH 1/2] optimize PushIndication --- .../Demo/Functions/GetFunEventsFn.cs | 18 +------- .../Demo/Functions/GetWeatherFn.cs | 16 +------ .../MessageHub/MessageHubExtensions.cs | 45 +++++++++++++++++++ .../Routing/Executor/MCPToolExecutor.cs | 10 +---- .../Routing/RoutingService.InvokeFunction.cs | 16 ++----- 5 files changed, 53 insertions(+), 52 deletions(-) create mode 100644 src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs index 818d4e5a1..018ecb9cc 100644 --- a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs +++ b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs @@ -18,28 +18,14 @@ public GetFunEventsFn(IServiceProvider services) public async Task Execute(RoleDialogModel message) { var args = JsonSerializer.Deserialize(message.FunctionArgs); - var conv = _services.GetRequiredService(); - var messageHub = _services.GetRequiredService>>(); await Task.Delay(1000); - message.Indication = $"Start querying event data in {args?.City}"; - messageHub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = message, - RefId = conv.ConversationId - }); + _services.PushIndication(message, $"Start querying event data in {args?.City}"); await Task.Delay(1500); - message.Indication = $"Still searching events in {args?.City}"; - messageHub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = message, - RefId = conv.ConversationId - }); + _services.PushIndication(message, $"Still searching events in {args?.City}"); await Task.Delay(1500); diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs index 95ed84594..5378d4941 100644 --- a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs +++ b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs @@ -28,13 +28,7 @@ public async Task Execute(RoleDialogModel message) await Task.Delay(1000); - message.Indication = $"Start querying weather data in {args?.City}"; - messageHub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = message, - RefId = conv.ConversationId - }); + _services.PushIndication(message, $"Start querying weather data in {args?.City}"); await Task.Delay(1500); @@ -49,13 +43,7 @@ public async Task Execute(RoleDialogModel message) }); #endif - message.Indication = $"Still working on it... Hold on, {args?.City}"; - messageHub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = message, - RefId = conv.ConversationId - }); + _services.PushIndication(message, $"Still working on it... Hold on, {args?.City}"); await Task.Delay(1500); diff --git a/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs b/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs new file mode 100644 index 000000000..2ed3510ca --- /dev/null +++ b/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs @@ -0,0 +1,45 @@ +namespace BotSharp.Core.MessageHub; + +/// +/// Raising the "working on it" line in chat. +/// +public static class MessageHubExtensions +{ + /// + /// Announce in the conversation this call is running in. + /// is pushed as-is and its Indication overwritten, so clone it + /// when the original is still in use by a call in flight. + /// + public static void PushIndication(this IServiceProvider services, RoleDialogModel message, string? indication) + { + var conv = services.GetRequiredService(); + var hub = services.GetRequiredService>>(); + hub.PushIndication(message, indication, conv.ConversationId); + } + + /// + /// The same, for callers already holding both -- a thread where resolving a scoped service to + /// ask for the conversation id would be a race. + /// + public static void PushIndication( + this MessageHub> hub, + RoleDialogModel message, + string? indication, + string? conversationId) + { + // Nothing to announce, or no one to announce it to: subscribers match on RefId. + if (string.IsNullOrWhiteSpace(indication) || string.IsNullOrWhiteSpace(conversationId)) + { + return; + } + + message.Indication = indication; + + hub.Push(new() + { + EventName = ChatEvent.OnIndicationReceived, + Data = message, + RefId = conversationId + }); + } +} diff --git a/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs b/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs index 06d345ff1..47b2313db 100644 --- a/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs +++ b/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs @@ -145,15 +145,7 @@ public void Report(ProgressNotificationValue value) // Cloned: this is pushed to observers that read it, and the function's own message is // still being used by the call in flight. Its indication is not ours to overwrite. - var indication = RoleDialogModel.From(_message); - indication.Indication = value.Message; - - _hub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = indication, - RefId = _conversationId - }); + _hub.PushIndication(RoleDialogModel.From(_message), value.Message, _conversationId); } } diff --git a/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs b/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs index 5d08098f6..f7e2e8c89 100644 --- a/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs +++ b/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs @@ -26,20 +26,10 @@ public async Task InvokeFunction(string name, RoleDialogModel message, Inv // Clone message var clonedMessage = RoleDialogModel.From(message); clonedMessage.FunctionName = name; + // Assigned even when empty: the hooks below log this field, and the value From() copied + // would otherwise linger there. clonedMessage.Indication = await funcExecutor.GetIndicatorAsync(message); - - // An empty indicator means the callee has nothing worth announcing for this call. - if (!string.IsNullOrEmpty(clonedMessage.Indication)) - { - var conv = _services.GetRequiredService(); - var messageHub = _services.GetRequiredService>>(); - messageHub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = clonedMessage, - RefId = conv.ConversationId - }); - } + _services.PushIndication(clonedMessage, clonedMessage.Indication); var hooks = _services.GetHooksOrderByPriority(clonedMessage.CurrentAgentId); foreach (var hook in hooks) From 2103a20c2ad1b86a718026414d77bcaae42d1b3f Mon Sep 17 00:00:00 2001 From: "nick.yi" Date: Thu, 24 Sep 2026 10:10:44 +0800 Subject: [PATCH 2/2] optimize ConversationHub --- .../Demo/Functions/GetFunEventsFn.cs | 4 +- .../Demo/Functions/GetWeatherFn.cs | 4 +- .../MessageHub/ConversationHub.cs | 42 +++++++++++++++++++ .../MessageHub/MessageHubExtensions.cs | 37 ++-------------- .../Routing/Executor/MCPToolExecutor.cs | 27 ++++-------- .../Routing/RoutingService.InvokeFunction.cs | 2 +- 6 files changed, 58 insertions(+), 58 deletions(-) create mode 100644 src/Infrastructure/BotSharp.Core/MessageHub/ConversationHub.cs diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs index 018ecb9cc..9c6fe7dd6 100644 --- a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs +++ b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs @@ -21,11 +21,11 @@ public async Task Execute(RoleDialogModel message) await Task.Delay(1000); - _services.PushIndication(message, $"Start querying event data in {args?.City}"); + _services.GetHub().PushIndication(message, $"Start querying event data in {args?.City}"); await Task.Delay(1500); - _services.PushIndication(message, $"Still searching events in {args?.City}"); + _services.GetHub().PushIndication(message, $"Still searching events in {args?.City}"); await Task.Delay(1500); diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs index 5378d4941..c0d980046 100644 --- a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs +++ b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs @@ -28,7 +28,7 @@ public async Task Execute(RoleDialogModel message) await Task.Delay(1000); - _services.PushIndication(message, $"Start querying weather data in {args?.City}"); + _services.GetHub().PushIndication(message, $"Start querying weather data in {args?.City}"); await Task.Delay(1500); @@ -43,7 +43,7 @@ public async Task Execute(RoleDialogModel message) }); #endif - _services.PushIndication(message, $"Still working on it... Hold on, {args?.City}"); + _services.GetHub().PushIndication(message, $"Still working on it... Hold on, {args?.City}"); await Task.Delay(1500); diff --git a/src/Infrastructure/BotSharp.Core/MessageHub/ConversationHub.cs b/src/Infrastructure/BotSharp.Core/MessageHub/ConversationHub.cs new file mode 100644 index 000000000..488a59cfd --- /dev/null +++ b/src/Infrastructure/BotSharp.Core/MessageHub/ConversationHub.cs @@ -0,0 +1,42 @@ +namespace BotSharp.Core.MessageHub; + +/// +/// Pushing events into one conversation. The id is captured when the hub is taken, on the thread +/// that asked for it -- a reporter running on a transport's own thread can keep this and push +/// without resolving a scoped service again. Don't hold one past the conversation it was taken for. +/// +public sealed class ConversationHub +{ + private readonly MessageHub> _hub; + + internal ConversationHub(MessageHub> hub, string? conversationId) + { + _hub = hub; + ConversationId = conversationId; + } + + /// Where every push goes; null when taken outside a conversation. + public string? ConversationId { get; } + + /// + /// Raise the "working on it" line in chat. is pushed as-is and its + /// Indication overwritten, so clone it when the original is still in use by a call in flight. + /// + public void PushIndication(RoleDialogModel message, string? indication) + { + // Nothing to announce, or no one to announce it to: subscribers match on RefId. + if (string.IsNullOrWhiteSpace(indication) || string.IsNullOrWhiteSpace(ConversationId)) + { + return; + } + + message.Indication = indication; + + _hub.Push(new() + { + EventName = ChatEvent.OnIndicationReceived, + Data = message, + RefId = ConversationId + }); + } +} diff --git a/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs b/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs index 2ed3510ca..bf07cc3fa 100644 --- a/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs +++ b/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs @@ -1,45 +1,14 @@ namespace BotSharp.Core.MessageHub; -/// -/// Raising the "working on it" line in chat. -/// public static class MessageHubExtensions { /// - /// Announce in the conversation this call is running in. - /// is pushed as-is and its Indication overwritten, so clone it - /// when the original is still in use by a call in flight. + /// The hub bound to the conversation this call is running in. /// - public static void PushIndication(this IServiceProvider services, RoleDialogModel message, string? indication) + public static ConversationHub GetHub(this IServiceProvider services) { var conv = services.GetRequiredService(); var hub = services.GetRequiredService>>(); - hub.PushIndication(message, indication, conv.ConversationId); - } - - /// - /// The same, for callers already holding both -- a thread where resolving a scoped service to - /// ask for the conversation id would be a race. - /// - public static void PushIndication( - this MessageHub> hub, - RoleDialogModel message, - string? indication, - string? conversationId) - { - // Nothing to announce, or no one to announce it to: subscribers match on RefId. - if (string.IsNullOrWhiteSpace(indication) || string.IsNullOrWhiteSpace(conversationId)) - { - return; - } - - message.Indication = indication; - - hub.Push(new() - { - EventName = ChatEvent.OnIndicationReceived, - Data = message, - RefId = conversationId - }); + return new ConversationHub(hub, conv.ConversationId); } } diff --git a/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs b/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs index 47b2313db..33ec616f5 100644 --- a/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs +++ b/src/Infrastructure/BotSharp.Core/Routing/Executor/MCPToolExecutor.cs @@ -87,17 +87,12 @@ public Task GetIndicatorAsync(RoleDialogModel message) /// private sealed class ToolProgressIndicator : IProgress { - private readonly MessageHub> _hub; - private readonly string _conversationId; + private readonly ConversationHub _hub; private readonly RoleDialogModel _message; - private ToolProgressIndicator( - MessageHub> hub, - string conversationId, - RoleDialogModel message) + private ToolProgressIndicator(ConversationHub hub, RoleDialogModel message) { _hub = hub; - _conversationId = conversationId; _message = message; } @@ -106,22 +101,16 @@ private ToolProgressIndicator( /// tool invoked outside one, from a task or a test. Null is the right answer there rather /// than a reporter that drops everything: it also tells the server not to bother sending. /// - /// The conversation id is read HERE, on the thread that starts the call, and captured. - /// runs on whichever thread the MCP transport is reading on, and - /// resolving a scoped service from there to ask again would be a race for a value that + /// The hub is taken HERE, on the thread that starts the call, so it carries the conversation + /// id with it. runs on whichever thread the MCP transport is reading on, + /// and resolving a scoped service from there to ask again would be a race for a value that /// cannot change during the call. /// /// public static ToolProgressIndicator? For(IServiceProvider services, RoleDialogModel message) { - var conversationId = services.GetRequiredService().ConversationId; - if (string.IsNullOrWhiteSpace(conversationId)) - { - return null; - } - - var hub = services.GetRequiredService>>(); - return new ToolProgressIndicator(hub, conversationId, message); + var hub = services.GetHub(); + return string.IsNullOrWhiteSpace(hub.ConversationId) ? null : new ToolProgressIndicator(hub, message); } /// @@ -145,7 +134,7 @@ public void Report(ProgressNotificationValue value) // Cloned: this is pushed to observers that read it, and the function's own message is // still being used by the call in flight. Its indication is not ours to overwrite. - _hub.PushIndication(RoleDialogModel.From(_message), value.Message, _conversationId); + _hub.PushIndication(RoleDialogModel.From(_message), value.Message); } } diff --git a/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs b/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs index f7e2e8c89..b932f358b 100644 --- a/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs +++ b/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs @@ -29,7 +29,7 @@ public async Task InvokeFunction(string name, RoleDialogModel message, Inv // Assigned even when empty: the hooks below log this field, and the value From() copied // would otherwise linger there. clonedMessage.Indication = await funcExecutor.GetIndicatorAsync(message); - _services.PushIndication(clonedMessage, clonedMessage.Indication); + _services.GetHub().PushIndication(clonedMessage, clonedMessage.Indication); var hooks = _services.GetHooksOrderByPriority(clonedMessage.CurrentAgentId); foreach (var hook in hooks)