using BotSharp.Plugin.OpenAI.Models.Realtime; using OpenAI.Chat; namespace BotSharp.Plugin.OpenAI.Providers.Realtime; /// /// Reference to https://platform.openai.com/docs/api-reference/realtime-server-events /// public class RealTimeCompletionProvider : IRealTimeCompletion { public string Provider => "openai"; public string Model => _model; private readonly IServiceProvider _services; private readonly ILogger _logger; private readonly BotSharpOptions _botsharpOptions; protected string _model = "gpt-4o-mini-realtime-preview"; private LlmRealtimeSession _session; public RealTimeCompletionProvider( IServiceProvider services, ILogger logger, BotSharpOptions botsharpOptions) { _logger = logger; _services = services; _botsharpOptions = botsharpOptions; } public async Task Connect( RealtimeHubConnection conn, Func onModelReady, Func onModelAudioDeltaReceived, Func onModelAudioResponseDone, Func onModelAudioTranscriptDone, Func, Task> onModelResponseDone, Func onConversationItemCreated, Func onInputAudioTranscriptionDone, Func onInterruptionDetected) { var settingsService = _services.GetRequiredService(); var realtimeSettings = _services.GetRequiredService(); _model = realtimeSettings.Model; var settings = settingsService.GetSetting(Provider, _model); if (_session != null) { _session.Dispose(); } _session = new LlmRealtimeSession(_services, new ChatSessionOptions { JsonOptions = _botsharpOptions.JsonSerializerOptions }); await _session.ConnectAsync( uri: new Uri($"wss://api.openai.com/v1/realtime?model={_model}"), headers: new Dictionary { {"Authorization", $"Bearer {settings.ApiKey}"}, {"OpenAI-Beta", "realtime=v1"} }, cancellationToken: CancellationToken.None); _ = ReceiveMessage( realtimeSettings, conn, onModelReady, onModelAudioDeltaReceived, onModelAudioResponseDone, onModelAudioTranscriptDone, onModelResponseDone, onConversationItemCreated, onInputAudioTranscriptionDone, onInterruptionDetected); } public async Task Disconnect() { if (_session != null) { await _session.DisconnectAsync(); _session.Dispose(); } } public async Task AppenAudioBuffer(string message) { var audioAppend = new { type = "input_audio_buffer.append", audio = message }; await SendEventToModel(audioAppend); } public async Task AppenAudioBuffer(ArraySegment data, int length) { var message = Convert.ToBase64String(data.AsSpan(0, length).ToArray()); await AppenAudioBuffer(message); } public async Task TriggerModelInference(string? instructions = null) { // Triggering model inference if (!string.IsNullOrEmpty(instructions)) { await SendEventToModel(new { type = "response.create", response = new { instructions } }); } else { await SendEventToModel(new { type = "response.create" }); } } public async Task CancelModelResponse() { await SendEventToModel(new { type = "response.cancel" }); } public async Task RemoveConversationItem(string itemId) { await SendEventToModel(new { type = "conversation.item.delete", item_id = itemId }); } private async Task ReceiveMessage( RealtimeModelSettings realtimeSettings, RealtimeHubConnection conn, Func onModelReady, Func onModelAudioDeltaReceived, Func onModelAudioResponseDone, Func onModelAudioTranscriptDone, Func, Task> onModelResponseDone, Func onConversationItemCreated, Func onInputAudioTranscriptionDone, Func onInterruptionDetected) { DateTime? startTime = null; await foreach (ChatSessionUpdate update in _session.ReceiveUpdatesAsync(CancellationToken.None)) { var receivedText = update?.RawResponse; if (string.IsNullOrEmpty(receivedText)) { continue; } var response = JsonSerializer.Deserialize(receivedText); if (realtimeSettings?.ModelResponseTimeoutSeconds > 0 && !string.IsNullOrWhiteSpace(realtimeSettings?.ModelResponseTimeoutEndEvent) && startTime.HasValue && (DateTime.UtcNow - startTime.Value).TotalSeconds >= realtimeSettings.ModelResponseTimeoutSeconds && response.Type != realtimeSettings.ModelResponseTimeoutEndEvent) { startTime = null; await TriggerModelInference("Responsd to user immediately"); continue; } if (response.Type == "error") { _logger.LogError($"{response.Type}: {receivedText}"); var error = JsonSerializer.Deserialize(receivedText); if (error?.Body.Type == "server_error") { break; } } else if (response.Type == "session.created") { _logger.LogInformation($"{response.Type}: {receivedText}"); await onModelReady(); } else if (response.Type == "session.updated") { _logger.LogInformation($"{response.Type}: {receivedText}"); } else if (response.Type == "response.audio_transcript.delta") { _logger.LogDebug($"{response.Type}: {receivedText}"); } else if (response.Type == "response.audio_transcript.done") { _logger.LogInformation($"{response.Type}: {receivedText}"); var data = JsonSerializer.Deserialize(receivedText); await onModelAudioTranscriptDone(data.Transcript); } else if (response.Type == "response.audio.delta") { var audio = JsonSerializer.Deserialize(receivedText); if (audio?.Delta != null) { _logger.LogDebug($"{response.Type}: {receivedText}"); await onModelAudioDeltaReceived(audio.Delta, audio.ItemId); } } else if (response.Type == "response.audio.done") { _logger.LogInformation($"{response.Type}: {receivedText}"); await onModelAudioResponseDone(); } else if (response.Type == "response.done") { _logger.LogInformation($"{response.Type}: {receivedText}"); var data = JsonSerializer.Deserialize(receivedText).Body; if (data.Status != "completed") { if (data.StatusDetails.Type == "incomplete" && data.StatusDetails.Reason == "max_output_tokens") { await onInterruptionDetected(); await TriggerModelInference("Response user concisely"); } } else { var messages = await OnResponsedDone(conn, receivedText); await onModelResponseDone(messages); } } else if (response.Type == "conversation.item.created") { _logger.LogInformation($"{response.Type}: {receivedText}"); var data = JsonSerializer.Deserialize(receivedText); if (data?.Item?.Role == "user") { startTime = DateTime.UtcNow; } await onConversationItemCreated(receivedText); } else if (response.Type == "conversation.item.input_audio_transcription.completed") { _logger.LogInformation($"{response.Type}: {receivedText}"); var message = await OnUserAudioTranscriptionCompleted(conn, receivedText); if (!string.IsNullOrEmpty(message.Content)) { await onInputAudioTranscriptionDone(message); } } else if (response.Type == "input_audio_buffer.speech_started") { _logger.LogInformation($"{response.Type}: {receivedText}"); // Handle user interuption await onInterruptionDetected(); } else if (response.Type == "input_audio_buffer.speech_stopped") { _logger.LogInformation($"{response.Type}: {receivedText}"); } else if (response.Type == "input_audio_buffer.committed") { _logger.LogInformation($"{response.Type}: {receivedText}"); } } _session.Dispose(); } public async Task SendEventToModel(object message) { if (_session == null) return; await _session.SendEventToModelAsync(message); } public async Task UpdateSession(RealtimeHubConnection conn, bool isInit = false) { var convService = _services.GetRequiredService(); var conv = await convService.GetConversation(conn.ConversationId); var agentService = _services.GetRequiredService(); var agent = await agentService.LoadAgent(conn.CurrentAgentId); var (prompt, messages, options) = PrepareOptions(agent, []); var instruction = messages.FirstOrDefault()?.Content.FirstOrDefault()?.Text ?? agent?.Description ?? string.Empty; var functions = options.Tools.Select(x => { var fn = new FunctionDef { Name = x.FunctionName, Description = x.FunctionDescription }; fn.Parameters = JsonSerializer.Deserialize(x.FunctionParameters); return fn; }).ToArray(); var realtimeModelSettings = _services.GetRequiredService(); var sessionUpdate = new { type = "session.update", session = new RealtimeSessionUpdateRequest { InputAudioFormat = realtimeModelSettings.InputAudioFormat, OutputAudioFormat = realtimeModelSettings.OutputAudioFormat, Voice = realtimeModelSettings.Voice, Instructions = instruction, ToolChoice = "auto", Tools = functions, Modalities = realtimeModelSettings.Modalities, Temperature = Math.Max(options.Temperature ?? realtimeModelSettings.Temperature, 0.6f), MaxResponseOutputTokens = realtimeModelSettings.MaxResponseOutputTokens, TurnDetection = new RealtimeSessionTurnDetection { InterruptResponse = realtimeModelSettings.InterruptResponse/*, Threshold = realtimeModelSettings.TurnDetection.Threshold, PrefixPadding = realtimeModelSettings.TurnDetection.PrefixPadding, SilenceDuration = realtimeModelSettings.TurnDetection.SilenceDuration*/ }, InputAudioNoiseReduction = new InputAudioNoiseReduction { Type = "near_field" } } }; if (realtimeModelSettings.InputAudioTranscribe) { var words = new List(); HookEmitter.Emit(_services, hook => words.AddRange(hook.OnModelTranscriptPrompt(agent)), agent.Id); sessionUpdate.session.InputAudioTranscription = new InputAudioTranscription { Model = realtimeModelSettings.InputAudioTranscription.Model, Language = realtimeModelSettings.InputAudioTranscription.Language, Prompt = string.Join(", ", words.Select(x => x.ToLower().Trim()).Distinct()).SubstringMax(1024) }; } await HookEmitter.Emit(_services, async hook => { await hook.OnSessionUpdated(agent, instruction, functions, isInit); }, agent.Id); await SendEventToModel(sessionUpdate); await Task.Delay(300); return instruction; } public async Task InsertConversationItem(RoleDialogModel message) { if (message.Role == AgentRole.Function) { var functionConversationItem = new { type = "conversation.item.create", item = new { call_id = message.ToolCallId, type = "function_call_output", output = message.Content } }; await SendEventToModel(functionConversationItem); } else if (message.Role == AgentRole.Assistant) { var conversationItem = new { type = "conversation.item.create", item = new { type = "message", role = message.Role, content = new object[] { new { type = "text", text = message.Content } } } }; await SendEventToModel(conversationItem); } else if (message.Role == AgentRole.User) { var conversationItem = new { type = "conversation.item.create", item = new { type = "message", role = message.Role, content = new object[] { new { type = "input_text", text = message.Content } } } }; await SendEventToModel(conversationItem); } else { throw new NotImplementedException($"Unrecognized role {message.Role}."); } } public void SetModelName(string model) { _model = model; } #region Private methods private async Task> OnResponsedDone(RealtimeHubConnection conn, string response) { var outputs = new List(); var data = JsonSerializer.Deserialize(response).Body; if (data.Status != "completed") { _logger.LogError(data.StatusDetails.ToString()); /*if (data.StatusDetails.Type == "incomplete" && data.StatusDetails.Reason == "max_output_tokens") { await TriggerModelInference("Response user concisely"); }*/ return []; } var prompts = new List(); var inputTokenDetails = data.Usage?.InputTokenDetails; var outputTokenDetails = data.Usage?.OutputTokenDetails; foreach (var output in data.Outputs) { if (output.Type == "function_call") { outputs.Add(new RoleDialogModel(AgentRole.Assistant, output.Arguments) { CurrentAgentId = conn.CurrentAgentId, FunctionName = output.Name, FunctionArgs = output.Arguments, ToolCallId = output.CallId, MessageId = output.Id, MessageType = MessageTypeName.FunctionCall }); prompts.Add($"{output.Name}({output.Arguments})"); } else if (output.Type == "message") { var content = output.Content.FirstOrDefault()?.Transcript ?? string.Empty; outputs.Add(new RoleDialogModel(output.Role, content) { CurrentAgentId = conn.CurrentAgentId, MessageId = output.Id, MessageType = MessageTypeName.Plain }); prompts.Add(content); } } // After chat completion hook var text = string.Join("\r\n", prompts); var contentHooks = _services.GetServices(); foreach (var hook in contentHooks) { await hook.AfterGenerated(new RoleDialogModel(AgentRole.Assistant, text) { CurrentAgentId = conn.CurrentAgentId }, new TokenStatsModel { Provider = Provider, Model = _model, Prompt = text, TextInputTokens = inputTokenDetails?.TextTokens ?? 0 - inputTokenDetails?.CachedTokenDetails?.TextTokens ?? 0, CachedTextInputTokens = data.Usage?.InputTokenDetails?.CachedTokenDetails?.TextTokens ?? 0, AudioInputTokens = inputTokenDetails?.AudioTokens ?? 0 - inputTokenDetails?.CachedTokenDetails?.AudioTokens ?? 0, CachedAudioInputTokens = inputTokenDetails?.CachedTokenDetails?.AudioTokens ?? 0, TextOutputTokens = outputTokenDetails?.TextTokens ?? 0, AudioOutputTokens = outputTokenDetails?.AudioTokens ?? 0 }); } return outputs; } private async Task OnUserAudioTranscriptionCompleted(RealtimeHubConnection conn, string response) { var data = JsonSerializer.Deserialize(response); return new RoleDialogModel(AgentRole.User, data.Transcript) { CurrentAgentId = conn.CurrentAgentId }; } private (string, IEnumerable, ChatCompletionOptions) PrepareOptions(Agent agent, List conversations) { var agentService = _services.GetRequiredService(); var state = _services.GetRequiredService(); var fileStorage = _services.GetRequiredService(); var settingsService = _services.GetRequiredService(); var settings = settingsService.GetSetting(Provider, _model); var allowMultiModal = settings != null && settings.MultiModal; var messages = new List(); var temperature = float.Parse(state.GetState("temperature", "0.0")); var maxTokens = int.TryParse(state.GetState("max_tokens"), out var tokens) ? tokens : agent.LlmConfig?.MaxOutputTokens ?? LlmConstant.DEFAULT_MAX_OUTPUT_TOKEN; var options = new ChatCompletionOptions() { ToolChoice = ChatToolChoice.CreateAutoChoice(), Temperature = temperature, MaxOutputTokenCount = maxTokens }; var functions = agent.Functions.Concat(agent.SecondaryFunctions ?? []); foreach (var function in functions) { if (!agentService.RenderFunction(agent, function)) continue; var property = agentService.RenderFunctionProperty(agent, function); options.Tools.Add(ChatTool.CreateFunctionTool( functionName: function.Name, functionDescription: function.Description, functionParameters: BinaryData.FromObjectAsJson(property))); } if (!string.IsNullOrEmpty(agent.Instruction) || !agent.SecondaryInstructions.IsNullOrEmpty()) { var text = agentService.RenderedInstruction(agent); messages.Add(new SystemChatMessage(text)); } if (!string.IsNullOrEmpty(agent.Knowledges)) { messages.Add(new SystemChatMessage(agent.Knowledges)); } var samples = ProviderHelper.GetChatSamples(agent.Samples); foreach (var sample in samples) { messages.Add(sample.Role == AgentRole.User ? new UserChatMessage(sample.Content) : new AssistantChatMessage(sample.Content)); } var filteredMessages = conversations.Select(x => x).ToList(); var firstUserMsgIdx = filteredMessages.FindIndex(x => x.Role == AgentRole.User); if (firstUserMsgIdx > 0) { filteredMessages = filteredMessages.Where((_, idx) => idx >= firstUserMsgIdx).ToList(); } foreach (var message in filteredMessages) { if (message.Role == AgentRole.Function) { messages.Add(new AssistantChatMessage(new List { ChatToolCall.CreateFunctionToolCall(message.ToolCallId, message.FunctionName, BinaryData.FromString(message.FunctionArgs ?? string.Empty)) })); messages.Add(new ToolChatMessage(message.ToolCallId, message.Content)); } else if (message.Role == AgentRole.User) { var text = !string.IsNullOrWhiteSpace(message.Payload) ? message.Payload : message.Content; var textPart = ChatMessageContentPart.CreateTextPart(text); var contentParts = new List { textPart }; if (allowMultiModal && !message.Files.IsNullOrEmpty()) { foreach (var file in message.Files) { if (!string.IsNullOrEmpty(file.FileData)) { var (contentType, bytes) = FileUtility.GetFileInfoFromData(file.FileData); var contentPart = ChatMessageContentPart.CreateImagePart(BinaryData.FromBytes(bytes), contentType, ChatImageDetailLevel.Auto); contentParts.Add(contentPart); } else if (!string.IsNullOrEmpty(file.FileStorageUrl)) { var contentType = FileUtility.GetFileContentType(file.FileStorageUrl); var bytes = fileStorage.GetFileBytes(file.FileStorageUrl); var contentPart = ChatMessageContentPart.CreateImagePart(BinaryData.FromBytes(bytes), contentType, ChatImageDetailLevel.Auto); contentParts.Add(contentPart); } else if (!string.IsNullOrEmpty(file.FileUrl)) { var uri = new Uri(file.FileUrl); var contentPart = ChatMessageContentPart.CreateImagePart(uri, ChatImageDetailLevel.Auto); contentParts.Add(contentPart); } } } messages.Add(new UserChatMessage(contentParts) { ParticipantName = message.FunctionName }); } else if (message.Role == AgentRole.Assistant) { messages.Add(new AssistantChatMessage(message.Content)); } } var prompt = GetPrompt(messages, options); return (prompt, messages, options); } private string GetPrompt(IEnumerable messages, ChatCompletionOptions options) { var prompt = string.Empty; if (!messages.IsNullOrEmpty()) { // System instruction var verbose = string.Join("\r\n", messages .Select(x => x as SystemChatMessage) .Where(x => x != null) .Select(x => { if (!string.IsNullOrEmpty(x.ParticipantName)) { // To display Agent name in log return $"[{x.ParticipantName}]: {x.Content.FirstOrDefault()?.Text ?? string.Empty}"; } return $"{AgentRole.System}: {x.Content.FirstOrDefault()?.Text ?? string.Empty}"; })); prompt += $"{verbose}\r\n"; verbose = string.Join("\r\n", messages .Where(x => x as SystemChatMessage == null) .Select(x => { var fnMessage = x as ToolChatMessage; if (fnMessage != null) { return $"{AgentRole.Function}: {fnMessage.Content.FirstOrDefault()?.Text ?? string.Empty}"; } var userMessage = x as UserChatMessage; if (userMessage != null) { var content = x.Content.FirstOrDefault()?.Text ?? string.Empty; return !string.IsNullOrEmpty(userMessage.ParticipantName) && userMessage.ParticipantName != "route_to_agent" ? $"{userMessage.ParticipantName}: {content}" : $"{AgentRole.User}: {content}"; } var assistMessage = x as AssistantChatMessage; if (assistMessage != null) { var toolCall = assistMessage.ToolCalls?.FirstOrDefault(); return toolCall != null ? $"{AgentRole.Assistant}: Call function {toolCall?.FunctionName}({toolCall?.FunctionArguments})" : $"{AgentRole.Assistant}: {assistMessage.Content.FirstOrDefault()?.Text ?? string.Empty}"; } return string.Empty; })); if (!string.IsNullOrEmpty(verbose)) { prompt += $"\r\n[CONVERSATION]\r\n{verbose}\r\n"; } } if (!options.Tools.IsNullOrEmpty()) { var functions = string.Join("\r\n", options.Tools.Select(fn => { return $"\r\n{fn.FunctionName}: {fn.FunctionDescription}\r\n{fn.FunctionParameters}"; })); prompt += $"\r\n[FUNCTIONS]{functions}\r\n"; } return prompt; } #endregion }