Files
voice-cat/src/VoiceCat.Core/VoiceCatClient.cs
T
Talon f3ac779bf4
Build and test / test (macos-latest) (push) Canceled after 0s
Build and test / test (ubuntu-24.04) (push) Canceled after 0s
Build and test / test (windows-latest) (push) Canceled after 0s
Build and test / apple-client (push) Canceled after 0s
Survive changing networks and deepen the receive buffer
Media died silently whenever a client's source address changed. The relay bound a
peer's endpoint once and refused to move it, and the client stopped offering its
binding token after the first bind, so a Wi-Fi/cellular handover stranded the
session in both directions. Add an authenticated Rebind media frame: the binding
token travels in the clear for peer lookup only, and the AEAD tag over header and
token plus the peer's existing replay window are what authorize the move, so a
captured rebind cannot be replayed to redirect someone else's downlink. The client
rebuilds its UDP socket instead of retrying on one still pinned to a vanished
interface.

Nothing judged the control connection live: pings were sent and pongs ignored, so a
blackholed TCP path went unnoticed for minutes while the UI showed a live session.
Treat any server traffic as liveness and fail the connection when it stops, which
drives the existing reconnect.

The receive jitter buffer had lost its depth floor, so a channel without FEC or
DRED played out with no buffer at all and ordinary reordering became concealment.
Restore a one-frame floor, observe every arrival rather than only accepted ones —
a shallow buffer was rejecting the late arrivals that should have deepened it —
and allow playout to hold a frame so depth can follow a degrading link. A stalled
consumer now sheds the oldest queued packet instead of refusing the live talkspurt.

Add a deterministic network-impairment simulation covering bursty loss, jitter,
reordering, duplication, outages and a stalled consumer, a handover test against a
real relay, a replay test for the rebind path, and a blackholed control connection
driven through a freezable TCP proxy.
2026-09-24 19:15:16 +02:00

353 lines
20 KiB
C#

