Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
136 changes: 136 additions & 0 deletions Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
extern alias Eventing;

using System;
using System.Collections.Generic;
using System.Security.Claims;
using System.Threading;
using System.Threading.Tasks;
using FluentAssertions;
using Microsoft.AspNetCore.SignalR;
using Moq;
using NUnit.Framework;
using Resgrid.Model;
using Resgrid.Model.Services;
using EventingHubs = Eventing::Resgrid.Web.Eventing.Hubs;
using ServicesHubs = Resgrid.Web.Services.Hubs;

namespace Resgrid.Tests.Web.Eventing
{
/// <summary>
/// SignalR does not wait for in-flight hub calls before disconnecting, and on the Redis backplane a group
/// change for a connection no server holds waits 30 seconds for an ack and then throws. A hub call that
/// finds its connection closed skips the group change.
/// </summary>
[TestFixture]
public class ClosedConnectionHubTests
{
private const int DepartmentId = 7;
private const int LinkId = 3;
private const int CallId = 5;
private const string ChannelId = "ch-1";
private const string UserId = "b0000000-0000-0000-0000-000000000001";

private static readonly Func<Harness, Task> NoArrange = _ => Task.CompletedTask;

private static readonly TestCaseData[] GroupChangingCalls =
{
Case("Eventing.Connect", h => h.Eventing.Connect(DepartmentId)),
Case("Eventing.SubscribeToDepartmentLink", h => h.Eventing.SubscribeToDepartmentLink(LinkId)),
Case("Eventing.UnsubscribeToDepartmentLink", h => h.Eventing.UnsubscribeToDepartmentLink(LinkId)),
Case("Eventing.SubscribeToCall", h => h.Eventing.SubscribeToCall(CallId)),
Case("Eventing.UnsubscribeToCall", h => h.Eventing.UnsubscribeToCall(CallId)),
Case("Services.Connect", h => h.Services.Connect(DepartmentId)),
Case("Services.SubscribeToDepartmentLink", h => h.Services.SubscribeToDepartmentLink(LinkId)),
Case("Services.UnsubscribeToDepartmentLink", h => h.Services.UnsubscribeToDepartmentLink(LinkId)),
Case("Chat.Connect", h => h.Chat.Connect()),
Case("Chat.JoinChannel", h => h.Chat.JoinChannel(ChannelId)),
Case("Chat.LeaveChannel", h => h.Chat.LeaveChannel(ChannelId), arrange: h => h.Chat.JoinChannel(ChannelId))
};

[TestCaseSource(nameof(GroupChangingCalls))]
public async Task Open_connection_changes_groups(Func<Harness, Task> arrange, Func<Harness, Task> act)
{
var harness = new Harness();
await arrange(harness);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

kody code-review Kody Rules high

Unhandled task rejection in await arrange(harness) prevents arrangement failures from being reported with test context. Catch Exception and call Assert.Fail with the failure details in Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs:57, :66, and :70, and Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs:192 and :194.

Kody rule violation: Handle async operations with proper error handling

try { await arrange(harness); } catch (Exception exception) { Assert.Fail($"Arrange failed: {exception}"); }
Prompt for LLM

File Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs:

Line 54:

Unhandled task rejection in `await arrange(harness)` prevents arrangement failures from being reported with test context. Catch `Exception` and call `Assert.Fail` with the failure details in `Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs:57`, `:66`, and `:70`, and `Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs:192` and `:194`.

Suggested Code:

try { await arrange(harness); } catch (Exception exception) { Assert.Fail($"Arrange failed: {exception}"); }

Talk to Kody by mentioning @kody

Was this suggestion helpful? React with 👍 or 👎 to help Kody learn from this interaction.

​

​

harness.Groups.Invocations.Clear();

await act(harness);

harness.Groups.Invocations.Should().NotBeEmpty("otherwise the closed-connection case proves nothing");
}

[TestCaseSource(nameof(GroupChangingCalls))]
public async Task Closed_connection_leaves_groups_alone(Func<Harness, Task> arrange, Func<Harness, Task> act)
{
var harness = new Harness();
await arrange(harness);
harness.Groups.Invocations.Clear();
harness.CloseConnection();

await act(harness);

harness.Groups.Invocations.Should().BeEmpty();
}

private static TestCaseData Case(string name, Func<Harness, Task> act, Func<Harness, Task> arrange = null) =>
new TestCaseData(arrange ?? NoArrange, act).SetArgDisplayNames(name);

public sealed class Harness
{
private readonly CancellationTokenSource _connectionAborted = new CancellationTokenSource();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

kody code-review Kody Rules high

Resource leak occurs because the CancellationTokenSource stored in _connectionAborted remains undisposed after the test completes. Implement IDisposable on Harness and dispose _connectionAborted during teardown.

Kody rule violation: Use using statements for disposable resources

private readonly CancellationTokenSource _connectionAborted = new CancellationTokenSource();

public void Dispose() => _connectionAborted.Dispose();
Prompt for LLM

File Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs:

Line 80:

Resource leak occurs because the `CancellationTokenSource` stored in `_connectionAborted` remains undisposed after the test completes. Implement `IDisposable` on `Harness` and dispose `_connectionAborted` during teardown.

Suggested Code:

private readonly CancellationTokenSource _connectionAborted = new CancellationTokenSource();

public void Dispose() => _connectionAborted.Dispose();

Talk to Kody by mentioning @kody

Was this suggestion helpful? React with 👍 or 👎 to help Kody learn from this interaction.

​

​


public Harness()
{
var context = new Mock<HubCallerContext>();
context.SetupGet(x => x.ConnectionId).Returns("conn-1");
context.SetupGet(x => x.ConnectionAborted).Returns(_connectionAborted.Token);
context.SetupGet(x => x.Items).Returns(new Dictionary<object, object>());
context.SetupGet(x => x.User).Returns(new ClaimsPrincipal(new ClaimsIdentity(new[]
{
new Claim(ClaimTypes.PrimarySid, UserId),
new Claim(ClaimTypes.PrimaryGroupSid, DepartmentId.ToString())
}, "test")));

var clients = new Mock<IHubCallerClients>();
clients.SetupGet(x => x.Caller).Returns(Mock.Of<ISingleClientProxy>());
clients.Setup(x => x.Group(It.IsAny<string>())).Returns(Mock.Of<IClientProxy>());

var links = new Mock<IDepartmentLinksService>();
links.Setup(x => x.GetLinkByIdAsync(LinkId))
.ReturnsAsync(new DepartmentLink { DepartmentId = DepartmentId, LinkedDepartmentId = 8, LinkEnabled = true });

var calls = new Mock<ICallsService>();
calls.Setup(x => x.GetCallByIdAsync(CallId, It.IsAny<bool>())).ReturnsAsync(new Call { DepartmentId = DepartmentId });

var channels = new Mock<IChatChannelService>();
channels.Setup(x => x.GetChannelByIdAsync(ChannelId))
.ReturnsAsync(new ChatChannel { ChatChannelId = ChannelId, DepartmentId = DepartmentId });

var permissions = new Mock<IChatPermissionService>();
permissions.Setup(x => x.IsActiveDepartmentUserAsync(DepartmentId, UserId)).ReturnsAsync(true);
permissions.Setup(x => x.CanAccessChannelAsync(It.IsAny<ChatChannel>(), UserId, null)).ReturnsAsync(true);
permissions.Setup(x => x.GetChannelAccessVersionAsync(ChannelId)).ReturnsAsync("v1");

Eventing = Attach(new EventingHubs.EventingHub(links.Object, calls.Object), context.Object, clients.Object);
Services = Attach(new ServicesHubs.EventingHub(links.Object), context.Object, clients.Object);
Chat = Attach(new EventingHubs.ChatHub(channels.Object, permissions.Object, Mock.Of<IChatMessageService>(),
Mock.Of<IChatPresenceService>()), context.Object, clients.Object);
}

public Mock<IGroupManager> Groups { get; } = new Mock<IGroupManager>();
public EventingHubs.EventingHub Eventing { get; }
public ServicesHubs.EventingHub Services { get; }
public EventingHubs.ChatHub Chat { get; }

public void CloseConnection() => _connectionAborted.Cancel();

private T Attach<T>(T hub, HubCallerContext context, IHubCallerClients clients) where T : Hub
{
hub.Context = context;
hub.Groups = Groups.Object;
hub.Clients = clients;
return hub;
}
}
}
}
22 changes: 18 additions & 4 deletions Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -142,8 +142,8 @@ public async Task Periodic_sync_applies_matrix_changes_to_tracked_connections()
var hubContext = new Mock<IHubContext<GeolocationHub>>();
hubContext.SetupGet(x => x.Groups).Returns(groups.Object);
var sync = new GeolocationVisibilitySync(tracker, new GeolocationMembership(_visibility.Object), hubContext.Object);
tracker.GetOrAdd("conn-1", DepartmentId, Viewer).SubscribeToDepartmentMap();
tracker.GetOrAdd("conn-2", DepartmentId, Subject).SubscribeToDepartmentMap();
tracker.Track("conn-1", DepartmentId, Viewer, CancellationToken.None).SubscribeToDepartmentMap();
tracker.Track("conn-2", DepartmentId, Subject, CancellationToken.None).SubscribeToDepartmentMap();
tracker.Remove("conn-2");

