diff --git a/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs b/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs index 5acd3c3..025d17e 100644 --- a/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs +++ b/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs @@ -22,22 +22,24 @@ namespace OpenAC.MosswartMassacre.Backend; /// Incoming frames are reassembled across receive calls; the original read a /// single 4 KiB buffer and lost any larger frame. /// +/// +/// Each begins a run of its own. ends +/// the run without blocking the caller: the run stops reconnecting, sends +/// what was already queued (a farewell frame, a logoff inventory) for up to +/// ten seconds, and only then closes. Whatever that run still receives goes +/// nowhere, so a relog never hears the previous character's frames. +/// /// internal sealed class BackendClient : IDisposable { public const string SecretHeader = "X-Plugin-Secret"; private static readonly TimeSpan ReconnectDelay = TimeSpan.FromSeconds(2); - private static readonly TimeSpan FlushTimeout = TimeSpan.FromMilliseconds(500); + private static readonly TimeSpan FlushTimeout = TimeSpan.FromSeconds(10); private readonly IPluginLogger _log; - private readonly ConcurrentQueue _inbound = new(); - private readonly ConcurrentQueue _connections = new(); private readonly object _gate = new(); - private CancellationTokenSource? _cts; - private Task? _loop; - private Channel? _outbound; - private Task? _sendLoop; + private Run? _run; public BackendClient(IPluginLogger log) { @@ -50,7 +52,7 @@ internal sealed class BackendClient : IDisposable get { lock (_gate) - return _cts is not null; + return _run is not null; } } @@ -60,71 +62,66 @@ internal sealed class BackendClient : IDisposable get { lock (_gate) - return _outbound is not null; + return _run?.Outbound is not null; } } + /// The last run stopped, finishing in the background; tests wait on it. + internal Task LastStop { get; private set; } = Task.CompletedTask; + public void Start(Uri endpoint, string secret, string playerName, Func utcNow) { lock (_gate) { - if (_cts is not null) + if (_run is not null) return; - _cts = new CancellationTokenSource(); - CancellationToken token = _cts.Token; + var run = new Run(); + _run = run; _log.Info("[WebSocket] connecting"); - _loop = Task.Run(() => RunAsync(endpoint, secret, playerName, utcNow, token)); + run.Loop = Task.Run(() => RunAsync(run, endpoint, secret, playerName, utcNow)); } } public void Stop() { - Task? loop; - CancellationTokenSource? cts; - Channel? outbound; - Task? sendLoop; + Run? run; lock (_gate) { - if (_cts is null) - return; - cts = _cts; - _cts = null; - outbound = _outbound; - _outbound = null; - sendLoop = _sendLoop; - loop = _loop; - _loop = null; + run = _run; + _run = null; } + if (run is null) + return; - // Let what was already queued go out (a farewell frame is sent just - // before stopping), briefly, then cut the socket. - if (outbound is not null) + Channel? outbound; + Task? sendLoop; + lock (run) { - outbound.Writer.TryComplete(); - try + run.Stopping = true; + outbound = run.Outbound; + run.Outbound = null; + sendLoop = run.SendLoop; + } + outbound?.Writer.TryComplete(); + + LastStop = Task.Run(async () => + { + if (sendLoop is not null) + await Task.WhenAny(sendLoop, Task.Delay(FlushTimeout)).ConfigureAwait(false); + await run.Cts.CancelAsync().ConfigureAwait(false); + if (run.Loop is { } loop) { - (sendLoop ?? outbound.Reader.Completion).Wait(FlushTimeout); + try + { + await loop.ConfigureAwait(false); + } + catch (Exception ex) when (ex is OperationCanceledException or WebSocketException) + { + } } - catch (AggregateException) - { - } - } - cts.Cancel(); - try - { - loop?.Wait(TimeSpan.FromSeconds(2)); - } - catch (AggregateException) - { - } - cts.Dispose(); - while (_inbound.TryDequeue(out _)) - { - } - while (_connections.TryDequeue(out _)) - { - } - _log.Info("[WebSocket] DISABLED"); + run.Cts.Dispose(); + _log.Info("[WebSocket] DISABLED"); + }); } /// Queues a frame. Dropped when no socket is open. @@ -132,26 +129,43 @@ internal sealed class BackendClient : IDisposable { Channel? outbound; lock (_gate) - outbound = _outbound; + outbound = _run?.Outbound; return outbound is not null && outbound.Writer.TryWrite(json); } /// The next frame the backend sent, trimmed, if any. - public bool TryReceive(out string frame) => _inbound.TryDequeue(out frame!); + public bool TryReceive(out string frame) + { + Run? run; + lock (_gate) + run = _run; + if (run is not null) + return run.Inbound.TryDequeue(out frame!); + frame = string.Empty; + return false; + } /// /// The next completed connect-and-register, if any. The tick thread uses it /// to run what the original ran right after registering. /// - public bool TryTakeConnection(out BackendConnection connection) => - _connections.TryDequeue(out connection); + public bool TryTakeConnection(out BackendConnection connection) + { + Run? run; + lock (_gate) + run = _run; + if (run is not null) + return run.Connections.TryDequeue(out connection); + connection = default; + return false; + } public void Dispose() => Stop(); - private async Task RunAsync( - Uri endpoint, string secret, string playerName, Func utcNow, CancellationToken token) + private async Task RunAsync(Run run, Uri endpoint, string secret, string playerName, Func utcNow) { - while (!token.IsCancellationRequested) + CancellationToken token = run.Cts.Token; + while (!token.IsCancellationRequested && !run.Stopping) { var socket = new ClientWebSocket(); Channel? outbound = null; @@ -168,18 +182,20 @@ internal sealed class BackendClient : IDisposable outbound = Channel.CreateUnbounded( new UnboundedChannelOptions { SingleReader = true, SingleWriter = false }); - lock (_gate) - { - if (token.IsCancellationRequested) - break; - _outbound = outbound; - } - _connections.Enqueue(new BackendConnection(connectedAt)); - - Task receive = ReceiveLoopAsync(socket, token); Task send = SendLoopAsync(socket, outbound.Reader, token); - lock (_gate) - _sendLoop = send; + lock (run) + { + if (run.Stopping) + { + outbound.Writer.TryComplete(); + break; + } + run.Outbound = outbound; + run.SendLoop = send; + } + run.Connections.Enqueue(new BackendConnection(connectedAt)); + + Task receive = ReceiveLoopAsync(socket, run.Inbound, token); await Task.WhenAny(receive, send).ConfigureAwait(false); } catch (OperationCanceledException) when (token.IsCancellationRequested) @@ -192,17 +208,17 @@ internal sealed class BackendClient : IDisposable } finally { - lock (_gate) + lock (run) { - if (ReferenceEquals(_outbound, outbound)) - _outbound = null; + if (ReferenceEquals(run.Outbound, outbound)) + run.Outbound = null; } outbound?.Writer.TryComplete(); socket.Abort(); socket.Dispose(); } - if (token.IsCancellationRequested) + if (token.IsCancellationRequested || run.Stopping) break; _log.Info("[WebSocket] Reconnecting in 2 seconds..."); try @@ -216,7 +232,7 @@ internal sealed class BackendClient : IDisposable } } - private async Task ReceiveLoopAsync(ClientWebSocket socket, CancellationToken token) + private async Task ReceiveLoopAsync(ClientWebSocket socket, ConcurrentQueue inbound, CancellationToken token) { var buffer = new byte[8192]; using var message = new MemoryStream(); @@ -247,7 +263,7 @@ internal sealed class BackendClient : IDisposable string text = Encoding.UTF8.GetString(message.GetBuffer(), 0, (int)message.Length).Trim(); message.SetLength(0); if (result.MessageType == WebSocketMessageType.Text && text.Length > 0) - _inbound.Enqueue(text); + inbound.Enqueue(text); } } @@ -275,6 +291,24 @@ internal sealed class BackendClient : IDisposable byte[] bytes = Encoding.UTF8.GetBytes(text); return socket.SendAsync(new ArraySegment(bytes), WebSocketMessageType.Text, true, token); } + + /// One started connection loop and everything that belongs to it alone. + private sealed class Run + { + public CancellationTokenSource Cts { get; } = new(); + + public ConcurrentQueue Inbound { get; } = new(); + + public ConcurrentQueue Connections { get; } = new(); + + public volatile bool Stopping; + + public Task? Loop { get; set; } + + public Channel? Outbound { get; set; } + + public Task? SendLoop { get; set; } + } } /// One completed connect-and-register, stamped with the UTC time it connected. diff --git a/src/OpenAC.MosswartMassacre/FeatureCatalog.cs b/src/OpenAC.MosswartMassacre/FeatureCatalog.cs index d56f699..5ace75f 100644 --- a/src/OpenAC.MosswartMassacre/FeatureCatalog.cs +++ b/src/OpenAC.MosswartMassacre/FeatureCatalog.cs @@ -22,7 +22,7 @@ internal static class FeatureCatalog var kills = new KillTracker(context.Clock); var rares = new RareTracker(context, delayed); var combat = new CombatStatsTracker(context); - var stats = new SessionStats(context, kills, rares); + var stats = new SessionStats(context, kills, rares, delayed); var telemetry = new TelemetryStream(context, stats); features.Add(stats); features.Add(combat); @@ -42,12 +42,14 @@ internal static class FeatureCatalog context.MetaState = () => HostReadings.MetaState(context.Host); var mainWindow = new Views.MainWindow(context, stats); features.Add(mainWindow); - session.AddConnectionListener(telemetry); // Cross-computer vital and debuff sharing. var sharing = new Sharing.VitalSharing(context); features.Add(sharing); + // On connect the original re-subscribed to sharing before its first + // telemetry frame. session.AddConnectionListener(sharing); + session.AddConnectionListener(telemetry); session.AddShareHandler(sharing); features.Add(new Sharing.VitalSharingOverlay(context, sharing)); diff --git a/src/OpenAC.MosswartMassacre/Inventory/InventoryFeature.cs b/src/OpenAC.MosswartMassacre/Inventory/InventoryFeature.cs index 7df03bc..823aa18 100644 --- a/src/OpenAC.MosswartMassacre/Inventory/InventoryFeature.cs +++ b/src/OpenAC.MosswartMassacre/Inventory/InventoryFeature.cs @@ -312,6 +312,11 @@ internal sealed class InventoryFeature : IMmFeature { if (!_active || !InventoryLog) return; + // The host may already have let go of the session's objects by the + // time it reports the logoff. An empty list then is not the inventory: + // sending it would empty the backend's copy and overwrite the saved one. + if (!_settled || _reader.CaptureOwned().Count == 0) + return; try { DumpInventory(); diff --git a/src/OpenAC.MosswartMassacre/Radar/RadarStream.cs b/src/OpenAC.MosswartMassacre/Radar/RadarStream.cs index dc4f00e..c86ea37 100644 --- a/src/OpenAC.MosswartMassacre/Radar/RadarStream.cs +++ b/src/OpenAC.MosswartMassacre/Radar/RadarStream.cs @@ -50,9 +50,7 @@ internal sealed class RadarStream : IMmFeature _context.Verbose("[Radar] Stopped nearby objects tracker"); } - public void OnLogin() - { - } + public void OnLogin() => _mapsSent.Clear(); public void OnLogoff() => Stop(); diff --git a/src/OpenAC.MosswartMassacre/Streams/SessionStats.cs b/src/OpenAC.MosswartMassacre/Streams/SessionStats.cs index 3b9aaee..ddacdc2 100644 --- a/src/OpenAC.MosswartMassacre/Streams/SessionStats.cs +++ b/src/OpenAC.MosswartMassacre/Streams/SessionStats.cs @@ -15,12 +15,14 @@ internal sealed class SessionStats : IMmFeature public const uint NumDeathsProperty = 43; private readonly MmContext _context; + private readonly DelayedCommands _delayed; private readonly Action _onDied; private IDisposable? _statsTimer; - public SessionStats(MmContext context, KillTracker kills, RareTracker rares) + public SessionStats(MmContext context, KillTracker kills, RareTracker rares, DelayedCommands delayed) { _context = context; + _delayed = delayed; Kills = kills; Rares = rares; _onDied = _ => OnDeath(); @@ -47,6 +49,8 @@ internal sealed class SessionStats : IMmFeature public void OnLogoff() { + // A rare announcement still waiting belongs to the character leaving. + _delayed.Clear(); _context.Host.Events.LocalPlayerDied -= _onDied; _statsTimer?.Dispose(); _statsTimer = null; diff --git a/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendClientTests.cs b/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendClientTests.cs index 36aff2a..bd468f7 100644 --- a/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendClientTests.cs +++ b/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendClientTests.cs @@ -93,6 +93,44 @@ public sealed class BackendClientTests Assert.Equal(1, server.Connections); } + [Fact] + public async Task Stop_sends_what_was_queued_before_it_closes_without_blocking() + { + using var server = new TestSocketServer(); + using var client = new BackendClient(new RecordingLogger()); + client.Start(server.Uri, "", "Mossy", () => Connected); + server.NextFrame(); + Assert.True(await WaitAsync(() => client.IsConnected)); + + string big = "{\"type\":\"full_inventory\",\"pad\":\"" + new string('x', 400_000) + "\"}"; + client.Send(big); + client.Send("{\"type\":\"share_unsubscribe\"}"); + var watch = System.Diagnostics.Stopwatch.StartNew(); + client.Stop(); + Assert.True(watch.ElapsedMilliseconds < 200, $"Stop blocked for {watch.ElapsedMilliseconds} ms"); + + Assert.Equal(big, server.NextFrame()); + Assert.Equal("{\"type\":\"share_unsubscribe\"}", server.NextFrame()); + await client.LastStop.WaitAsync(TimeSpan.FromSeconds(15)); + } + + [Fact] + public async Task A_new_start_right_after_stop_hears_only_its_own_connection() + { + using var server = new TestSocketServer(); + using var client = new BackendClient(new RecordingLogger()); + client.Start(server.Uri, "", "Mossy", () => Connected); + server.NextFrame(); + Assert.True(await WaitAsync(() => client.IsConnected)); + client.Stop(); + + client.Start(server.Uri, "", "Horan", () => Connected); + + Assert.Equal("{\"type\":\"register\",\"player_name\":\"Horan\"}", server.NextFrame(8000)); + Assert.True(await WaitAsync(() => client.TryTakeConnection(out _))); + Assert.False(client.TryTakeConnection(out _)); + } + internal static async Task WaitAsync(Func condition, int timeoutMs = 5000) { for (int waited = 0; waited < timeoutMs; waited += 20) diff --git a/tests/OpenAC.MosswartMassacre.Tests/Inventory/InventoryFeatureTests.cs b/tests/OpenAC.MosswartMassacre.Tests/Inventory/InventoryFeatureTests.cs index a5cabb6..8a12a4e 100644 --- a/tests/OpenAC.MosswartMassacre.Tests/Inventory/InventoryFeatureTests.cs +++ b/tests/OpenAC.MosswartMassacre.Tests/Inventory/InventoryFeatureTests.cs @@ -223,6 +223,29 @@ public sealed class InventoryFeatureTests Assert.Single(harness.FramesOfTypeRaw("full_inventory")); } + [Fact] + public void Logoff_after_the_host_let_go_of_the_items_sends_nothing_and_keeps_the_saved_copy() + { + // The host may reset the session before it reports the logoff; an + // empty list then would wipe the backend's inventory. + using var harness = new PluginHarness("Mossy", arrange: host => Own(host, Tapers())); + LoginAndSettle(harness); + harness.Tick(0.1); + harness.Tick(0.1); + harness.Plugin.OnLogoff(); + string? saved = harness.Host.Storage.ReadText("inventory.json"); + Assert.NotNull(saved); + harness.Plugin.OnLogin(); + LoginAndSettle(harness); + harness.Frames.Clear(); + harness.Host.Automation.Items.Owned.Clear(); + + harness.Plugin.OnLogoff(); + + Assert.Empty(harness.FramesOfTypeRaw("full_inventory")); + Assert.Equal(saved, harness.Host.Storage.ReadText("inventory.json")); + } + [Fact] public void The_taper_count_follows_the_stacks_in_the_packs() { diff --git a/tests/OpenAC.MosswartMassacre.Tests/ReviewFixTests.cs b/tests/OpenAC.MosswartMassacre.Tests/ReviewFixTests.cs new file mode 100644 index 0000000..c9bb609 --- /dev/null +++ b/tests/OpenAC.MosswartMassacre.Tests/ReviewFixTests.cs @@ -0,0 +1,37 @@ + +namespace OpenAC.MosswartMassacre.Tests; + +/// Pins for the fixes a review of the port asked for. +public sealed class ReviewFixTests +{ + [Fact] + public void On_connect_the_sharing_subscribe_goes_out_before_the_first_telemetry() + { + using var harness = new PluginHarness("Mossy"); + harness.Login(); + harness.Context.Settings.Current.VitalSharingEnabled = true; + harness.Frames.Clear(); + + harness.SimulateConnected(); + + string[] types = [.. harness.Frames.Select(f => (string)Newtonsoft.Json.Linq.JObject.Parse(f)["type"]!)]; + Assert.Equal(["share_subscribe", "telemetry"], types.Take(2)); + } + + [Fact] + public void A_rare_announcement_still_waiting_at_logoff_is_dropped() + { + using var harness = new PluginHarness("Mossy"); + harness.Login(); + harness.Host.Automation.Chat.Deliver("Mossy has discovered the Pearl of Blood Drinking!", kind: 4); + harness.Host.Automation.Chat.Submitted.Clear(); + + harness.Plugin.OnLogoff(); + harness.Host.Automation.Character.Name = "Horan"; + harness.Plugin.OnLogin(); + harness.Host.Automation.Chat.Submitted.Clear(); + harness.Run(4); + + Assert.DoesNotContain(harness.Host.Automation.Chat.Submitted, line => line.StartsWith("/a ", StringComparison.Ordinal)); + } +}