Skip to content

Commit 964578c

Browse files
HandyS11claude
andcommitted
feat(connections): carry IsConnected/WasConnected on ConnectionStatusChangedEvent
WasConnected is computed from in-process published history, not the persisted store (which would still say Connected right after boot). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent d2cb081 commit 964578c

10 files changed

Lines changed: 107 additions & 12 deletions

File tree

src/RustPlusBot.Abstractions/Events/ConnectionStatusChangedEvent.cs

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,4 +3,14 @@ namespace RustPlusBot.Abstractions.Events;
33
/// <summary>Published when a server's live-connection state changes, so #info can re-render.</summary>
44
/// <param name="GuildId">The owning guild snowflake.</param>
55
/// <param name="ServerId">The server whose connection state changed.</param>
6-
public sealed record ConnectionStatusChangedEvent(ulong GuildId, Guid ServerId);
6+
/// <param name="IsConnected">True when the new status is Connected.</param>
7+
/// <param name="WasConnected">
8+
/// True when the previous status published in this process was Connected. Deliberately
9+
/// in-process (not store-derived): the persisted status survives restarts and would still
10+
/// read Connected right after boot, re-triggering unreachable sweeps on every startup.
11+
/// </param>
12+
public sealed record ConnectionStatusChangedEvent(
13+
ulong GuildId,
14+
Guid ServerId,
15+
bool IsConnected,
16+
bool WasConnected);

src/RustPlusBot.Features.Connections/Supervisor/ConnectionSupervisor.cs

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,10 @@ internal sealed partial class ConnectionSupervisor(
4646
private readonly SemaphoreSlim _gate = new(1, 1);
4747
private readonly ConcurrentDictionary<(ulong Guild, Guid Server), LiveSocket> _liveSockets = new();
4848
private readonly ConnectionOptions _options = options.Value;
49+
50+
/// <summary>Last status published per key IN THIS PROCESS — the store's persisted status survives restarts and would falsely report Connected at boot.</summary>
51+
private readonly ConcurrentDictionary<(ulong Guild, Guid Server), ConnectionStatus> _publishedStatuses = new();
52+
4953
private readonly CancellationTokenSource _shutdown = new();
5054
private bool _disposed;
5155

@@ -953,7 +957,12 @@ private async Task PublishStatusAsync(
953957

954958
if (changed)
955959
{
956-
await eventBus.PublishAsync(new ConnectionStatusChangedEvent(key.Guild, key.Server), ct)
960+
var wasConnected = _publishedStatuses.TryGetValue(key, out var previous)
961+
&& previous == ConnectionStatus.Connected;
962+
_publishedStatuses[key] = status;
963+
await eventBus.PublishAsync(
964+
new ConnectionStatusChangedEvent(key.Guild, key.Server,
965+
status == ConnectionStatus.Connected, wasConnected), ct)
957966
.ConfigureAwait(false);
958967
}
959968
}

tests/RustPlusBot.Features.Alarms.Tests/AlarmStateRelayTests.cs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -380,7 +380,8 @@ public async Task ConnectionStatus_not_connected_refreshes_all_alarms_unreachabl
380380
]);
381381

382382
await h.Relay.HandleConnectionStatusAsync(
383-
new ConnectionStatusChangedEvent(10UL, serverId), CancellationToken.None);
383+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: false, WasConnected: true),
384+
CancellationToken.None);
384385

385386
await h.Refresher.Received(1).RefreshAsync(
386387
Arg.Is<SmartAlarm>(a => a.EntityId == 42UL), unreachable: true, Arg.Any<CancellationToken>());
@@ -402,7 +403,8 @@ public async Task ConnectionStatus_connected_does_nothing()
402403
});
403404

404405
await h.Relay.HandleConnectionStatusAsync(
405-
new ConnectionStatusChangedEvent(10UL, serverId), CancellationToken.None);
406+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: true, WasConnected: true),
407+
CancellationToken.None);
406408

407409
await h.Refresher.DidNotReceive().RefreshAsync(
408410
Arg.Any<ulong>(), Arg.Any<Guid>(), Arg.Any<ulong>(), Arg.Any<bool>(), Arg.Any<CancellationToken>());

