diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetFunEventsFn.cs index 818d4e5a1..9c6fe7dd6 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.GetHub().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.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 95ed84594..c0d980046 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.GetHub().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.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 new file mode 100644 index 000000000..bf07cc3fa --- /dev/null +++ b/src/Infrastructure/BotSharp.Core/MessageHub/MessageHubExtensions.cs @@ -0,0 +1,14 @@ +namespace BotSharp.Core.MessageHub; + +public static class MessageHubExtensions +{ + /// + /// The hub bound to the conversation this call is running in. + /// + public static ConversationHub GetHub(this IServiceProvider services) + { + var conv = services.GetRequiredService(); + var hub = services.GetRequiredService>>(); + 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 06d345ff1..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,15 +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. - 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); } } diff --git a/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs b/src/Infrastructure/BotSharp.Core/Routing/RoutingService.InvokeFunction.cs index 5d08098f6..b932f358b 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.GetHub().PushIndication(clonedMessage, clonedMessage.Indication); var hooks = _services.GetHooksOrderByPriority(clonedMessage.CurrentAgentId); foreach (var hook in hooks)