BotSharp/src/Infrastructure/BotSharp.Core/Session/BotSharpRealtimeSession.cs

88 lines
2.6 KiB
C#
Raw Normal View History

2025-05-01 17:47:30 +00:00
using System.ClientModel;
2025-05-07 21:34:39 +00:00
using System.Net.WebSockets;
2025-05-01 17:47:30 +00:00
using System.Runtime.CompilerServices;
2025-05-07 21:34:39 +00:00
namespace BotSharp.Core.Session;
2025-05-01 17:47:30 +00:00
public class BotSharpRealtimeSession : IDisposable
{
private readonly IServiceProvider _services;
private readonly WebSocket _websocket;
private readonly ChatSessionOptions? _sessionOptions;
private readonly object _singleReceiveLock = new();
private AsyncWebsocketDataCollectionResult _receivedCollectionResult;
2025-05-19 19:54:40 +00:00
private bool _disposed = false;
2025-05-01 17:47:30 +00:00
public BotSharpRealtimeSession(
IServiceProvider services,
WebSocket websocket,
ChatSessionOptions? sessionOptions)
{
_services = services;
_websocket = websocket;
_sessionOptions = sessionOptions;
}
public async IAsyncEnumerable<ChatSessionUpdate> ReceiveUpdatesAsync([EnumeratorCancellation] CancellationToken cancellationToken = default)
{
await foreach (ClientResult result in ReceiveInnerUpdatesAsync(cancellationToken))
{
var update = HandleSessionResult(result);
yield return update;
}
}
private async IAsyncEnumerable<ClientResult> ReceiveInnerUpdatesAsync([EnumeratorCancellation] CancellationToken cancellationToken = default)
{
lock (_singleReceiveLock)
{
_receivedCollectionResult ??= new(_websocket, _sessionOptions, cancellationToken);
}
await foreach (var result in _receivedCollectionResult)
{
yield return result;
}
}
private ChatSessionUpdate HandleSessionResult(ClientResult result)
{
using var response = result.GetRawResponse();
var bytes = response.Content.ToArray();
var text = Encoding.UTF8.GetString(bytes, 0, bytes.Length);
return new ChatSessionUpdate
{
RawResponse = text
};
}
2025-05-14 23:14:01 +00:00
public async Task SendEventAsync(string message)
2025-05-01 17:47:30 +00:00
{
2025-05-19 19:54:40 +00:00
if (_disposed || _websocket.State != WebSocketState.Open)
2025-05-01 17:47:30 +00:00
{
2025-05-19 19:54:40 +00:00
return;
2025-05-01 17:47:30 +00:00
}
2025-05-19 19:54:40 +00:00
var buffer = Encoding.UTF8.GetBytes(message);
await _websocket.SendAsync(new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None);
2025-05-01 17:47:30 +00:00
}
2025-05-14 23:14:01 +00:00
public async Task DisconnectAsync()
2025-05-01 17:47:30 +00:00
{
2025-05-19 19:54:40 +00:00
if (_disposed || _websocket.State != WebSocketState.Open)
2025-05-01 17:47:30 +00:00
{
2025-05-19 19:54:40 +00:00
return;
2025-05-01 17:47:30 +00:00
}
2025-05-19 19:54:40 +00:00
2025-09-23 18:11:31 +00:00
await _websocket.CloseAsync(WebSocketCloseStatus.NormalClosure, "Normal Closure", CancellationToken.None);
2025-05-01 17:47:30 +00:00
}
public void Dispose()
{
2025-05-19 19:54:40 +00:00
if (_disposed) return;
_disposed = true;
2025-05-01 17:47:30 +00:00
_websocket.Dispose();
}
}