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-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-11 23:27:07 +00:00
|
|
|
conn.EntryAgentId = agent.Id;
|
|
|
|
|
|
2025-02-09 23:31:46 +00:00
|
|
|
var routing = _services.GetRequiredService<IRoutingService>();
|
|
|
|
|
var dialogs = convService.GetDialogHistory();
|
|
|
|
|
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-02-11 23:27:07 +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-02-27 04:40:50 +00:00
|
|
|
if (message.FunctionName == "route_to_agent")
|
|
|
|
|
{
|
|
|
|
|
var routedAgentId = routing.Context.GetCurrentAgentId();
|
|
|
|
|
if (conn.EntryAgentId != routedAgentId)
|
|
|
|
|
{
|
|
|
|
|
conn.EntryAgentId = routedAgentId;
|
|
|
|
|
await completer.UpdateSession(conn);
|
|
|
|
|
await completer.TriggerModelInference("Reply based on the function's output.");
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
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
|
|
|
|
|
{
|
|
|
|
|
// append transcript to conversation
|
|
|
|
|
storage.Append(conn.ConversationId, message);
|
|
|
|
|
dialogs.Add(message);
|
|
|
|
|
|
|
|
|
|
foreach (var hook in hookProvider.HooksOrderByPriority)
|
|
|
|
|
{
|
|
|
|
|
hook.SetAgent(agent)
|
|
|
|
|
.SetConversation(conversation);
|
|
|
|
|
|
|
|
|
|
if (!string.IsNullOrEmpty(message.Content))
|
|
|
|
|
{
|
2025-02-26 19:16:57 +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 =>
|
|
|
|
|
{
|
|
|
|
|
// append transcript to conversation
|
|
|
|
|
storage.Append(conn.ConversationId, message);
|
|
|
|
|
dialogs.Add(message);
|
|
|
|
|
},
|
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);
|
|
|
|
|
}
|
|
|
|
|
}
|