diff --git a/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs b/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs index 025d17e..eda77b6 100644 --- a/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs +++ b/src/OpenAC.MosswartMassacre/Backend/BackendClient.cs @@ -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. /// +/// +/// 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 +/// . +/// /// 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 _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? outbound; Task? sendLoop; @@ -120,10 +130,16 @@ internal sealed class BackendClient : IDisposable } } run.Cts.Dispose(); - _log.Info("[WebSocket] DISABLED"); }); } + /// + /// 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. + /// + public bool TryTakeStatus(out string line) => _status.TryDequeue(out line!); + /// Queues a frame. Dropped when no socket is open. 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( 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); } } + /// Queues a line for chat, then logs it; queued first so a line seen in the log is already drainable. + 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); diff --git a/src/OpenAC.MosswartMassacre/Backend/BackendSession.cs b/src/OpenAC.MosswartMassacre/Backend/BackendSession.cs index d740abe..2699bce 100644 --- a/src/OpenAC.MosswartMassacre/Backend/BackendSession.cs +++ b/src/OpenAC.MosswartMassacre/Backend/BackendSession.cs @@ -17,6 +17,12 @@ namespace OpenAC.MosswartMassacre.Backend; /// stop_radar 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. +/// +/// 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. +/// /// 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 /// Drains the socket. Public so tests can pump without a host tick. public void Pump() { + DrainStatus(); + while (_context.Backend.TryTakeConnection(out BackendConnection connection)) OnConnected(connection.ConnectedAtUtc); @@ -93,6 +103,13 @@ internal sealed class BackendSession : IMmFeature } } + /// Writes the socket's queued status lines to chat. Tick thread only. + private void DrainStatus() + { + while (_context.Backend.TryTakeStatus(out string line)) + _context.Chat.Write(line); + } + /// Commands queued from the backend and not yet run. 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) diff --git a/src/OpenAC.MosswartMassacre/Sharing/VitalSharing.cs b/src/OpenAC.MosswartMassacre/Sharing/VitalSharing.cs index 6141a54..6baf92b 100644 --- a/src/OpenAC.MosswartMassacre/Sharing/VitalSharing.cs +++ b/src/OpenAC.MosswartMassacre/Sharing/VitalSharing.cs @@ -140,16 +140,23 @@ internal sealed class VitalSharing : IMmFeature, IBackendConnectionListener, ISh /// public void OnBackendConnected() { - if (!_context.Settings.Current.VitalSharingEnabled || !_context.IsLoggedIn) - return; - _context.Send(new + try { - type = "share_subscribe", - timestamp = Wire.Timestamp(_context.Clock), - character_name = _context.CharacterName, - tags = Tags, - }); - _context.Verbose("[WebSocket] Vital sharing re-subscribed"); + if (!_context.Settings.Current.VitalSharingEnabled || !_context.IsLoggedIn) + return; + _context.Send(new + { + type = "share_subscribe", + timestamp = Wire.Timestamp(_context.Clock), + character_name = _context.CharacterName, + tags = Tags, + }); + _context.Chat.Write("[WebSocket] Vital sharing re-subscribed"); + } + catch (Exception ex) + { + _context.Chat.Write($"[WebSocket] Re-subscribe error: {ex.Message}"); + } } // ── Outgoing ── diff --git a/src/OpenAC.MosswartMassacre/Streams/TelemetryStream.cs b/src/OpenAC.MosswartMassacre/Streams/TelemetryStream.cs index 02f3a2d..a01a826 100644 --- a/src/OpenAC.MosswartMassacre/Streams/TelemetryStream.cs +++ b/src/OpenAC.MosswartMassacre/Streams/TelemetryStream.cs @@ -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); diff --git a/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendChatStatusTests.cs b/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendChatStatusTests.cs new file mode 100644 index 0000000..e9f8c39 --- /dev/null +++ b/tests/OpenAC.MosswartMassacre.Tests/Backend/BackendChatStatusTests.cs @@ -0,0 +1,122 @@ +using OpenAC.MosswartMassacre.Tests.Fakes; + +namespace OpenAC.MosswartMassacre.Tests.Backend; + +/// +/// 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. +/// +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)); + } + + /// Every chat line from on was written on , the thread that ran the tick. + 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))]; +} diff --git a/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeAutomation.cs b/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeAutomation.cs index 9315a5d..8817ef2 100644 --- a/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeAutomation.cs +++ b/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeAutomation.cs @@ -431,7 +431,14 @@ public sealed class RecordingChat : IPluginChat return message; } - public void PostSystemMessage(string text) => Posted.Add(text); + /// The managed thread each ran on. + public List PostedThreadIds { get; } = []; + + public void PostSystemMessage(string text) + { + Posted.Add(text); + PostedThreadIds.Add(Environment.CurrentManagedThreadId); + } public void PostMessage(string text, int logTextType) { diff --git a/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeHost.cs b/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeHost.cs index 16ed618..88e35ce 100644 --- a/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeHost.cs +++ b/tests/OpenAC.MosswartMassacre.Tests/Fakes/FakeHost.cs @@ -295,12 +295,25 @@ public sealed class RecordingLogger : IPluginLogger { public List 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); + + /// True when a logged line matches; safe while a background thread logs. + public bool Any(Func match) + { + lock (Messages) + return Messages.Any(match); + } + + private void Add(string line) + { + lock (Messages) + Messages.Add(line); + } } public sealed class FakeGameState : IGameState