From 7959fbf9b54c22d8f337ed5e135bb0255e4e7ee8 Mon Sep 17 00:00:00 2001 From: Jicheng Lu <103353@smsassist.com> Date: Wed, 23 Apr 2025 18:07:54 -0500 Subject: [PATCH] init --- .../BotSharp.Core.Crontab/CrontabPlugin.cs | 4 +- .../Controllers/ConversationController.cs | 2 +- ...ketsMiddleware.cs => ChatHubMiddleware.cs} | 4 +- .../ChatStreamMiddleware.cs | 155 ++++++++++++++++++ .../Models/Stream/ChatStreamEventResponse.cs | 15 ++ src/Plugins/BotSharp.Plugin.ChatHub/Using.cs | 3 +- src/WebStarter/Program.cs | 2 +- 7 files changed, 178 insertions(+), 7 deletions(-) rename src/Plugins/BotSharp.Plugin.ChatHub/{WebSocketsMiddleware.cs => ChatHubMiddleware.cs} (94%) create mode 100644 src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs create mode 100644 src/Plugins/BotSharp.Plugin.ChatHub/Models/Stream/ChatStreamEventResponse.cs diff --git a/src/Infrastructure/BotSharp.Core.Crontab/CrontabPlugin.cs b/src/Infrastructure/BotSharp.Core.Crontab/CrontabPlugin.cs index feef2546..edfca3ea 100644 --- a/src/Infrastructure/BotSharp.Core.Crontab/CrontabPlugin.cs +++ b/src/Infrastructure/BotSharp.Core.Crontab/CrontabPlugin.cs @@ -36,7 +36,7 @@ public class CrontabPlugin : IBotSharpPlugin services.AddScoped(); services.AddScoped(); - services.AddHostedService(); - services.AddHostedService(); + //services.AddHostedService(); + //services.AddHostedService(); } } diff --git a/src/Infrastructure/BotSharp.OpenAPI/Controllers/ConversationController.cs b/src/Infrastructure/BotSharp.OpenAPI/Controllers/ConversationController.cs index a76d197c..ebbda722 100644 --- a/src/Infrastructure/BotSharp.OpenAPI/Controllers/ConversationController.cs +++ b/src/Infrastructure/BotSharp.OpenAPI/Controllers/ConversationController.cs @@ -654,4 +654,4 @@ public class ConversationController : ControllerBase return jsonOption; } #endregion -} +} \ No newline at end of file diff --git a/src/Plugins/BotSharp.Plugin.ChatHub/WebSocketsMiddleware.cs b/src/Plugins/BotSharp.Plugin.ChatHub/ChatHubMiddleware.cs similarity index 94% rename from src/Plugins/BotSharp.Plugin.ChatHub/WebSocketsMiddleware.cs rename to src/Plugins/BotSharp.Plugin.ChatHub/ChatHubMiddleware.cs index 47c646b7..0c63a092 100644 --- a/src/Plugins/BotSharp.Plugin.ChatHub/WebSocketsMiddleware.cs +++ b/src/Plugins/BotSharp.Plugin.ChatHub/ChatHubMiddleware.cs @@ -3,11 +3,11 @@ using System.Text.RegularExpressions; namespace BotSharp.Plugin.ChatHub; -public class WebSocketsMiddleware +public class ChatHubMiddleware { private readonly RequestDelegate _next; - public WebSocketsMiddleware(RequestDelegate next) + public ChatHubMiddleware(RequestDelegate next) { _next = next; } diff --git a/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs b/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs new file mode 100644 index 00000000..a9411ed3 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs @@ -0,0 +1,155 @@ +using Azure; +using BotSharp.Abstraction.Realtime; +using BotSharp.Abstraction.Realtime.Models; +using Microsoft.AspNetCore.Http; +using System.Net.WebSockets; + +namespace BotSharp.Plugin.ChatHub; + +public class ChatStreamMiddleware +{ + private readonly RequestDelegate _next; + private readonly ILogger _logger; + + public ChatStreamMiddleware( + RequestDelegate next, + ILogger logger) + { + _next = next; + _logger = logger; + } + + public async Task Invoke(HttpContext httpContext) + { + var request = httpContext.Request; + + if (request.Path.StartsWithSegments("/chat/stream")) + { + if (httpContext.WebSockets.IsWebSocketRequest) + { + try + { + var services = httpContext.RequestServices; + var segments = request.Path.Value.Split("/"); + var agentId = segments[segments.Length - 2]; + var conversationId = segments[segments.Length - 1]; + + using var webSocket = await httpContext.WebSockets.AcceptWebSocketAsync(); + await HandleWebSocket(services, agentId, conversationId, webSocket); + } + catch (Exception ex) + { + _logger.LogError(ex, $"Error when connecting Chat stream. ({ex.Message})"); + } + return; + } + } + + await _next(httpContext); + } + + private async Task HandleWebSocket(IServiceProvider services, string agentId, string conversationId, WebSocket webSocket) + { + var hub = services.GetRequiredService(); + var conn = hub.SetHubConnection(conversationId); + + // load conversation and state + var convService = services.GetRequiredService(); + convService.SetConversationId(conversationId, []); + await convService.GetConversationRecordOrCreateNew(agentId); + + var buffer = new byte[1024 * 1024 * 8]; + WebSocketReceiveResult result; + + do + { + result = await webSocket.ReceiveAsync(new(buffer), CancellationToken.None); + if (result.MessageType != WebSocketMessageType.Text) + { + continue; + } + + var receivedText = Encoding.UTF8.GetString(buffer, 0, result.Count); + if (string.IsNullOrEmpty(receivedText)) + { + continue; + } + + var (eventType, data) = MapEvents(conn, receivedText); + + if (eventType == "start") + { + await ConnectToModel(hub, webSocket); + } + else if (eventType == "media") + { + if (!string.IsNullOrEmpty(data)) + { + await hub.Completer.AppenAudioBuffer(data); + } + } + else if (eventType == "disconnect") + { + await hub.Completer.Disconnect(); + } + } + while (!webSocket.CloseStatus.HasValue); + + await webSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None); + } + + private async Task ConnectToModel(IRealtimeHub hub, WebSocket webSocket) + { + await hub.ConnectToModel(async data => + { + await SendEventToUser(webSocket, data); + }); + } + + private async Task SendEventToUser(WebSocket webSocket, string message) + { + var buffer = Encoding.UTF8.GetBytes(message); + await webSocket.SendAsync(new ArraySegment(buffer), WebSocketMessageType.Text, true, CancellationToken.None); + } + + private (string, string) MapEvents(RealtimeHubConnection conn, string receivedText) + { + var response = JsonSerializer.Deserialize(receivedText); + string data = string.Empty; + + switch (response.Event) + { + case "start": + conn.ResetStreamState(); + break; + case "media": + var mediaResponse = JsonSerializer.Deserialize(receivedText); + data = mediaResponse?.Payload ?? string.Empty; + break; + case "disconnect": + break; + } + + conn.OnModelMessageReceived = message => + JsonSerializer.Serialize(new + { + @event = "media", + media = new { payload = message } + }); + + conn.OnModelAudioResponseDone = () => + JsonSerializer.Serialize(new + { + @event = "mark", + mark = new { name = "responsePart" } + }); + + conn.OnModelUserInterrupted = () => + JsonSerializer.Serialize(new + { + @event = "clear" + }); + + return (response.Event, data); + } +} diff --git a/src/Plugins/BotSharp.Plugin.ChatHub/Models/Stream/ChatStreamEventResponse.cs b/src/Plugins/BotSharp.Plugin.ChatHub/Models/Stream/ChatStreamEventResponse.cs new file mode 100644 index 00000000..817c86b1 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.ChatHub/Models/Stream/ChatStreamEventResponse.cs @@ -0,0 +1,15 @@ +using System.Text.Json.Serialization; + +namespace BotSharp.Plugin.ChatHub.Models.Stream; + +internal class ChatStreamEventResponse +{ + [JsonPropertyName("event")] + public string Event { get; set; } +} + +internal class ChatStreamMediaEventResponse : ChatStreamEventResponse +{ + [JsonPropertyName("payload")] + public string Payload { get; set; } +} diff --git a/src/Plugins/BotSharp.Plugin.ChatHub/Using.cs b/src/Plugins/BotSharp.Plugin.ChatHub/Using.cs index bc50d257..73240fbb 100644 --- a/src/Plugins/BotSharp.Plugin.ChatHub/Using.cs +++ b/src/Plugins/BotSharp.Plugin.ChatHub/Using.cs @@ -33,4 +33,5 @@ global using BotSharp.Abstraction.Messaging.Enums; global using BotSharp.Abstraction.Messaging.Models.RichContent; global using BotSharp.Abstraction.Templating; global using BotSharp.Plugin.ChatHub.Settings; -global using BotSharp.Plugin.ChatHub.Enums; \ No newline at end of file +global using BotSharp.Plugin.ChatHub.Enums; +global using BotSharp.Plugin.ChatHub.Models.Stream; \ No newline at end of file diff --git a/src/WebStarter/Program.cs b/src/WebStarter/Program.cs index 21599a68..390c7e5d 100644 --- a/src/WebStarter/Program.cs +++ b/src/WebStarter/Program.cs @@ -44,7 +44,7 @@ var app = builder.Build(); // Enable SignalR app.MapHub("/chatHub"); -app.UseMiddleware(); +app.UseMiddleware(); // Use BotSharp app.UseBotSharp()