Show the backend socket's status and errors in chat, as the original did

Live, /mm ws enable printed nothing about the connection: the socket's
connecting / CONNECTED / REGISTERED / error / reconnect / DISABLED lines went
only to the host log, which the windowed client does not show. The original
printed every one of them to chat.

BackendClient now queues each status and error line (original wording) and
still logs it; BackendSession drains the queue to chat on the tick, and right
after a start or stop on the tick thread so "[WebSocket] connecting…" and
"[WebSocket] DISABLED" come before the command's own reply, as in the
original. The socket thread never writes to chat. The original's telemetry-loop
and cleanup lines are emitted at the matching points, "Starting telemetry
loop" comes from the telemetry listener, and "Vital sharing re-subscribed" is
no longer verbose-only. No frame on the wire changes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Erik 2026-09-26 07:30:30 +02:00
parent 336c5cfdf5
commit 28794299cb
7 changed files with 226 additions and 21 deletions

View file

@ -29,6 +29,13 @@ namespace OpenAC.MosswartMassacre.Backend;
/// ten seconds, and only then closes. Whatever that run still receives goes
/// nowhere, so a relog never hears the previous character's frames.
/// </para>
/// <para>
/// Every status and error line (connecting, connected, registered, errors,
/// reconnects, disabled) goes to the host log and is also queued for the
/// player's chat, worded as the original printed it there. The socket thread
/// never writes to chat; the tick thread drains the queue through
/// <see cref="TryTakeStatus"/>.
/// </para>
/// </remarks>
internal sealed class BackendClient : IDisposable
{
@ -39,6 +46,7 @@ internal sealed class BackendClient : IDisposable
private readonly IPluginLogger _log;
private readonly object _gate = new();
private readonly ConcurrentQueue<string> _status = new();
private Run? _run;
public BackendClient(IPluginLogger log)
@ -77,7 +85,7 @@ internal sealed class BackendClient : IDisposable
return;
var run = new Run();
_run = run;
_log.Info("[WebSocket] connecting");
Status("[WebSocket] connecting…");
run.Loop = Task.Run(() => RunAsync(run, endpoint, secret, playerName, utcNow));
}
}
@ -92,6 +100,8 @@ internal sealed class BackendClient : IDisposable
}
if (run is null)
return;
// Said at once, as the original did, while the run finishes behind it.
Status("[WebSocket] DISABLED");
Channel<string>? outbound;
Task? sendLoop;
@ -120,10 +130,16 @@ internal sealed class BackendClient : IDisposable
}
}
run.Cts.Dispose();
_log.Info("[WebSocket] DISABLED");
});
}
/// <summary>
/// The next status or error line for the player's chat, oldest first. Lines
/// of a stopped run still come out, as the original printed them after it
/// said DISABLED.
/// </summary>
public bool TryTakeStatus(out string line) => _status.TryDequeue(out line!);
/// <summary>Queues a frame. Dropped when no socket is open.</summary>
public bool Send(string json)
{
@ -173,12 +189,12 @@ internal sealed class BackendClient : IDisposable
{
socket.Options.SetRequestHeader(SecretHeader, secret);
await socket.ConnectAsync(endpoint, token).ConfigureAwait(false);
_log.Info("[WebSocket] CONNECTED");
Status("[WebSocket] CONNECTED");
DateTime connectedAt = utcNow();
string register = Wire.Serialize(new { type = "register", player_name = playerName });
await SendFrameAsync(socket, register, token).ConfigureAwait(false);
_log.Info("[WebSocket] REGISTERED");
Status("[WebSocket] REGISTERED");
outbound = Channel.CreateUnbounded<string>(
new UnboundedChannelOptions { SingleReader = true, SingleWriter = false });
@ -197,17 +213,27 @@ internal sealed class BackendClient : IDisposable
Task receive = ReceiveLoopAsync(socket, run.Inbound, token);
await Task.WhenAny(receive, send).ConfigureAwait(false);
// The original's telemetry loop ran for as long as the socket
// did, and said why it ended. A stop here is the stop flag
// first and the token later, so both count as cancelled.
bool cancelled = token.IsCancellationRequested || run.Stopping;
if (cancelled)
Status("[WebSocket] Telemetry loop cancelled");
Status($"[WebSocket] Telemetry loop ended - State: {socket.State}, Cancelled: {cancelled}");
}
catch (OperationCanceledException) when (token.IsCancellationRequested)
{
Status("[WebSocket] Connection cancelled");
break;
}
catch (Exception ex)
{
_log.Warn($"[WebSocket] Connection error: {ex.Message}");
Status($"[WebSocket] Connection error: {ex.Message}", warn: true);
}
finally
{
Status($"[WebSocket] Cleaning up connection - Final state: {socket.State}");
lock (run)
{
if (ReferenceEquals(run.Outbound, outbound))
@ -220,7 +246,7 @@ internal sealed class BackendClient : IDisposable
if (token.IsCancellationRequested || run.Stopping)
break;
_log.Info("[WebSocket] Reconnecting in 2 seconds...");
Status("[WebSocket] Reconnecting in 2 seconds...");
try
{
await Task.Delay(ReconnectDelay, token).ConfigureAwait(false);
@ -249,7 +275,7 @@ internal sealed class BackendClient : IDisposable
}
catch (Exception ex)
{
_log.Warn($"[WebSocket] receive error: {ex.Message}");
Status($"[WebSocket] receive error: {ex.Message}", warn: true);
return;
}
@ -282,10 +308,20 @@ internal sealed class BackendClient : IDisposable
}
catch (Exception ex)
{
_log.Warn($"[WebSocket] Send error: {ex.Message}");
Status($"[WebSocket] Send error: {ex.Message}", warn: true);
}
}
/// <summary>Queues a line for chat, then logs it; queued first so a line seen in the log is already drainable.</summary>
private void Status(string line, bool warn = false)
{
_status.Enqueue(line);
if (warn)
_log.Warn(line);
else
_log.Info(line);
}
private static Task SendFrameAsync(ClientWebSocket socket, string text, CancellationToken token)
{
byte[] bytes = Encoding.UTF8.GetBytes(text);

View file

@ -17,6 +17,12 @@ namespace OpenAC.MosswartMassacre.Backend;
/// <c>stop_radar</c> toggle the radar; any other text is run as if typed into
/// the chat bar. Commands run one per tick, as the original ran one per timer
/// tick.
/// <para>
/// The socket's status and error lines reach the chat here, on the tick
/// thread: each tick drains them before anything else, and a start or stop
/// made on this thread drains at once so its line comes out in the order the
/// original printed it.
/// </para>
/// </remarks>
internal sealed class BackendSession : IMmFeature
{
@ -47,11 +53,13 @@ internal sealed class BackendSession : IMmFeature
{
if (_context.Settings.Current.WebSocketEnabled)
StartSocket();
DrainStatus();
}
public void OnLogoff()
{
_context.Backend.Stop();
DrainStatus();
_pendingCommands.Clear();
}
@ -64,6 +72,8 @@ internal sealed class BackendSession : IMmFeature
/// <summary>Drains the socket. Public so tests can pump without a host tick.</summary>
public void Pump()
{
DrainStatus();
while (_context.Backend.TryTakeConnection(out BackendConnection connection))
OnConnected(connection.ConnectedAtUtc);
@ -93,6 +103,13 @@ internal sealed class BackendSession : IMmFeature
}
}
/// <summary>Writes the socket's queued status lines to chat. Tick thread only.</summary>
private void DrainStatus()
{
while (_context.Backend.TryTakeStatus(out string line))
_context.Chat.Write(line);
}
/// <summary>Commands queued from the backend and not yet run.</summary>
internal int PendingCommandCount => _pendingCommands.Count;
@ -188,6 +205,7 @@ internal sealed class BackendSession : IMmFeature
_context.Settings.Save();
if (_context.IsLoggedIn)
StartSocket();
DrainStatus();
_context.Chat.Write("WS streaming ENABLED.");
}
else if (sub.Equals("disable", StringComparison.OrdinalIgnoreCase))
@ -195,6 +213,7 @@ internal sealed class BackendSession : IMmFeature
_context.Settings.Current.WebSocketEnabled = false;
_context.Settings.Save();
_context.Backend.Stop();
DrainStatus();
_context.Chat.Write("WS streaming DISABLED.");
}
else if (sub.Equals("url", StringComparison.OrdinalIgnoreCase) && args.Length > 2)

