2025-02-07 22:40:57 +00:00
using BotSharp.Abstraction.Realtime ;
using System.Net.WebSockets ;
using BotSharp.Abstraction.Realtime.Models ;
2025-02-09 23:31:46 +00:00
using BotSharp.Abstraction.MLTasks ;
2025-02-26 17:41:52 +00:00
using BotSharp.Abstraction.Conversations.Enums ;
2025-03-01 00:19:44 +00:00
using BotSharp.Abstraction.Routing.Models ;
2025-03-05 22:57:34 +00:00
using NetTopologySuite.Index.HPRtree ;
using BotSharp.Abstraction.Agents.Models ;
using Microsoft.Identity.Client.Extensions.Msal ;
using Microsoft.AspNetCore.Cors.Infrastructure ;
2025-02-07 22:40:57 +00:00
namespace BotSharp.Core.Realtime ;
public class RealtimeHub : IRealtimeHub
{
private readonly IServiceProvider _services ;
private readonly ILogger _logger ;
2025-02-27 04:40:50 +00:00
2025-02-07 22:40:57 +00:00
public RealtimeHub ( IServiceProvider services , ILogger < RealtimeHub > logger )
{
_services = services ;
_logger = logger ;
}
public async Task Listen ( WebSocket userWebSocket ,
Func < string , RealtimeHubConnection > onUserMessageReceived )
{
2025-03-03 21:30:06 +00:00
var buffer = new byte [ 1024 * 16 ] ;
2025-02-07 22:40:57 +00:00
WebSocketReceiveResult result ;
2025-02-09 23:31:46 +00:00
var completer = _services . GetServices < IRealTimeCompletion > ( ) . First ( x = > x . Provider = = "openai" ) ;
2025-02-07 22:40:57 +00:00
do
{
result = await userWebSocket . ReceiveAsync ( new ArraySegment < byte > ( buffer ) , CancellationToken . None ) ;
string receivedText = Encoding . UTF8 . GetString ( buffer , 0 , result . Count ) ;
2025-03-03 21:30:06 +00:00
2025-02-07 22:40:57 +00:00
if ( string . IsNullOrEmpty ( receivedText ) )
{
continue ;
}
var conn = onUserMessageReceived ( receivedText ) ;
2025-02-09 23:31:46 +00:00
if ( conn . Event = = "user_connected" )
2025-02-07 22:40:57 +00:00
{
2025-02-09 23:31:46 +00:00
await ConnectToModel ( completer , userWebSocket , conn ) ;
2025-02-07 22:40:57 +00:00
}
2025-02-09 23:31:46 +00:00
else if ( conn . Event = = "user_data_received" )
2025-02-07 22:40:57 +00:00
{
2025-02-09 23:31:46 +00:00
await completer . AppenAudioBuffer ( conn . Data ) ;
2025-02-07 22:40:57 +00:00
}
2025-03-05 22:57:34 +00:00
else if ( conn . Event = = "user_dtmf_received" )
{
await HandleUserDtmfReceived ( completer , conn ) ;
}
2025-02-09 23:31:46 +00:00
else if ( conn . Event = = "user_disconnected" )
2025-02-07 22:40:57 +00:00
{
2025-02-09 23:31:46 +00:00
await completer . Disconnect ( ) ;
2025-03-05 22:57:34 +00:00
await HandleUserDisconnected ( conn ) ;
2025-02-07 22:40:57 +00:00
}
} while ( ! result . CloseStatus . HasValue ) ;
await userWebSocket . CloseAsync ( result . CloseStatus . Value , result . CloseStatusDescription , CancellationToken . None ) ;
}
2025-02-09 23:31:46 +00:00
private async Task ConnectToModel ( IRealTimeCompletion completer , WebSocket userWebSocket , RealtimeHubConnection conn )
2025-02-07 22:40:57 +00:00
{
2025-02-09 23:31:46 +00:00
var hookProvider = _services . GetRequiredService < ConversationHookProvider > ( ) ;
var convService = _services . GetRequiredService < IConversationService > ( ) ;
convService . SetConversationId ( conn . ConversationId , [ ] ) ;
var conversation = await convService . GetConversation ( conn . ConversationId ) ;
2025-02-11 23:27:07 +00:00
2025-02-09 23:31:46 +00:00
var agentService = _services . GetRequiredService < IAgentService > ( ) ;
var agent = await agentService . LoadAgent ( conversation . AgentId ) ;
2025-02-28 18:44:59 +00:00
conn . CurrentAgentId = agent . Id ;
2025-02-11 23:27:07 +00:00
2025-03-03 21:30:06 +00:00
// Set model
var model = agent . LlmConfig . Model ;
if ( ! model . Contains ( "-realtime-" ) )
{
var llmProviderService = _services . GetRequiredService < ILlmProviderService > ( ) ;
model = llmProviderService . GetProviderModel ( "openai" , "gpt-4" , realTime : true ) . Name ;
}
completer . SetModelName ( model ) ;
conn . Model = model ;
2025-02-09 23:31:46 +00:00
var routing = _services . GetRequiredService < IRoutingService > ( ) ;
2025-02-28 18:44:59 +00:00
routing . Context . Push ( agent . Id ) ;
2025-02-09 23:31:46 +00:00
var dialogs = convService . GetDialogHistory ( ) ;
2025-02-28 03:19:20 +00:00
if ( dialogs . Count = = 0 )
{
dialogs . Add ( new RoleDialogModel ( AgentRole . User , "Hi" ) ) ;
}
2025-02-09 23:31:46 +00:00
routing . Context . SetDialogs ( dialogs ) ;
await completer . Connect ( conn ,
onModelReady : async ( ) = >
{
2025-03-05 22:57:58 +00:00
// Control initial session, prevent initial response interruption
await completer . UpdateSession ( conn , turnDetection : false ) ;
2025-02-10 23:28:03 +00:00
// Add dialog history
foreach ( var item in dialogs )
{
2025-02-11 23:27:07 +00:00
await completer . InsertConversationItem ( item ) ;
2025-02-10 23:28:03 +00:00
}
if ( dialogs . LastOrDefault ( ) ? . Role = = AgentRole . Assistant )
{
2025-03-05 22:57:58 +00:00
await completer . TriggerModelInference ( $"Rephase your last response:\r\n{dialogs.LastOrDefault()?.Content}" ) ;
2025-02-10 23:28:03 +00:00
}
else
{
await completer . TriggerModelInference ( "Reply based on the conversation context." ) ;
}
2025-03-05 22:57:58 +00:00
// Start turn detection
await Task . Delay ( 1000 * 8 ) ;
await completer . UpdateSession ( conn , turnDetection : true ) ;
2025-02-09 23:31:46 +00:00
} ,
onModelAudioDeltaReceived : async audioDeltaData = >
{
2025-03-05 18:29:28 +00:00
// If this is the first delta of a new response, set the start timestamp
if ( ! conn . ResponseStartTimestamp . HasValue )
{
conn . ResponseStartTimestamp = conn . LatestMediaTimestamp ;
_logger . LogDebug ( $"Setting start timestamp for new response: {conn.ResponseStartTimestamp}ms" ) ;
}
2025-02-09 23:31:46 +00:00
var data = conn . OnModelMessageReceived ( audioDeltaData ) ;
await SendEventToUser ( userWebSocket , data ) ;
2025-03-05 18:29:28 +00:00
// Send mark messages to Media Streams so we know if and when AI response playback is finished
if ( ! string . IsNullOrEmpty ( conn . StreamId ) )
{
var markEvent = new
{
@event = "mark" ,
streamSid = conn . StreamId ,
mark = new { name = "responsePart" }
} ;
await SendEventToUser ( userWebSocket , markEvent ) ;
conn . MarkQueue . Enqueue ( "responsePart" ) ;
}
2025-02-09 23:31:46 +00:00
} ,
onModelAudioResponseDone : async ( ) = >
{
var data = conn . OnModelAudioResponseDone ( ) ;
await SendEventToUser ( userWebSocket , data ) ;
} ,
onAudioTranscriptDone : async transcript = >
{
} ,
2025-02-11 23:27:07 +00:00
onModelResponseDone : async messages = >
2025-02-09 23:31:46 +00:00
{
foreach ( var message in messages )
{
// Invoke function
2025-02-26 17:41:52 +00:00
if ( message . MessageType = = MessageTypeName . FunctionCall )
2025-02-09 23:31:46 +00:00
{
await routing . InvokeFunction ( message . FunctionName , message ) ;
2025-02-11 23:27:07 +00:00
message . Role = AgentRole . Function ;
2025-03-01 00:19:44 +00:00
if ( message . FunctionName = = "route_to_agent" )
2025-02-27 04:40:50 +00:00
{
2025-03-01 00:19:44 +00:00
var inst = JsonSerializer . Deserialize < RoutingArgs > ( message . FunctionArgs ? ? "{}" ) ;
message . Content = $"Connected to agent of {inst.AgentName}" ;
conn . CurrentAgentId = routing . Context . GetCurrentAgentId ( ) ;
await completer . UpdateSession ( conn ) ;
await completer . InsertConversationItem ( message ) ;
2025-03-01 02:42:03 +00:00
await completer . TriggerModelInference ( $"Guide the user through the next steps of the process as this Agent ({inst.AgentName}), following its instructions and operational procedures." ) ;
2025-02-27 04:40:50 +00:00
}
2025-03-01 00:19:44 +00:00
else if ( message . FunctionName = = "util-routing-fallback_to_router" )
{
var inst = JsonSerializer . Deserialize < FallbackArgs > ( message . FunctionArgs ? ? "{}" ) ;
message . Content = $"Returned to Router due to {inst.Reason}" ;
conn . CurrentAgentId = routing . Context . GetCurrentAgentId ( ) ;
2025-02-28 18:44:59 +00:00
2025-03-01 00:19:44 +00:00
await completer . UpdateSession ( conn ) ;
await completer . InsertConversationItem ( message ) ;
2025-03-01 02:42:03 +00:00
await completer . TriggerModelInference ( $"Check with user whether to proceed the new request: {inst.Reason}" ) ;
2025-03-01 00:19:44 +00:00
}
else
{
await completer . InsertConversationItem ( message ) ;
await completer . TriggerModelInference ( "Reply based on the function's output." ) ;
}
2025-02-09 23:31:46 +00:00
}
2025-03-05 18:29:28 +00:00
else
2025-03-05 13:48:06 +00:00
{
2025-03-05 18:29:28 +00:00
// append output audio transcript to conversation
dialogs . Add ( message ) ;
foreach ( var hook in hookProvider . HooksOrderByPriority )
{
hook . SetAgent ( agent )
. SetConversation ( conversation ) ;
2025-02-11 23:27:07 +00:00
2025-03-05 18:29:28 +00:00
await hook . OnResponseGenerated ( message ) ;
}
2025-02-11 23:27:07 +00:00
}
2025-02-09 23:31:46 +00:00
}
} ,
2025-02-11 23:27:07 +00:00
onConversationItemCreated : async response = >
{
} ,
onInputAudioTranscriptionCompleted : async message = >
{
2025-03-01 02:42:03 +00:00
// append input audio transcript to conversation
2025-02-11 23:27:07 +00:00
dialogs . Add ( message ) ;
2025-03-01 02:42:03 +00:00
foreach ( var hook in hookProvider . HooksOrderByPriority )
{
hook . SetAgent ( agent )
. SetConversation ( conversation ) ;
await hook . OnMessageReceived ( message ) ;
}
2025-02-11 23:27:07 +00:00
} ,
2025-02-09 23:31:46 +00:00
onUserInterrupted : async ( ) = >
{
2025-03-05 18:29:28 +00:00
// Reset states
conn . MarkQueue . Clear ( ) ;
conn . LastAssistantItem = null ;
conn . ResponseStartTimestamp = null ;
2025-02-09 23:31:46 +00:00
var data = conn . OnModelUserInterrupted ( ) ;
await SendEventToUser ( userWebSocket , data ) ;
} ) ;
2025-02-07 22:40:57 +00:00
}
2025-03-05 22:57:34 +00:00
private async Task HandleUserDtmfReceived ( IRealTimeCompletion completer , RealtimeHubConnection conn )
{
var routing = _services . GetRequiredService < IRoutingService > ( ) ;
var hookProvider = _services . GetRequiredService < ConversationHookProvider > ( ) ;
var agentService = _services . GetRequiredService < IAgentService > ( ) ;
var agent = await agentService . LoadAgent ( conn . CurrentAgentId ) ;
var dialogs = routing . Context . GetDialogs ( ) ;
var convService = _services . GetRequiredService < IConversationService > ( ) ;
var conversation = await convService . GetConversation ( conn . ConversationId ) ;
var message = new RoleDialogModel ( AgentRole . User , conn . Data )
{
CurrentAgentId = routing . Context . GetCurrentAgentId ( )
} ;
dialogs . Add ( message ) ;
foreach ( var hook in hookProvider . HooksOrderByPriority )
{
hook . SetAgent ( agent )
. SetConversation ( conversation ) ;
await hook . OnMessageReceived ( message ) ;
}
await completer . InsertConversationItem ( message ) ;
await completer . TriggerModelInference ( "Reply based on the user input" ) ;
}
private async Task HandleUserDisconnected ( RealtimeHubConnection conn )
{
// Save dialog history
var routing = _services . GetRequiredService < IRoutingService > ( ) ;
var storage = _services . GetRequiredService < IConversationStorage > ( ) ;
var dialogs = routing . Context . GetDialogs ( ) ;
foreach ( var item in dialogs )
{
storage . Append ( conn . ConversationId , item ) ;
}
}
2025-02-09 23:31:46 +00:00
private async Task SendEventToUser ( WebSocket webSocket , object message )
2025-02-07 22:40:57 +00:00
{
var data = JsonSerializer . Serialize ( message ) ;
var buffer = Encoding . UTF8 . GetBytes ( data ) ;
await webSocket . SendAsync ( new ArraySegment < byte > ( buffer ) , WebSocketMessageType . Text , true , CancellationToken . None ) ;
}
}