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-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 )
{
var buffer = new byte [ 1024 * 4 ] ;
WebSocketReceiveResult result ;
2025-02-09 23:31:46 +00:00
var llmProviderService = _services . GetRequiredService < ILlmProviderService > ( ) ;
var model = llmProviderService . GetProviderModel ( "openai" , "gpt-4" ,
realTime : true ) . Name ;
var completer = _services . GetServices < IRealTimeCompletion > ( ) . First ( x = > x . Provider = = "openai" ) ;
completer . SetModelName ( model ) ;
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 ) ;
_logger . LogDebug ( $"Received from user: {receivedText}" ) ;
if ( string . IsNullOrEmpty ( receivedText ) )
{
continue ;
}
var conn = onUserMessageReceived ( receivedText ) ;
2025-02-09 23:31:46 +00:00
conn . Model = model ;
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-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-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 storage = _services . GetRequiredService < IConversationStorage > ( ) ;
2025-02-11 23:27:07 +00:00
2025-02-09 23:31:46 +00:00
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-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 ( ) = >
{
// Control initial session
2025-02-27 04:40:50 +00:00
await completer . UpdateSession ( conn ) ;
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-01 00:19:44 +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-02-09 23:31:46 +00:00
} ,
onModelAudioDeltaReceived : async audioDeltaData = >
{
var data = conn . OnModelMessageReceived ( audioDeltaData ) ;
await SendEventToUser ( userWebSocket , data ) ;
} ,
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-02-11 23:27:07 +00:00
else
{
2025-03-01 02:42:03 +00:00
// append output audio transcript to conversation
2025-02-11 23:27:07 +00:00
storage . Append ( conn . ConversationId , message ) ;
dialogs . Add ( message ) ;
foreach ( var hook in hookProvider . HooksOrderByPriority )
{
hook . SetAgent ( agent )
. SetConversation ( conversation ) ;
2025-03-01 02:42:03 +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
storage . Append ( conn . ConversationId , message ) ;
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 ( ) = >
{
var data = conn . OnModelUserInterrupted ( ) ;
await SendEventToUser ( userWebSocket , data ) ;
} ) ;
2025-02-07 22:40:57 +00:00
}
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 ) ;
}
}