using System.Text.Json.Nodes; using backend.Conversations; using backend.Tools; using Microsoft.Extensions.Logging; namespace backend.Realtime; public interface ISessionOutput { Task WriteEnvelopeAsync(JsonObject envelope, CancellationToken ct); Task WriteBinaryAsync(ReadOnlyMemory bytes, CancellationToken ct); } public class RealtimeSession { private readonly Guid _deviceId; private readonly IRealtimeUpstream _upstream; private readonly ISessionOutput _output; private readonly ConversationLog _log; private readonly ToolRegistry _registry; private readonly RealtimeSettings _cfg; private readonly IDeviceChannel _channel; private readonly ILogger _logger; private Conversation? _conversation; private EndReason _endReason = EndReason.Unknown; private bool _endRequested; private DateTime _lastSpeechAt = DateTime.UtcNow; private readonly List<(string CallId, string Name, string ArgsJson)> _pendingCalls = new(); private readonly System.Threading.Channels.Channel _uplink = System.Threading.Channels.Channel.CreateBounded( new System.Threading.Channels.BoundedChannelOptions(64) { FullMode = System.Threading.Channels.BoundedChannelFullMode.DropOldest, }); public ValueTask WriteUplinkAsync(byte[] frame, CancellationToken ct) => _uplink.Writer.WriteAsync(frame, ct); public RealtimeSession( Guid deviceId, IRealtimeUpstream upstream, ISessionOutput output, ConversationLog log, ToolRegistry registry, RealtimeSettings cfg, IDeviceChannel channel, ILogger logger) { _deviceId = deviceId; _upstream = upstream; _output = output; _log = log; _registry = registry; _cfg = cfg; _channel = channel; _logger = logger; } public Guid? ConversationId => _conversation?.Id; public async Task RunAsync(CancellationToken ct) { try { var enabled = _registry.EnabledFor(_cfg.EnabledTools).ToList(); await _upstream.SendJsonAsync(RealtimeEvents.SessionUpdate(_cfg, enabled), ct); var ack = await WaitForUpstreamTypeAsync("session.updated", ct); if (ack is null) { _endReason = EndReason.Error; return; } _conversation = await _log.StartAsync(_deviceId, ct); await _output.WriteEnvelopeAsync(new JsonObject { ["type"] = "session_started", ["conversation_id"] = _conversation.Id.ToString(), }, ct); _lastSpeechAt = DateTime.UtcNow; var uplinkTask = Task.Run(() => PumpUplinkAsync(ct), ct); var idleTask = Task.Run(() => WatchdogAsync(ct), ct); await PumpAsync(ct); _uplink.Writer.TryComplete(); await uplinkTask; await idleTask; } catch (OperationCanceledException) { if (_endReason == EndReason.Unknown) _endReason = EndReason.Error; } catch (Exception ex) { _logger.LogError(ex, "RealtimeSession crashed"); _endReason = EndReason.Error; } finally { await _upstream.DisposeAsync(); if (_conversation is not null) { await _log.EndAsync(_conversation.Id, _endReason, CancellationToken.None); } var endEnv = new JsonObject { ["type"] = "session_ended", ["reason"] = _endReason.ToString().ToLowerInvariant(), }; await _output.WriteEnvelopeAsync(endEnv, CancellationToken.None); } } private async Task PumpAsync(CancellationToken ct) { while (!ct.IsCancellationRequested && !_endRequested) { var evt = await ReceiveOrIdleAsync(ct); if (evt is null && _endRequested) return; if (evt is null && ct.IsCancellationRequested) return; if (evt is null) continue; var type = (string?)evt["type"]; switch (type) { case "input_audio_buffer.speech_started": _lastSpeechAt = DateTime.UtcNow; break; case "response.audio.delta": await ForwardAudioDeltaAsync(evt, ct); break; case "conversation.item.input_audio_transcription.completed": { var text = (string?)evt["transcript"]; if (_conversation is not null && !string.IsNullOrEmpty(text)) await _log.AppendUserAsync(_conversation.Id, text!, ct); break; } case "response.audio_transcript.done": { var text = (string?)evt["transcript"]; if (_conversation is not null && !string.IsNullOrEmpty(text)) await _log.AppendAssistantAsync(_conversation.Id, text!, ct); break; } case "response.function_call_arguments.done": { var callId = (string?)evt["call_id"] ?? ""; var name = (string?)evt["name"] ?? ""; var args = (string?)evt["arguments"] ?? "{}"; _pendingCalls.Add((callId, name, args)); break; } case "response.done": { await _output.WriteEnvelopeAsync( new JsonObject { ["type"] = "assistant_done" }, ct); if (_pendingCalls.Count == 0) break; var calls = _pendingCalls.ToList(); _pendingCalls.Clear(); if (calls.Any(c => c.Name == "end_session")) { _endReason = EndReason.Tool; _endRequested = true; break; } foreach (var c in calls) { await ExecuteToolAsync(c.CallId, c.Name, c.ArgsJson, ct); } break; } } } } private async Task ExecuteToolAsync(string callId, string name, string argsJson, CancellationToken ct) { var tool = _registry.Get(name); string outputJson; if (tool is null) { outputJson = $$"""{"ok":false,"error":"unknown tool: {{name}}"}"""; } else { try { using var argsDoc = System.Text.Json.JsonDocument.Parse(argsJson); var ctx = new DeviceContext(_deviceId, _conversation!.Id, _channel); var res = await tool.ExecuteAsync( new ToolInvocation(callId, argsDoc.RootElement.Clone()), ctx, ct); outputJson = res.Output.GetRawText(); } catch (Exception ex) { outputJson = System.Text.Json.JsonSerializer.Serialize(new { ok = false, error = ex.Message, }); } } if (_conversation is not null) { await _log.AppendToolAsync(_conversation.Id, name, argsJson, outputJson, ct); } await _upstream.SendJsonAsync(RealtimeEvents.FunctionCallOutput(callId, outputJson), ct); await _upstream.SendJsonAsync(RealtimeEvents.ResponseCreate(), ct); } private async Task ForwardAudioDeltaAsync(JsonObject evt, CancellationToken ct) { var b64 = (string?)evt["delta"]; if (string.IsNullOrEmpty(b64)) return; var bytes = Convert.FromBase64String(b64); await _output.WriteBinaryAsync(bytes, ct); } private async Task WatchdogAsync(CancellationToken ct) { try { var timeout = TimeSpan.FromSeconds(_cfg.IdleTimeoutSeconds); while (!ct.IsCancellationRequested && !_endRequested) { await Task.Delay(250, ct); if (DateTime.UtcNow - _lastSpeechAt > timeout) { _endReason = EndReason.Idle; _endRequested = true; return; } } } catch (OperationCanceledException) { } } private async Task ReceiveOrIdleAsync(CancellationToken ct) { using var timeout = CancellationTokenSource.CreateLinkedTokenSource(ct); var receive = _upstream.ReceiveJsonAsync(timeout.Token); var delay = Task.Delay(250, timeout.Token); var done = await Task.WhenAny(receive, delay); if (done == receive) return await receive; timeout.Cancel(); try { await receive; } catch { /* ignore */ } return null; } private async Task PumpUplinkAsync(CancellationToken ct) { try { await foreach (var frame in _uplink.Reader.ReadAllAsync(ct)) { await _upstream.SendJsonAsync(RealtimeEvents.InputAudioBufferAppend(frame), ct); } } catch (OperationCanceledException) { } catch (Exception ex) { _logger.LogWarning(ex, "uplink pump errored"); } } private async Task WaitForUpstreamTypeAsync(string type, CancellationToken ct) { while (true) { var evt = await _upstream.ReceiveJsonAsync(ct); if (evt is null) return null; if ((string?)evt["type"] == type) return evt; } } }