diff --git a/src/Infrastructure/BotSharp.Abstraction/Agents/Models/AgentLlmConfig.cs b/src/Infrastructure/BotSharp.Abstraction/Agents/Models/AgentLlmConfig.cs index b069de92b..4684b4b90 100644 --- a/src/Infrastructure/BotSharp.Abstraction/Agents/Models/AgentLlmConfig.cs +++ b/src/Infrastructure/BotSharp.Abstraction/Agents/Models/AgentLlmConfig.cs @@ -80,6 +80,16 @@ public AgentLlmConfig(AgentTemplateLlmConfig templateLlmConfig) [JsonPropertyName("realtime")] [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] public LlmRealtimeConfig? Realtime { get; set; } + + /// + /// Live config, kept apart from because the two are different kinds + /// of voice session: a live model runs the conversation itself and delegates the thinking to + /// a backend model, so its provider and model are not interchangeable with a realtime one. + /// Setting this is what makes a conversation a live session. + /// + [JsonPropertyName("live")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LlmLiveConfig? Live { get; set; } } public class LlmImageCompositionConfig : LlmProviderModel @@ -92,4 +102,8 @@ public class LlmAudioTranscriptionConfig : LlmProviderModel public class LlmRealtimeConfig : LlmConfigBase { +} + +public class LlmLiveConfig : LlmConfigBase +{ } \ No newline at end of file diff --git a/src/Infrastructure/BotSharp.Abstraction/MLTasks/ILiveCompletion.cs b/src/Infrastructure/BotSharp.Abstraction/MLTasks/ILiveCompletion.cs new file mode 100644 index 000000000..d0d317596 --- /dev/null +++ b/src/Infrastructure/BotSharp.Abstraction/MLTasks/ILiveCompletion.cs @@ -0,0 +1,16 @@ +namespace BotSharp.Abstraction.MLTasks; + +/// +/// A full duplex voice model that owns the conversation itself: it listens and speaks at the +/// same time, arbitrates interruptions internally, and delegates reasoning and tool calls to a +/// separate backend model. +/// +/// It drives the same hub as and is used wherever a realtime +/// completer is expected, so the interface adds nothing to that contract. It exists to keep the +/// two families apart in DI: a live provider is never returned by +/// GetServices<IRealTimeCompletion>(), so it is reachable only through the agent's +/// live config, never as a realtime fallback. +/// +public interface ILiveCompletion : IRealTimeCompletion +{ +} diff --git a/src/Infrastructure/BotSharp.Abstraction/MLTasks/Settings/LlmModelSetting.cs b/src/Infrastructure/BotSharp.Abstraction/MLTasks/Settings/LlmModelSetting.cs index 815b3b44a..111442367 100644 --- a/src/Infrastructure/BotSharp.Abstraction/MLTasks/Settings/LlmModelSetting.cs +++ b/src/Infrastructure/BotSharp.Abstraction/MLTasks/Settings/LlmModelSetting.cs @@ -185,7 +185,13 @@ public enum LlmModelType Embedding = 4, Audio = 5, Realtime = 6, - Web = 7 + Web = 7, + + /// + /// Full duplex voice model that runs the conversation itself and delegates the thinking to + /// a backend model. Kept apart from : the two are not interchangeable. + /// + Live = 8 } public enum LlmModelCapability @@ -203,5 +209,6 @@ public enum LlmModelCapability AudioGeneration = 10, Realtime = 11, WebSearch = 12, - PdfReading = 13 + PdfReading = 13, + Live = 14 } diff --git a/src/Infrastructure/BotSharp.Abstraction/Realtime/IRealtimeHook.cs b/src/Infrastructure/BotSharp.Abstraction/Realtime/IRealtimeHook.cs index b4fed5143..ff9b7147f 100644 --- a/src/Infrastructure/BotSharp.Abstraction/Realtime/IRealtimeHook.cs +++ b/src/Infrastructure/BotSharp.Abstraction/Realtime/IRealtimeHook.cs @@ -6,6 +6,20 @@ namespace BotSharp.Abstraction.Realtime; public interface IRealtimeHook : IHookBase { + /// + /// Establishes the ambient identity for a stream connection, before anything runs on it. + /// + /// Synchronous on purpose. Ambient identity is AsyncLocal-backed, and a write made inside an + /// `async` method does not flow back out to its caller - so an awaited hook cannot establish + /// identity for the connection that follows it, only for itself. + /// + /// + /// Who the client says it is. A browser cannot set headers on a WebSocket handshake, so this + /// arrives as a query parameter and is UNVERIFIED: treat it as a claim to be checked, not as + /// proof. Null when the client sent none. + /// + void OnAuthenticate(string agentId, string conversationId, string? userId) { } + Task OnModelReady(Agent agent, IRealTimeCompletion completer) => Task.CompletedTask; diff --git a/src/Infrastructure/BotSharp.Abstraction/Realtime/Models/RealtimeHubConnection.cs b/src/Infrastructure/BotSharp.Abstraction/Realtime/Models/RealtimeHubConnection.cs index c0dac6d5b..08f16e63b 100644 --- a/src/Infrastructure/BotSharp.Abstraction/Realtime/Models/RealtimeHubConnection.cs +++ b/src/Infrastructure/BotSharp.Abstraction/Realtime/Models/RealtimeHubConnection.cs @@ -16,6 +16,18 @@ public class RealtimeHubConnection public Func OnModelUserInterrupted { get; set; } = null!; public Func OnUserSpeechDetected { get; set; } = () => string.Empty; + /// + /// Serializes a partial transcript (role, delta) for the user stream. Display only: a turn + /// still reaches conversation storage only once it is complete. + /// + public Func? OnModelTranscriptDelta { get; set; } + + /// + /// Set by the hub so a provider can push display-only events straight to the user stream, + /// without going through the conversation hooks that write history. + /// + public Func? SendEventToUser { get; set; } + public void ResetResponseState() { LastAssistantItemId = null; diff --git a/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs b/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs index d5182f9b7..b47ae4205 100644 --- a/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs +++ b/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs @@ -52,12 +52,15 @@ public async Task ConnectToModel( routing.Context.SetDialogs(dialogs); routing.Context.SetMessageId(_conn.ConversationId, Guid.Empty.ToString()); - var (provider, model) = GetLlmProviderModel(agent); + var (provider, model, isLive) = GetLlmProviderModel(agent); - _completer = _services.GetServices().First(x => x.Provider == provider); + _completer = ResolveCompleter(provider, isLive); _completer.SetModelName(model); _completer.SetOptions(options); + // Lets a provider stream display-only events, such as partial transcripts. + _conn.SendEventToUser = responseToUser; + await _completer.Connect( conn: _conn, onModelReady: async () => @@ -189,27 +192,54 @@ public RealtimeHubConnection SetHubConnection(string conversationId) return _conn; } - private (string, string) GetLlmProviderModel(Agent agent) + /// + /// Picks the voice model for this conversation, and which family it belongs to. + /// + /// The agent's live config wins when it is set, because a live session is a deliberate + /// choice rather than a variation of a realtime one. Everything else resolves to realtime: + /// RealtimeModelSettings describes a realtime provider, so a live provider is never reached + /// by fallback, only by an agent asking for it. + /// + private (string provider, string model, bool isLive) GetLlmProviderModel(Agent agent) { - var provider = agent?.LlmConfig?.Realtime?.Provider; - var model = agent?.LlmConfig?.Realtime?.Model; - var settingService = _services.GetRequiredService(); + var live = agent?.LlmConfig?.Live; + if (live?.IsValid == true) + { + return (live.Provider!, live.Model!, true); + } - if (!string.IsNullOrEmpty(provider) && !string.IsNullOrEmpty(model)) + var realtime = agent?.LlmConfig?.Realtime; + if (realtime?.IsValid == true) { - return (provider, model); + return (realtime.Provider!, realtime.Model!, false); } - provider = _settings.Provider; - model = _settings.Model; + if (!string.IsNullOrEmpty(_settings.Provider) && !string.IsNullOrEmpty(_settings.Model)) + { + return (_settings.Provider, _settings.Model, false); + } + + return ("openai", "gpt-realtime", false); + } + + /// + /// Live providers register as only, so the two families are + /// resolved from separate DI lists and a provider name is looked up in one of them, never + /// both. A live completer still drives the hub as a realtime one - the interface derives + /// from - so nothing downstream has to know the difference. + /// + private IRealTimeCompletion ResolveCompleter(string provider, bool isLive) + { + var completer = isLive + ? _services.GetServices().FirstOrDefault(x => x.Provider == provider) + : _services.GetServices().FirstOrDefault(x => x.Provider == provider); - if (!string.IsNullOrEmpty(provider) && !string.IsNullOrEmpty(model)) + if (completer == null) { - return (provider, model); + var family = isLive ? "live" : "realtime"; + throw new InvalidOperationException($"Can't resolve the {family} completion provider by \"{provider}\"."); } - provider = "openai"; - model = "gpt-realtime"; - return (provider, model); + return completer; } } diff --git a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs index 95ed84594..a5db2031c 100644 --- a/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs +++ b/src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs @@ -38,16 +38,16 @@ public async Task Execute(RoleDialogModel message) await Task.Delay(1500); -#if DEBUG - var intermediateMsg = RoleDialogModel.From(message, AgentRole.Assistant, $"Here is your weather in {args?.City}"); - messageHub.Push(new() - { - EventName = ChatEvent.OnIntermediateMessageReceivedFromAssistant, - Data = intermediateMsg, - RefId = conv.ConversationId, - SaveDataToDb = true - }); -#endif +//#if DEBUG +// var intermediateMsg = RoleDialogModel.From(message, AgentRole.Assistant, $"Here is your weather in {args?.City}"); +// messageHub.Push(new() +// { +// EventName = ChatEvent.OnIntermediateMessageReceivedFromAssistant, +// Data = intermediateMsg, +// RefId = conv.ConversationId, +// SaveDataToDb = true +// }); +//#endif message.Indication = $"Still working on it... Hold on, {args?.City}"; messageHub.Push(new() @@ -62,20 +62,20 @@ public async Task Execute(RoleDialogModel message) message.Content = $"It is a sunny day!"; message.StopCompletion = false; -#if DEBUG - var sidecar = _services.GetService(); - if (sidecar != null) - { - var text = $"I want to know fun events in {args?.City}"; - var states = new List - { - new() { Key = StateConst.CHANNEL, Value = ConversationChannel.Email } - }; - - var msg = await sidecar.SendMessage(message.CurrentAgentId, text, states: states); - message.Content = $"{message.Content} {msg.Content}"; - } -#endif +//#if DEBUG +// var sidecar = _services.GetService(); +// if (sidecar != null) +// { +// var text = $"I want to know fun events in {args?.City}"; +// var states = new List +// { +// new() { Key = StateConst.CHANNEL, Value = ConversationChannel.Email } +// }; + +// var msg = await sidecar.SendMessage(message.CurrentAgentId, text, states: states); +// message.Content = $"{message.Content} {msg.Content}"; +// } +//#endif return true; } diff --git a/src/Infrastructure/BotSharp.Core/Infrastructures/CompletionProvider.cs b/src/Infrastructure/BotSharp.Core/Infrastructures/CompletionProvider.cs index 6373d99c6..58018ca4a 100644 --- a/src/Infrastructure/BotSharp.Core/Infrastructures/CompletionProvider.cs +++ b/src/Infrastructure/BotSharp.Core/Infrastructures/CompletionProvider.cs @@ -40,6 +40,10 @@ public static object GetCompletion( { return GetRealTimeCompletion(services, provider: provider, model: model); } + else if (settings.Type == LlmModelType.Live) + { + return GetLiveCompletion(services, provider: provider, model: model); + } else { return GetChatCompletion(services, provider: provider, model: model, agentConfig: agentConfig); @@ -205,6 +209,35 @@ public static IRealTimeCompletion GetRealTimeCompletion( return completer; } + /// + /// Live providers are registered as only, so they are resolved + /// from their own list: a provider name is never shared between the two voice families here. + /// + public static ILiveCompletion GetLiveCompletion( + IServiceProvider services, + string? provider = null, + string? model = null, + string? modelId = null, + bool? multiModal = null, + AgentLlmConfig? agentConfig = null) + { + var completions = services.GetServices(); + (provider, model) = GetProviderAndModel(services, provider: provider, model: model, modelId: modelId, + multiModal: multiModal, + modelType: LlmModelType.Live, + agentConfig: agentConfig); + + var completer = completions.FirstOrDefault(x => x.Provider == provider); + if (completer == null) + { + var logger = services.GetRequiredService>(); + logger.LogError($"Can't resolve live completion provider by {provider}"); + } + + completer?.SetModelName(model); + return completer; + } + private static (string, string) GetProviderAndModel( IServiceProvider services, string? provider = null, diff --git a/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs b/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs index 91f2db767..93754610e 100644 --- a/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs +++ b/src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs @@ -1,5 +1,7 @@ +using BotSharp.Abstraction.Realtime; using BotSharp.Abstraction.Realtime.Options; using BotSharp.Abstraction.Realtime.Sessions; +using BotSharp.Core.Infrastructures; using BotSharp.Core.Session; using Microsoft.AspNetCore.Http; using System.Net.WebSockets; @@ -34,6 +36,23 @@ public async Task Invoke(HttpContext httpContext) var agentId = segments[segments.Length - 2]; var conversationId = segments[segments.Length - 1]; + /* + * Identity, before anything else touches the connection. + * + * The handshake carries no Authorization header - a browser cannot set one on + * a WebSocket - so the signed-in user arrives as a query parameter. It stays + * a query parameter rather than a path segment because agentId and + * conversationId are read from the END of the path above, so an extra segment + * would silently shift both. + * + * Emitted through the SYNCHRONOUS Emit overload: ambient identity is + * AsyncLocal-backed and a write inside an async method does not escape it, so + * an awaited hook would establish nothing for the work awaited below. + */ + var userId = httpContext.Request.Query["user-id"].FirstOrDefault(); + HookEmitter.Emit(services, + hook => hook.OnAuthenticate(agentId, conversationId, userId), agentId); + using var webSocket = await httpContext.WebSockets.AcceptWebSocketAsync(); await HandleWebSocket(services, agentId, conversationId, webSocket); } @@ -96,11 +115,29 @@ private async Task HandleWebSocket(IServiceProvider services, string agentId, st else if (eventType == "disconnect") { _logger.LogDebug($"Disconnecting chat stream connection for conversation ({conversationId})"); - await hub.Completer.Disconnect(); break; } } + /* + * Runs however the loop ended, not just on an explicit disconnect. The browser sends + * "disconnect" and closes the socket in the same breath, so the close often wins the + * race; closing the tab sends nothing at all. Either way the loop simply ends, and + * leaving the model session undisconnected stranded the final turn in the provider's + * buffer with nothing left to flush it. + */ + if (hub.Completer != null) + { + try + { + await hub.Completer.Disconnect(); + } + catch (Exception ex) + { + _logger.LogError(ex, $"Error when disconnecting the model for conversation ({conversationId})"); + } + } + await convService.SaveStates(); await session.DisconnectAsync(); } @@ -160,6 +197,13 @@ private void InitEvents(RealtimeHubConnection conn) mark = new { name = "responsePart" } }); + conn.OnModelTranscriptDelta = (role, delta) => + JsonSerializer.Serialize(new + { + @event = "transcript", + transcript = new { role, delta } + }); + conn.OnModelUserInterrupted = () => JsonSerializer.Serialize(new { diff --git a/src/Plugins/BotSharp.Plugin.MongoStorage/Models/AgentLlmConfigMongoModel.cs b/src/Plugins/BotSharp.Plugin.MongoStorage/Models/AgentLlmConfigMongoModel.cs index 465740343..c984c2cdd 100644 --- a/src/Plugins/BotSharp.Plugin.MongoStorage/Models/AgentLlmConfigMongoModel.cs +++ b/src/Plugins/BotSharp.Plugin.MongoStorage/Models/AgentLlmConfigMongoModel.cs @@ -15,6 +15,7 @@ public class AgentLlmConfigMongoModel public LlmImageCompositionConfigMongoModel? ImageComposition { get; set; } public LlmAudioTranscriptionConfigMongoModel? AudioTranscription { get; set; } public LlmRealtimeConfigMongoModel? Realtime { get; set; } + public LlmLiveConfigMongoModel? Live { get; set; } @@ -35,7 +36,8 @@ public class AgentLlmConfigMongoModel ReasoningEffortLevel = config.ReasoningEffortLevel, ImageComposition = LlmImageCompositionConfigMongoModel.ToMongoModel(config.ImageComposition), AudioTranscription = LlmAudioTranscriptionConfigMongoModel.ToMongoModel(config.AudioTranscription), - Realtime = LlmRealtimeConfigMongoModel.ToMongoModel(config.Realtime) + Realtime = LlmRealtimeConfigMongoModel.ToMongoModel(config.Realtime), + Live = LlmLiveConfigMongoModel.ToMongoModel(config.Live) }; } @@ -56,7 +58,8 @@ public class AgentLlmConfigMongoModel ReasoningEffortLevel = config.ReasoningEffortLevel, ImageComposition = LlmImageCompositionConfigMongoModel.ToDomainModel(config.ImageComposition), AudioTranscription = LlmAudioTranscriptionConfigMongoModel.ToDomainModel(config.AudioTranscription), - Realtime = LlmRealtimeConfigMongoModel.ToDomainModel(config.Realtime) + Realtime = LlmRealtimeConfigMongoModel.ToDomainModel(config.Realtime), + Live = LlmLiveConfigMongoModel.ToDomainModel(config.Live) }; } } @@ -159,4 +162,40 @@ public class LlmRealtimeConfigMongoModel : LlmProviderModelMongoModel ReasoningEffortLevel = config.ReasoningEffortLevel }; } -} \ No newline at end of file +} + +[BsonIgnoreExtraElements(Inherited = true)] +public class LlmLiveConfigMongoModel : LlmProviderModelMongoModel +{ + public string? ReasoningEffortLevel { get; set; } + + public static LlmLiveConfig? ToDomainModel(LlmLiveConfigMongoModel? config) + { + if (config == null) + { + return null; + } + + return new LlmLiveConfig + { + Provider = config.Provider, + Model = config.Model, + ReasoningEffortLevel = config.ReasoningEffortLevel + }; + } + + public static LlmLiveConfigMongoModel? ToMongoModel(LlmLiveConfig? config) + { + if (config == null) + { + return null; + } + + return new LlmLiveConfigMongoModel + { + Provider = config.Provider, + Model = config.Model, + ReasoningEffortLevel = config.ReasoningEffortLevel + }; + } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveEventTypes.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveEventTypes.cs new file mode 100644 index 000000000..cd394d026 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveEventTypes.cs @@ -0,0 +1,117 @@ +namespace BotSharp.Plugin.OpenAI.Models.Live; + +/// +/// Event type strings sent by the client to the Live endpoint. +/// Reference to https://developers.openai.com/api/docs/guides/voice-websockets +/// +public static class LiveClientEventType +{ + public const string SessionStart = "session.start"; + public const string SessionUpdate = "session.update"; + public const string SessionClose = "session.close"; + public const string InputAudioAppend = "session.input_audio.append"; + public const string InputAudioMute = "session.input_audio.mute"; + public const string InputAudioUnmute = "session.input_audio.unmute"; + + /// + /// Trusted behavior instructions injected into the live conversation. + /// + public const string InstructionsAppend = "session.instructions.append"; + + /// + /// Factual context the model may use but does not speak immediately. + /// + public const string ThinkingAppend = "session.thinking.append"; + + /// + /// Content the model paraphrases and speaks aloud. + /// + public const string CommentaryAppend = "session.commentary.append"; + + /// + /// Appends an item (function_call_output or typed message) to the backend Responses turn. + /// + public const string ResponseItemCreate = "response.item.create"; + + /// + /// Asks the backend handler to continue once all required tool results are submitted. + /// + public const string ResponseCreate = "response.create"; +} + +/// +/// Event type strings received from the Live endpoint. +/// +public static class LiveServerEventType +{ + public const string SessionStarted = "session.started"; + public const string SessionUpdated = "session.updated"; + public const string SessionClosed = "session.closed"; + public const string UsageUpdated = "session.usage.updated"; + public const string Error = "error"; + + public const string OutputAudioDelta = "session.output_audio.delta"; + public const string InputTranscriptDelta = "session.input_transcript.delta"; + public const string OutputTranscriptDelta = "session.output_transcript.delta"; + + public const string DelegationCreated = "session.delegation.created"; + + public const string InputAudioMuted = "session.input_audio.muted"; + public const string InputAudioUnmuted = "session.input_audio.unmuted"; + + public const string InstructionsAppended = "session.instructions.appended"; + public const string ThinkingAppended = "session.thinking.appended"; + public const string CommentaryAppended = "session.commentary.appended"; + + /// + /// Envelope carrying a nested Responses API event from the backend handler. + /// + public const string ResponseEvent = "response.event"; +} + +/// +/// Nested event types carried inside a envelope. +/// +public static class LiveResponseInnerEventType +{ + public const string OutputItemDone = "response.output_item.done"; + public const string OutputTextDelta = "response.output_text.delta"; + public const string Completed = "response.completed"; + public const string Failed = "response.failed"; + public const string Incomplete = "response.incomplete"; +} + +/// +/// Output item types carried by . +/// +public static class LiveResponseItemType +{ + public const string FunctionCall = "function_call"; + public const string Message = "message"; +} + +public static class LiveResponseItemStatus +{ + public const string Completed = "completed"; +} + +/// +/// Stage a backend message item belongs to. +/// +public static class LiveResponsePhase +{ + /// + /// The answer the backend handler settled on, and the only message worth recording: + /// earlier phases are drafts and working notes the caller never hears. + /// + public const string FinalAnswer = "final_answer"; +} + +public static class LiveDelegationType +{ + /// + /// Live calls the backend model itself. The only mode this provider supports; "client", + /// where the application answers delegations on its own, is not implemented. + /// + public const string Responses = "responses"; +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveModelConstants.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveModelConstants.cs new file mode 100644 index 000000000..b195f3e03 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveModelConstants.cs @@ -0,0 +1,11 @@ +namespace BotSharp.Plugin.OpenAI.Models.Live; + +public static class LiveModelConstants +{ + /// + /// Full duplex voice model served from the v1/live/sessions endpoint. + /// + public const string GPT_Live_1 = "gpt-live-1"; + + public const string Endpoint = "wss://api.openai.com/v1/live/sessions"; +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LivePromptConstants.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LivePromptConstants.cs new file mode 100644 index 000000000..d7db67fdb --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LivePromptConstants.cs @@ -0,0 +1,74 @@ +namespace BotSharp.Plugin.OpenAI.Models.Live; + +public static class LivePromptConstants +{ + /// + /// Agent channel whose instruction becomes the voice model prompt. An agent that has no + /// instruction for this channel falls back to . + /// + public const string VoiceInstructionChannel = "live"; + + /// + /// Spoken to open the call when the agent has no ".welcome" template of its own. Short and + /// channel neutral on purpose: it is heard, not read, and it is handed to the model as + /// commentary, so it is paraphrased rather than recited. + /// + public const string DefaultGreeting = "Hello! How can I help you today?"; + + /// + /// Default prompt for the voice model, following the structure OpenAI recommends in + /// https://developers.openai.com/api/docs/guides/live-prompting + /// + /// This governs how the conversation sounds and when work is handed to the backend. It is + /// deliberately short: reasoning, business rules and procedures belong in the backend prompt, + /// which is where the BotSharp agent instruction goes. + /// + /// The prompting guide leaves what to do *during* backend work to the product, so the + /// "while the backend is working" rules below are ours rather than OpenAI's. + /// + /// Override per agent with a "live" channel instruction (instruction.live.liquid). + /// + public const string DefaultVoiceInstruction = + """ + You are a calm, friendly voice assistant. Speak warmly and naturally, at an unhurried pace. + + # Backchannel policy + Use moderate backchannels. Acknowledge naturally without competing with the main response. + + # Interruption policy + Stop speaking when the user interrupts. Listen to what they say. + + # Delegation policy + The backend holds the tools, the records and the business knowledge. It does the thinking. + + Delegate to the backend when: + - The request needs a backend capability. + - The answer depends on information you do not already have in this conversation. + - The user asks you to do something rather than just talk. + + Do not delegate when: + - You can answer from the conversation or from a result the backend already gave you. + - The user is making small talk, or clarifying something you just said. + + Delegate before giving an answer that depends on backend work. Do not guess the result + while waiting. + + # While the backend is working + Keep the conversation going. The work happens behind you, not instead of you. + + - Acknowledge the request in a few words, then stay present. Do not fall silent. + - Never mention tools, systems, lookups or "the backend". From where the user sits you + are simply thinking about it. + - Go on answering whatever you already know: small talk, a clarifying question, something + said earlier in the call. Only the delegated answer has to wait. + - Keep listening while you wait. Anything the user adds may change the request, and you + can pass it on. + - If it really is taking a while, say so once. Repeating filler is worse than a short pause. + - When the result arrives, simply give it. Do not announce that it arrived, and do not + recap what you were doing. + + # Response style + Keep replies short and easy to listen to. Ask one question at a time. + Never read out identifiers, code or raw data unless the user asks for them. + """; +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveResponseEvent.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveResponseEvent.cs new file mode 100644 index 000000000..2359bf571 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveResponseEvent.cs @@ -0,0 +1,189 @@ +namespace BotSharp.Plugin.OpenAI.Models.Live; + +/// +/// response.event - envelope wrapping a Responses API lifecycle event produced by the +/// backend handler. Dispatch on Event.Type and keep the outer DelegationId. +/// +public class LiveResponseEventEnvelope : LiveServerEvent +{ + [JsonPropertyName("delegation_id")] + public string? DelegationId { get; set; } + + [JsonPropertyName("event")] + public LiveInnerResponseEvent? Event { get; set; } +} + +public class LiveInnerResponseEvent +{ + [JsonPropertyName("type")] + public string Type { get; set; } = null!; + + [JsonPropertyName("sequence_number")] + public int? SequenceNumber { get; set; } + + [JsonPropertyName("item_id")] + public string? ItemId { get; set; } + + [JsonPropertyName("output_index")] + public int? OutputIndex { get; set; } + + [JsonPropertyName("content_index")] + public int? ContentIndex { get; set; } + + [JsonPropertyName("delta")] + public string? Delta { get; set; } + + [JsonPropertyName("item")] + public LiveResponseOutputItem? Item { get; set; } + + [JsonPropertyName("response")] + public LiveResponseBody? Response { get; set; } + + #region Flattened function call fields + // response.output_item.done reports the finished tool call either nested under "item" + // or flattened onto the event itself, so both shapes are read. + [JsonPropertyName("call_id")] + public string? CallId { get; set; } + + [JsonPropertyName("name")] + public string? Name { get; set; } + + [JsonPropertyName("arguments")] + public string? Arguments { get; set; } + #endregion + + public LiveResponseOutputItem? ResolveItem() + { + if (Item != null) + { + return Item; + } + + if (string.IsNullOrEmpty(CallId) && string.IsNullOrEmpty(Name)) + { + return null; + } + + return new LiveResponseOutputItem + { + Id = ItemId, + Type = LiveResponseItemType.FunctionCall, + CallId = CallId, + Name = Name, + Arguments = Arguments + }; + } +} + +public class LiveResponseBody +{ + [JsonPropertyName("id")] + public string? Id { get; set; } + + [JsonPropertyName("status")] + public string? Status { get; set; } + + [JsonPropertyName("usage")] + public LiveResponseUsage? Usage { get; set; } +} + +public class LiveResponseUsage +{ + [JsonPropertyName("input_tokens")] + public long InputTokens { get; set; } + + [JsonPropertyName("output_tokens")] + public long OutputTokens { get; set; } + + [JsonPropertyName("total_tokens")] + public long TotalTokens { get; set; } + + [JsonPropertyName("input_tokens_details")] + public LiveResponseInputTokenDetail? InputTokenDetails { get; set; } + + [JsonPropertyName("output_tokens_details")] + public LiveResponseOutputTokenDetail? OutputTokenDetails { get; set; } +} + +public class LiveResponseInputTokenDetail +{ + [JsonPropertyName("cached_tokens")] + public long CachedTokens { get; set; } +} + +public class LiveResponseOutputTokenDetail +{ + [JsonPropertyName("reasoning_tokens")] + public long ReasoningTokens { get; set; } +} + +public class LiveResponseOutputItem +{ + [JsonPropertyName("id")] + public string? Id { get; set; } + + /// + /// function_call, message, reasoning, etc. + /// + [JsonPropertyName("type")] + public string? Type { get; set; } + + [JsonPropertyName("status")] + public string? Status { get; set; } + + [JsonPropertyName("role")] + public string? Role { get; set; } + + [JsonPropertyName("name")] + public string? Name { get; set; } + + [JsonPropertyName("call_id")] + public string? CallId { get; set; } + + [JsonPropertyName("arguments")] + public string? Arguments { get; set; } + + /// + /// Where a message item sits in the backend turn: final_answer is the answer the voice + /// model is given to speak, earlier phases are working notes on the way to it. + /// + [JsonPropertyName("phase")] + public string? Phase { get; set; } + + [JsonPropertyName("content")] + public LiveResponseOutputContent[]? Content { get; set; } + + public bool IsCompleted => string.IsNullOrEmpty(Status) + || Status.Equals(LiveResponseItemStatus.Completed, StringComparison.OrdinalIgnoreCase); + + /// + /// Joins the text parts of a message item. A message carries its text in content[], and + /// the parts are fragments of one utterance rather than alternatives, so they concatenate. + /// + public string? GetOutputText() + { + if (Content == null || Content.Length == 0) + { + return null; + } + + var parts = Content + .Select(x => x.Text ?? x.Transcript) + .Where(x => !string.IsNullOrWhiteSpace(x)); + + var text = string.Join(string.Empty, parts).Trim(); + return string.IsNullOrEmpty(text) ? null : text; + } +} + +public class LiveResponseOutputContent +{ + [JsonPropertyName("type")] + public string? Type { get; set; } + + [JsonPropertyName("text")] + public string? Text { get; set; } + + [JsonPropertyName("transcript")] + public string? Transcript { get; set; } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveServerEvent.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveServerEvent.cs new file mode 100644 index 000000000..3f9cf4d63 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveServerEvent.cs @@ -0,0 +1,141 @@ +namespace BotSharp.Plugin.OpenAI.Models.Live; + +/// +/// Common envelope of every event received from the Live endpoint. +/// +public class LiveServerEvent +{ + [JsonPropertyName("type")] + public string Type { get; set; } = null!; + + [JsonPropertyName("event_id")] + public string? EventId { get; set; } + + /// + /// Echoes the event_id of the client event this one answers. + /// + [JsonPropertyName("client_event_id")] + public string? ClientEventId { get; set; } +} + +public class LiveErrorEvent : LiveServerEvent +{ + [JsonPropertyName("error")] + public LiveErrorBody Body { get; set; } = new(); +} + +public class LiveErrorBody +{ + [JsonPropertyName("type")] + public string? Type { get; set; } + + /// + /// e.g. immutable_field_update. + /// + [JsonPropertyName("code")] + public string? Code { get; set; } + + [JsonPropertyName("message")] + public string? Message { get; set; } + + public override string ToString() => $"{Type}: {Message} ({Code})"; +} + +public class LiveSessionStartedEvent : LiveServerEvent +{ + [JsonPropertyName("session")] + public LiveSessionConfig? Session { get; set; } +} + +public class LiveSessionClosedEvent : LiveServerEvent +{ + /// + /// close_requested, expired, content, remote_hangup or connection_lost. + /// + [JsonPropertyName("reason")] + public string? Reason { get; set; } + + [JsonPropertyName("usage")] + public LiveUsage? Usage { get; set; } +} + +public class LiveUsageUpdatedEvent : LiveServerEvent +{ + [JsonPropertyName("usage")] + public LiveUsage? Usage { get; set; } +} + +/// +/// Live voice sessions are billed by duration rather than tokens. +/// Backend Responses usage is reported separately through the response.event envelope. +/// +public class LiveUsage +{ + [JsonPropertyName("seconds")] + public double Seconds { get; set; } + + [JsonPropertyName("context_window")] + public LiveContextWindowUsage? ContextWindow { get; set; } +} + +public class LiveContextWindowUsage +{ + [JsonPropertyName("usage_ratio")] + public double UsageRatio { get; set; } +} + +/// +/// session.output_audio.delta - a chunk of base64 encoded model audio. +/// +public class LiveOutputAudioDelta : LiveServerEvent +{ + [JsonPropertyName("delta")] + public string? Delta { get; set; } + + [JsonPropertyName("item_id")] + public string? ItemId { get; set; } +} + +/// +/// session.input_transcript.delta and session.output_transcript.delta. +/// Live emits no turn completion marker, and intervals may overlap. +/// +public class LiveTranscriptDelta : LiveServerEvent +{ + [JsonPropertyName("delta")] + public string? Delta { get; set; } + + [JsonPropertyName("start_ms")] + public long? StartMs { get; set; } + + [JsonPropertyName("end_ms")] + public long? EndMs { get; set; } +} + +/// +/// session.delegation.created - raised when delegation type is "client" and the model +/// hands a task to the application backend. +/// +public class LiveDelegationCreatedEvent : LiveServerEvent +{ + [JsonPropertyName("offset_ms")] + public long? OffsetMs { get; set; } + + [JsonPropertyName("delegation")] + public LiveDelegation? Delegation { get; set; } +} + +public class LiveDelegation +{ + [JsonPropertyName("id")] + public string Id { get; set; } = null!; + + [JsonPropertyName("type")] + public string? Type { get; set; } + + /// + /// client or responses. + /// + [JsonPropertyName("target")] + public string? Target { get; set; } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveSessionConfig.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveSessionConfig.cs new file mode 100644 index 000000000..2aeecacff --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Models/Live/LiveSessionConfig.cs @@ -0,0 +1,172 @@ +using BotSharp.Abstraction.Functions.Models; + +namespace BotSharp.Plugin.OpenAI.Models.Live; + +/// +/// Session object sent with session.start, and partially with session.update. +/// Reference to https://developers.openai.com/api/docs/guides/live +/// +public class LiveSessionConfig +{ + /// + /// Required at session.start, immutable afterwards. + /// + [JsonPropertyName("model")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Model { get; set; } + + /// + /// Conversation style guidance, up to 16,384 tokens. Detailed workflows belong in + /// instead. + /// + [JsonPropertyName("instructions")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Instructions { get; set; } + + /// + /// Prior conversation turns to seed the session with. + /// + [JsonPropertyName("input")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveSessionInputItem[]? Input { get; set; } + + [JsonPropertyName("audio")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveAudioConfig? Audio { get; set; } + + [JsonPropertyName("delegation")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveDelegationConfig? Delegation { get; set; } + + /// + /// Enables session recording so it can be downloaded or forked later. + /// + [JsonPropertyName("store")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public bool? Store { get; set; } +} + +public class LiveSessionInputItem +{ + /// + /// developer, user or assistant. + /// + [JsonPropertyName("role")] + public string Role { get; set; } = null!; + + [JsonPropertyName("content")] + public string Content { get; set; } = null!; +} + +public class LiveAudioConfig +{ + /// + /// Wire format of the audio carried over the WebSocket. Negotiated automatically on WebRTC. + /// + [JsonPropertyName("format")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveAudioFormat? Format { get; set; } + + [JsonPropertyName("output")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveOutputAudioConfig? Output { get; set; } +} + +/// +/// audio/pcm at 24000 or 16000, audio/pcmu or audio/pcma at 8000. +/// +public class LiveAudioFormat +{ + [JsonPropertyName("type")] + public string Type { get; set; } = "audio/pcm"; + + [JsonPropertyName("rate")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public int? Rate { get; set; } +} + +public class LiveOutputAudioConfig +{ + [JsonPropertyName("voice")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Voice { get; set; } +} + +/// +/// Decides who runs the reasoning and tool calls behind the voice conversation. +/// The mode cannot be changed on a live session; start a new one instead. +/// +public class LiveDelegationConfig +{ + /// + /// responses or client. + /// + [JsonPropertyName("type")] + public string Type { get; set; } = LiveDelegationType.Responses; + + [JsonPropertyName("responses")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveResponsesDelegationConfig? Responses { get; set; } +} + +/// +/// Backend handler configuration used when delegation type is "responses". +/// This is the only part of the session that session.update may change. +/// +public class LiveResponsesDelegationConfig +{ + /// + /// Backend reasoning model, required when the delegation is created. + /// + [JsonPropertyName("model")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Model { get; set; } + + [JsonPropertyName("instructions")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Instructions { get; set; } + + [JsonPropertyName("tools")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public FunctionDef[]? Tools { get; set; } + + /// + /// auto, required, none or a named function. + /// + [JsonPropertyName("tool_choice")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? ToolChoice { get; set; } + + [JsonPropertyName("parallel_tool_calls")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public bool? ParallelToolCalls { get; set; } + + /// + /// Minimum 16 when set. + /// + [JsonPropertyName("max_output_tokens")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public int? MaxOutputTokens { get; set; } + + /// + /// auto, default, flex, or priority (also accepted as "fast"). + /// Applies to the backend Responses call, not to the billed voice minutes. + /// + [JsonPropertyName("service_tier")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? ServiceTier { get; set; } + + [JsonPropertyName("reasoning")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public LiveReasoningConfig? Reasoning { get; set; } +} + +public class LiveReasoningConfig +{ + /// + /// minimal, low, medium, high or xhigh. + /// + [JsonPropertyName("effort")] + [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)] + public string? Effort { get; set; } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/OpenAiPlugin.cs b/src/Plugins/BotSharp.Plugin.OpenAI/OpenAiPlugin.cs index daade7a94..9bacc99b6 100644 --- a/src/Plugins/BotSharp.Plugin.OpenAI/OpenAiPlugin.cs +++ b/src/Plugins/BotSharp.Plugin.OpenAI/OpenAiPlugin.cs @@ -7,6 +7,7 @@ using BotSharp.Plugin.OpenAI.Providers.Audio; using Microsoft.Extensions.Configuration; using BotSharp.Plugin.OpenAI.Providers.Realtime; +using BotSharp.Plugin.OpenAI.Providers.Live; namespace BotSharp.Plugin.OpenAI; @@ -35,5 +36,6 @@ public void RegisterDI(IServiceCollection services, IConfiguration config) services.AddScoped(); services.AddScoped(); services.AddScoped(); + services.AddScoped(); } } \ No newline at end of file diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/IdleFlushTimer.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/IdleFlushTimer.cs new file mode 100644 index 000000000..5c2395196 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/IdleFlushTimer.cs @@ -0,0 +1,77 @@ +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +/// +/// Fires a callback once a stream of events has been quiet for a while. The Live endpoint +/// never marks the end of a turn, so audio and transcript boundaries are inferred from silence. +/// +internal sealed class IdleFlushTimer : IDisposable +{ + private readonly TimeSpan _idle; + private readonly Func _onIdle; + private readonly ILogger? _logger; + private readonly string _name; + private readonly object _lock = new(); + + private Timer? _timer; + private bool _disposed; + + public IdleFlushTimer(string name, TimeSpan idle, Func onIdle, ILogger? logger = null) + { + _name = name; + _idle = idle; + _onIdle = onIdle; + _logger = logger; + } + + /// + /// Restarts the idle countdown. Call on every delta. + /// + public void Touch() + { + lock (_lock) + { + if (_disposed) return; + + _timer ??= new Timer(OnElapsed, null, Timeout.Infinite, Timeout.Infinite); + _timer.Change(_idle, Timeout.InfiniteTimeSpan); + } + } + + /// + /// Stops the countdown without invoking the callback, e.g. after an explicit flush. + /// + public void Cancel() + { + lock (_lock) + { + if (_disposed) return; + _timer?.Change(Timeout.Infinite, Timeout.Infinite); + } + } + + private void OnElapsed(object? state) => _ = FireAsync(); + + private async Task FireAsync() + { + try + { + await _onIdle(); + } + catch (Exception ex) + { + _logger?.LogError(ex, "Error while flushing idle {Name} buffer.", _name); + } + } + + public void Dispose() + { + lock (_lock) + { + if (_disposed) return; + + _disposed = true; + _timer?.Dispose(); + _timer = null; + } + } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveAppendBudget.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveAppendBudget.cs new file mode 100644 index 000000000..12c551409 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveAppendBudget.cs @@ -0,0 +1,140 @@ +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +/// +/// Keeps session.*.append content inside the server's per-append budget. +/// +/// The limit is stated in tokens, not characters, and the two diverge sharply by script: +/// Latin prose runs around four characters per token, while CJK is closer to one. A plain +/// character cap is therefore far too generous for Chinese, Japanese and Korean, which is +/// exactly where an over-long append would be rejected. +/// +internal static class LiveAppendBudget +{ + /// + /// Upper bound on how many appends one piece of content may be split across, so a runaway + /// answer cannot flood the session. + /// + public const int MaxChunks = 5; + + /// + /// Marks content that had to be dropped, so the model can tell it saw only part of it. + /// + public const string ElisionMarker = " [...]"; + + /// + /// Rough token count. Wide (CJK and full-width) characters are charged a whole token each, + /// everything else a quarter. Deliberately pessimistic: overshooting the limit is rejected + /// by the server, undershooting only costs a little context. + /// + public static int EstimateTokens(string? text) + { + if (string.IsNullOrEmpty(text)) return 0; + + var wide = 0; + var narrow = 0; + + foreach (var ch in text) + { + if (IsWide(ch)) wide++; + else narrow++; + } + + return wide + (narrow + 3) / 4; + } + + /// + /// Splits content on sentence boundaries so a long append is delivered in full instead of + /// being amputated. Appends concatenate on the server, so the pieces reassemble into the + /// original text - which is why this, and not truncation, is the default treatment. + /// Returns whether anything still had to be dropped at the chunk ceiling. + /// + public static (List Chunks, bool Truncated) SplitIntoChunks(string text, int maxTokens, int maxChunks) + { + if (EstimateTokens(text) <= maxTokens) + { + return ([text], false); + } + + var chunks = new List(); + var remaining = text.AsSpan(); + + while (!remaining.IsEmpty && chunks.Count < maxChunks) + { + if (EstimateTokens(remaining.ToString()) <= maxTokens) + { + chunks.Add(remaining.ToString().Trim()); + remaining = default; + break; + } + + var cut = FindCutIndex(remaining, maxTokens); + var boundary = FindSentenceBoundary(remaining, cut); + + chunks.Add(remaining[..boundary].ToString().Trim()); + remaining = remaining[boundary..].TrimStart(); + } + + chunks.RemoveAll(string.IsNullOrWhiteSpace); + return (chunks, !remaining.IsEmpty); + } + + /// + /// Largest prefix length that still fits the budget, never landing between a surrogate pair. + /// + private static int FindCutIndex(ReadOnlySpan text, int maxTokens) + { + var wide = 0; + var narrow = 0; + + for (var i = 0; i < text.Length; i++) + { + if (IsWide(text[i])) wide++; + else narrow++; + + if (wide + (narrow + 3) / 4 > maxTokens) + { + // Stepping back off a low surrogate keeps the pair intact. + return char.IsLowSurrogate(text[i]) ? Math.Max(0, i - 1) : i; + } + } + + return text.Length; + } + + /// + /// Walks back from the budget limit to the last sentence end, then to the last space, + /// so a chunk does not stop mid-word. + /// + private static int FindSentenceBoundary(ReadOnlySpan text, int limit) + { + // Do not give back more than half the chunk chasing a boundary. + var floor = limit / 2; + + for (var i = limit - 1; i > floor; i--) + { + var ch = text[i]; + if (ch is '.' or '!' or '?' or '\n' or '。' or '!' or '?') + { + return i + 1; + } + } + + for (var i = limit - 1; i > floor; i--) + { + if (char.IsWhiteSpace(text[i])) + { + return i + 1; + } + } + + return limit; + } + + private static bool IsWide(char c) + => (c >= 0x1100 && c <= 0x11FF) // Hangul Jamo + || (c >= 0x2E80 && c <= 0xA4CF) // CJK radicals, kana, CJK unified ideographs + || (c >= 0xAC00 && c <= 0xD7AF) // Hangul syllables + || (c >= 0xF900 && c <= 0xFAFF) // CJK compatibility ideographs + || (c >= 0xFF00 && c <= 0xFF60) // Full-width forms + || (c >= 0xFFE0 && c <= 0xFFE6); // Full-width signs +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Events.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Events.cs new file mode 100644 index 000000000..1229fa02f --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Events.cs @@ -0,0 +1,590 @@ +using BotSharp.Abstraction.Realtime.Settings; + +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +public partial class LiveCompletionProvider +{ + private readonly StringBuilder _outputTranscript = new(); + private readonly StringBuilder _inputTranscript = new(); + private readonly object _transcriptLock = new(); + + private IdleFlushTimer? _outputTranscriptTimer; + private IdleFlushTimer? _inputTranscriptTimer; + private IdleFlushTimer? _audioIdleTimer; + + /// + /// Everything that touches the conversation - running a tool, recording a turn - runs here + /// rather than on the receive loop, which has to stay free to drain audio. See + /// for why it is serial. + /// + private SerialWorkQueue? _conversationWork; + + /// + /// Who spoke most recently, so the final flush can emit the turns in the order they happened. + /// + private volatile string? _lastTranscriptRole; + + #region Receive loop + private async Task ReceiveMessage(RealtimeModelSettings realtimeSettings) + { + var session = _session; + if (session == null) return; + + // Captured rather than read from the fields on the way out. A reconnect replaces both + // while this loop is unwinding, and tearing down whatever the fields point at by then + // would close the session that just replaced this one. + var work = _conversationWork; + + await foreach (ChatSessionUpdate update in session.ReceiveUpdatesAsync(CancellationToken.None)) + { + var receivedText = update?.RawResponse; + if (string.IsNullOrEmpty(receivedText)) + { + continue; + } + + LiveServerEvent? response; + try + { + response = JsonSerializer.Deserialize(receivedText); + } + catch (JsonException ex) + { + _logger.LogError(ex, "Unable to parse {Provider} event: {Payload}", Provider, receivedText); + continue; + } + + if (response == null || string.IsNullOrEmpty(response.Type)) + { + continue; + } + + if (await HandleServerEvent(response.Type, receivedText)) + { + break; + } + } + + // The stream is over, so whatever is still buffered is a finished turn. + await FlushPendingTurns(work); + + if (ReferenceEquals(_session, session)) + { + DisposeSessionWorkers(); + } + + session.Dispose(); + } + + /// + /// Returns true when the receive loop should stop. + /// + private async Task HandleServerEvent(string type, string receivedText) + { + switch (type) + { + case LiveServerEventType.Error: + var error = JsonSerializer.Deserialize(receivedText); + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + // Command level failures leave the session usable; only transport errors end it. + return error?.Body?.Type == "server_error"; + + case LiveServerEventType.SessionStarted: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + _isBlocking = false; + await _onModelReady(); + + // After _onModelReady, so the backend handler is configured before the model + // opens its mouth and the greeting cannot arrive ahead of the session setup. + await SendGreeting(); + return false; + + case LiveServerEventType.SessionUpdated: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + return false; + + case LiveServerEventType.SessionClosed: + var closed = JsonSerializer.Deserialize(receivedText); + _logger.Log(TraceLevel, "{Type}: reason {Reason}, billed {Seconds}s", + type, closed?.Reason, closed?.Usage?.Seconds); + // Deliberately not flushed. A turn cut short by the session ending is not a + // completed turn, and its partial text has already been shown live. + return true; + + case LiveServerEventType.UsageUpdated: + var usage = JsonSerializer.Deserialize(receivedText); + _logger.Log(TraceLevel, "{Type}: {Seconds}s, context {Ratio}", + type, usage?.Usage?.Seconds, usage?.Usage?.ContextWindow?.UsageRatio); + return false; + + case LiveServerEventType.OutputAudioDelta: + await OnOutputAudioDelta(receivedText); + return false; + + case LiveServerEventType.OutputTranscriptDelta: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + await OnTranscriptDelta(receivedText, _outputTranscript, _outputTranscriptTimer, AgentRole.Assistant); + return false; + + case LiveServerEventType.InputTranscriptDelta: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + await OnTranscriptDelta(receivedText, _inputTranscript, _inputTranscriptTimer, AgentRole.User); + return false; + + case LiveServerEventType.ResponseEvent: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + OnBackendResponseEvent(receivedText); + return false; + + case LiveServerEventType.DelegationCreated: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + await OnDelegationCreated(receivedText); + return false; + + default: + _logger.Log(TraceLevel, "{Type}: {Payload}", type, receivedText); + return false; + } + } + #endregion + + #region Audio + private async Task OnOutputAudioDelta(string receivedText) + { + var audio = JsonSerializer.Deserialize(receivedText); + if (string.IsNullOrEmpty(audio?.Delta)) + { + return; + } + + await _onModelAudioDeltaReceived(audio.Delta, audio.ItemId ?? string.Empty); + + // Audio going quiet is the end of the model's turn, and OnModelAudioIdle closes it. + // Audio deliberately does not arm the transcript flush any more: on a full duplex session + // that holds the line open with a continuous stream, that kept the turn open forever. + _audioIdleTimer?.Touch(); + } + #endregion + + /// + /// The model's audio stream has gone quiet. + /// + /// That is the truest end-of-turn signal available - it says the model stopped speaking, + /// rather than inferring it from a gap between transcript deltas - so the turn is closed here. + /// Whether it ever runs is the open question: a full duplex connection that holds the line + /// open with a continuous stream would never go idle, and the turn would then be closed by + /// the transcript timer or by the user speaking. Logged so it is clear which one did it. + /// + private async Task OnModelAudioIdle() + { + _logger.LogInformation("{Provider} model audio idle for {Ms}ms; closing the turn.", + Provider, LiveSettings.AudioIdleMs); + + await _onModelAudioResponseDone(); + FlushOutputTranscript(); + } + + #region Transcripts + private async Task OnTranscriptDelta(string receivedText, StringBuilder buffer, IdleFlushTimer? timer, string role) + { + var data = JsonSerializer.Deserialize(receivedText); + if (string.IsNullOrEmpty(data?.Delta)) + { + return; + } + + // Live marks no end of turn, but the speakers take it in turns: the other side starting + // to talk is what closes the one before it. Flushing here rather than on a timer means a + // turn is written the moment it is actually over, and a thinking pause never splits it. + // + // Both flushes no-op on an empty buffer, so a run of deltas from the same speaker costs + // nothing. The flush runs before the append so the turns come out in the order spoken. + if (role == AgentRole.Assistant) + { + FlushInputTranscript(); + } + else + { + FlushOutputTranscript(); + } + + lock (_transcriptLock) + { + buffer.Append(data.Delta); + } + + _lastTranscriptRole = role; + + timer?.Touch(); + + // Show the words as they are spoken. Storage still waits for the turn to complete, so a + // half-spoken utterance never becomes a conversation record. + await PublishTranscriptDelta(role, data.Delta); + } + + /// + /// Pushes a partial transcript straight to the user stream, bypassing the conversation hooks + /// that write history. Awaited rather than fired off, so deltas arrive in the order spoken. + /// + private async Task PublishTranscriptDelta(string role, string delta) + { + var serialize = _conn.OnModelTranscriptDelta; + var send = _conn.SendEventToUser; + + if (serialize == null || send == null) + { + return; + } + + try + { + await send(serialize(role, delta)); + } + catch (Exception ex) + { + // A dead UI socket must not take the voice session down with it. + _logger.LogError(ex, "Failed to stream a {Role} transcript delta to the user.", role); + } + } + + private string TakeBuffer(StringBuilder buffer) + { + lock (_transcriptLock) + { + var text = buffer.ToString().Trim(); + buffer.Clear(); + return text; + } + } + + /// + /// Closes the model's turn. The buffer is taken here, on the caller's thread, so the turn + /// boundary lands where the caller decided it should; recording it is queued, because that + /// reaches conversation storage and hooks and must not hold up the receive loop. + /// + private void FlushOutputTranscript() + { + _outputTranscriptTimer?.Cancel(); + + var text = TakeBuffer(_outputTranscript); + if (string.IsNullOrEmpty(text)) + { + return; + } + + _logger.LogInformation("{Provider} model transcript: {Transcript}", Provider, text); + + // Read now rather than inside the queued work: by the time that runs the model may + // already be speaking its next turn, and this turn would be filed under its item id. + var messageId = _conn.LastAssistantItemId ?? Guid.NewGuid().ToString(); + + _conversationWork?.Enqueue(() => DeliverOutputTranscript(text, messageId)); + } + + private async Task DeliverOutputTranscript(string text, string messageId) + { + await _onModelAudioTranscriptDone(text); + + var message = new RoleDialogModel(AgentRole.Assistant, text) + { + CurrentAgentId = _conn.CurrentAgentId, + MessageId = messageId, + MessageType = MessageTypeName.Plain + }; + + await _onModelResponseDone([message]); + } + + /// + /// Writes out whatever is still buffered when the conversation ends. + /// + /// Alternation closes every turn but the last one: nobody speaks after it, so nothing + /// triggers its flush. Once the stream is over no further delta can arrive, which makes the + /// buffer a complete turn rather than a half-finished one - so this records it instead of + /// discarding it. Safe to call more than once; both flushes no-op on an empty buffer. + /// + /// Queues the turns like any other flush and then drains, so the last words of the call + /// reach storage before the session is torn down. + /// + /// + /// The queue to drain, for a caller that may no longer own the one in the field. + /// + private async Task FlushPendingTurns(SerialWorkQueue? work = null) + { + // Emit in the order they were spoken, so whoever spoke last is written last. Each is + // guarded separately: this runs during teardown, where one failing turn must neither + // take the other down with it nor break the caller's shutdown path. + var flushes = _lastTranscriptRole == AgentRole.Assistant + ? new Action[] { FlushInputTranscript, FlushOutputTranscript } + : [FlushOutputTranscript, FlushInputTranscript]; + + foreach (var flush in flushes) + { + try + { + flush(); + } + catch (Exception ex) + { + _logger.LogError(ex, "Failed to flush the final {Provider} transcript.", Provider); + } + } + + var queue = work ?? _conversationWork; + if (queue != null) + { + await queue.DrainAsync(); + } + } + + /// + /// Closes the caller's turn. Split the same way as : + /// the boundary is taken here, the recording is queued. + /// + private void FlushInputTranscript() + { + _inputTranscriptTimer?.Cancel(); + + var text = TakeBuffer(_inputTranscript); + if (string.IsNullOrEmpty(text)) + { + return; + } + + _logger.LogInformation("{Provider} user transcript: {Transcript}", Provider, text); + + _conversationWork?.Enqueue(() => _onInputAudioTranscriptionDone(new RoleDialogModel(AgentRole.User, text) + { + CurrentAgentId = _conn.CurrentAgentId + })); + } + + /// + /// Reports what a backend turn cost. Live voice time is billed per second rather than per + /// token; token counts come from the backend handler, which is the only thing here that + /// produces any. + /// + /// Filed against the backend model rather than the voice model that owns the session. These + /// tokens were spent by the delegated Responses call, and the stats are priced by looking the + /// reported model up in LlmProviders - so naming the voice model here would charge a + /// reasoning turn at the rates of a model that never saw it. + /// + /// Called from the lifecycle event that carries the usage rather than held in a field until + /// the spoken transcript is flushed. Those two run on different threads, so the field went + /// stale both ways: a second response completing first overwrote the usage nobody had read + /// yet, and a transcript flushed before the backend finished reported zeros. + /// + /// Prompt is left empty deliberately: the loggers that listen on this hook treat a non-empty + /// Prompt as the text sent to the model and write it to the content log and the completion + /// log on their own. A live turn has no such text - the session instruction went out once, at + /// session update - so filling it with the spoken transcript put the assistant's own words in + /// the log a second time, beside the response entry + /// already produced. + /// + private async Task ReportGenerated(LiveResponseUsage? usage, string messageId) + { + var contentHooks = _services.GetHooks(_conn.CurrentAgentId); + foreach (var hook in contentHooks) + { + // The cost belongs to the backend turn, not to anything spoken, so the message + // carries no content: it is here for the agent id the stats are filed under. + await hook.AfterGenerated(new RoleDialogModel(AgentRole.Assistant, string.Empty) + { + CurrentAgentId = _conn.CurrentAgentId, + MessageId = messageId + }, + new TokenStatsModel + { + Provider = _backendProvider.IfNullOrEmptyAs(Provider)!, + Model = _backendModel ?? string.Empty, + Prompt = string.Empty, + TextInputTokens = (usage?.InputTokens ?? 0) - (usage?.InputTokenDetails?.CachedTokens ?? 0), + CachedTextInputTokens = usage?.InputTokenDetails?.CachedTokens ?? 0, + TextOutputTokens = usage?.OutputTokens ?? 0 + }); + } + } + #endregion + + #region Delegation + /// + /// Unwraps a backend Responses event. Completed tool calls are surfaced to the hub so + /// the existing routing pipeline executes them. + /// + private void OnBackendResponseEvent(string receivedText) + { + var envelope = JsonSerializer.Deserialize(receivedText); + var inner = envelope?.Event; + if (inner == null) + { + return; + } + + if (!string.IsNullOrEmpty(envelope?.DelegationId)) + { + _currentDelegationId = envelope.DelegationId; + } + + switch (inner.Type) + { + case LiveResponseInnerEventType.OutputItemDone: + OnBackendOutputItemDone(inner); + return; + + case LiveResponseInnerEventType.Completed: + case LiveResponseInnerEventType.Incomplete: + case LiveResponseInnerEventType.Failed: + _logger.LogInformation("{Type}: backend response {Status}", inner.Type, inner.Response?.Status); + + // A turn that came back incomplete or failed was still billed for what it did + // produce, so all three report. Queued rather than awaited: this reaches the + // stats store, and the receive loop has to stay free to drain audio. + var usage = inner.Response?.Usage; + if (usage != null) + { + var responseId = inner.Response?.Id ?? Guid.NewGuid().ToString(); + _conversationWork?.Enqueue(() => ReportGenerated(usage, responseId)); + } + return; + + default: + _logger.LogDebug("{Type}: {Payload}", inner.Type, receivedText); + return; + } + } + + /// + /// A finished item from the backend turn. Two shapes matter: a function_call, which goes to + /// the existing routing pipeline, and a message, the answer the backend settled on, which is + /// only logged - the assistant turn is still recorded from the spoken transcript. + /// + private void OnBackendOutputItemDone(LiveInnerResponseEvent inner) + { + var item = inner.ResolveItem(); + if (item == null || !item.IsCompleted) + { + return; + } + + var messageId = item.Id ?? inner.ItemId ?? Guid.NewGuid().ToString(); + + switch (item.Type) + { + case LiveResponseItemType.FunctionCall: + if (string.IsNullOrEmpty(item.Name)) + { + return; + } + + _logger.Log(TraceLevel, "{Provider} tool call {Name}({Arguments})", Provider, item.Name, item.Arguments); + + var call = new RoleDialogModel(AgentRole.Assistant, item.Arguments ?? "{}") + { + CurrentAgentId = _conn.CurrentAgentId, + FunctionName = item.Name, + FunctionArgs = item.Arguments, + ToolCallId = item.CallId, + MessageId = messageId, + MessageType = MessageTypeName.FunctionCall + }; + + // Queued, never awaited here. Running the function inline would stop the socket + // being read for as long as it takes - and on a full duplex call the model is + // still speaking through that, so the caller hears the reply cut in half. + _conversationWork?.Enqueue(() => _onModelResponseDone([call])); + return; + + case LiveResponseItemType.Message: + // Logged, not recorded. The conversation still takes its assistant turn from the + // spoken transcript, and this is the answer the voice model is about to speak, + // so storing it here would put the same turn in the history twice. + var text = item.GetOutputText(); + if (string.IsNullOrEmpty(text)) + { + return; + } + + _logger.Log(TraceLevel, "{Provider} backend {Phase} answer: {Text}", + Provider, item.Phase ?? LiveResponsePhase.FinalAnswer, text); + return; + + default: + return; + } + } + + /// + /// Raised only under client delegation: the model has handed a task to this application. + /// Results are returned with session.thinking.append or session.commentary.append. + /// + private async Task OnDelegationCreated(string receivedText) + { + var data = JsonSerializer.Deserialize(receivedText); + var delegationId = data?.Delegation?.Id; + _currentDelegationId = delegationId; + + _logger.LogInformation("{Type}: delegation {Id} to {Target}", + LiveServerEventType.DelegationCreated, delegationId, data?.Delegation?.Target); + + //await _onConversationItemCreated(receivedText); + } + #endregion + + #region Turn boundary timers and background work + private void ResetSessionState() + { + DisposeSessionWorkers(); + + lock (_transcriptLock) + { + _outputTranscript.Clear(); + _inputTranscript.Clear(); + } + + _lastTranscriptRole = null; + + _conversationWork = new SerialWorkQueue("conversation", _logger); + + var settings = LiveSettings; + _outputTranscriptTimer = new IdleFlushTimer( + "model transcript", TimeSpan.FromMilliseconds(settings.TranscriptIdleMs), FlushOutputTranscriptAsync, _logger); + // The user's buffer has nothing keeping it armed the way model audio arms the other one, + // so it waits longer: a speaker pausing mid-sentence must not be mistaken for a finished + // turn now that the model replying is what normally closes it. + _inputTranscriptTimer = new IdleFlushTimer( + "user transcript", TimeSpan.FromMilliseconds(settings.InputTranscriptIdleMs), FlushInputTranscriptAsync, _logger); + _audioIdleTimer = new IdleFlushTimer( + "model audio", TimeSpan.FromMilliseconds(settings.AudioIdleMs), OnModelAudioIdle, _logger); + } + + // The flushes only take a buffer and queue its delivery, so there is nothing left to await; + // the timer still wants a Func. + private Task FlushOutputTranscriptAsync() + { + FlushOutputTranscript(); + return Task.CompletedTask; + } + + private Task FlushInputTranscriptAsync() + { + FlushInputTranscript(); + return Task.CompletedTask; + } + + private void DisposeSessionWorkers() + { + _outputTranscriptTimer?.Dispose(); + _inputTranscriptTimer?.Dispose(); + _audioIdleTimer?.Dispose(); + + _outputTranscriptTimer = null; + _inputTranscriptTimer = null; + _audioIdleTimer = null; + + // Already drained by FlushPendingTurns on every path that gets here; this only closes + // the queue so a late event cannot start work on a session that is gone. + _conversationWork?.Dispose(); + _conversationWork = null; + } + #endregion +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Greeting.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Greeting.cs new file mode 100644 index 000000000..5c80d7059 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Greeting.cs @@ -0,0 +1,100 @@ +using BotSharp.Abstraction.Messaging; +using BotSharp.Abstraction.Templating; +using BotSharp.Abstraction.Users; + +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +public partial class LiveCompletionProvider +{ + /// + /// Opens the call, so the caller hears something before they have said anything. + /// + /// The greeting goes out as commentary, the channel the model speaks from. Putting it in + /// session.instructions would work once and then stay there, leaving a standing order to + /// greet for the rest of the call. Nothing is written to the conversation here either: the + /// model speaks it, the transcript comes back, and the normal turn flush records it like + /// any other thing it said. + /// + private async Task SendGreeting() + { + if (!LiveSettings.GreetOnStart) + { + return; + } + + string? greeting; + try + { + greeting = await ResolveGreeting(); + } + catch (Exception ex) + { + // A missing template or a template that fails to render is not worth losing the call. + _logger.LogError(ex, "Failed to resolve the {Provider} greeting.", Provider); + return; + } + + if (string.IsNullOrWhiteSpace(greeting)) + { + return; + } + + _logger.LogInformation("{Provider} opening the call with: {Greeting}", Provider, greeting); + await AppendSessionContext(LiveClientEventType.CommentaryAppend, greeting); + } + + /// + /// The agent's own ".welcome" template wins, rendered the way the text chat renders it, so + /// each agent keeps its own opening line. The configured greeting is the fallback for agents + /// that have no template of their own. + /// + private async Task ResolveGreeting() + { + var agentService = _services.GetRequiredService(); + var agent = await agentService.LoadAgent(_conn.CurrentAgentId); + + var template = agent?.Templates?.FirstOrDefault(x => x.Name == ".welcome"); + if (template != null && !string.IsNullOrWhiteSpace(template.Content)) + { + var render = _services.GetRequiredService(); + var user = _services.GetService(); + var content = render.Render(template.Content, new Dictionary + { + { "user", user! } + }); + + var spoken = ToSpeakableText(content); + if (!string.IsNullOrWhiteSpace(spoken)) + { + return spoken; + } + } + + return LiveSettings.Greeting; + } + + /// + /// A welcome template may hold rich content rather than plain prose. Only the text of it can + /// be spoken, so buttons and cards are reduced to their wording and the rest dropped. + /// + private string ToSpeakableText(string content) + { + try + { + var richContentService = _services.GetRequiredService(); + var messages = richContentService.ConvertToMessages(content); + + var spoken = string.Join(" ", messages + .Select(x => x?.Text) + .Where(x => !string.IsNullOrWhiteSpace(x))) + .Trim(); + + return !string.IsNullOrWhiteSpace(spoken) ? spoken : content; + } + catch (Exception ex) + { + _logger.LogWarning(ex, "Could not read the welcome template as rich content; speaking it as plain text."); + return content; + } + } +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Session.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Session.cs new file mode 100644 index 000000000..9fecdf3bd --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.Session.cs @@ -0,0 +1,326 @@ +using BotSharp.Abstraction.Realtime.Settings; +using BotSharp.Abstraction.Templating; +using BotSharp.Abstraction.Templating.Constants; +using OpenAI.Chat; + +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +public partial class LiveCompletionProvider +{ + /// + /// Pushes the agent's current instruction and tools to the backend handler. + /// Only delegation.responses may change once a session is running, so the rest of the + /// session object is sent once with session.start and left alone afterwards. + /// + public async Task UpdateSession(RealtimeHubConnection conn, bool isInit = false) + { + var (sessionConfig, instruction) = await BuildSessionConfig(conn, isInit); + + // session.start already carried this configuration. + if (isInit) + { + return instruction; + } + + if (sessionConfig.Delegation?.Responses != null) + { + await SendEventToModel(new + { + type = LiveClientEventType.SessionUpdate, + event_id = NewEventId("update"), + session = new + { + delegation = new + { + type = sessionConfig.Delegation.Type, + responses = sessionConfig.Delegation.Responses + } + } + }); + } + + // The agent instruction reaches the backend through delegation.responses.instructions + // above, so nothing is appended to the voice model here. + // + // No settle delay before returning, unlike the realtime provider this was modelled on. + // The caller's next act is response.create on the same socket, and the server applies + // what it is sent in order, so there is nothing for a delay to win. + return instruction; + } + + #region Session config + private async Task<(LiveSessionConfig, string)> BuildSessionConfig(RealtimeHubConnection conn, bool isInit) + { + var agentService = _services.GetRequiredService(); + var realtimeModelSettings = _services.GetRequiredService(); + var liveSettings = LiveSettings; + + var agent = await agentService.LoadAgent(conn.CurrentAgentId); + var (_, messages, options) = PrepareOptions(agent, []); + + var instruction = messages.FirstOrDefault()?.Content.FirstOrDefault()?.Text ?? agent?.Description ?? string.Empty; + var functions = options.Tools.Select(x => new FunctionDef + { + Name = x.FunctionName, + Description = x.FunctionDescription, + Parameters = JsonSerializer.Deserialize(x.FunctionParameters) + }).ToArray(); + + _lastAgentInstruction = instruction; + + var backendModel = agent?.LlmConfig?.Model; + if (string.IsNullOrEmpty(backendModel)) + { + _logger.LogWarning("No llm_config model on agent {AgentId}; the {Provider} delegation has no backend to think with.", + agent?.Id, Provider); + } + + _backendProvider = agent?.LlmConfig?.Provider; + _backendModel = backendModel; + + var config = new LiveSessionConfig + { + Model = _model, + // The voice model gets conversation style and a delegation policy; the agent + // instruction is workflow detail and belongs to the backend handler instead. + Instructions = GetVoiceInstruction(agent), + Store = liveSettings.Store ? true : null, + Audio = BuildAudioConfig(realtimeModelSettings, liveSettings), + Delegation = BuildDelegationConfig(agent, instruction, functions, realtimeModelSettings, backendModel) + }; + + await HookEmitter.Emit(_services, async hook => + { + await hook.OnSessionUpdated(agent, instruction, functions, isInit); + }, agent.Id); + + return (config, instruction); + } + + /// + /// Prompt for the voice model: the agent's "live" channel instruction when it has one + /// (instruction.live.liquid), otherwise . + /// It is a liquid template like any other instruction, so it is rendered before being sent. + /// The agent's own instruction is workflow detail and goes to the backend handler instead. + /// + private string GetVoiceInstruction(Agent? agent) + { + var instruction = agent?.ChannelInstructions? + .FirstOrDefault(x => x.Channel.IsEqualTo(LivePromptConstants.VoiceInstructionChannel))? + .Instruction; + + if (string.IsNullOrWhiteSpace(instruction)) + { + return LivePromptConstants.DefaultVoiceInstruction; + } + + var agentService = _services.GetRequiredService(); + var render = _services.GetRequiredService(); + var renderData = new Dictionary(agentService.CollectRenderData(agent!)) + { + [TemplateRenderConstant.RENDER_AGENT] = agent! + }; + + return render.Render(instruction, renderData); + } + + private LiveAudioConfig BuildAudioConfig(RealtimeModelSettings realtimeModelSettings, LiveSettings liveSettings) + { + // A Live WebSocket carries a single negotiated format in both directions, so the + // output format wins when the two realtime settings disagree. + var format = _realtimeOptions?.OutputAudioFormat + ?? _realtimeOptions?.InputAudioFormat + ?? realtimeModelSettings.OutputAudioFormat + ?? realtimeModelSettings.InputAudioFormat; + + // Live accepts only "format" and "output" here. It has no audio.input section: + // noise reduction and turn detection are handled inside the model. + return new LiveAudioConfig + { + Format = ConvertAudioFormat(format), + Output = new LiveOutputAudioConfig + { + Voice = liveSettings.Voice + } + }; + } + + private LiveDelegationConfig BuildDelegationConfig( + Agent? agent, + string instruction, + FunctionDef[] functions, + RealtimeModelSettings realtimeModelSettings, + string? backendModel) + { + var liveSettings = LiveSettings; + var reasoningEffort = GetReasoningEffort(agent, backendModel); + + return new LiveDelegationConfig + { + Type = LiveDelegationType.Responses, + Responses = new LiveResponsesDelegationConfig + { + Model = backendModel, + Instructions = instruction, + Tools = functions, + ToolChoice = "auto", + ParallelToolCalls = liveSettings.ParallelToolCalls, + MaxOutputTokens = realtimeModelSettings.MaxResponseOutputTokens >= 16 + ? realtimeModelSettings.MaxResponseOutputTokens + : null, + ServiceTier = liveSettings.ServiceTier, + Reasoning = !string.IsNullOrWhiteSpace(reasoningEffort) + ? new LiveReasoningConfig { Effort = reasoningEffort } + : null + } + }; + } + + /// + /// Maps the BotSharp audio format names onto the media types Live accepts: + /// audio/pcm at 24k or 16k, audio/pcmu and audio/pcma at 8k. + /// + private LiveAudioFormat ConvertAudioFormat(string? format) + { + if (string.IsNullOrWhiteSpace(format)) + { + return new LiveAudioFormat { Type = "audio/pcm", Rate = 24000 }; + } + + switch (format.ToLowerInvariant()) + { + case "pcm16": + case "audio/pcm": + return new LiveAudioFormat { Type = "audio/pcm", Rate = 24000 }; + case "pcm16_16k": + return new LiveAudioFormat { Type = "audio/pcm", Rate = 16000 }; + case "g711_ulaw": + case "audio/pcmu": + return new LiveAudioFormat { Type = "audio/pcmu", Rate = 8000 }; + case "g711_alaw": + case "audio/pcma": + return new LiveAudioFormat { Type = "audio/pcma", Rate = 8000 }; + default: + _logger.LogWarning("Unknown audio format {Format} for {Provider}, falling back to audio/pcm 24kHz.", format, Provider); + return new LiveAudioFormat { Type = "audio/pcm", Rate = 24000 }; + } + } + + /// + /// Reasoning effort for the delegated backend model, which is the half of the session that + /// does the thinking. It comes from the agent's own llm_config, not from llm_config.live: + /// that one names the voice model, and the voice model has no reasoning to configure. + /// + private string? GetReasoningEffort(Agent? agent, string? backendModel) + { + var state = _services.GetRequiredService(); + var reasoningEffort = state.GetState("reasoning_effort_level"); + + if (string.IsNullOrEmpty(reasoningEffort)) + { + reasoningEffort = agent?.LlmConfig?.ReasoningEffortLevel; + } + + if (string.IsNullOrEmpty(reasoningEffort)) + { + // Settings for the backend model too, for the same reason. + var settings = GetModelSetting(backendModel)?.Reasoning; + + reasoningEffort = settings?.EffortLevel; + if (settings?.Parameters != null + && settings.Parameters.TryGetValue("EffortLevel", out var settingValue) + && !string.IsNullOrEmpty(settingValue?.Default)) + { + reasoningEffort = settingValue.Default; + } + } + + return reasoningEffort; + } + #endregion + + #region Instruction and tool preparation + private (string, IEnumerable, ChatCompletionOptions) PrepareOptions(Agent agent, List conversations) + { + var agentService = _services.GetRequiredService(); + var state = _services.GetRequiredService(); + + 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 + }; + + // Prepare instruction and functions + var renderData = agentService.CollectRenderData(agent); + var (instruction, functions) = agentService.PrepareInstructionAndFunctions(agent, renderData); + if (!string.IsNullOrWhiteSpace(instruction)) + { + messages.Add(new SystemChatMessage(instruction)); + } + + foreach (var function in functions) + { + if (!agentService.RenderFunction(agent, function, renderData)) + { + continue; + } + + var property = agentService.RenderFunctionProperty(agent, function, renderData); + + options.Tools.Add(ChatTool.CreateFunctionTool( + functionName: function.Name, + functionDescription: function.Description, + functionParameters: BinaryData.FromObjectAsJson(property))); + } + + 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.IfNullOrEmptyAs(message.FunctionName), message.FunctionName, BinaryData.FromString(message.FunctionArgs ?? "{}")) + })); + + messages.Add(new ToolChatMessage(message.ToolCallId.IfNullOrEmptyAs(message.FunctionName), message.LlmContent)); + } + else if (message.Role == AgentRole.User) + { + messages.Add(new UserChatMessage(message.LlmContent)); + } + else if (message.Role == AgentRole.Assistant) + { + messages.Add(new AssistantChatMessage(message.LlmContent)); + } + } + + return (instruction, messages, options); + } + #endregion +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.cs new file mode 100644 index 000000000..1b65bb008 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/LiveCompletionProvider.cs @@ -0,0 +1,438 @@ +using BotSharp.Abstraction.Realtime.Options; +using BotSharp.Abstraction.Realtime.Settings; +using System.Text.RegularExpressions; + +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +/// +/// Full duplex voice provider for gpt-live-1, served from wss://api.openai.com/v1/live/sessions. +/// Reference to https://developers.openai.com/api/docs/guides/live +/// +/// Live differs from the realtime endpoint in three ways that shape this implementation: +/// the model listens and speaks at the same time and arbitrates interruptions itself, +/// transcripts arrive as deltas with no turn completion marker, and reasoning plus tool +/// calls are delegated to a separate backend model. +/// +public partial class LiveCompletionProvider : ILiveCompletion +{ + public string Provider => "openai"; + public string Model => _model; + + private readonly IServiceProvider _services; + private readonly ILogger _logger; + private readonly BotSharpOptions _botsharpOptions; + private readonly OpenAiSettings _openAiSettings; + + /// + /// Level for the event tracing below. A local debug build raises it to Critical so the live + /// traffic stands out in the console; anywhere else it stays at Information, where a running + /// call does not read as a fault. + /// +#if DEBUG + private const LogLevel TraceLevel = LogLevel.Critical; +#else + private const LogLevel TraceLevel = LogLevel.Information; +#endif + + private string _model = LiveModelConstants.GPT_Live_1; + private LlmRealtimeSession? _session; + private RealtimeOptions? _realtimeOptions; + private bool _isBlocking = false; + + /// + /// Last delegation handed out by the model, used to attribute appended context. + /// + private string? _currentDelegationId; + + /// + /// Agent instruction most recently sent to the backend handler. Callers hand it back to + /// TriggerModelInference to mean "carry on", which must not be appended to the voice model. + /// + private string? _lastAgentInstruction; + + /// + /// Provider and model of the backend handler, as resolved for the agent when the session was + /// last configured. Kept so the token stats can be filed against the model that actually + /// spent them, from the response event that reports the usage - which carries no agent. + /// + private string? _backendProvider; + private string? _backendModel; + + private RealtimeHubConnection _conn = null!; + private Func _onModelReady = null!; + private Func _onModelAudioDeltaReceived = null!; + private Func _onModelAudioResponseDone = null!; + private Func _onModelAudioTranscriptDone = null!; + private Func, Task> _onModelResponseDone = null!; + private Func _onConversationItemCreated = null!; + private Func _onInputAudioTranscriptionDone = null!; + private Func _onInterruptionDetected = null!; + + public LiveCompletionProvider( + IServiceProvider services, + ILogger logger, + BotSharpOptions botsharpOptions, + OpenAiSettings openAiSettings) + { + _services = services; + _logger = logger; + _botsharpOptions = botsharpOptions; + _openAiSettings = openAiSettings; + } + + private LiveSettings LiveSettings => _openAiSettings.Live ?? new LiveSettings(); + + public async Task Connect( + RealtimeHubConnection conn, + Func onModelReady, + Func onModelAudioDeltaReceived, + Func onModelAudioResponseDone, + Func onModelAudioTranscriptDone, + Func, Task> onModelResponseDone, + Func onConversationItemCreated, + Func onInputAudioTranscriptionDone, + Func onInterruptionDetected) + { + _logger.LogInformation($"Connecting {Provider} live server..."); + + _conn = conn; + _onModelReady = onModelReady; + _onModelAudioDeltaReceived = onModelAudioDeltaReceived; + _onModelAudioResponseDone = onModelAudioResponseDone; + _onModelAudioTranscriptDone = onModelAudioTranscriptDone; + _onModelResponseDone = onModelResponseDone; + _onConversationItemCreated = onConversationItemCreated; + _onInputAudioTranscriptionDone = onInputAudioTranscriptionDone; + _onInterruptionDetected = onInterruptionDetected; + + _currentDelegationId = null; + _lastAgentInstruction = null; + _backendProvider = null; + _backendModel = null; + ResetSessionState(); + + var realtimeSettings = _services.GetRequiredService(); + + _session = new LlmRealtimeSession(_services, new ChatSessionOptions + { + Provider = Provider, + JsonOptions = _botsharpOptions.JsonSerializerOptions, + Logger = _logger + }); + + var modelSetting = GetModelSetting(); + + // The Live endpoint carries the model in session.start rather than in the query string. + await _session.ConnectAsync( + uri: GetEndpoint(modelSetting), + headers: BuildHeaders(modelSetting), + cancellationToken: CancellationToken.None); + + // session.start must be the first message on the socket. + var (sessionConfig, _) = await BuildSessionConfig(conn, isInit: true); + var startEvent = new + { + type = LiveClientEventType.SessionStart, + event_id = NewEventId("start"), + session = sessionConfig + }; + + if (_logger.IsEnabled(LogLevel.Debug)) + { + // An unknown field makes the server reject the whole session, so log what was sent. + _logger.LogDebug("{Type}: {Payload}", LiveClientEventType.SessionStart, + JsonSerializer.Serialize(startEvent, _botsharpOptions.JsonSerializerOptions)); + } + + await SendEventToModel(startEvent); + + _ = ReceiveMessage(realtimeSettings); + } + + public async Task Reconnect(RealtimeHubConnection conn) + { + _logger.LogInformation($"Reconnecting {Provider} live server..."); + + _isBlocking = true; + _conn = conn; + await Disconnect(); + await Task.Delay(500); + await Connect( + _conn, + _onModelReady, + _onModelAudioDeltaReceived, + _onModelAudioResponseDone, + _onModelAudioTranscriptDone, + _onModelResponseDone, + _onConversationItemCreated, + _onInputAudioTranscriptionDone, + _onInterruptionDetected); + } + + public async Task Disconnect() + { + _logger.LogInformation($"Disconnecting {Provider} live server..."); + + // The conversation is ending, so record the turn nobody spoke after. Done before the + // timers go, since disposing them would drop the only other route to that last turn. + await FlushPendingTurns(); + + DisposeSessionWorkers(); + + if (_session != null) + { + // Asking the server to close yields a final session.closed carrying billed duration. + await SendEventToModel(new + { + type = LiveClientEventType.SessionClose, + event_id = NewEventId("close") + }); + + await _session.DisconnectAsync(); + _session.Dispose(); + _session = null; + } + } + + public async Task AppenAudioBuffer(string message) + { + if (_isBlocking) return; + + await SendEventToModel(new + { + type = LiveClientEventType.InputAudioAppend, + audio = message + }); + } + + public async Task AppenAudioBuffer(ArraySegment data, int length) + { + if (_isBlocking) return; + + var message = Convert.ToBase64String(data.AsSpan(0, length).ToArray()); + await AppenAudioBuffer(message); + } + + public async Task SendEventToModel(object message) + { + if (_session == null) return; + + await _session.SendEventToModelAsync(message); + } + + /// + /// Live has no manual turn taking: the model decides when to speak. Text handed in here is + /// something to say now, so it goes to the channel the model speaks from rather than to + /// session.instructions, which would make a one-off line a permanent standing directive. + /// + public async Task TriggerModelInference(string? instructions = null) + { + // RealtimeConversationHook passes the agent instruction back purely to mean "continue". + // Speaking that aloud would read the whole agent prompt to the caller. + if (!string.IsNullOrWhiteSpace(instructions) + && !string.Equals(instructions, _lastAgentInstruction, StringComparison.Ordinal)) + { + await AppendSessionContext(LiveClientEventType.CommentaryAppend, ExtractSpokenText(instructions)); + } + + await SendEventToModel(new + { + type = LiveClientEventType.ResponseCreate, + event_id = NewEventId("continue") + }); + } + + /// + /// Live arbitrates barge-in inside the model, so there is no response to cancel. + /// + public Task CancelModelResponse() + { + _logger.LogDebug("{Provider} handles interruption in-model; cancel request ignored.", Provider); + return Task.CompletedTask; + } + + /// + /// The live conversation history is owned by the server and cannot be edited item by item. + /// + public Task RemoveConversationItem(string itemId) + { + _logger.LogWarning("{Provider} does not support removing conversation item {ItemId}.", Provider, itemId); + return Task.CompletedTask; + } + + public async Task InsertConversationItem(RoleDialogModel message) + { + if (message.Role == AgentRole.Function) + { + // Tool results go back to the backend handler, not to the voice model. + await SendEventToModel(new + { + type = LiveClientEventType.ResponseItemCreate, + event_id = NewEventId("tool_result"), + item = new + { + type = "function_call_output", + call_id = message.ToolCallId, + output = message.Content ?? string.Empty + } + }); + } + else if (message.Role == AgentRole.Assistant) + { + // Commentary is content the model paraphrases and speaks aloud. + await AppendSessionContext(LiveClientEventType.CommentaryAppend, message.Content); + } + else if (message.Role == AgentRole.User) + { + await SendEventToModel(new + { + type = LiveClientEventType.ResponseItemCreate, + event_id = NewEventId("typed_input"), + item = new + { + type = "message", + role = AgentRole.User, + content = new object[] + { + new + { + type = "input_text", + text = message.Content ?? string.Empty + } + } + } + }); + } + } + + public void SetModelName(string model) + { + _model = model; + } + + public void SetOptions(RealtimeOptions? options) + { + _realtimeOptions = options; + } + + #region Private methods + /// + /// Callers wrap what they want said, as in: Say to user: "your order shipped". Commentary is + /// paraphrased aloud, so the wrapper has to come off or the model reads the instruction out. + /// Anything that is not a short prefix followed by a fully quoted line is passed through. + /// + private static string ExtractSpokenText(string instructions) + { + var match = QuotedPayload.Match(instructions.Trim()); + return match.Success ? match.Groups[1].Value : instructions; + } + + private static readonly Regex QuotedPayload = + new(@"^[^""]{0,40}:\s*""(.+)""$", RegexOptions.Singleline | RegexOptions.Compiled); + + /// + /// Sends one of the session.*.append events. Each append is capped at 500 tokens + /// server side, so the content is truncated before it leaves. + /// + private async Task AppendSessionContext(string eventType, string? content) + { + if (string.IsNullOrWhiteSpace(content)) return; + + // Instructions are session wide directives, so they are not attributed to a delegation. + // Thinking and commentary answer a specific delegation the model handed out. + var delegationId = eventType == LiveClientEventType.InstructionsAppend + ? null + : _currentDelegationId; + + // A misconfigured budget would otherwise drop the content entirely. + var maxTokens = Math.Max(16, LiveSettings.MaxAppendTokens); + // Appends concatenate on the server, so splitting long content across several of them + // reassembles the original rather than corrupting it. Only the chunk ceiling loses text. + var (chunks, truncated) = LiveAppendBudget.SplitIntoChunks(content, maxTokens, LiveAppendBudget.MaxChunks); + + if (truncated) + { + var estimated = LiveAppendBudget.EstimateTokens(content); + + if (eventType == LiveClientEventType.InstructionsAppend) + { + // Half a directive can mean the opposite of the whole - "do not mention the + // discount to the user" cut short still reads as an instruction to obey. Leaving + // the model on its existing instructions is the safer failure. + _logger.LogError( + "Dropping a {Tokens} token instruction: it exceeds {Chunks} appends of {Max} tokens and " + + "cannot be delivered without risking a partial directive.", + estimated, LiveAppendBudget.MaxChunks, maxTokens); + return; + } + + _logger.LogWarning( + "{EventType} content of about {Tokens} tokens exceeded {Chunks} appends of {Max} tokens and was trimmed.", + eventType, estimated, LiveAppendBudget.MaxChunks, maxTokens); + + if (chunks.Count > 0) + { + chunks[^1] += LiveAppendBudget.ElisionMarker; + } + } + + foreach (var chunk in chunks) + { + await SendEventToModel(new + { + type = eventType, + event_id = NewEventId("ctx"), + delegation_id = delegationId, + content = chunk + }); + } + } + + private Dictionary BuildHeaders(LlmModelSetting? settings) + { + return new Dictionary + { + { "Authorization", $"Bearer {settings?.ApiKey}" } + }; + } + + /// + /// Socket address from the model's entry in LlmProviders, so a deployment can point the call + /// at a proxy or a regional host without a code change, falling back to the public Live + /// endpoint. A configured value that is not a usable absolute URI is logged and ignored + /// rather than thrown: a typo in settings should not take the call down. + /// + private Uri GetEndpoint(LlmModelSetting? settings) + { + var endpoint = settings?.Endpoint; + + if (string.IsNullOrWhiteSpace(endpoint)) + { + return new Uri(LiveModelConstants.Endpoint); + } + + if (!Uri.TryCreate(endpoint, UriKind.Absolute, out var uri)) + { + _logger.LogWarning("Ignoring the endpoint configured for {Provider}.{Model}, which is not an absolute URI: {Endpoint}.", + Provider, _model, endpoint); + return new Uri(LiveModelConstants.Endpoint); + } + + return uri; + } + + /// + /// Looks up a model's settings under the OpenAI provider key, which is shared with the + /// realtime and chat providers: the live models are listed alongside them in LlmProviders. + /// + private LlmModelSetting? GetModelSetting(string? model = null) + { + var llmProviderService = _services.GetRequiredService(); + model = model.IfNullOrEmptyAs(_model); + + return llmProviderService.GetSetting(Provider, model); + } + + private static string NewEventId(string prefix) => $"{prefix}_{Guid.NewGuid():N}"; + #endregion +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/README.md b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/README.md new file mode 100644 index 000000000..e17caf5ce --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/README.md @@ -0,0 +1,214 @@ +# OpenAI Live provider (`gpt-live-1`) + +Implements `ILiveCompletion` against OpenAI's Live endpoint, `wss://api.openai.com/v1/live/sessions`. +It shares the **`openai`** provider key with the realtime and chat providers; what separates it is +the interface, not the name. + +`ILiveCompletion` derives from `IRealTimeCompletion`, so a live session drives the same hub, +middleware and hooks as a realtime one. The separate interface is what keeps the two families +apart in DI: the live provider registers as `ILiveCompletion` only, so +`GetServices()` never returns it and `gpt-realtime` cannot be replaced by +accident. A live session happens only when an agent asks for one. + +Reference: https://developers.openai.com/api/docs/guides/live + +## Enabling it + +Live is chosen per agent, through `llm_config.live` - a sibling of `llm_config.realtime`, not a +variation of it: + +```jsonc +// agents//agent.json +"llm_config": { + "live": { + "provider": "openai", + "model": "gpt-live-1" + } +} +``` + +`RealtimeHub` resolves it like this: + +| Agent config | Family | Resolved from | +| --- | --- | --- | +| `llm_config.live` set | live | `GetServices()` | +| `llm_config.realtime` set | realtime | `GetServices()` | +| neither | realtime | `RealtimeModel` settings, else `openai`/`gpt-realtime` | + +Both families answer to the provider name `openai`, so the name alone decides nothing: the family +comes from which agent config is set, and the name is then looked up in that family's list only. +`RealtimeModel` therefore always means realtime, and a name with no provider in that list throws +with the family and provider named. + +Deployment-wide defaults for the live session itself: + +```jsonc +// appsettings.json +"OpenAi": { + "Live": { + "Voice": "marin", // Live has its own voice set; "alloy" is rejected + "BackendModel": "gpt-5.6-luna", // model that does the reasoning and runs tools + "TranscriptIdleMs": 800, + "AudioIdleMs": 500 + } +} +``` + +Audio format still comes from `RealtimeModel` (`g711_ulaw` for telephony, `pcm16` for browser +audio), as does `MaxResponseOutputTokens`; those settings are shared by both families. + +Credentials come from the `gpt-live-1` model entry under the `openai` provider in `LlmProviders`, +next to the realtime and chat models, so an existing API key is not duplicated. The entry carries +`"Type": "live"` and the `Live` capability, which keeps it out of realtime model lookups: + +```jsonc +{ "Id": "gpt-live", "Name": "gpt-live-1", "ApiKey": "", "Type": "live", "Capabilities": [ "Live" ] } +``` + +Add `"Endpoint"` to that entry to send the socket somewhere else - a proxy, or a regional host. +Left out, it is `wss://api.openai.com/v1/live/sessions`; a value that is not an absolute URI is +logged and ignored rather than failing the call. + +## How the Live protocol maps onto the completion interface + +| Interface member | Live wire event | +| --- | --- | +| `Connect` | `session.start` (carries model, instructions, audio, delegation) | +| `AppenAudioBuffer` | `session.input_audio.append` | +| `UpdateSession` | `session.update` (delegation only) + `session.instructions.append` | +| `InsertConversationItem` (function) | `response.item.create` with `function_call_output` | +| `InsertConversationItem` (assistant) | `session.commentary.append` | +| `InsertConversationItem` (user) | `response.item.create` with a typed `message` | +| `TriggerModelInference` | `session.commentary.append` + `response.create` | +| `Disconnect` | `session.close` | + +| Callback | Live wire event | +| --- | --- | +| `onModelReady` | `session.started` | +| `onModelAudioDeltaReceived` | `session.output_audio.delta` | +| `onModelAudioTranscriptDone` / `onModelResponseDone` | `session.output_transcript.delta`, flushed on silence | +| `onInputAudioTranscriptionDone` | `session.input_transcript.delta`, flushed on silence | +| `onModelResponseDone` (tool call) | `response.event` → `response.output_item.done` | +| `onConversationItemCreated` | `session.delegation.created` | + +## Greeting the caller + +The model opens the call rather than waiting to be spoken to. On `session.started`, once the +backend handler is configured, the greeting is sent as `session.commentary.append` - the channel +the model speaks from, and a one-shot one, unlike `session.instructions.append` which would leave +a standing order to greet for the rest of the call. + +The agent's own `.welcome` template wins, rendered exactly as the text chat renders it, so an +agent keeps the opening line it already has. Rich content is reduced to its wording, since only +text can be spoken. Agents without a template fall back to `OpenAi:Live:Greeting`, and +`GreetOnStart: false` turns the whole thing off. + +Nothing is written to the conversation here: the model speaks the greeting, its transcript comes +back, and the normal turn flush records it like anything else it said. + +## Two prompts, not one + +Live splits the prompt across its two halves, and so does this provider: + +| Prompt | Goes to | Content | +| --- | --- | --- | +| `session.instructions` | voice model | how the conversation sounds, and when to delegate | +| `delegation.responses.instructions` | backend model | the BotSharp agent instruction and its tools | + +The voice prompt defaults to `LivePromptConstants.DefaultVoiceInstruction`, which follows the +structure in OpenAI's [prompting guide](https://developers.openai.com/api/docs/guides/live-prompting): +personality, backchannel policy, interruption policy, delegation policy, response style. Keep it +short - the guide is explicit that reasoning and procedures belong in the backend prompt. + +Override it per agent with a `live` channel instruction - the same mechanism as any other +channel override, so it needs no code and no setting: + +``` +agents//instructions/instruction.live.liquid +``` + +It is rendered as a liquid template like the default instruction. An agent with no such file +keeps `DefaultVoiceInstruction`. + +The agent instruction is never appended to the voice model. It reaches the backend through +`session.start`, and is refreshed with `session.update` on `delegation.responses` whenever +`UpdateSession` runs - after a function call or an agent transfer. + +## TriggerModelInference speaks, it does not re-instruct + +Every caller of `TriggerModelInference(text)` passes a one-off line for the current turn, never +standing policy - `Say to user: "..."` from the conversation hook, `Response based on the user +input: ...` from Twilio. `session.instructions.append` would make those permanent, so they go to +`session.commentary.append`, which the model paraphrases aloud and then leaves behind. + +A wrapper such as `Say to user: "your order shipped"` is unwrapped first, otherwise the model +reads the instruction out. Anything that is not a short prefix plus a fully quoted line passes +through untouched. + +The agent instruction is the one thing filtered out here: the hook hands it back to mean +"continue", and speaking it aloud would read the whole agent prompt to the caller. + +## Append budget + +Every `session.*.append` is capped at 500 tokens server side, while the session instructions +they build up hold 16,384. Appends concatenate, so over-budget content is **split on sentence +boundaries** across up to 5 appends and reassembles into the original - nothing is truncated in +the normal case. `MaxAppendTokens` (default 450) sets the per-append budget. + +Only content too large for 5 appends loses anything, and the channels differ there: + +- **commentary** and **thinking** keep what fits, marked with `[...]` so the model can see it + is partial. +- **instructions** are dropped entirely with an error logged. Half a directive still reads as a + directive - "do not mention the discount to the user" cut short inverts its meaning - so + leaving the model on its existing instructions is the safer failure. + +The budget is counted in tokens rather than characters because the two diverge by script: Latin +prose is roughly four characters per token, CJK closer to one. A character cap sized for English +overshoots the server limit badly for Chinese, Japanese and Korean. + +## Nothing runs on the receive loop + +The socket has a single consumer, and anything awaited inside it stops the socket being read. +On a half duplex session that costs nothing, because the model is silent while a tool runs. On +Live it is the whole problem: the model keeps talking through a delegation, so audio frames pile +up unread and the caller hears the reply cut in half - "Okay, checking" ... ten seconds ... "that +now." + +So the receive loop parses an event and moves on. Conversation work - invoking a tool, recording +a turn - goes to a `SerialWorkQueue`, which runs it on one background worker in the order it was +queued. Serial rather than fire and forget, because two tool calls must not interleave their +session updates, and because the conversation state this work touches is not thread safe. + +What stays on the loop is what must not queue behind a tool call: audio deltas and live +transcript deltas, both of which are only a write to the caller's socket. + +Turn boundaries are still decided on the loop. The flushes take the buffer there, synchronously, +and queue only its delivery - otherwise a turn queued behind a slow tool would pick up the words +the model spoke after it. + +## Behaviour that differs from the realtime provider + +- **No turn completion event.** Live streams transcript deltas and never marks the end of a turn, + so turns are closed by **alternation**: each speaker's buffer is flushed when the other one + starts talking. That writes a turn the moment it is genuinely over, and a thinking pause mid + sentence no longer splits one reply into two. + Idle timers remain only as a backstop for the last turn of a session, when nobody speaks after + it - `TranscriptIdleMs` for the model, the longer `InputTranscriptIdleMs` for the user, since + model audio arms the former but nothing arms the latter between words. `AudioIdleMs` still + drives `onModelAudioResponseDone`. +- **No interruption callback.** Live is full duplex and arbitrates barge-in inside the model, + so `onInterruptionDetected` is never raised and `CancelModelResponse` is a no-op. Client-side + playback buffers are deliberately not cleared. +- **`RemoveConversationItem` is unsupported.** The server owns the live history. +- **`audio` has no `input` section.** Live accepts only `audio.format` and `audio.output.voice`; + noise reduction and turn detection are handled inside the model, unlike the realtime API. + The server rejects the entire session on an unknown field, so keep this object minimal. +- **Only responses delegation is supported.** Live also offers a "client" mode, where the + application answers delegations itself; it is not implemented here. +- **Only `delegation.responses` is mutable.** Everything else in the session object is fixed at + `session.start`; changing the delegation mode requires a new session. +- **Billing is per second, not per token.** `TokenStatsModel` reports the backend Responses + usage; the voice minutes are reported separately by `session.closed`. +- `RealtimeModelSettings.ModelResponseTimeoutSeconds` has no effect here, since it is defined in + terms of the realtime `response.done` lifecycle that Live does not have. diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/SerialWorkQueue.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/SerialWorkQueue.cs new file mode 100644 index 000000000..ef15d5022 --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Providers/Live/SerialWorkQueue.cs @@ -0,0 +1,106 @@ +using System.Threading.Channels; + +namespace BotSharp.Plugin.OpenAI.Providers.Live; + +/// +/// Runs queued work on a single background worker, one item at a time and in the order it was +/// queued. +/// +/// The Live socket has one consumer, and anything awaited on it stops the socket being read. +/// That is fatal on a full duplex call: the model keeps speaking while a tool runs, so audio +/// frames pile up unread and the caller hears a gap where the reply should be. Conversation +/// work is handed here instead, so the receive loop only ever parses an event and moves on. +/// +/// Serial rather than fire and forget for two reasons: two tool calls must not interleave their +/// session updates, and the conversation state this work touches is not thread safe. +/// +internal sealed class SerialWorkQueue : IDisposable +{ + private readonly Channel> _queue = Channel.CreateUnbounded>( + new UnboundedChannelOptions { SingleReader = true }); + + /// + /// The queue whose worker the current call is running on, if any. Flows into the work's own + /// continuations, which is how spots being called from inside itself. + /// + private static readonly AsyncLocal _running = new(); + + private readonly Task _worker; + private readonly ILogger? _logger; + private readonly string _name; + + /// + /// The name is suffixed with a fresh id. There is one queue per call, so without it the log + /// lines of concurrent sessions all read the same and none of them can be followed. + /// + public SerialWorkQueue(string name, ILogger? logger = null) + { + _name = $"{name}-{Guid.NewGuid():N}"; + _logger = logger; + _worker = Task.Run(RunAsync); + } + + /// + /// Name as it appears in this queue's log lines, id included. + /// + public string Name => _name; + + /// + /// Hands work to the worker and returns at once. Never throws: a queue that fails to accept + /// work must not take the caller's receive loop down with it. + /// + public void Enqueue(Func work) + { + if (_queue.Writer.TryWrite(work)) return; + + // Only happens once the queue is closed, i.e. the session is already tearing down. + _logger?.LogWarning("Dropped {Name} work: the queue is closed.", _name); + } + + /// + /// Stops accepting work and waits for what is already queued to finish. Called during + /// teardown, where the final turns still have to reach storage before the session goes. + /// + public async Task DrainAsync() + { + _queue.Writer.TryComplete(); + + // Queued work can end up here: a tool call runs on this worker, and handling it can + // tear the session down - a reconnect, say. Waiting would be waiting on ourselves. + // Closing the queue is still right; what is already in it runs as this worker unwinds. + if (ReferenceEquals(_running.Value, this)) + { + _logger?.LogDebug("Not waiting on the {Name} queue from inside its own worker.", _name); + return; + } + + try + { + await _worker; + } + catch (Exception ex) + { + _logger?.LogError(ex, "Failed to drain the {Name} queue.", _name); + } + } + + private async Task RunAsync() + { + _running.Value = this; + + await foreach (var work in _queue.Reader.ReadAllAsync()) + { + try + { + await work(); + } + catch (Exception ex) + { + // One failed item must not end the worker; the rest of the call still needs it. + _logger?.LogError(ex, "Error while running queued {Name} work.", _name); + } + } + } + + public void Dispose() => _queue.Writer.TryComplete(); +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Settings/LiveSettings.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Settings/LiveSettings.cs new file mode 100644 index 000000000..5c1cadd7a --- /dev/null +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Settings/LiveSettings.cs @@ -0,0 +1,73 @@ +namespace BotSharp.Plugin.OpenAI.Settings; + +/// +/// Defaults for the OpenAI Live endpoint (gpt-live-1). Agent level realtime config and +/// conversation state take precedence over these values at runtime. +/// +public class LiveSettings +{ + /// + /// Voice model that owns the conversation. + /// + public string Model { get; set; } = LiveModelConstants.GPT_Live_1; + + /// + /// Live only accepts its own voice set; realtime voice names such as "alloy" are rejected. + /// + public string Voice { get; set; } = "marin"; + + /// + /// Whether the model opens the call rather than waiting to be spoken to. + /// + public bool GreetOnStart { get; set; } = true; + + /// + /// Fallback greeting for agents with no ".welcome" template of their own; an agent that has + /// one greets with that instead, so it keeps the opening line it already uses in text chat. + /// Either way the model paraphrases it rather than reading it out word for word, which is + /// what keeps a spoken greeting sounding natural. + /// + public string? Greeting { get; set; } = LivePromptConstants.DefaultGreeting; + + /// + /// auto, default, flex, or priority (also accepted as "fast"). + /// Applies to the backend Responses call, not to the billed voice minutes. + /// + public string? ServiceTier { get; set; } = "auto"; + + public bool? ParallelToolCalls { get; set; } + + /// + /// Records the session so it can be downloaded or forked later. + /// + public bool Store { get; set; } + + /// + /// A turn is normally closed by the other speaker starting, or by the model's audio going + /// idle, since Live emits no completion marker. This is the backstop for the model's last + /// turn: armed by its transcript deltas alone, so a pause longer than this mid reply will + /// split it into two messages. + /// + public int TranscriptIdleMs { get; set; } = 500; + + /// + /// Backstop for the user's last turn. Longer than because + /// nothing arms the buffer between words: at 500ms a speaker pausing mid-sentence would be + /// filed as a finished turn before the model ever got to reply and close it properly. + /// + public int InputTranscriptIdleMs { get; set; } = 4000; + + /// + /// Idle gap after which the model audio stream is treated as finished. + /// + public int AudioIdleMs { get; set; } = 1000; + + /// + /// Budget for a single session.*.append. The server caps each append at 500 tokens and + /// rejects anything larger, so this leaves headroom for the estimate being approximate. + /// Spoken content over the budget is split across appends; a directive or a tool result is + /// trimmed instead, because splitting those would change what they mean. + /// + public int MaxAppendTokens { get; set; } = 450; + +} diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Settings/OpenAiSettings.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Settings/OpenAiSettings.cs index 7f9ddeeb3..24387a443 100644 --- a/src/Plugins/BotSharp.Plugin.OpenAI/Settings/OpenAiSettings.cs +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Settings/OpenAiSettings.cs @@ -12,4 +12,9 @@ public class OpenAiSettings /// Conversation state keys take precedence over these values at runtime. /// public WebSearchSettings? WebSearch { get; set; } + + /// + /// Defaults for the Live endpoint (gpt-live-1): voice, backend delegation and turn boundaries. + /// + public LiveSettings? Live { get; set; } } diff --git a/src/Plugins/BotSharp.Plugin.OpenAI/Using.cs b/src/Plugins/BotSharp.Plugin.OpenAI/Using.cs index 4cb466d09..ed424c797 100644 --- a/src/Plugins/BotSharp.Plugin.OpenAI/Using.cs +++ b/src/Plugins/BotSharp.Plugin.OpenAI/Using.cs @@ -38,6 +38,7 @@ global using BotSharp.Plugin.OpenAI.Models.Text; global using BotSharp.Plugin.OpenAI.Models.Realtime; +global using BotSharp.Plugin.OpenAI.Models.Live; global using BotSharp.Plugin.OpenAI.Models.Image; global using BotSharp.Plugin.OpenAI.Models.Web; global using BotSharp.Plugin.OpenAI.Settings; \ No newline at end of file diff --git a/src/WebStarter/appsettings.json b/src/WebStarter/appsettings.json index da8ea1926..823cc3110 100644 --- a/src/WebStarter/appsettings.json +++ b/src/WebStarter/appsettings.json @@ -42,7 +42,16 @@ }, "OpenAi": { - "UseResponseApi": true + "UseResponseApi": true, + "Live": { + "Model": "gpt-live-1", + "Voice": "marin", + "GreetOnStart": true, + "Greeting": "Hello! How can I help you today?", + "Store": false, + "TranscriptIdleMs": 800, + "AudioIdleMs": 500 + } }, "LlmProviders": [ @@ -213,6 +222,16 @@ "AudioOutputCost": 0 } }, + { + "Id": "gpt-live", + "Name": "gpt-live-1", + "ApiKey": "", + "Type": "live", + "MultiModal": true, + "Capabilities": [ + "Live" + ] + }, { "Id": "text-embedding-3", "Name": "text-embedding-3-small",