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) <noreply@anthropic.com>
This commit is contained in:
Erik 2026-09-25 13:06:50 +02:00
parent 71f14473c6
commit aec0ddb8d7
8 changed files with 223 additions and 82 deletions

View file

@ -22,22 +22,24 @@ namespace OpenAC.MosswartMassacre.Backend;
/// Incoming frames are reassembled across receive calls; the original read a /// Incoming frames are reassembled across receive calls; the original read a
/// single 4 KiB buffer and lost any larger frame. /// single 4 KiB buffer and lost any larger frame.
/// </para> /// </para>
/// <para>
/// Each <see cref="Start"/> begins a run of its own. <see cref="Stop"/> 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.
/// </para>
/// </remarks> /// </remarks>
internal sealed class BackendClient : IDisposable internal sealed class BackendClient : IDisposable
{ {
public const string SecretHeader = "X-Plugin-Secret"; public const string SecretHeader = "X-Plugin-Secret";
private static readonly TimeSpan ReconnectDelay = TimeSpan.FromSeconds(2); 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 IPluginLogger _log;
private readonly ConcurrentQueue<string> _inbound = new();
private readonly ConcurrentQueue<BackendConnection> _connections = new();
private readonly object _gate = new(); private readonly object _gate = new();
private CancellationTokenSource? _cts; private Run? _run;
private Task? _loop;
private Channel<string>? _outbound;
private Task? _sendLoop;
public BackendClient(IPluginLogger log) public BackendClient(IPluginLogger log)
{ {
@ -50,7 +52,7 @@ internal sealed class BackendClient : IDisposable
get get
{ {
lock (_gate) lock (_gate)
return _cts is not null; return _run is not null;
} }
} }
@ -60,71 +62,66 @@ internal sealed class BackendClient : IDisposable
get get
{ {
lock (_gate) lock (_gate)
return _outbound is not null; return _run?.Outbound is not null;
} }
} }
/// <summary>The last run stopped, finishing in the background; tests wait on it.</summary>
internal Task LastStop { get; private set; } = Task.CompletedTask;
public void Start(Uri endpoint, string secret, string playerName, Func<DateTime> utcNow) public void Start(Uri endpoint, string secret, string playerName, Func<DateTime> utcNow)
{ {
lock (_gate) lock (_gate)
{ {
if (_cts is not null) if (_run is not null)
return; return;
_cts = new CancellationTokenSource(); var run = new Run();
CancellationToken token = _cts.Token; _run = run;
_log.Info("[WebSocket] connecting"); _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() public void Stop()
{ {
Task? loop; Run? run;
CancellationTokenSource? cts;
Channel<string>? outbound;
Task? sendLoop;
lock (_gate) lock (_gate)
{ {
if (_cts is null) run = _run;
return; _run = null;
cts = _cts;
_cts = null;
outbound = _outbound;
_outbound = null;
sendLoop = _sendLoop;
loop = _loop;
_loop = null;
} }
if (run is null)
return;
// Let what was already queued go out (a farewell frame is sent just Channel<string>? outbound;
// before stopping), briefly, then cut the socket. Task? sendLoop;
if (outbound is not null) lock (run)
{ {
outbound.Writer.TryComplete(); run.Stopping = true;
try 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) run.Cts.Dispose();
{ _log.Info("[WebSocket] DISABLED");
} });
}
cts.Cancel();
try
{
loop?.Wait(TimeSpan.FromSeconds(2));
}
catch (AggregateException)
{
}
cts.Dispose();
while (_inbound.TryDequeue(out _))
{
}
while (_connections.TryDequeue(out _))
{
}
_log.Info("[WebSocket] DISABLED");
} }
/// <summary>Queues a frame. Dropped when no socket is open.</summary> /// <summary>Queues a frame. Dropped when no socket is open.</summary>
@ -132,26 +129,43 @@ internal sealed class BackendClient : IDisposable
{ {
Channel<string>? outbound; Channel<string>? outbound;
lock (_gate) lock (_gate)
outbound = _outbound; outbound = _run?.Outbound;
return outbound is not null && outbound.Writer.TryWrite(json); return outbound is not null && outbound.Writer.TryWrite(json);
} }
/// <summary>The next frame the backend sent, trimmed, if any.</summary> /// <summary>The next frame the backend sent, trimmed, if any.</summary>
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;
}
/// <summary> /// <summary>
/// The next completed connect-and-register, if any. The tick thread uses it /// The next completed connect-and-register, if any. The tick thread uses it
/// to run what the original ran right after registering. /// to run what the original ran right after registering.
/// </summary> /// </summary>
public bool TryTakeConnection(out BackendConnection connection) => public bool TryTakeConnection(out BackendConnection connection)
_connections.TryDequeue(out 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(); public void Dispose() => Stop();
private async Task RunAsync( private async Task RunAsync(Run run, Uri endpoint, string secret, string playerName, Func<DateTime> utcNow)
Uri endpoint, string secret, string playerName, Func<DateTime> utcNow, CancellationToken token)
{ {
while (!token.IsCancellationRequested) CancellationToken token = run.Cts.Token;
while (!token.IsCancellationRequested && !run.Stopping)
{ {
var socket = new ClientWebSocket(); var socket = new ClientWebSocket();
Channel<string>? outbound = null; Channel<string>? outbound = null;
@ -168,18 +182,20 @@ internal sealed class BackendClient : IDisposable
outbound = Channel.CreateUnbounded<string>( outbound = Channel.CreateUnbounded<string>(
new UnboundedChannelOptions { SingleReader = true, SingleWriter = false }); 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); Task send = SendLoopAsync(socket, outbound.Reader, token);
lock (_gate) lock (run)
_sendLoop = send; {
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); await Task.WhenAny(receive, send).ConfigureAwait(false);
} }
catch (OperationCanceledException) when (token.IsCancellationRequested) catch (OperationCanceledException) when (token.IsCancellationRequested)
@ -192,17 +208,17 @@ internal sealed class BackendClient : IDisposable
} }
finally finally
{ {
lock (_gate) lock (run)
{ {
if (ReferenceEquals(_outbound, outbound)) if (ReferenceEquals(run.Outbound, outbound))
_outbound = null; run.Outbound = null;
} }
outbound?.Writer.TryComplete(); outbound?.Writer.TryComplete();
socket.Abort(); socket.Abort();
socket.Dispose(); socket.Dispose();
} }
if (token.IsCancellationRequested) if (token.IsCancellationRequested || run.Stopping)
break; break;
_log.Info("[WebSocket] Reconnecting in 2 seconds..."); _log.Info("[WebSocket] Reconnecting in 2 seconds...");
try 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<string> inbound, CancellationToken token)
{ {
var buffer = new byte[8192]; var buffer = new byte[8192];
using var message = new MemoryStream(); 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(); string text = Encoding.UTF8.GetString(message.GetBuffer(), 0, (int)message.Length).Trim();
message.SetLength(0); message.SetLength(0);
if (result.MessageType == WebSocketMessageType.Text && text.Length > 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); byte[] bytes = Encoding.UTF8.GetBytes(text);
return socket.SendAsync(new ArraySegment<byte>(bytes), WebSocketMessageType.Text, true, token); return socket.SendAsync(new ArraySegment<byte>(bytes), WebSocketMessageType.Text, true, token);
} }
/// <summary>One started connection loop and everything that belongs to it alone.</summary>
private sealed class Run
{
public CancellationTokenSource Cts { get; } = new();
public ConcurrentQueue<string> Inbound { get; } = new();
public ConcurrentQueue<BackendConnection> Connections { get; } = new();
public volatile bool Stopping;
public Task? Loop { get; set; }
public Channel<string>? Outbound { get; set; }
public Task? SendLoop { get; set; }
}
} }
/// <summary>One completed connect-and-register, stamped with the UTC time it connected.</summary> /// <summary>One completed connect-and-register, stamped with the UTC time it connected.</summary>

