add custom stream mode

This commit is contained in:
Jicheng Lu 2025-05-02 13:57:03 -05:00
parent d703477b4b
commit ed39786539
9 changed files with 459 additions and 132 deletions

View file

@ -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<IRealtimeHook>(_services, async hook => await hook.OnModelReady(agent, _completer));
await (init?.Invoke(data) ?? Task.CompletedTask);
},
onModelAudioDeltaReceived: async (audioDeltaData, itemId) =>
{

View file

@ -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);
}
}

View file

@ -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();
}
}

View file

@ -0,0 +1,7 @@
namespace BotSharp.Test.RealtimeVoice.Enums;
internal enum SessionMode
{
StreamChannel = 1,
CustomStream = 2
}

View file

@ -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<IStreamChannel>();
Console.WriteLine("PCM-16 Microphone Capture (24kHz Sample Rate)");
Console.WriteLine("-----------------------------------------------");
var convService = services.GetRequiredService<IConversationService>();
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<IRealtimeHub>();
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<ModelResponseEvent>(data);
if (response.Event == "speech_detected")
{
channel.ClearBuffer();
}
else if (response.Event == "media")
{
var message = JsonSerializer.Deserialize<ModelResponseMediaEvent>(data);
await channel.SendAsync(Convert.FromBase64String(message.Media), CancellationToken.None);
}
});
StreamReceiveResult result;
var buffer = new byte[1024 * 32];
do
{
var seg = new ArraySegment<byte>(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);

View file

@ -0,0 +1,46 @@
using System.Text.Json;
namespace BotSharp.Test.RealtimeVoice.Session;
internal partial class ConsoleChatSession
{
/// <summary>
/// Start a new chat session via custom stream
/// </summary>
/// <param name="agentId"></param>
/// <returns></returns>
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<ModelResponseEvent>(data);
if (response.Event == "speech_detected")
{
audioOut.ClearBuffer();
}
else if (response.Event == "media")
{
var message = JsonSerializer.Deserialize<ModelResponseMediaEvent>(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) { }
}
}

View file

@ -0,0 +1,50 @@
using System.Text.Json;
namespace BotSharp.Test.RealtimeVoice.Session;
internal partial class ConsoleChatSession
{
/// <summary>
/// Start a new chat session via stream channel
/// </summary>
/// <param name="agentId"></param>
/// <returns></returns>
private async Task StartStreamChannelAsync(string agentId)
{
DisplayRemarks();
var (hub, conversationId) = await Setup(agentId);
var channel = _services.GetRequiredService<IStreamChannel>();
await channel.ConnectAsync(conversationId);
await hub.ConnectToModel(async data =>
{
var response = JsonSerializer.Deserialize<ModelResponseEvent>(data);
if (response.Event == "speech_detected")
{
channel.ClearBuffer();
}
else if (response.Event == "media")
{
var message = JsonSerializer.Deserialize<ModelResponseMediaEvent>(data);
await channel.SendAsync(Convert.FromBase64String(message.Media), CancellationToken.None);
}
});
StreamReceiveResult result;
var buffer = new byte[1024 * 32];
do
{
var seg = new ArraySegment<byte>(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);
}
}

View file

@ -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);
}
/// <summary>
/// Start a new session
/// </summary>
/// <param name="agentId"></param>
/// <param name="mode"></param>
/// <returns></returns>
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<byte>.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<byte>.Shared.Return(buffer);
}
}
/// <summary>
/// Create a new conversation and set up the events
/// </summary>
/// <param name="agentId"></param>
/// <returns></returns>
private async Task<(IRealtimeHub, string)> Setup(string agentId)
{
var convService = _services.GetRequiredService<IConversationService>();
var hub = _services.GetRequiredService<IRealtimeHub>();
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");
}
}

View file

@ -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;