From aec0ddb8d7d337c8a5bc9fa3c57de9007c89f97b Mon Sep 17 00:00:00 2001 From: Erik Date: Fri, 25 Sep 2026 13:06:50 +0200 Subject: [PATCH] Review fixes: logoff inventory, socket stop, connect order, relog leftovers F2: the logoff inventory is skipped when the host already let go of the items (it may reset the session before reporting the logoff): an empty list would have emptied the backend's inventory and overwritten the saved copy. F3/F4: each socket start owns its run; Stop no longer blocks the game thread and gives the run up to ten seconds to send what was queued (a large logoff inventory) before closing, flushing that run's own send loop; a relog never hears the previous run's frames. F7: on connect, share_subscribe goes out before the first telemetry frame, as the original sent them. F8: a rare announcement still waiting at logoff is dropped rather than sent as the next character. F9: each dungeon map is sent once per login, as the original (reloaded per login) did. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Backend/BackendClient.cs | 186 +++++++++++------- src/OpenAC.MosswartMassacre/FeatureCatalog.cs | 6 +- .../Inventory/InventoryFeature.cs | 5 + .../Radar/RadarStream.cs | 4 +- .../Streams/SessionStats.cs | 6 +- .../Backend/BackendClientTests.cs | 38 ++++ .../Inventory/InventoryFeatureTests.cs | 23 +++ .../ReviewFixTests.cs | 37 ++++ 8 files changed, 223 insertions(+), 82 deletions(-) create mode 100644 tests/OpenAC.MosswartMassacre.Tests/ReviewFixTests.cs 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)); + } +}