using Cysharp.Threading.Tasks; using NUnit.Framework; using ShrinkNetwork; public class NetworkQueueTests { [Test] public async Task StateReplacementIsScopedToSessionAndMovesAfterReliableWork() { using var queue = new ShrinkNetworkWorkQueue(10, 100); var output = new List(); UniTask Write(int i) { output.Add(i); return UniTask.CompletedTask; } var old = queue.EnqueueAsync(1, "position", 10, () => Write(1), "player"); var reliable = queue.EnqueueAsync(1, "position", 10, () => Write(2)); var other = queue.EnqueueAsync(2, "position", 10, () => Write(3), "player"); var latest = queue.EnqueueAsync(1, "position", 20, () => Write(4), "player"); Assert.That(await old, Is.EqualTo(ShrinkNetworkQueueResult.Replaced)); Assert.That(queue.CaptureDiagnostics().PendingBytes, Is.EqualTo(40)); Assert.That(await queue.PumpAsync(10), Is.EqualTo(3)); Assert.That(output, Is.EqualTo(new[] { 2, 3, 4 })); foreach (var result in new[] { await reliable, await other, await latest }) Assert.That(result, Is.EqualTo(ShrinkNetworkQueueResult.Completed)); } [Test] public async Task CapacityRejectionDoesNotLoseExistingReliableWorkOrState() { using var queue = new ShrinkNetworkWorkQueue(2, 20, 1); var a = queue.EnqueueAsync(1, "a", 10, () => UniTask.CompletedTask, "x"); Assert.That(await queue.EnqueueAsync(1, "a", 1, () => UniTask.CompletedTask), Is.EqualTo(ShrinkNetworkQueueResult.Rejected)); Assert.That(await queue.EnqueueAsync(1, "a", 21, () => UniTask.CompletedTask, "x"), Is.EqualTo(ShrinkNetworkQueueResult.Rejected)); await queue.PumpAsync(1); Assert.That(await a, Is.EqualTo(ShrinkNetworkQueueResult.Completed)); Assert.That(queue.CaptureDiagnostics().PendingBytes, Is.Zero); } [Test] public async Task BudgetDoesNotPreemptHandlerAndConcurrentPumpsDoNotRunHandlersInParallel() { using var queue = new ShrinkNetworkWorkQueue(4, 100); var completion = new UniTaskCompletionSource(); var first = queue.EnqueueAsync(1, "a", 10, () => completion.Task); var second = queue.EnqueueAsync(1, "a", 10, () => UniTask.CompletedTask); var pump = queue.PumpAsync(4, 100, TimeSpan.FromMilliseconds(1)).AsTask(); await Task.Delay(10); Assert.That(pump.IsCompleted, Is.False); Assert.That(await queue.PumpAsync(4), Is.Zero); completion.TrySetResult(); Assert.That(await pump, Is.EqualTo(1)); Assert.That(await first, Is.EqualTo(ShrinkNetworkQueueResult.Completed)); queue.Dispose(); Assert.That(await second, Is.EqualTo(ShrinkNetworkQueueResult.Canceled)); } [Test] public async Task ByteBudgetMakesProgressForLargePacketsAndHandlerFailureDoesNotLeakQueue() { using var queue = new ShrinkNetworkWorkQueue(4, 100); var failed = queue.EnqueueAsync(1, "a", 30, () => UniTask.FromException(new IOException("handler"))).AsTask(); var next = queue.EnqueueAsync(1, "a", 10, () => UniTask.CompletedTask); Assert.That(await queue.PumpAsync(4, 1), Is.EqualTo(1)); Assert.ThrowsAsync(async () => await failed); Assert.That(await queue.PumpAsync(4), Is.EqualTo(1)); Assert.That(await next, Is.EqualTo(ShrinkNetworkQueueResult.Completed)); Assert.That(queue.CaptureDiagnostics().PendingCount, Is.Zero); } }