Skip to content
Merged
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,16 @@ public AgentLlmConfig(AgentTemplateLlmConfig templateLlmConfig)
[JsonPropertyName("realtime")]
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public LlmRealtimeConfig? Realtime { get; set; }

/// <summary>
/// Live config, kept apart from <see cref="Realtime"/> 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.
/// </summary>
[JsonPropertyName("live")]
[JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
public LlmLiveConfig? Live { get; set; }
}

public class LlmImageCompositionConfig : LlmProviderModel
Expand All @@ -92,4 +102,8 @@ public class LlmAudioTranscriptionConfig : LlmProviderModel

public class LlmRealtimeConfig : LlmConfigBase
{
}

public class LlmLiveConfig : LlmConfigBase
{
}
16 changes: 16 additions & 0 deletions src/Infrastructure/BotSharp.Abstraction/MLTasks/ILiveCompletion.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
namespace BotSharp.Abstraction.MLTasks;

/// <summary>
/// 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 <see cref="IRealTimeCompletion"/> 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
/// <c>GetServices&lt;IRealTimeCompletion&gt;()</c>, so it is reachable only through the agent's
/// live config, never as a realtime fallback.
/// </summary>
public interface ILiveCompletion : IRealTimeCompletion
{
}
Original file line number Diff line number Diff line change
Expand Up @@ -185,7 +185,13 @@ public enum LlmModelType
Embedding = 4,
Audio = 5,
Realtime = 6,
Web = 7
Web = 7,

/// <summary>
/// Full duplex voice model that runs the conversation itself and delegates the thinking to
/// a backend model. Kept apart from <see cref="Realtime"/>: the two are not interchangeable.
/// </summary>
Live = 8
}

public enum LlmModelCapability
Expand All @@ -203,5 +209,6 @@ public enum LlmModelCapability
AudioGeneration = 10,
Realtime = 11,
WebSearch = 12,
PdfReading = 13
PdfReading = 13,
Live = 14
}
14 changes: 14 additions & 0 deletions src/Infrastructure/BotSharp.Abstraction/Realtime/IRealtimeHook.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,20 @@ namespace BotSharp.Abstraction.Realtime;

public interface IRealtimeHook : IHookBase
{
/// <summary>
/// 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.
/// </summary>
/// <param name="userId">
/// 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.
/// </param>
void OnAuthenticate(string agentId, string conversationId, string? userId) { }

Task OnModelReady(Agent agent, IRealTimeCompletion completer)
=> Task.CompletedTask;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,18 @@ public class RealtimeHubConnection
public Func<string> OnModelUserInterrupted { get; set; } = null!;
public Func<string> OnUserSpeechDetected { get; set; } = () => string.Empty;

/// <summary>
/// Serializes a partial transcript (role, delta) for the user stream. Display only: a turn
/// still reaches conversation storage only once it is complete.
/// </summary>
public Func<string, string, string>? OnModelTranscriptDelta { get; set; }

/// <summary>
/// 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.
/// </summary>
public Func<string, Task>? SendEventToUser { get; set; }

public void ResetResponseState()
{
LastAssistantItemId = null;
Expand Down
60 changes: 45 additions & 15 deletions src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<IRealTimeCompletion>().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 () =>
Expand Down Expand Up @@ -189,27 +192,54 @@ public RealtimeHubConnection SetHubConnection(string conversationId)
return _conn;
}

private (string, string) GetLlmProviderModel(Agent agent)
/// <summary>
/// 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.
/// </summary>
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<ISettingService>();
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);
}

/// <summary>
/// Live providers register as <see cref="ILiveCompletion"/> 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 <see cref="IRealTimeCompletion"/> - so nothing downstream has to know the difference.
/// </summary>
private IRealTimeCompletion ResolveCompleter(string provider, bool isLive)
{
var completer = isLive
? _services.GetServices<ILiveCompletion>().FirstOrDefault(x => x.Provider == provider)
: _services.GetServices<IRealTimeCompletion>().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;
}
}
48 changes: 24 additions & 24 deletions src/Infrastructure/BotSharp.Core/Demo/Functions/GetWeatherFn.cs
Original file line number Diff line number Diff line change
Expand Up @@ -38,16 +38,16 @@ public async Task<bool> 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()
Expand All @@ -62,20 +62,20 @@ public async Task<bool> Execute(RoleDialogModel message)
message.Content = $"It is a sunny day!";
message.StopCompletion = false;

#if DEBUG
var sidecar = _services.GetService<IConversationSideCar>();
if (sidecar != null)
{
var text = $"I want to know fun events in {args?.City}";
var states = new List<MessageState>
{
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<IConversationSideCar>();
// if (sidecar != null)
// {
// var text = $"I want to know fun events in {args?.City}";
// var states = new List<MessageState>
// {
// 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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -205,6 +209,35 @@ public static IRealTimeCompletion GetRealTimeCompletion(
return completer;
}

/// <summary>
/// Live providers are registered as <see cref="ILiveCompletion"/> only, so they are resolved
/// from their own list: a provider name is never shared between the two voice families here.
/// </summary>
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<ILiveCompletion>();
(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<ILogger<CompletionProvider>>();
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,
Expand Down
46 changes: 45 additions & 1 deletion src/Plugins/BotSharp.Plugin.ChatHub/ChatStreamMiddleware.cs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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<IRealtimeHook>(services,
hook => hook.OnAuthenticate(agentId, conversationId, userId), agentId);

using var webSocket = await httpContext.WebSockets.AcceptWebSocketAsync();
await HandleWebSocket(services, agentId, conversationId, webSocket);
}
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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
{
Expand Down
Loading
Loading