From 274b85025c005beec4696ae367baca7d5e10225d Mon Sep 17 00:00:00 2001 From: Talon Date: Tue, 15 Sep 2026 22:58:16 +0200 Subject: [PATCH] Add configurable media-aware managed session reaping --- CLAUDE.md | 2 +- PROGRESS.md | 20 ++- docs/api-dotnet.md | 17 ++- docs/porting-to-dotnet.md | 8 +- docs/protocol.md | 4 + docs/roadmap.md | 6 +- dotnet/README.md | 8 +- .../VoiceCat.Server/Transport/MediaRelay.cs | 10 +- .../Transport/SessionActivity.cs | 18 +++ .../Transport/TlsControlConnection.cs | 6 +- dotnet/src/VoiceCat.Server/VoiceServer.cs | 70 ++++++++-- .../src/VoiceCat.Server/VoiceServerOptions.cs | 21 +++ dotnet/tests/VoiceCat.Tests/ReaperTests.cs | 129 ++++++++++++++++++ dotnet/tests/VoiceCat.Tests/ServerTests.cs | 4 +- 14 files changed, 291 insertions(+), 32 deletions(-) create mode 100644 dotnet/src/VoiceCat.Server/Transport/SessionActivity.cs create mode 100644 dotnet/src/VoiceCat.Server/VoiceServerOptions.cs create mode 100644 dotnet/tests/VoiceCat.Tests/ReaperTests.cs diff --git a/CLAUDE.md b/CLAUDE.md index 6c08448..c3a5642 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -38,7 +38,7 @@ dotnet test dotnet/VoiceCat.slnx -c Release --no-build See `dotnet/README.md` for C# conventions and required native voice/CLI conformance, and `docs/api-dotnet.md` for managed interfaces. Phase 4 remains in progress; -server administration/reaping and audio/client/UI phases are still pending. +media-aware reaping is implemented; server administration and audio/client/UI phases remain pending. The default development preset is **`dev`** — it builds everything (server + tools + tests) with real vcpkg deps. The `skeleton` preset (no deps, stubs only) is a fast smoke check; see diff --git a/PROGRESS.md b/PROGRESS.md index 7d18f56..b6a047c 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -10,6 +10,24 @@ up instantly. Newest status at the top. ## ▶ Where we left off / next action +- **Done (2026-09-15): Managed media-aware reaper.** Voice checkpoint committed as + `05eacb3`. Added `VoiceServerOptions` (name/guests/capacity, handshake deadline, + idle timeout and sweep interval), preserving the previous constructor overload. + Default expiry/sweep are 45 s / 15 s; zero idle timeout disables reaping. Parsed TCP + envelopes, authenticated owned-stream voice and exact bound-endpoint UDP keepalives + refresh one monotonic session timestamp. Invalid media does not refresh it. Removed + the independent 60-second TCP-only timeout so media-active sessions remain connected. + Reaping sends a fatal disconnect and retires presence/media routing with one LEFT + event. Shutdown awaits active control/media/reaper loops and unfinished handshakes. + Tests inject `TimeProvider` timestamps to verify TCP/UDP activity, rejected media, + single departure events, disabled reaping and shutdown. **Verified:** 160/160 managed + tests with all native conformance enabled; managed Release build has zero warnings; + native dev build and 29/29 CTest tests green; `git diff --check` passes. + **Next:** protected channel joins and channel CRUD, permissions/moderation/account + administration, then production configuration/publishing. Phase 4 remains in progress. + Follow with audio/core/managed CLI, switch Windows to the managed library, then C# + AppKit/UIKit clients. Keep the Swift broadcast extension and frozen shared-ring boundary. + - **Done (2026-09-15): Phase 4 encrypted voice checkpoint.** Existing pending codec/DSP and initial server work committed as `4067bab`. Managed server now binds UDP on the TCP port number, issues 16-byte session tokens, supports subscription and multi-stream @@ -136,7 +154,7 @@ dotnet test dotnet/VoiceCat.slnx -c Release --no-restore ./dotnet/check-licenses.ps1 ``` -Last results: **154/154 managed tests, no skips with those variables set; 29/29 native +Last results: **160/160 managed tests, no skips with those variables set; 29/29 native CTest tests; warning-free Release build; 22 package licenses approved; locked restore and `git diff --check` passed.** Without the variables, native interoperability tests skip; that is not equivalent verification. Desktop CI stages codec/DSP bindings on diff --git a/docs/api-dotnet.md b/docs/api-dotnet.md index 8904f20..008afc5 100644 --- a/docs/api-dotnet.md +++ b/docs/api-dotnet.md @@ -181,7 +181,14 @@ owns each `TlsSession`; handlers exchange envelopes through bounded queues. This checkpoint caps connections at 64, queued input at 32 envelopes, queued output at 64 envelopes, and each control payload at 64 KiB (stricter than the shared framer's 16 MiB limit). Queue exhaustion disconnects slow consumers. Handshake timeout is -15 seconds; completed TLS connections have a 60-second receive-idle timeout. +15 seconds by default. Completed TLS connections use the server's media-aware reaper. + +The existing `VoiceServer(directory, endpoint, allowGuests, name)` constructor remains +available. An overload accepts `VoiceServerOptions` and an optional `TimeProvider`. +Options configure server name, guest access, connection limit (default 64), handshake +timeout (15 seconds), idle timeout (45 seconds) and reaper interval (15 seconds). +Zero idle timeout disables reaping; active reaping requires a positive interval. +Invalid options fail before creating credentials, databases or sockets. Authentication starts users in unprotected Lobby (id 1), subject to its capacity. Success returns permissions, then a cloned snapshot; peers receive joined/updated/left @@ -218,8 +225,12 @@ Control handlers publish immutable routing snapshots. Crypto is created within t TLS owner loop and transferred once. A coalesced notification wakes retired-key cleanup. The synchronous fan-out core allocates zero managed bytes with platform ChaCha20-Poly1305; socket scheduling and the allocating BouncyCastle fallback are outside that guarantee. -UDP keepalives are echoed for bound endpoints. The full media-aware reaper remains pending; -the existing 60-second TLS receive-idle timeout still applies. +UDP keepalives are echoed for bound endpoints. Any parsed control envelope, authenticated +voice from an active owned stream, or exact header-only keepalive from a bound endpoint +refreshes a shared monotonic activity timestamp. Invalid media does not refresh it. +The reaper sends a fatal disconnect, removes presence/routing and broadcasts one LEFT +event. Valid UDP activity keeps a TCP-idle client alive. Shutdown cancels and awaits +the accept, reaper, control and media loops before disposing credentials/storage. `AccountStore(path)` retains the C++ schema version 2, accepts version 1 migration, and rejects unknown revisions. Opening an existing channel table does not reseed it. diff --git a/docs/porting-to-dotnet.md b/docs/porting-to-dotnet.md index 04e83ae..82653bb 100644 --- a/docs/porting-to-dotnet.md +++ b/docs/porting-to-dotnet.md @@ -745,8 +745,12 @@ BouncyCastle crypto fallback are excluded from its zero-allocation guarantee. Two real C++ `vccli` processes also pass join/text/bidirectional voice tests using finite `--test-tone-ms` external capture/playback. The transport load test delivers all 2,500 recipient packets from a sender paced at 50 pps to 50 subscribers. -Administration, protected joins, production configuration and media-aware reaping -still remain before Phase 4 completion. +**Reaper checkpoint:** configurable 45-second idle expiry / 15-second sweep replaces +the TCP-only idle timeout. Control envelopes, authenticated voice and bound-endpoint +keepalives refresh shared monotonic activity; invalid media does not. Reaping removes +presence and media routing, and can be disabled. Tests inject a clock to cover silent +clients, UDP-only activity, forged media, single departure events and disabled expiry. +Administration, protected joins and production configuration remain before Phase 4 completion. 1. `VoiceCat.Server`: accept loop, `ConnSession` protocol handling, session registry. 2. `Db` on `Microsoft.Data.Sqlite` — same schema. **Resolve the Argon2id hash-compat diff --git a/docs/protocol.md b/docs/protocol.md index a9716ee..8e4d2a1 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -305,6 +305,10 @@ message TextMessage { laptop sleep) that never produce a TCP EOF are cleaned up, and peers' audio engines `remove_stream` and stop PLC. The timeout and sweep interval are configurable via `server::Config::reaper_timeout_ms` / `reaper_sweep_ms` (set to 0 to disable). + The managed server uses `VoiceServerOptions.IdleTimeout` / `ReaperInterval` with the + same 45-second / 15-second defaults (zero idle timeout disables reaping). It refreshes + activity on parsed control envelopes, authenticated owned-stream voice, and exact + bound-endpoint keepalives; rejected media does not refresh activity. Timing is monotonic. - **UDP:** a separate lightweight keepalive on the media channel (voice.md §6) keeps NAT bindings alive and detects media-path failure independently of the control channel. - **Graceful disconnect.** A client ending its session sends `Disconnect { code = 0; diff --git a/docs/roadmap.md b/docs/roadmap.md index 45e8672..99e9b2a 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -16,8 +16,10 @@ packaging and TLS/server/client migration remain later checkpoints. desktop packaging, and managed control/UDP server slices are implemented. Two C++ `vccli` processes authenticate, join, chat and exchange mono/stereo voice through the managed server. The 50-subscriber fan-out core has an allocation regression test. -- **Next:** finish managed server administration, protected joins, production configuration - and media-aware reaping, then audio/client core, Windows cutover, C# AppKit and UIKit. +- **Media-aware reaping:** monotonic control/valid-UDP activity, configurable 45-second + idle timeout / 15-second sweep, and graceful shutdown are implemented and tested. +- **Next:** finish managed server administration, protected joins and production configuration, + then audio/client core, Windows cutover, C# AppKit and UIKit. Keep the Swift ReplayKit extension and its shared ring; defer C++ removal until parity. - See `docs/porting-to-dotnet.md` and `dotnet/README.md`. diff --git a/dotnet/README.md b/dotnet/README.md index 27e7f86..942b6c3 100644 --- a/dotnet/README.md +++ b/dotnet/README.md @@ -140,8 +140,12 @@ the optional `VOICECAT_BUILD_DOTNET_ORACLE=ON` configure flag and a real-deps bu The server also advertises UDP on the TCP port number, supports voice subscription and stream signaling, and reseals encoded audio for subscribers in the same channel. UDP binding fixes the first endpoint for the session; reconnect after endpoint changes. -Protected joins, administration, moderation, production configuration and full -media-aware reaping remain before Phase 4 completion. +Protected joins, administration, moderation and production configuration remain +before Phase 4 completion. The server's media-aware reaper defaults to 45 seconds +of inactivity with a 15-second sweep. Parsed control envelopes, valid encrypted +voice and keepalives from bound endpoints refresh activity; invalid media does not. +`VoiceServerOptions` configures timeouts and capacity; zero idle timeout disables +reaping. The constructor overload accepts `TimeProvider` for deterministic expiry tests. Enable deterministic native voice interoperability (no audio hardware required): diff --git a/dotnet/src/VoiceCat.Server/Transport/MediaRelay.cs b/dotnet/src/VoiceCat.Server/Transport/MediaRelay.cs index e6faf8a..01d06bc 100644 --- a/dotnet/src/VoiceCat.Server/Transport/MediaRelay.cs +++ b/dotnet/src/VoiceCat.Server/Transport/MediaRelay.cs @@ -7,10 +7,11 @@ using VoiceCat.Protocol; namespace VoiceCat.Server.Transport; -internal sealed class MediaPeer(byte[] token, MediaSessionCrypto crypto) +internal sealed class MediaPeer(byte[] token, MediaSessionCrypto crypto, SessionActivity? activity = null) { public byte[] Token { get; } = token; public MediaSessionCrypto Crypto { get; } = crypto; + public SessionActivity Activity { get; } = activity ?? new(TimeProvider.System); // Only the UDP loop reads or changes the endpoint and binding state. public SocketAddress? Endpoint { get; set; } public void Dispose() { Crypto.Dispose(); CryptographicOperations.ZeroMemory(Token); } @@ -101,10 +102,15 @@ internal sealed class MediaRelay : IAsyncDisposable if (source is null) continue; if (header.Type == MediaFrameType.Keepalive) { - if (length == VoiceFrameHeader.Size) await SendAsync(input.AsMemory(0, length), sender).ConfigureAwait(false); + if (length == VoiceFrameHeader.Size) + { + source.Peer.Activity.Touch(); + await SendAsync(input.AsMemory(0, length), sender).ConfigureAwait(false); + } continue; } if (!fanout.TryStart(input.AsSpan(0, length), source, current)) continue; + source.Peer.Activity.Touch(); while (fanout.TryNext(out ReadOnlyMemory packet, out SocketAddress? endpoint)) await SendAsync(packet, endpoint!).ConfigureAwait(false); } diff --git a/dotnet/src/VoiceCat.Server/Transport/SessionActivity.cs b/dotnet/src/VoiceCat.Server/Transport/SessionActivity.cs new file mode 100644 index 0000000..c8dea2a --- /dev/null +++ b/dotnet/src/VoiceCat.Server/Transport/SessionActivity.cs @@ -0,0 +1,18 @@ +namespace VoiceCat.Server.Transport; + +internal sealed class SessionActivity(TimeProvider clock) +{ + private long lastSeen = clock.GetTimestamp(); + public void Touch() + { + long now = clock.GetTimestamp(); + long previous = Volatile.Read(ref lastSeen); + while (now > previous) + { + long observed = Interlocked.CompareExchange(ref lastSeen, now, previous); + if (observed == previous) return; + previous = observed; + } + } + public bool IsExpired(TimeSpan timeout) => clock.GetElapsedTime(Volatile.Read(ref lastSeen)) >= timeout; +} diff --git a/dotnet/src/VoiceCat.Server/Transport/TlsControlConnection.cs b/dotnet/src/VoiceCat.Server/Transport/TlsControlConnection.cs index d4ad006..f41e765 100644 --- a/dotnet/src/VoiceCat.Server/Transport/TlsControlConnection.cs +++ b/dotnet/src/VoiceCat.Server/Transport/TlsControlConnection.cs @@ -27,12 +27,12 @@ internal sealed class TlsControlConnection : IAsyncDisposable public Task Completion { get; } public CancellationToken CancellationToken => lifetime.Token; - internal TlsControlConnection(Socket socket, TlsSession tls, CancellationToken cancellationToken) + internal TlsControlConnection(Socket socket, TlsSession tls, CancellationToken cancellationToken, TimeSpan? handshakeTimeout = null) { this.socket = socket; this.tls = tls; lifetime = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); - lifetime.CancelAfter(TimeSpan.FromSeconds(15)); + lifetime.CancelAfter(handshakeTimeout ?? TimeSpan.FromSeconds(15)); Completion = RunAsync(); } @@ -94,8 +94,8 @@ internal sealed class TlsControlConnection : IAsyncDisposable try { mediaCrypto = new(encryptor, tls.CreateMediaDecryptor()); } catch { encryptor.Dispose(); throw; } mediaReady.SetResult(); + lifetime.CancelAfter(Timeout.InfiniteTimeSpan); } - if (tls.IsReady) lifetime.CancelAfter(TimeSpan.FromSeconds(60)); while ((count = tls.ReadPlaintext(plaintext)) > 0) Parse(plaintext.AsSpan(0, count)); await FlushAsync(sendBuffer, cancellationToken).ConfigureAwait(false); receive = socket.ReceiveAsync(ciphertext, SocketFlags.None, cancellationToken).AsTask(); diff --git a/dotnet/src/VoiceCat.Server/VoiceServer.cs b/dotnet/src/VoiceCat.Server/VoiceServer.cs index 1bb71ae..0574f41 100644 --- a/dotnet/src/VoiceCat.Server/VoiceServer.cs +++ b/dotnet/src/VoiceCat.Server/VoiceServer.cs @@ -19,6 +19,8 @@ public sealed class VoiceServer : IAsyncDisposable private readonly IReadOnlyList channels; private readonly bool allowGuests; private readonly string name; + private readonly VoiceServerOptions options; + private readonly TimeProvider clock; private readonly CancellationTokenSource shutdown = new(); private readonly object gate = new(); private readonly Dictionary sessions = []; @@ -27,6 +29,7 @@ public sealed class VoiceServer : IAsyncDisposable private uint nextUser; private uint nextSsrc; private readonly Task accepting; + private readonly Task reaping; private int disposed; public IPEndPoint EndPoint => (IPEndPoint)listener.LocalEndPoint!; @@ -34,9 +37,16 @@ public sealed class VoiceServer : IAsyncDisposable public event Action? ConnectionFailed; public VoiceServer(string directory, IPEndPoint endpoint, bool allowGuests = true, string name = "VoiceCat Server") + : this(directory, endpoint, new VoiceServerOptions { AllowGuests = allowGuests, Name = name }) { } + + public VoiceServer(string directory, IPEndPoint endpoint, VoiceServerOptions options, TimeProvider? timeProvider = null) { - this.allowGuests = allowGuests; - this.name = name; + ArgumentNullException.ThrowIfNull(options); + options.Validate(); + this.options = options; + clock = timeProvider ?? TimeProvider.System; + allowGuests = options.AllowGuests; + name = options.Name; credentials = ServerCredentials.LoadOrCreate(directory, name); try { @@ -44,7 +54,7 @@ public sealed class VoiceServer : IAsyncDisposable channels = accounts.LoadChannels(); listener = new Socket(endpoint.AddressFamily, SocketType.Stream, ProtocolType.Tcp); listener.Bind(endpoint); - listener.Listen(64); + listener.Listen(options.MaximumConnections); media = new((IPEndPoint)listener.LocalEndPoint!); media.Failed += exception => ConnectionFailed?.Invoke(exception); } @@ -57,6 +67,7 @@ public sealed class VoiceServer : IAsyncDisposable throw; } accepting = AcceptAsync(); + reaping = ReapAsync(); } private async Task AcceptAsync() @@ -68,11 +79,11 @@ public sealed class VoiceServer : IAsyncDisposable Socket socket = await listener.AcceptAsync(shutdown.Token).ConfigureAwait(false); lock (gate) { - if (sessions.Count >= 64) { socket.Dispose(); continue; } + if (sessions.Count >= options.MaximumConnections) { socket.Dispose(); continue; } socket.NoDelay = true; string address = ((IPEndPoint)socket.RemoteEndPoint!).Address.ToString(); - var connection = new TlsControlConnection(socket, credentials.CreateTlsSession(), shutdown.Token); - var session = new Session(++nextSession, connection, address); + var connection = new TlsControlConnection(socket, credentials.CreateTlsSession(), shutdown.Token, options.HandshakeTimeout); + var session = new Session(++nextSession, connection, address, new(clock)); sessions.Add(session.Id, session); connections.RemoveAll(task => task.IsCompleted); connections.Add(HandleAsync(session)); @@ -88,6 +99,7 @@ public sealed class VoiceServer : IAsyncDisposable { await foreach (Envelope envelope in session.Connection.ReadAsync(shutdown.Token).ConfigureAwait(false)) { + session.Activity.Touch(); if (envelope.Ping is not null) { session.Connection.TrySend(new() { RequestId = envelope.RequestId, Pong = new() { Nonce = envelope.Ping.Nonce } }); @@ -101,7 +113,7 @@ public sealed class VoiceServer : IAsyncDisposable Reject(session, "Unsupported protocol version or banned address."); break; } - session.Media = new(RandomNumberGenerator.GetBytes(16), await session.Connection.TakeMediaCryptoAsync(shutdown.Token).ConfigureAwait(false)); + session.Media = new(RandomNumberGenerator.GetBytes(16), await session.Connection.TakeMediaCryptoAsync(shutdown.Token).ConfigureAwait(false), session.Activity); var hello = new ServerHello { ProtoVersion = 2, ServerName = name, ServerVersion = "0.1.0-dotnet", UdpPort = checked((uint)media.EndPoint.Port), ServerIdentityFingerprint = ByteString.CopyFrom(SHA256.HashData(credentials.Identity.PublicKey)) }; if (allowGuests) hello.AuthMethods.Add("guest"); hello.AuthMethods.Add("password"); @@ -159,6 +171,28 @@ public sealed class VoiceServer : IAsyncDisposable session.Connection.CompleteWrites(); } + private async Task ReapAsync() + { + if (options.IdleTimeout == TimeSpan.Zero) return; + using var timer = new PeriodicTimer(options.ReaperInterval, clock); + try + { + while (await timer.WaitForNextTickAsync(shutdown.Token).ConfigureAwait(false)) + { + lock (gate) + { + foreach (Session session in sessions.Values) + { + if (session.Closing || !session.Activity.IsExpired(options.IdleTimeout)) continue; + session.Closing = true; + Reject(session, "Receive idle timeout."); + } + } + } + } + catch (OperationCanceledException) when (shutdown.IsCancellationRequested) { } + } + private async Task AuthenticateAsync(Session session, ulong requestId, AuthRequest request) { User? user = null; @@ -321,23 +355,31 @@ public sealed class VoiceServer : IAsyncDisposable listener.Dispose(); try { - await accepting.ConfigureAwait(false); - Task[] pending; - lock (gate) pending = connections.ToArray(); - await Task.WhenAll(pending).ConfigureAwait(false); + await Task.WhenAll(accepting, reaping).ConfigureAwait(false); } finally { - try { await media.DisposeAsync().ConfigureAwait(false); } - finally { accounts.Dispose(); credentials.Dispose(); shutdown.Dispose(); } + try + { + Task[] pending; + lock (gate) pending = connections.ToArray(); + await Task.WhenAll(pending).ConfigureAwait(false); + } + finally + { + try { await media.DisposeAsync().ConfigureAwait(false); } + finally { accounts.Dispose(); credentials.Dispose(); shutdown.Dispose(); } + } } } - private sealed class Session(ulong id, TlsControlConnection connection, string address) + private sealed class Session(ulong id, TlsControlConnection connection, string address, SessionActivity activity) { public ulong Id { get; } = id; public TlsControlConnection Connection { get; } = connection; public string Address { get; } = address; + public SessionActivity Activity { get; } = activity; + public bool Closing { get; set; } public bool HelloReceived { get; set; } public User? User { get; set; } public MediaPeer? Media { get; set; } diff --git a/dotnet/src/VoiceCat.Server/VoiceServerOptions.cs b/dotnet/src/VoiceCat.Server/VoiceServerOptions.cs new file mode 100644 index 0000000..d5e51a2 --- /dev/null +++ b/dotnet/src/VoiceCat.Server/VoiceServerOptions.cs @@ -0,0 +1,21 @@ +namespace VoiceCat.Server; + +public sealed record VoiceServerOptions +{ + public string Name { get; init; } = "VoiceCat Server"; + public bool AllowGuests { get; init; } = true; + public int MaximumConnections { get; init; } = 64; + public TimeSpan HandshakeTimeout { get; init; } = TimeSpan.FromSeconds(15); + public TimeSpan IdleTimeout { get; init; } = TimeSpan.FromSeconds(45); + public TimeSpan ReaperInterval { get; init; } = TimeSpan.FromSeconds(15); + + internal void Validate() + { + ArgumentException.ThrowIfNullOrWhiteSpace(Name); + ArgumentOutOfRangeException.ThrowIfLessThan(MaximumConnections, 1); + if (HandshakeTimeout <= TimeSpan.Zero || HandshakeTimeout.TotalMilliseconds > uint.MaxValue - 1) throw new ArgumentOutOfRangeException(nameof(HandshakeTimeout)); + if (IdleTimeout < TimeSpan.Zero) throw new ArgumentOutOfRangeException(nameof(IdleTimeout)); + if (ReaperInterval < TimeSpan.Zero || ReaperInterval.TotalMilliseconds > uint.MaxValue - 1 || IdleTimeout > TimeSpan.Zero && ReaperInterval == TimeSpan.Zero) + throw new ArgumentOutOfRangeException(nameof(ReaperInterval)); + } +} diff --git a/dotnet/tests/VoiceCat.Tests/ReaperTests.cs b/dotnet/tests/VoiceCat.Tests/ReaperTests.cs new file mode 100644 index 0000000..4589ce5 --- /dev/null +++ b/dotnet/tests/VoiceCat.Tests/ReaperTests.cs @@ -0,0 +1,129 @@ +using VoiceCat.Protocol; +using System.Net.Sockets; +using VoiceCat.Server; +using Voicecat.V1; +using static VoiceCat.Tests.ServerTests; +using static VoiceCat.Tests.MediaRelayTests; + +namespace VoiceCat.Tests; + +public sealed class ReaperTests +{ + private static readonly VoiceServerOptions Options = new() + { + IdleTimeout = TimeSpan.FromSeconds(10), ReaperInterval = TimeSpan.FromMilliseconds(20) + }; + + [Fact] + public async Task SilentPeerIsReapedWhileTcpActivityKeepsObserverAlive() + { + var clock = new ManualClock(); + await using var fixture = new ServerFixture(options: Options, timeProvider: clock); + await using var alice = await fixture.ConnectAsync(); + await alice.LoginAsync("Alice"); + await using var bob = await fixture.ConnectAsync(); + User self = await bob.LoginAsync("Bob"); + clock.Advance(9); + alice.Send(new() { Ping = new() { Nonce = 99 } }); + await alice.ReadUntilAsync(e => e.Pong?.Nonce == 99); + clock.Advance(2); + Assert.Equal(self.Id, (await alice.ReadUntilAsync(e => e.UserEvent?.Kind == UserEvent.Types.Kind.Left)).UserEvent.LeftId); + Assert.Equal("Receive idle timeout.", (await bob.ReadUntilAsync(e => e.Disconnect is not null)).Disconnect.Reason); + alice.Send(new() { Subscribe = new() }); + int additionalDepartures = 0; + var snapshot = await alice.ReadUntilAsync(e => + { + if (e.UserEvent?.Kind == UserEvent.Types.Kind.Left) additionalDepartures++; + return e.ServerState is not null; + }); + Assert.Equal(0, additionalDepartures); + Assert.DoesNotContain(snapshot.ServerState.Users, user => user.Id == self.Id); + alice.Send(new() { Ping = new() { Nonce = 100 } }); + await alice.ReadUntilAsync(e => e.Pong?.Nonce == 100); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task ValidUdpActivityKeepsTcpIdleClientAlive(bool voice) + { + var clock = new ManualClock(); + await using var fixture = new ServerFixture(options: Options, timeProvider: clock); + await using var alice = await VoicePeer.ConnectAsync(fixture, "Alice"); + await using var bob = await VoicePeer.ConnectAsync(fixture, "Bob"); + uint ssrc = voice ? (await alice.AnnounceAsync(StreamKind.StreamMic)).Ssrc : 0; + clock.Advance(9); + if (voice) + { + await alice.SendAsync(alice.Seal(ssrc, [1, 2, 3])); + await bob.ReceiveVoiceAsync(); + } + else + { + byte[] keepalive = new byte[VoiceFrameHeader.Size]; + new VoiceFrameHeader(MediaFrameType.Keepalive, 0, 0, 0, 0, 0).Write(keepalive); + await alice.SendAsync(keepalive); + Assert.Equal(keepalive, await alice.ReceivePacketAsync()); + } + clock.Advance(2); + Assert.Equal(bob.Client.Authentication!.Self.Id, + (await alice.Client.ReadUntilAsync(e => e.UserEvent?.Kind == UserEvent.Types.Kind.Left)).UserEvent.LeftId); + alice.Client.Send(new() { Ping = new() { Nonce = 42 } }); + await alice.Client.ReadUntilAsync(e => e.Pong?.Nonce == 42); + } + + [Fact] + public async Task InvalidVoiceCannotKeepSilentSessionAlive() + { + var clock = new ManualClock(); + await using var fixture = new ServerFixture(options: Options, timeProvider: clock); + await using var alice = await VoicePeer.ConnectAsync(fixture, "Alice"); + await using var bob = await VoicePeer.ConnectAsync(fixture, "Bob"); + var stream = await bob.AnnounceAsync(StreamKind.StreamMic); + clock.Advance(9); + byte[] forged = bob.Seal(stream.Ssrc, [1]); + forged[^1] ^= 1; + await bob.SendAsync(forged); + alice.Client.Send(new() { Ping = new() { Nonce = 1 } }); + await alice.Client.ReadUntilAsync(e => e.Pong is not null); + clock.Advance(2); + Assert.Equal(bob.Client.Authentication!.Self.Id, + (await alice.Client.ReadUntilAsync(e => e.UserEvent?.Kind == UserEvent.Types.Kind.Left)).UserEvent.LeftId); + } + + [Fact] + public async Task ShutdownAwaitsActiveVoiceAndUnfinishedHandshake() + { + await using var fixture = new ServerFixture(options: Options); + await using var alice = await VoicePeer.ConnectAsync(fixture, "Alice"); + await using var bob = await VoicePeer.ConnectAsync(fixture, "Bob"); + var stream = await alice.AnnounceAsync(StreamKind.StreamMic); + await alice.SendAsync(alice.Seal(stream.Ssrc, [1, 2])); + await bob.ReceiveVoiceAsync(); + using var unfinished = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + await unfinished.ConnectAsync(fixture.Server.EndPoint); + await fixture.Server.DisposeAsync().AsTask().WaitAsync(TimeSpan.FromSeconds(10)); + await fixture.Server.DisposeAsync(); + } + + [Fact] + public async Task ReaperCanBeDisabled() + { + var clock = new ManualClock(); + await using var fixture = new ServerFixture(options: Options with { IdleTimeout = TimeSpan.Zero, ReaperInterval = TimeSpan.Zero }, timeProvider: clock); + await using var client = await fixture.ConnectAsync(); + await client.LoginAsync("Alice"); + clock.Advance(1000); + await Task.Delay(100); + client.Send(new() { Ping = new() { Nonce = 1 } }); + await client.ReadUntilAsync(e => e.Pong is not null); + } + + private sealed class ManualClock : TimeProvider + { + private long timestamp; + public override long TimestampFrequency => TimeSpan.TicksPerSecond; + public override long GetTimestamp() => Volatile.Read(ref timestamp); + public void Advance(int seconds) => Interlocked.Add(ref timestamp, seconds * TimeSpan.TicksPerSecond); + } +} diff --git a/dotnet/tests/VoiceCat.Tests/ServerTests.cs b/dotnet/tests/VoiceCat.Tests/ServerTests.cs index 30b066a..55360e4 100644 --- a/dotnet/tests/VoiceCat.Tests/ServerTests.cs +++ b/dotnet/tests/VoiceCat.Tests/ServerTests.cs @@ -136,10 +136,10 @@ public sealed class ServerTests public string Directory { get; } = Path.Combine(Path.GetTempPath(), "voicecat-server-" + Guid.NewGuid().ToString("N")); public VoiceServer Server { get; } private readonly string fingerprint; - public ServerFixture(bool guests = true) + public ServerFixture(bool guests = true, VoiceServerOptions? options = null, TimeProvider? timeProvider = null) { System.IO.Directory.CreateDirectory(Directory); - Server = new(Directory, new(IPAddress.Loopback, 0), guests); + Server = new(Directory, new(IPAddress.Loopback, 0), options ?? new() { AllowGuests = guests }, timeProvider); using var credentials = ServerCredentials.LoadOrCreate(Directory, "VoiceCat Server"); fingerprint = credentials.CertificateFingerprint; }