tests/RustPlusBot.Features.Alarms.Tests/Hosting/AlarmsHostedServiceTests.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -197,7 +197,8 @@ public async Task ConnectionStatusChangedEvent_non_connected_routes_to_relay_and
197197
&& !h.Refresher.ReceivedCalls().Any(c =>
198198
c.GetMethodInfo().Name == nameof(IAlarmRefresher.RefreshAsync)))
199199
{
200-
await h.Bus.PublishAsync(new ConnectionStatusChangedEvent(10UL, serverId));
200+
await h.Bus.PublishAsync(
201+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: false, WasConnected: true));
201202
await Task.Delay(20);
202203
}
203204

tests/RustPlusBot.Features.Connections.Tests/ConnectionSupervisorTests.cs

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -231,6 +231,72 @@ public async Task Heartbeat_Unreachable_ReconnectsAndRecovers()
231231
Assert.True(source.CreateCount >= 2);
232232
}
233233

234+
/// <summary>
235+
/// WasConnected must be computed from in-process state (the DB status survives restarts and
236+
/// would claim Connected at boot): statuses before the first Connected carry false; the drop
237+
/// after a Connected carries true.
238+
/// </summary>
239+
[Fact]
240+
public async Task StatusEvents_CarryWasConnected_OnlyAfterAConnectedDrop()
241+
{
242+
var source = new FakeRustSocketSource();
243+
source.EnqueueConnect(SocketConnectOutcome.Connected);
244+
source.EnqueueHeartbeat(HeartbeatResult.Ok(2)); // first heartbeat -> Connected
245+
source.EnqueueHeartbeat(HeartbeatResult.Unreachable); // next heartbeat -> drop
246+
source.EnqueueConnect(SocketConnectOutcome.Connected); // reconnect
247+
source.EnqueueHeartbeat(HeartbeatResult.Ok(4));
248+
await using var h = CreateHarness(source);
249+
var (serverId, _, _) = await SeedAsync(h.Provider);
250+
251+
using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30));
252+
var events = new List<ConnectionStatusChangedEvent>();
253+
_ = Task.Run(async () =>
254+
{
255+
await foreach (var e in h.Bus.SubscribeAsync<ConnectionStatusChangedEvent>(cts.Token))
256+
{
257+
lock (events)
258+
{
259+
events.Add(e);
260+
}
261+
}
262+
}, cts.Token);
263+
264+
await h.Supervisor.EnsureConnectionAsync(10UL, serverId);
265+
var recovered = await WaitForStateAsync(
266+
h.Provider, serverId, s => s.Status == ConnectionStatus.Connected && s.PlayerCount == 4);
267+
Assert.NotNull(recovered);
268+
269+
// The bus delivers asynchronously; wait until the collector has seen the drop.
270+
var deadline = DateTimeOffset.UtcNow.AddSeconds(10);
271+
while (DateTimeOffset.UtcNow < deadline)
272+
{
273+
lock (events)
274+
{
275+
if (events.Any(e => !e.IsConnected && e.WasConnected))
276+
{
277+
break;
278+
}
279+
}
280+
281+
await Task.Delay(15);
282+
}
283+
284+
await cts.CancelAsync();
285+
ConnectionStatusChangedEvent[] snapshot;
286+
lock (events)
287+
{
288+
snapshot = [.. events];
289+
}
290+
291+
var firstConnected = Array.FindIndex(snapshot, e => e.IsConnected);
292+
Assert.True(firstConnected >= 0, "expected a Connected status event");
293+
Assert.All(snapshot.Take(firstConnected), e => Assert.False(e.WasConnected));
294+
var drop = Array.FindIndex(
295+
snapshot, firstConnected, snapshot.Length - firstConnected, e => !e.IsConnected);
296+
Assert.True(drop > firstConnected, "expected a drop event after Connected");
297+
Assert.True(snapshot[drop].WasConnected);
298+
}
299+
234300
[Fact]
235301
public async Task StartAll_StartsAConnectionPerConnectableServer()
236302
{

tests/RustPlusBot.Features.StorageMonitors.Tests/Hosting/StorageMonitorsHostedServiceTests.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -160,7 +160,8 @@ public async Task ConnectionStatusChangedEvent_non_connected_routes_to_relay_and
160160
&& !h.Poster.ReceivedCalls().Any(c =>
161161
c.GetMethodInfo().Name == nameof(IStorageMonitorChannelPoster.EnsureAsync)))
162162
{
163-
await h.Bus.PublishAsync(new ConnectionStatusChangedEvent(Guild, serverId));
163+
await h.Bus.PublishAsync(
164+
new ConnectionStatusChangedEvent(Guild, serverId, IsConnected: false, WasConnected: true));
164165
await Task.Delay(20);
165166
}
166167