View file

@ -139,6 +139,8 @@ internal sealed class VitalSharing : IMmFeature, IBackendConnectionListener, ISh
/// no player id, unlike the first; the backend accepts both.
/// </summary>
public void OnBackendConnected()
{
try
{
if (!_context.Settings.Current.VitalSharingEnabled || !_context.IsLoggedIn)
return;
@ -149,7 +151,12 @@ internal sealed class VitalSharing : IMmFeature, IBackendConnectionListener, ISh
character_name = _context.CharacterName,
tags = Tags,
});
_context.Verbose("[WebSocket] Vital sharing re-subscribed");
_context.Chat.Write("[WebSocket] Vital sharing re-subscribed");
}
catch (Exception ex)
{
_context.Chat.Write($"[WebSocket] Re-subscribe error: {ex.Message}");
}
}
// ── Outgoing ──

View file

@ -31,6 +31,7 @@ internal sealed class TelemetryStream : IMmFeature, IBackendConnectionListener
{
if (!_context.IsLoggedIn)
return;
_context.Chat.Write("[WebSocket] Starting telemetry loop");
SendNow();
_timer?.Dispose();
_timer = _context.Scheduler.Every(Interval, SendNow);

View file

@ -0,0 +1,122 @@
using OpenAC.MosswartMassacre.Tests.Fakes;
namespace OpenAC.MosswartMassacre.Tests.Backend;
/// <summary>
/// The socket's status and error lines reach the player's chat, worded as the
/// original worded them, and only through the tick: the socket thread never
/// writes to chat itself.
/// </summary>
public sealed class BackendChatStatusTests
{
private const string P = "[Mosswart Massacre] ";
[Fact]
public async Task Enable_shows_connecting_then_connected_and_registered_on_the_next_tick()
{
using var server = new TestSocketServer();
using var harness = new PluginHarness("Mossy");
harness.Login();
harness.Context.Settings.Current.BackendUrl = server.Uri.ToString();
harness.Context.Commands.Dispatch("ws enable");
Assert.Equal([P + "[WebSocket] connecting…", P + "WS streaming ENABLED."], Tail(harness, 2));
server.NextFrame();
Assert.True(await BackendClientTests.WaitAsync(() => harness.Context.Backend.IsConnected));
Assert.DoesNotContain(P + "[WebSocket] CONNECTED", harness.ChatLines);
int before = harness.ChatLines.Count;
int tickThread = Environment.CurrentManagedThreadId;
harness.Tick(0.01);
Assert.Equal(
[P + "[WebSocket] CONNECTED", P + "[WebSocket] REGISTERED", P + "[WebSocket] Starting telemetry loop"],
Tail(harness, 3));
AssertPostedOn(harness, before, tickThread);
}
[Fact]
public async Task Streaming_saved_on_connects_at_login_and_says_so()
{
using var server = new TestSocketServer();
using var harness = new PluginHarness("Mossy");
harness.Context.Settings.Current.WebSocketEnabled = true;
harness.Context.Settings.Current.BackendUrl = server.Uri.ToString();
harness.Context.Settings.Save();
harness.Host.Automation.Character.IsInWorld = true;
harness.Plugin.OnLogin();
Assert.Contains(P + "[WebSocket] connecting…", harness.ChatLines);
Assert.Equal("{\"type\":\"register\",\"player_name\":\"Mossy\"}", server.NextFrame());
Assert.True(await BackendClientTests.WaitAsync(() => harness.Context.Backend.IsConnected));
harness.Tick(0.01);
Assert.Contains(P + "[WebSocket] REGISTERED", harness.ChatLines);
}
[Fact]
public async Task A_failed_connect_shows_the_error_cleanup_and_reconnect_on_the_next_tick()
{
using var harness = new PluginHarness("Mossy");
harness.Login();
// Nothing listens on port 1: the connect is refused.
harness.Context.Settings.Current.BackendUrl = "ws://127.0.0.1:1/websocket/";
var log = (RecordingLogger)harness.Host.Log;
harness.Context.Commands.Dispatch("ws enable");
Assert.True(await BackendClientTests.WaitAsync(
() => log.Any(line => line.Contains("[WebSocket] Reconnecting in 2 seconds...")), 15000));
Assert.DoesNotContain(harness.ChatLines, line => line.Contains("Connection error"));
int before = harness.ChatLines.Count;
int tickThread = Environment.CurrentManagedThreadId;
harness.Tick(0.01);
string[] tail = Tail(harness, 3);
Assert.StartsWith(P + "[WebSocket] Connection error: ", tail[0]);
Assert.StartsWith(P + "[WebSocket] Cleaning up connection - Final state: ", tail[1]);
Assert.Equal(P + "[WebSocket] Reconnecting in 2 seconds...", tail[2]);
AssertPostedOn(harness, before, tickThread);
}
[Fact]
public async Task Disable_shows_disabled_before_the_command_reply()
{
using var server = new TestSocketServer();
using var harness = new PluginHarness("Mossy");
harness.Login();
harness.Context.Settings.Current.BackendUrl = server.Uri.ToString();
harness.Context.Commands.Dispatch("ws enable");
server.NextFrame();
Assert.True(await BackendClientTests.WaitAsync(() => harness.Context.Backend.IsConnected));
harness.Tick(0.01);
harness.Context.Commands.Dispatch("ws disable");
Assert.Equal([P + "[WebSocket] DISABLED", P + "WS streaming DISABLED."], Tail(harness, 2));
}
[Fact]
public void A_reconnect_with_sharing_on_says_it_re_subscribed_even_without_verbose_logging()
{
using var harness = new PluginHarness("Mossy");
harness.Login();
harness.Context.Settings.Current.VitalSharingEnabled = true;
harness.Context.Settings.Current.VerboseLogging = false;
harness.SimulateConnected();
Assert.Equal(
[P + "[WebSocket] Vital sharing re-subscribed", P + "[WebSocket] Starting telemetry loop"],
Tail(harness, 2));
}
/// <summary>Every chat line from <paramref name="from"/> on was written on <paramref name="threadId"/>, the thread that ran the tick.</summary>
private static void AssertPostedOn(PluginHarness harness, int from, int threadId) =>
Assert.All(harness.Host.Automation.Chat.PostedThreadIds.Skip(from), id => Assert.Equal(threadId, id));
private static string[] Tail(PluginHarness harness, int count) =>
[.. harness.ChatLines.Skip(Math.Max(0, harness.ChatLines.Count - count))];
}

View file

@ -431,7 +431,14 @@ public sealed class RecordingChat : IPluginChat
return message;
}
public void PostSystemMessage(string text) => Posted.Add(text);
/// <summary>The managed thread each <see cref="PostSystemMessage"/> ran on.</summary>
public List<int> PostedThreadIds { get; } = [];
public void PostSystemMessage(string text)
{
Posted.Add(text);
PostedThreadIds.Add(Environment.CurrentManagedThreadId);
}
public void PostMessage(string text, int logTextType)
{

View file

@ -295,12 +295,25 @@ public sealed class RecordingLogger : IPluginLogger
{
public List<string> Messages { get; } = [];
public void Info(string message) => Messages.Add("info: " + message);
public void Info(string message) => Add("info: " + message);
public void Warn(string message) => Messages.Add("warn: " + message);
public void Warn(string message) => Add("warn: " + message);
public void Error(string message, Exception? exception = null)
=> Messages.Add("error: " + message);
=> Add("error: " + message);
/// <summary>True when a logged line matches; safe while a background thread logs.</summary>
public bool Any(Func<string, bool> match)
{
lock (Messages)
return Messages.Any(match);
}
private void Add(string line)
{
lock (Messages)
Messages.Add(line);
}
}
public sealed class FakeGameState : IGameState