View file

@ -22,7 +22,7 @@ internal static class FeatureCatalog
var kills = new KillTracker(context.Clock); var kills = new KillTracker(context.Clock);
var rares = new RareTracker(context, delayed); var rares = new RareTracker(context, delayed);
var combat = new CombatStatsTracker(context); 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); var telemetry = new TelemetryStream(context, stats);
features.Add(stats); features.Add(stats);
features.Add(combat); features.Add(combat);
@ -42,12 +42,14 @@ internal static class FeatureCatalog
context.MetaState = () => HostReadings.MetaState(context.Host); context.MetaState = () => HostReadings.MetaState(context.Host);
var mainWindow = new Views.MainWindow(context, stats); var mainWindow = new Views.MainWindow(context, stats);
features.Add(mainWindow); features.Add(mainWindow);
session.AddConnectionListener(telemetry);
// Cross-computer vital and debuff sharing. // Cross-computer vital and debuff sharing.
var sharing = new Sharing.VitalSharing(context); var sharing = new Sharing.VitalSharing(context);
features.Add(sharing); features.Add(sharing);
// On connect the original re-subscribed to sharing before its first
// telemetry frame.
session.AddConnectionListener(sharing); session.AddConnectionListener(sharing);
session.AddConnectionListener(telemetry);
session.AddShareHandler(sharing); session.AddShareHandler(sharing);
features.Add(new Sharing.VitalSharingOverlay(context, sharing)); features.Add(new Sharing.VitalSharingOverlay(context, sharing));

