From ed39786539f5c9a23176d110c0aa32db9dcf4ef6 Mon Sep 17 00:00:00 2001 From: Jicheng Lu <103353@smsassist.com> Date: Fri, 2 May 2025 13:57:03 -0500 Subject: [PATCH] add custom stream mode --- .../Services/RealtimeHub.cs | 2 +- .../Audio/AudioInStream.cs | 122 +++++++++++++ .../Audio/AudioOut.cs | 51 ++++++ .../Enums/SessionMode.cs | 7 + tests/BotSharp.Test.RealtimeVoice/Program.cs | 134 +------------- .../ConsoleChatSession.CustomStream.cs | 46 +++++ .../ConsoleChatSession.StreamChannel.cs | 50 +++++ .../Session/ConsoleChatSession.cs | 171 ++++++++++++++++++ tests/BotSharp.Test.RealtimeVoice/Using.cs | 8 + 9 files changed, 459 insertions(+), 132 deletions(-) create mode 100644 tests/BotSharp.Test.RealtimeVoice/Audio/AudioInStream.cs create mode 100644 tests/BotSharp.Test.RealtimeVoice/Audio/AudioOut.cs create mode 100644 tests/BotSharp.Test.RealtimeVoice/Enums/SessionMode.cs create mode 100644 tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.CustomStream.cs create mode 100644 tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.StreamChannel.cs create mode 100644 tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.cs diff --git a/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs b/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs index b81fe0f5..60729758 100644 --- a/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs +++ b/src/Infrastructure/BotSharp.Core.Realtime/Services/RealtimeHub.cs @@ -49,8 +49,8 @@ public class RealtimeHub : IRealtimeHub // Not TriggerModelInference, waiting for user utter. var instruction = await _completer.UpdateSession(_conn, isInit: true); var data = _conn.OnModelReady(); - await (init?.Invoke(data) ?? Task.CompletedTask); await HookEmitter.Emit(_services, async hook => await hook.OnModelReady(agent, _completer)); + await (init?.Invoke(data) ?? Task.CompletedTask); }, onModelAudioDeltaReceived: async (audioDeltaData, itemId) => { diff --git a/tests/BotSharp.Test.RealtimeVoice/Audio/AudioInStream.cs b/tests/BotSharp.Test.RealtimeVoice/Audio/AudioInStream.cs new file mode 100644 index 00000000..0531904d --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Audio/AudioInStream.cs @@ -0,0 +1,122 @@ +using NAudio.Wave; + +namespace BotSharp.Test.RealtimeVoice.Audio; + +internal class AudioInStream : Stream +{ + private const int SAMPLE_RATE = 16000; + private const int BYTES_PER_SAMPLE = 2; + private const int CHANNELS = 1; + private const int SAMPLING_SECONDS = 10; + private const int TIMEOUT_SECONDS = 100; + + private readonly byte[] _buffer = new byte[SAMPLE_RATE * BYTES_PER_SAMPLE * CHANNELS * SAMPLING_SECONDS]; + private readonly object _lock = new(); + private int _bufferReadPtr = 0; + private int _bufferWritePtr = 0; + private readonly WaveInEvent _waveInEvent; + + private AudioInStream() + { + _waveInEvent = new WaveInEvent + { + WaveFormat = new WaveFormat(SAMPLE_RATE, BYTES_PER_SAMPLE * 8, CHANNELS), + DeviceNumber = 0 + }; + + _waveInEvent.DataAvailable += (_, e) => + { + lock (_lock) + { + var bytesToCopy = e.BytesRecorded; + if (_bufferWritePtr + bytesToCopy >= _buffer.Length) + { + var chunkLength = _buffer.Length - _bufferWritePtr; + Array.Copy(e.Buffer, 0, _buffer, _bufferWritePtr, chunkLength); + bytesToCopy -= chunkLength; + _bufferWritePtr = 0; + } + Array.Copy(e.Buffer, e.BytesRecorded - bytesToCopy, _buffer, _bufferWritePtr, bytesToCopy); + _bufferWritePtr += bytesToCopy; + } + }; + + _waveInEvent.StartRecording(); + } + + public static AudioInStream Init() => new(); + + public override bool CanRead => true; + + public override bool CanSeek => false; + + public override bool CanWrite => false; + + public override long Length => throw new NotImplementedException(); + + public override long Position { get => throw new NotImplementedException(); set => throw new NotImplementedException(); } + + public override int Read(byte[] buffer, int offset, int count) + { + var total = count; + + while (GetAvailableBytes() < count) + { + Thread.Sleep(TIMEOUT_SECONDS); + } + + lock (_lock) + { + if (_bufferReadPtr + count >= _buffer.Length) + { + var chunkLength = _buffer.Length - _bufferReadPtr; + Array.Copy(_buffer, _bufferReadPtr, buffer, offset, chunkLength); + _bufferReadPtr = 0; + count -= chunkLength; + offset += chunkLength; + } + Array.Copy(_buffer, _bufferReadPtr, buffer, offset, count); + _bufferReadPtr += count; + } + + return total; + } + + private int GetAvailableBytes() + { + if (_bufferWritePtr >= _bufferReadPtr) + { + return _bufferWritePtr - _bufferReadPtr; + } + else + { + return _buffer.Length - _bufferReadPtr + _bufferWritePtr; + } + } + + public override void Flush() + { + throw new NotImplementedException(); + } + + public override long Seek(long offset, SeekOrigin origin) + { + throw new NotImplementedException(); + } + + public override void SetLength(long value) + { + throw new NotImplementedException(); + } + + public override void Write(byte[] buffer, int offset, int count) + { + throw new NotImplementedException(); + } + + protected override void Dispose(bool disposing = true) + { + _waveInEvent?.Dispose(); + base.Dispose(disposing); + } +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Audio/AudioOut.cs b/tests/BotSharp.Test.RealtimeVoice/Audio/AudioOut.cs new file mode 100644 index 00000000..ea5d71ad --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Audio/AudioOut.cs @@ -0,0 +1,51 @@ +using NAudio.Wave; + +namespace BotSharp.Test.RealtimeVoice.Audio; + +internal class AudioOut : IDisposable +{ + private const int SAMPLE_RATE = 24000; + private const int BYTES_PER_SAMPLE = 2; + private const int CHANNELS = 1; + private const int BUFFER_MINS = 10; + + private readonly BufferedWaveProvider _waveProvider; + private readonly WaveOutEvent _waveOutEvent; + + public AudioOut() + { + var audioFormat = new WaveFormat( + rate: SAMPLE_RATE, + bits: BYTES_PER_SAMPLE * 8, + channels: CHANNELS); + + _waveProvider = new BufferedWaveProvider(audioFormat) + { + BufferDuration = TimeSpan.FromMinutes(BUFFER_MINS), + DiscardOnBufferOverflow = true + }; + + _waveOutEvent = new WaveOutEvent() + { + DeviceNumber = 0 + }; + _waveOutEvent.Init(_waveProvider); + _waveOutEvent.Play(); + } + + public void Enqueue(BinaryData data) + { + var buffer = data?.ToArray() ?? []; + _waveProvider.AddSamples(buffer, 0, buffer.Length); + } + + public void ClearBuffer() + { + _waveProvider.ClearBuffer(); + } + + public void Dispose() + { + _waveOutEvent?.Dispose(); + } +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Enums/SessionMode.cs b/tests/BotSharp.Test.RealtimeVoice/Enums/SessionMode.cs new file mode 100644 index 00000000..dfd25c9f --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Enums/SessionMode.cs @@ -0,0 +1,7 @@ +namespace BotSharp.Test.RealtimeVoice.Enums; + +internal enum SessionMode +{ + StreamChannel = 1, + CustomStream = 2 +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Program.cs b/tests/BotSharp.Test.RealtimeVoice/Program.cs index 8a3ad997..5d52b992 100644 --- a/tests/BotSharp.Test.RealtimeVoice/Program.cs +++ b/tests/BotSharp.Test.RealtimeVoice/Program.cs @@ -1,136 +1,8 @@ -using BotSharp.Abstraction.Conversations.Enums; -using BotSharp.Abstraction.Conversations.Models; -using BotSharp.Abstraction.Conversations; using BotSharp.OpenAPI; -using System.Text.Json; using System.Reflection; var services = ServiceBuilder.CreateHostBuilder(Assembly.GetExecutingAssembly()); -var channel = services.GetRequiredService(); -Console.WriteLine("PCM-16 Microphone Capture (24kHz Sample Rate)"); -Console.WriteLine("-----------------------------------------------"); - -var convService = services.GetRequiredService(); -var conv = new Conversation -{ - AgentId = "01e2fc5c-2c89-4ec7-8470-7688608b496c", - Channel = ConversationChannel.Phone, - Title = $"Test", - Tags = [], -}; -conv = await convService.NewConversation(conv); - -await channel.ConnectAsync(conv.Id); - -var hub = services.GetRequiredService(); -var conn = hub.SetHubConnection(conv.Id); -conn.CurrentAgentId = conv.AgentId; - -conn.OnModelReady = () => - JsonSerializer.Serialize(new - { - @event = "init" - }); - -conn.OnModelMessageReceived = message => - JsonSerializer.Serialize(new - { - @event = "media", - media = message - }); - -conn.OnModelAudioResponseDone = () => - JsonSerializer.Serialize(new - { - @event = "mark", - mark = new { name = "responsePart" } - }); - -conn.OnModelUserInterrupted = () => - JsonSerializer.Serialize(new - { - @event = "interrupted" - }); - -conn.OnUserSpeechDetected = () => - JsonSerializer.Serialize(new - { - @event = "speech_detected" - }); - - -await hub.ConnectToModel(async data => -{ - var response = JsonSerializer.Deserialize(data); - if (response.Event == "speech_detected") - { - channel.ClearBuffer(); - } - else if (response.Event == "media") - { - var message = JsonSerializer.Deserialize(data); - await channel.SendAsync(Convert.FromBase64String(message.Media), CancellationToken.None); - } -}); - -StreamReceiveResult result; -var buffer = new byte[1024 * 32]; - -do -{ - var seg = new ArraySegment(buffer); - result = await channel.ReceiveAsync(seg, CancellationToken.None); - - await hub.Completer.AppenAudioBuffer(seg, result.Count); - - // Display the audio level - int audioLevel = CalculateAudioLevel(buffer, result.Count); - DisplayAudioLevel(audioLevel); -} while (result.Status == StreamChannelStatus.Open); - - -int CalculateAudioLevel(byte[] buffer, int bytesRecorded) -{ - // Simple audio level calculation (RMS) - int bytesPerSample = 2; // 16-bit PCM = 2 bytes per sample - int sampleCount = bytesRecorded / bytesPerSample; - if (sampleCount == 0) return 0; - - double sum = 0; - for (int i = 0; i < bytesRecorded; i += 2) - { - if (i + 1 < bytesRecorded) - { - short sample = (short)((buffer[i + 1] << 8) | buffer[i]); - double normalized = sample / (short.MaxValue * 1.0 + 1); - sum += normalized * normalized; - } - } - - double rms = Math.Sqrt(sum / sampleCount); - double db = 20 * Math.Log10(rms); - - if (double.IsInfinity(db) || double.IsNaN(db)) - { - return 0; - } - - db = Math.Clamp(db, -100, 0); - return (int)((db + 100) * 1); -} - -void DisplayAudioLevel(int level) -{ - const int sep = 50; - // Normalize level to 0-50 range for display - int displayLevel = (level * sep) / 100; - - // Clear the current line - Console.Write("\r" + new string(' ', 60)); - - // Display audio level as a bar - Console.Write("\rMicrophone: ["); - Console.Write(new string('#', displayLevel).PadRight(sep, ' ')); - Console.Write("]\r"); -} +var agentId = BuiltInAgentId.Chatbot; +var session = ConsoleChatSession.Init(services); +await session.StartAsync(agentId, SessionMode.StreamChannel); \ No newline at end of file diff --git a/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.CustomStream.cs b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.CustomStream.cs new file mode 100644 index 00000000..4d05f47f --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.CustomStream.cs @@ -0,0 +1,46 @@ +using System.Text.Json; + +namespace BotSharp.Test.RealtimeVoice.Session; + +internal partial class ConsoleChatSession +{ + /// + /// Start a new chat session via custom stream + /// + /// + /// + private async Task StartCustomStreamAsync(string agentId) + { + DisplayRemarks(); + + var (hub, conversationId) = await Setup(agentId); + var audioOut = new AudioOut(); + + await hub.ConnectToModel( + responseToUser: async data => + { + var response = JsonSerializer.Deserialize(data); + if (response.Event == "speech_detected") + { + audioOut.ClearBuffer(); + } + else if (response.Event == "media") + { + var message = JsonSerializer.Deserialize(data); + var binaryData = BinaryData.FromBytes(Convert.FromBase64String(message.Media)); + audioOut.Enqueue(binaryData); + } + }, + init: async data => + { + _ = Task.Run(async () => + { + using var audioIn = AudioInStream.Init(); + Console.WriteLine("\r\nListening microphone...\r\n"); + await SendAudio(hub, audioIn); + }); + }); + + while (true) { } + } +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.StreamChannel.cs b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.StreamChannel.cs new file mode 100644 index 00000000..311764b9 --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.StreamChannel.cs @@ -0,0 +1,50 @@ +using System.Text.Json; + +namespace BotSharp.Test.RealtimeVoice.Session; + +internal partial class ConsoleChatSession +{ + /// + /// Start a new chat session via stream channel + /// + /// + /// + private async Task StartStreamChannelAsync(string agentId) + { + DisplayRemarks(); + + var (hub, conversationId) = await Setup(agentId); + + var channel = _services.GetRequiredService(); + await channel.ConnectAsync(conversationId); + + await hub.ConnectToModel(async data => + { + var response = JsonSerializer.Deserialize(data); + if (response.Event == "speech_detected") + { + channel.ClearBuffer(); + } + else if (response.Event == "media") + { + var message = JsonSerializer.Deserialize(data); + await channel.SendAsync(Convert.FromBase64String(message.Media), CancellationToken.None); + } + }); + + StreamReceiveResult result; + var buffer = new byte[1024 * 32]; + + do + { + var seg = new ArraySegment(buffer); + result = await channel.ReceiveAsync(seg, CancellationToken.None); + + await hub.Completer.AppenAudioBuffer(seg, result.Count); + + // Display the audio level + int audioLevel = CalculateAudioLevel(buffer, result.Count); + DisplayAudioLevel(audioLevel); + } while (result.Status == StreamChannelStatus.Open); + } +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.cs b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.cs new file mode 100644 index 00000000..be5499b6 --- /dev/null +++ b/tests/BotSharp.Test.RealtimeVoice/Session/ConsoleChatSession.cs @@ -0,0 +1,171 @@ +using System.Buffers; +using System.Text.Json; + +namespace BotSharp.Test.RealtimeVoice.Session; + +internal partial class ConsoleChatSession +{ + private readonly IServiceProvider _services; + + private ConsoleChatSession( + IServiceProvider services) + { + _services = services; + } + + public static ConsoleChatSession Init(IServiceProvider services) + { + return new(services); + } + + /// + /// Start a new session + /// + /// + /// + /// + public async Task StartAsync(string agentId, SessionMode mode) + { + switch (mode) + { + case SessionMode.StreamChannel: + await StartStreamChannelAsync(agentId); + break; + case SessionMode.CustomStream: + await StartCustomStreamAsync(agentId); + break; + } + } + + private void DisplayRemarks() + { + Console.WriteLine("PCM-16 Microphone Capture (24kHz Sample Rate)"); + Console.WriteLine("-----------------------------------------------"); + } + + private async Task SendAudio(IRealtimeHub hub, Stream stream) + { + var buffer = ArrayPool.Shared.Rent(1024 * 16); + + try + { + while (true) + { + var bytesNum = await stream.ReadAsync(buffer, 0, buffer.Length, CancellationToken.None); + if (bytesNum == 0) break; + + var audioBytes = buffer.AsMemory(0, bytesNum); + var data = BinaryData.FromBytes(audioBytes); + await hub.Completer.AppenAudioBuffer(data.ToArray(), data.Length); + } + } + finally + { + ArrayPool.Shared.Return(buffer); + } + } + + /// + /// Create a new conversation and set up the events + /// + /// + /// + private async Task<(IRealtimeHub, string)> Setup(string agentId) + { + + var convService = _services.GetRequiredService(); + var hub = _services.GetRequiredService(); + + var conv = new Conversation + { + AgentId = agentId, + Channel = ConversationChannel.Phone, + Title = $"Test", + Tags = [], + }; + conv = await convService.NewConversation(conv); + + var conn = hub.SetHubConnection(conv.Id); + conn.CurrentAgentId = conv.AgentId; + + conn.OnModelReady = () => + JsonSerializer.Serialize(new + { + @event = "init" + }); + + conn.OnModelMessageReceived = message => + JsonSerializer.Serialize(new + { + @event = "media", + media = message + }); + + conn.OnModelAudioResponseDone = () => + JsonSerializer.Serialize(new + { + @event = "mark", + mark = new { name = "responsePart" } + }); + + conn.OnModelUserInterrupted = () => + JsonSerializer.Serialize(new + { + @event = "interrupted" + }); + + conn.OnUserSpeechDetected = () => + JsonSerializer.Serialize(new + { + @event = "speech_detected" + }); + + return (hub, conv.Id); + } + + + private int CalculateAudioLevel(byte[] buffer, int bytesRecorded) + { + // Simple audio level calculation (RMS) + int bytesPerSample = 2; // 16-bit PCM = 2 bytes per sample + int sampleCount = bytesRecorded / bytesPerSample; + if (sampleCount == 0) return 0; + + double sum = 0; + for (int i = 0; i < bytesRecorded; i += 2) + { + if (i + 1 < bytesRecorded) + { + short sample = (short)((buffer[i + 1] << 8) | buffer[i]); + double normalized = sample / (short.MaxValue * 1.0 + 1); + sum += normalized * normalized; + } + } + + double rms = Math.Sqrt(sum / sampleCount); + double db = 20 * Math.Log10(rms); + + if (double.IsInfinity(db) || double.IsNaN(db)) + { + return 0; + } + + db = Math.Clamp(db, -100, 0); + return (int)((db + 100) * 1); + } + + private void DisplayAudioLevel(int level) + { + const int sep = 50; + // Normalize level to 0-50 range for display + int displayLevel = (level * sep) / 100; + + // Clear the current line + Console.Write("\r" + new string(' ', 60)); + + // Display audio level as a bar + Console.Write("\rMicrophone: ["); + Console.Write(new string('#', displayLevel).PadRight(sep, ' ')); + Console.Write("]\r"); + } +} diff --git a/tests/BotSharp.Test.RealtimeVoice/Using.cs b/tests/BotSharp.Test.RealtimeVoice/Using.cs index 2a04d907..39b71ee2 100644 --- a/tests/BotSharp.Test.RealtimeVoice/Using.cs +++ b/tests/BotSharp.Test.RealtimeVoice/Using.cs @@ -6,6 +6,14 @@ global using BotSharp.Core; global using BotSharp.Core.Infrastructures; global using BotSharp.Core.Plugins; global using BotSharp.Logger; +global using BotSharp.Abstraction.Agents.Enums; global using BotSharp.Abstraction.Realtime.Models; global using BotSharp.Abstraction.Realtime; global using BotSharp.Abstraction.Realtime.Enums; +global using BotSharp.Abstraction.Conversations.Enums; +global using BotSharp.Abstraction.Conversations.Models; +global using BotSharp.Abstraction.Conversations; + +global using BotSharp.Test.RealtimeVoice.Session; +global using BotSharp.Test.RealtimeVoice.Audio; +global using BotSharp.Test.RealtimeVoice.Enums; \ No newline at end of file