_visibility.Setup(x => x.GetVisibilitySetKeysForViewerAsync(DepartmentId, Viewer)).ReturnsAsync(new[] { "abc" });
Expand Down Expand Up @@ -184,6 +184,20 @@ public async Task Hub_geolocation_connect_joins_the_viewers_groups_and_acknowled
tracker.Contains("conn-1").Should().BeFalse();
}

[Test]
public async Task Hub_call_that_finishes_after_its_connection_closed_leaves_nothing_tracked()
{
// SignalR runs OnDisconnectedAsync without waiting for in-flight hub calls.
var (hub, groups, caller, tracker) = CreateHub(Mock.Of<IUnitsService>(), Viewer, new CancellationToken(canceled: true));
await hub.OnDisconnectedAsync(null);

await hub.GeolocationConnect();

tracker.Contains("conn-1").Should().BeFalse("nothing would remove it, and every periodic sync would wait on it");
groups.Verify(x => x.AddToGroupAsync(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<CancellationToken>()), Times.Never);
caller.Verify(x => x.SendCoreAsync(It.IsAny<string>(), It.IsAny<object[]>(), It.IsAny<CancellationToken>()), Times.Never);
}

[Test]
public void Hub_location_publish_methods_are_reserved_for_the_internal_publisher()
{
Expand Down Expand Up @@ -216,7 +230,7 @@ private static (Mock<IHubClients<IClientProxy>> Clients, List<IReadOnlyList<stri
}

private (GeolocationHub Hub, Mock<IGroupManager> Groups, Mock<ISingleClientProxy> Caller, GeolocationConnectionTracker Tracker) CreateHub(
IUnitsService unitsService, string userId)
IUnitsService unitsService, string userId, CancellationToken connectionAborted = default)
{
var tracker = new GeolocationConnectionTracker();
var groups = new Mock<IGroupManager>();
Expand All @@ -226,7 +240,7 @@ private static (Mock<IHubClients<IClientProxy>> Clients, List<IReadOnlyList<stri

var context = new Mock<HubCallerContext>();
context.SetupGet(x => x.ConnectionId).Returns("conn-1");
context.SetupGet(x => x.ConnectionAborted).Returns(CancellationToken.None);
context.SetupGet(x => x.ConnectionAborted).Returns(connectionAborted);
context.SetupGet(x => x.User).Returns(new ClaimsPrincipal(new ClaimsIdentity(new[]
{
new Claim(ClaimTypes.PrimarySid, userId),
Expand Down
15 changes: 14 additions & 1 deletion Web/Resgrid.Web.Eventing/Hubs/ChatHub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,13 @@ private string GetUserId()
return Context.User?.FindFirst(ClaimTypes.PrimarySid)?.Value ?? String.Empty;
}

/// <summary>
/// True when the connection closed while this call was running (SignalR does not wait for in-flight
/// calls before disconnecting). No server holds the connection any more, so the Redis backplane would
/// wait 30 seconds for a group ack that never comes and then throw.
/// </summary>
private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested;

public override async Task OnConnectedAsync()
{
var departmentId = GetDepartmentId();
Expand Down Expand Up @@ -170,6 +177,9 @@ public async Task Connect()
!await _chatPermissionService.IsActiveDepartmentUserAsync(departmentId, userId))
return;

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, $"chatuser:{departmentId}:{userId.ToLowerInvariant()}");
await Groups.AddToGroupAsync(Context.ConnectionId, $"chatdept:{departmentId}");

Expand All @@ -187,6 +197,9 @@ public async Task JoinChannel(string channelId, int? asUnitId = null)
if (groupName == null)
throw new HubException("Chat authorization is temporarily unavailable.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, groupName);
Context.Items[GetJoinedChannelGroupContextKey(channel.ChatChannelId)] = groupName;

Expand All @@ -202,7 +215,7 @@ public async Task LeaveChannel(string channelId)
var groupName = Context.Items.TryGetValue(contextKey, out var trackedGroupName)
? trackedGroupName as string
: null;
if (groupName != null)
if (groupName != null && !ConnectionClosed)
{
await Groups.RemoveFromGroupAsync(Context.ConnectionId, groupName);
Context.Items.Remove(contextKey);
Expand Down
20 changes: 18 additions & 2 deletions Web/Resgrid.Web.Eventing/Hubs/EventingHub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,9 @@ public async Task Connect(int departmentId)
if (authenticatedDepartmentId <= 0 || departmentId != authenticatedDepartmentId)
throw new HubException("Not authorized for this department.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, authenticatedDepartmentId.ToString());
await Clients.Caller.SendAsync("onConnected", Context.ConnectionId);
}
Expand All @@ -55,14 +58,17 @@ public async Task SubscribeToDepartmentLink(int linkId)
if (link == null || !link.LinkEnabled || !linkedDepartmentId.HasValue)
throw new HubException("Not authorized for this department link.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, linkedDepartmentId.Value.ToString());
}

public async Task UnsubscribeToDepartmentLink(int linkId)
{
var link = await _departmentLinksService.GetLinkByIdAsync(linkId);
var linkedDepartmentId = GetLinkedDepartmentForCaller(link);
if (linkedDepartmentId.HasValue)
if (linkedDepartmentId.HasValue && !ConnectionClosed)
await Groups.RemoveFromGroupAsync(Context.ConnectionId, linkedDepartmentId.Value.ToString());
}

Expand All @@ -72,11 +78,21 @@ public async Task SubscribeToCall(int callId)
if (call == null || call.DepartmentId != GetDepartmentId())
throw new HubException("Not authorized for this call.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, GetCallGroupName(callId));
}

public Task UnsubscribeToCall(int callId) =>
Groups.RemoveFromGroupAsync(Context.ConnectionId, GetCallGroupName(callId));
ConnectionClosed ? Task.CompletedTask : Groups.RemoveFromGroupAsync(Context.ConnectionId, GetCallGroupName(callId));

/// <summary>
/// True when the connection closed while this call was running (SignalR does not wait for in-flight
/// calls before disconnecting). No server holds the connection any more, so the Redis backplane would
/// wait 30 seconds for a group ack that never comes and then throw.
/// </summary>
private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested;

/// <summary>
/// Single source of truth for the per-call group name, so a publisher added later cannot drift from
Expand Down
12 changes: 11 additions & 1 deletion Web/Resgrid.Web.Eventing/Hubs/GeolocationHub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -65,9 +65,10 @@ private string GetUserId()
return Context.User?.FindFirst(ClaimTypes.PrimarySid)?.Value;
}

/// <summary>Null when the connection closed while this call was running.</summary>
private GeolocationConnection TrackConnection(int departmentId)
{
return _connectionTracker.GetOrAdd(Context.ConnectionId, departmentId, GetUserId());
return _connectionTracker.Track(Context.ConnectionId, departmentId, GetUserId(), Context.ConnectionAborted);
}

/// <summary>
Expand All @@ -81,6 +82,9 @@ public async Task GeolocationConnect()
if (departmentId > 0)
{
var connection = TrackConnection(departmentId);
if (connection == null)
return;

connection.SubscribeToDepartmentMap();
await _membership.SyncAsync(Groups, connection, Context.ConnectionAborted);

Expand Down Expand Up @@ -127,6 +131,9 @@ public async Task UnitLocationConnect(int unitId)
return;

var connection = TrackConnection(departmentId);
if (connection == null)
return;

connection.SubscribeToUnit(unitId);
await _membership.SyncAsync(Groups, connection, Context.ConnectionAborted);

Expand All @@ -148,6 +155,9 @@ public async Task PersonLocationConnect(string userId)
return;

var connection = TrackConnection(departmentId);
if (connection == null)
return;

connection.SubscribeToPerson(userId);
await _membership.SyncAsync(Groups, connection, Context.ConnectionAborted);

Expand Down
19 changes: 17 additions & 2 deletions Web/Resgrid.Web.Eventing/Services/GeolocationConnectionTracker.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,24 @@ public sealed class GeolocationConnectionTracker
private readonly ConcurrentDictionary<string, GeolocationConnection> _connections =
new ConcurrentDictionary<string, GeolocationConnection>(StringComparer.Ordinal);

public GeolocationConnection GetOrAdd(string connectionId, int departmentId, string userId)
/// <summary>
/// Returns the tracked connection, or null when it has already closed. SignalR does not wait for
/// in-flight hub calls before OnDisconnectedAsync, so a subscribe can land after the entry was removed;
/// nothing would remove it again, and every periodic sync would wait out a 30 second Redis group ack
/// for it. ConnectionAborted is cancelled before OnDisconnectedAsync runs, so checking it after the
/// add catches that case.
/// </summary>
public GeolocationConnection Track(string connectionId, int departmentId, string userId, CancellationToken connectionAborted)
{
return _connections.GetOrAdd(connectionId, id => new GeolocationConnection(id, departmentId, userId));
var connection = _connections.GetOrAdd(connectionId, id => new GeolocationConnection(id, departmentId, userId));

if (connectionAborted.IsCancellationRequested)
{
Remove(connectionId);
return null;
}

return connection;
}

public bool Contains(string connectionId) => _connections.ContainsKey(connectionId);
Expand Down
15 changes: 14 additions & 1 deletion Web/Resgrid.Web.Services/Hubs/EventingHub.cs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,9 @@ public async Task Connect(int departmentId)
if (authenticatedDepartmentId <= 0 || authenticatedDepartmentId != departmentId)
throw new HubException("Not authorized for this department.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, departmentId.ToString());
await Clients.Caller.SendAsync("onConnected", Context.ConnectionId);
}
Expand All @@ -42,17 +45,27 @@ public async Task SubscribeToDepartmentLink(int linkId)
if (link == null || !link.LinkEnabled || !linkedDepartmentId.HasValue)
throw new HubException("Not authorized for this department link.");

if (ConnectionClosed)
return;

await Groups.AddToGroupAsync(Context.ConnectionId, linkedDepartmentId.Value.ToString());
}

public async Task UnsubscribeToDepartmentLink(int linkId)
{
var link = await _departmentLinksService.GetLinkByIdAsync(linkId);
var linkedDepartmentId = GetLinkedDepartmentForCaller(link);
if (linkedDepartmentId.HasValue)
if (linkedDepartmentId.HasValue && !ConnectionClosed)
await Groups.RemoveFromGroupAsync(Context.ConnectionId, linkedDepartmentId.Value.ToString());
}

/// <summary>
/// True when the connection closed while this call was running (SignalR does not wait for in-flight
/// calls before disconnecting). No server holds the connection any more, so the Redis backplane would
/// wait 30 seconds for a group ack that never comes and then throw.
/// </summary>
private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested;

public Task PersonnelStatusUpdated(int departmentId, int id) =>
PublishAsync("personnelStatusUpdated", departmentId, id);

Expand Down
Loading
Loading