using System.Collections.Concurrent;
using System.Net;
using System.Net.Sockets;
using System.Threading.Channels;
using VoiceCat.Crypto;
using VoiceCat.Transport;
using VoiceCat.Protocol;
using VoiceCat.Audio;
using Voicecat.V1;
using Channel = Voicecat.V1.Channel;
namespace VoiceCat.Core;
public enum ClientConnectionState { Disconnected, Connecting, VerifyingIdentity, Authenticating, Connected }
public sealed record ServerIdentityChallenge(string Host, ushort Port, string CertificateFingerprint, TofuStatus Status);
public sealed partial class VoiceCatClient : IAsyncDisposable
{
private readonly string clientName;
private readonly string clientVersion;
private readonly TofuStore pins;
private readonly SemaphoreSlim lifecycle = new(1);
private readonly CancellationTokenSource disposed = new();
private readonly object stateGate = new();
private readonly ConcurrentDictionary<ulong, TaskCompletionSource<Envelope>> pending = new();
private readonly System.Threading.Channels.Channel<Envelope> events = System.Threading.Channels.Channel.CreateBounded<Envelope>(128);
private readonly Dictionary<uint, Channel> channels = [];
private readonly Dictionary<uint, User> users = [];
private readonly Dictionary<uint, StreamInfo> localStreams = [];
public AudioEngine Audio { get; }
public IReadOnlyList<StreamInfo> LocalStreams { get { lock (stateGate) return localStreams.Values.Select(s => s.Clone()).ToArray(); } }
private TlsControlConnection? control;
private MediaSessionCrypto? mediaCrypto;
private ClientMediaTransport? media;
private Task keepalive = Task.CompletedTask;
public event EncodedVoiceHandler? VoiceReceived;
private CancellationTokenSource? connectionLifetime;
private Task reader = Task.CompletedTask;
private long nextRequest;
private AuthResult? authentication;
private ServerHello? hello;
private uint adaptiveLossChannel;
private int adaptiveLossPercent = -1;
private ClientConnectionState state;
private long lastControlInbound;
// A control connection that stops answering is dead even though the socket still looks open.
// A phone that changes interface leaves TCP blackholed rather than reset, and the OS will not
// report it for minutes, so liveness is judged here instead.
private TimeSpan controlKeepaliveInterval = TimeSpan.FromSeconds(10);
private TimeSpan controlSilenceTimeout = TimeSpan.FromSeconds(30);
// Instance scoped so tests can shorten the window without disturbing parallel tests.
internal void SetControlLiveness(TimeSpan keepalive, TimeSpan silenceTimeout)
{ controlKeepaliveInterval = keepalive; controlSilenceTimeout = silenceTimeout; }
public event Action<ClientConnectionState>? ConnectionStateChanged;
public ClientConnectionState State { get { lock (stateGate) return state; } }
public Exception? ConnectionFailure { get; private set; }
public Task Completion => reader;
public AuthResult? Authentication { get { lock (stateGate) return authentication?.Clone(); } }
public ServerHello? ServerHello { get { lock (stateGate) return hello?.Clone(); } }
public bool SupportsAdaptivePacketLoss { get { lock (stateGate) return hello?.Features.Contains(ProtocolFeatures.AdaptivePacketLoss) == true; } }
public IReadOnlyList<Channel> Channels { get { lock (stateGate) return channels.Values.Select(c => c.Clone()).ToArray(); } }
public IReadOnlyList<User> Users { get { lock (stateGate) return users.Values.Select(u => u.Clone()).ToArray(); } }
public bool TryReadEvent(out Envelope? envelope) => events.Reader.TryRead(out envelope);
public IAsyncEnumerable<Envelope> ReadEventsAsync(CancellationToken cancellationToken = default) => events.Reader.ReadAllAsync(cancellationToken);
// deviceClockedAudio: the platform's audio callback drives Audio.RunCycle() itself (iOS
// keeps render callbacks running in the background while sleep-paced threads coalesce there),
// so the audio engine must not start its own pacing worker.
public VoiceCatClient(string clientName = "VoiceCat .NET", string clientVersion = "0.1.0", string? tofuStorePath = null, bool deviceClockedAudio = false)
{
this.clientName = clientName;
this.clientVersion = clientVersion;
pins = new(tofuStorePath ?? Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.LocalApplicationData), "VoiceCat", "tofu.txt"));
Audio = new(TrySendEncodedVoice, startWorker: !deviceClockedAudio);
VoiceReceived += Audio.Receive;
}
public async Task ConnectAsync(string host, ushort port, Func<ServerIdentityChallenge, CancellationToken, ValueTask<bool>>? confirmIdentity = null, CancellationToken cancellationToken = default)
{
ArgumentException.ThrowIfNullOrWhiteSpace(host);
ArgumentOutOfRangeException.ThrowIfZero(port);
await lifecycle.WaitAsync(cancellationToken).ConfigureAwait(false);
Socket? socket = null;
bool started = false;
try
{
ObjectDisposedException.ThrowIf(disposed.IsCancellationRequested, this);
if (control is not null) throw new InvalidOperationException("Disconnect before reconnecting.");
started = true;
ConnectionFailure = null;
connectionLifetime = CancellationTokenSource.CreateLinkedTokenSource(disposed.Token);
using var connecting = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, connectionLifetime.Token);
CancellationToken token = connecting.Token;
SetState(ClientConnectionState.Connecting);
socket = new(SocketType.Stream, ProtocolType.Tcp) { NoDelay = true };
await socket.ConnectAsync(host, port, token).ConfigureAwait(false);
string? fingerprint = null;
control = new(socket, TlsSession.CreateClient(value => { fingerprint = value; return true; }), connectionLifetime.Token);
socket = null; // Transport owns it from here.
mediaCrypto = await control.TakeMediaCryptoAsync(token).ConfigureAwait(false);
string certificatePin = fingerprint ?? throw new IOException("TLS did not report a certificate fingerprint.");
TofuStatus pinStatus = pins.Check(host, port, certificatePin);
if (pinStatus != TofuStatus.Matched)
{
SetState(ClientConnectionState.VerifyingIdentity);
if (confirmIdentity is null || !await confirmIdentity(new(host, port, certificatePin, pinStatus), token).ConfigureAwait(false))
throw new System.Security.Authentication.AuthenticationException("Server identity was rejected.");
pins.Pin(host, port, certificatePin);
}
reader = ReadAsync(control, connectionLifetime.Token);
var clientHello = new ClientHello { ProtoVersion = 2, ClientName = clientName, ClientVersion = clientVersion };
clientHello.Features.Add(ProtocolFeatures.AdaptivePacketLoss);
Envelope response = await RequestAsync(new() { ClientHello = clientHello }, token).ConfigureAwait(false);
if (response.ServerHello?.ProtoVersion != 2) throw new IOException("Unsupported server protocol.");
lock (stateGate) hello = response.ServerHello.Clone();
keepalive = KeepaliveAsync(connectionLifetime.Token);
SetState(ClientConnectionState.Authenticating);
}
catch
{
socket?.Dispose();
if (started) await CloseAsync().ConfigureAwait(false);
throw;
}
finally { lifecycle.Release(); }
}
public Task<AuthResult> AuthenticateGuestAsync(string nickname, CancellationToken cancellationToken = default) =>
AuthenticateAsync(new() { Guest = new() { Nickname = nickname } }, cancellationToken);
public Task<AuthResult> AuthenticateUserAsync(string username, string password, CancellationToken cancellationToken = default) =>
AuthenticateAsync(new() { Password = new() { Username = username, Password = password } }, cancellationToken);
private async Task<AuthResult> AuthenticateAsync(AuthRequest request, CancellationToken cancellationToken)
{
if (State != ClientConnectionState.Authenticating) throw new InvalidOperationException("Authentication requires a connected TLS session.");
Envelope response = await RequestAsync(new() { AuthRequest = request }, cancellationToken).ConfigureAwait(false);
AuthResult result = response.AuthResult ?? throw new IOException("Unexpected authentication response.");
if (result.Ok)
{
var endpoint = (IPEndPoint)control!.RemoteEndPoint;
IPAddress address = endpoint.Address.IsIPv4MappedToIPv6 ? endpoint.Address.MapToIPv4() : endpoint.Address;
media = new(new(address, checked((int)ServerHello!.UdpPort)), result.UdpToken.Span, mediaCrypto!, connectionLifetime!.Token);
media.Received += (header, packet) => VoiceReceived?.Invoke(header, packet);
SetState(ClientConnectionState.Connected);
}
return result;
}
public async Task<VoiceSubscriptionResult> SubscribeVoiceAsync(bool subscribe = true, CancellationToken cancellationToken = default)
{
if (subscribe && media is not null) await media.Bound.WaitAsync(TimeSpan.FromSeconds(5), cancellationToken).ConfigureAwait(false);
return (await RequestAsync(subscribe ? new() { SubscribeVoice = new() } : new() { UnsubscribeVoice = new() }, cancellationToken).ConfigureAwait(false)).VoiceSubscriptionResult;
}
public bool TrySendEncodedVoice(uint ssrc, uint timestamp, ReadOnlySpan<byte> payload, VoiceFrameFlags flags = VoiceFrameFlags.None) =>
media?.TrySend(new(MediaFrameType.Voice, flags, 0, ssrc, 0, timestamp), payload) == true;
public async Task<StreamInfo> StartStreamAsync(StreamKind kind, string label = "", int captureChannels = 1, CancellationToken cancellationToken = default)
{
if (State != ClientConnectionState.Connected) throw new InvalidOperationException("Client is disconnected.");
var response = (await RequestAsync(new() { StreamAnnounce = new() { Kind = kind, Label = label } }, cancellationToken).ConfigureAwait(false)).StreamAnnounceResult;
if (!response.Ok) throw new InvalidOperationException(response.Error);
var info = new StreamInfo { StreamId = response.StreamId, Ssrc = response.Ssrc, Kind = kind, Audio = response.EffectiveAudio.Clone(), Label = label };
try
{
lock (stateGate)
{
if (State != ClientConnectionState.Connected) throw new InvalidOperationException("Client disconnected during stream negotiation.");
Audio.AddLocalStream(info, captureChannels); localStreams[info.StreamId] = info;
User? self = authentication is null ? null : users.GetValueOrDefault(authentication.Self.Id, authentication.Self);
if (self is not null && self.ChannelId == adaptiveLossChannel && adaptiveLossPercent >= 0)
Audio.SetExpectedPacketLoss(adaptiveLossPercent);
}
return info.Clone();
}
catch { if (State == ClientConnectionState.Connected) Send(new() { StreamStop = new() { StreamId = info.StreamId } }); throw; }
}
public void StopStream(uint streamId)
{
lock (stateGate) { localStreams.Remove(streamId); Audio.RemoveLocalStream(streamId); }
Send(new() { StreamStop = new() { StreamId = streamId } });
}
private async Task KeepaliveAsync(CancellationToken cancellationToken)
{
try
{
Volatile.Write(ref lastControlInbound, Environment.TickCount64);
using var timer = new PeriodicTimer(controlKeepaliveInterval);
while (await timer.WaitForNextTickAsync(cancellationToken).ConfigureAwait(false))
{
// Every server reply counts as liveness, so a busy session never trips this; an
// unanswered ping is what exposes a path that has stopped carrying anything.
if (Environment.TickCount64 - Volatile.Read(ref lastControlInbound) > controlSilenceTimeout.TotalMilliseconds)
{
ConnectionFailure ??= new IOException("Server stopped responding on the control connection.");
connectionLifetime?.Cancel();
break;
}
Send(new() { Ping = new() { Nonce = checked((ulong)Environment.TickCount64) } });
}
}
catch (Exception exception) when (exception is OperationCanceledException or IOException or InvalidOperationException) { }
}
public async Task<Envelope> RequestAsync(Envelope request, CancellationToken cancellationToken = default)
{
TlsControlConnection connection = control ?? throw new InvalidOperationException("Client is disconnected.");
var completion = new TaskCompletionSource<Envelope>(TaskCreationOptions.RunContinuationsAsynchronously);
ulong id = checked((ulong)Interlocked.Increment(ref nextRequest));
Envelope outbound = request.Clone(); outbound.RequestId = id;
if (!pending.TryAdd(id, completion)) throw new InvalidOperationException("Request ids exhausted.");
try
{
if (!connection.TrySend(outbound)) throw new IOException("Control queue is full or closed.");
return await completion.Task.WaitAsync(TimeSpan.FromSeconds(15), cancellationToken).ConfigureAwait(false);
}
finally { pending.TryRemove(id, out _); }
}
public void Send(Envelope message)
{
TlsControlConnection connection = control ?? throw new InvalidOperationException("Client is disconnected.");
if (!connection.TrySend(message.Clone())) throw new IOException("Control queue is full or closed.");
}
private async Task ReadAsync(TlsControlConnection connection, CancellationToken cancellationToken)
{
Exception? failure = null;
try
{
await foreach (Envelope message in connection.ReadAsync(cancellationToken).ConfigureAwait(false))
{
Volatile.Write(ref lastControlInbound, Environment.TickCount64);
Apply(message);
if (message.RequestId != 0 && pending.TryRemove(message.RequestId, out var completion)) completion.TrySetResult(message.Clone());
if (!events.Writer.TryWrite(message.Clone())) throw new IOException("Client event queue exhausted; consume events regularly.");
if (message.Disconnect is not null) { connection.CompleteWrites(); break; }
}
}
// A liveness failure has already recorded the real cause and cancelled this read, so do
// not replace it with the cancellation it produced.
catch (Exception exception) { failure = exception; ConnectionFailure ??= exception; }
finally
{
connectionLifetime?.Cancel();
foreach (var operation in pending.Values) operation.TrySetException(failure ?? new IOException("Connection closed."));
SetState(ClientConnectionState.Disconnected);
}
}
private void Apply(Envelope message)
{
lock (stateGate)
{
if (message.AuthResult?.Ok == true) authentication = message.AuthResult.Clone();
if (message.ServerState is not null)
{
channels.Clear(); users.Clear();
foreach (var channel in message.ServerState.Channels) channels[channel.Id] = channel.Clone();
foreach (var user in message.ServerState.Users) users[user.Id] = user.Clone();
}
if (message.ChannelEvent is not null)
{
if (message.ChannelEvent.Kind == ChannelEvent.Types.Kind.Deleted) channels.Remove(message.ChannelEvent.DeletedId);
else if (message.ChannelEvent.Channel is not null)
{
PacketLossMode previousMode = channels.GetValueOrDefault(message.ChannelEvent.Channel.Id)?.Audio.PacketLossMode
?? PacketLossMode.PacketLossManual;
channels[message.ChannelEvent.Channel.Id] = message.ChannelEvent.Channel.Clone();
User? self = authentication is null ? null : users.GetValueOrDefault(authentication.Self.Id, authentication.Self);
if (self?.ChannelId == message.ChannelEvent.Channel.Id && previousMode != message.ChannelEvent.Channel.Audio.PacketLossMode)
{ adaptiveLossChannel = 0; adaptiveLossPercent = -1; }
}
}
if (message.UserEvent is not null)
{
if (message.UserEvent.Kind == UserEvent.Types.Kind.Left) users.Remove(message.UserEvent.LeftId);
else if (message.UserEvent.User is not null) users[message.UserEvent.User.Id] = message.UserEvent.User.Clone();
}
if (message.PacketLossUpdate is not null && authentication is not null)
{
User self = users.GetValueOrDefault(authentication.Self.Id, authentication.Self);
if (self.ChannelId == message.PacketLossUpdate.ChannelId)
{
adaptiveLossChannel = self.ChannelId;
adaptiveLossPercent = checked((int)message.PacketLossUpdate.AppliedPercent);
Audio.SetExpectedPacketLoss(checked((int)message.PacketLossUpdate.AppliedPercent));
}
}
if (authentication is not null && (message.ServerState is not null || message.UserEvent is not null))
{
User self = users.GetValueOrDefault(authentication.Self.Id, authentication.Self);
Audio.SetRemoteStreams(users.Values.ToArray(), self.Id, self.ChannelId);
foreach (var id in localStreams.Keys.Where(id => !self.Streams.Any(s => s.StreamId == id)).ToArray())
{ Audio.RemoveLocalStream(id); localStreams.Remove(id); }
}
}
}
private void SetState(ClientConnectionState value)
{
lock (stateGate) state = value;
ConnectionStateChanged?.Invoke(value);
}
public async Task DisconnectAsync()
{
connectionLifetime?.Cancel();
await lifecycle.WaitAsync().ConfigureAwait(false);
try { await CloseAsync().ConfigureAwait(false); }
finally { lifecycle.Release(); }
}
private async Task CloseAsync()
{
connectionLifetime?.Cancel();
try { await reader.ConfigureAwait(false); }
finally
{
try { await keepalive.ConfigureAwait(false); }
catch (Exception exception) when (exception is IOException or OperationCanceledException or SocketException or ObjectDisposedException) { }
try { if (media is not null) await media.DisposeAsync().ConfigureAwait(false); }
catch (Exception exception) when (exception is IOException or OperationCanceledException or SocketException or ObjectDisposedException) { }
try { if (control is not null) await control.DisposeAsync().ConfigureAwait(false); }
catch (Exception exception) when (exception is IOException or OperationCanceledException or SocketException or ObjectDisposedException) { }
control = null;
media = null;
mediaCrypto?.Dispose(); mediaCrypto = null;
connectionLifetime?.Dispose(); connectionLifetime = null;
lock (stateGate) { authentication = null; hello = null; channels.Clear(); users.Clear(); adaptiveLossChannel = 0; adaptiveLossPercent = -1; }
lock (stateGate)
{
foreach (var id in localStreams.Keys) Audio.RemoveLocalStream(id);
localStreams.Clear(); Audio.SetRemoteStreams([], 0, 0);
}
SetState(ClientConnectionState.Disconnected);
}
}
public async ValueTask DisposeAsync()
{
if (disposed.IsCancellationRequested) return;
disposed.Cancel();
await DisconnectAsync().ConfigureAwait(false);
events.Writer.TryComplete();
Audio.Dispose();
}
}