tests/RustPlusBot.Features.StorageMonitors.Tests/StorageMonitorStateRelayTests.cs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -114,7 +114,8 @@ public async Task HandleConnectionStatusAsync_NotConnected_PostsUnreachable()
114114
]);
115115

116116
await h.Relay.HandleConnectionStatusAsync(
117-
new ConnectionStatusChangedEvent(Guild, Server), CancellationToken.None);
117+
new ConnectionStatusChangedEvent(Guild, Server, IsConnected: false, WasConnected: true),
118+
CancellationToken.None);
118119

119120
await h.Poster.Received(1).EnsureAsync(555UL, Arg.Any<ulong?>(),
120121
Arg.Any<global::Discord.Embed>(), Arg.Any<global::Discord.MessageComponent>(),
@@ -132,7 +133,8 @@ public async Task HandleConnectionStatusAsync_Connected_DoesNothing()
132133
});
133134

134135
await h.Relay.HandleConnectionStatusAsync(
135-
new ConnectionStatusChangedEvent(Guild, Server), CancellationToken.None);
136+
new ConnectionStatusChangedEvent(Guild, Server, IsConnected: true, WasConnected: true),
137+
CancellationToken.None);
136138

137139
await h.Poster.DidNotReceive().EnsureAsync(Arg.Any<ulong>(), Arg.Any<ulong?>(),
138140
Arg.Any<global::Discord.Embed>(), Arg.Any<global::Discord.MessageComponent>(),

tests/RustPlusBot.Features.Switches.Tests/Hosting/SwitchesHostedServiceTests.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -152,7 +152,8 @@ public async Task ConnectionStatusChangedEvent_non_connected_routes_to_relay_and
152152
&& !h.Poster.ReceivedCalls().Any(c =>
153153
c.GetMethodInfo().Name == nameof(ISwitchChannelPoster.EnsureAsync)))
154154
{
155-
await h.Bus.PublishAsync(new ConnectionStatusChangedEvent(10UL, serverId));
155+
await h.Bus.PublishAsync(
156+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: false, WasConnected: true));
156157
await Task.Delay(20);
157158
}
158159

tests/RustPlusBot.Features.Switches.Tests/SwitchStateRelayTests.cs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -89,7 +89,8 @@ public async Task ConnectionStatus_not_connected_marks_switches_unreachable()
8989
]);
9090

9191
await h.Relay.HandleConnectionStatusAsync(
92-
new ConnectionStatusChangedEvent(10UL, serverId), CancellationToken.None);
92+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: false, WasConnected: true),
93+
CancellationToken.None);
9394

9495
await h.Poster.Received(1).EnsureAsync(777UL, 900UL, Arg.Any<global::Discord.Embed>(),
9596
Arg.Any<global::Discord.MessageComponent>(), Arg.Any<CancellationToken>());
@@ -107,7 +108,8 @@ public async Task ConnectionStatus_connected_does_nothing()
107108
});
108109

109110
await h.Relay.HandleConnectionStatusAsync(
110-
new ConnectionStatusChangedEvent(10UL, serverId), CancellationToken.None);
111+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: true, WasConnected: true),
112+
CancellationToken.None);
111113

112114
await h.Poster.DidNotReceive().EnsureAsync(Arg.Any<ulong>(), Arg.Any<ulong?>(),
113115
Arg.Any<global::Discord.Embed>(),

tests/RustPlusBot.Features.Workspace.Tests/Hosting/WorkspaceConnectionStatusTests.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,8 @@ public async Task ConnectionStatusChanged_ReconcilesThatServer()
3232
&& !reconciler.ReceivedCalls().Any(c =>
3333
c.GetMethodInfo().Name == nameof(IWorkspaceReconciler.ReconcileServerAsync)))
3434
{
35-
await bus.PublishAsync(new ConnectionStatusChangedEvent(10UL, serverId));
35+
await bus.PublishAsync(
36+
new ConnectionStatusChangedEvent(10UL, serverId, IsConnected: false, WasConnected: false));
3637
await Task.Delay(20);
3738
}
3839

0 commit comments

Comments
 (0)