View file

@ -312,6 +312,11 @@ internal sealed class InventoryFeature : IMmFeature
{ {
if (!_active || !InventoryLog) if (!_active || !InventoryLog)
return; 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 try
{ {
DumpInventory(); DumpInventory();

View file

@ -50,9 +50,7 @@ internal sealed class RadarStream : IMmFeature
_context.Verbose("[Radar] Stopped nearby objects tracker"); _context.Verbose("[Radar] Stopped nearby objects tracker");
} }
public void OnLogin() public void OnLogin() => _mapsSent.Clear();
{
}
public void OnLogoff() => Stop(); public void OnLogoff() => Stop();

View file

@ -15,12 +15,14 @@ internal sealed class SessionStats : IMmFeature
public const uint NumDeathsProperty = 43; public const uint NumDeathsProperty = 43;
private readonly MmContext _context; private readonly MmContext _context;
private readonly DelayedCommands _delayed;
private readonly Action<string> _onDied; private readonly Action<string> _onDied;
private IDisposable? _statsTimer; private IDisposable? _statsTimer;
public SessionStats(MmContext context, KillTracker kills, RareTracker rares) public SessionStats(MmContext context, KillTracker kills, RareTracker rares, DelayedCommands delayed)
{ {
_context = context; _context = context;
_delayed = delayed;
Kills = kills; Kills = kills;
Rares = rares; Rares = rares;
_onDied = _ => OnDeath(); _onDied = _ => OnDeath();
@ -47,6 +49,8 @@ internal sealed class SessionStats : IMmFeature
public void OnLogoff() public void OnLogoff()
{ {
// A rare announcement still waiting belongs to the character leaving.
_delayed.Clear();
_context.Host.Events.LocalPlayerDied -= _onDied; _context.Host.Events.LocalPlayerDied -= _onDied;
_statsTimer?.Dispose(); _statsTimer?.Dispose();
_statsTimer = null; _statsTimer = null;

View file

@ -93,6 +93,44 @@ public sealed class BackendClientTests
Assert.Equal(1, server.Connections); 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<bool> WaitAsync(Func<bool> condition, int timeoutMs = 5000) internal static async Task<bool> WaitAsync(Func<bool> condition, int timeoutMs = 5000)
{ {
for (int waited = 0; waited < timeoutMs; waited += 20) for (int waited = 0; waited < timeoutMs; waited += 20)

View file

@ -223,6 +223,29 @@ public sealed class InventoryFeatureTests
Assert.Single(harness.FramesOfTypeRaw("full_inventory")); 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] [Fact]
public void The_taper_count_follows_the_stacks_in_the_packs() public void The_taper_count_follows_the_stacks_in_the_packs()
{ {

View file

@ -0,0 +1,37 @@
namespace OpenAC.MosswartMassacre.Tests;
/// <summary>Pins for the fixes a review of the port asked for.</summary>
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));
}
}