diff --git a/Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs b/Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs new file mode 100644 index 000000000..dd516054c --- /dev/null +++ b/Tests/Resgrid.Tests/Web/Eventing/ClosedConnectionHubTests.cs @@ -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 +{ + /// + /// 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. + /// + [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 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 arrange, Func act) + { + var harness = new Harness(); + await arrange(harness); + 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 arrange, Func 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 act, Func arrange = null) => + new TestCaseData(arrange ?? NoArrange, act).SetArgDisplayNames(name); + + public sealed class Harness + { + private readonly CancellationTokenSource _connectionAborted = new CancellationTokenSource(); + + public Harness() + { + var context = new Mock(); + context.SetupGet(x => x.ConnectionId).Returns("conn-1"); + context.SetupGet(x => x.ConnectionAborted).Returns(_connectionAborted.Token); + context.SetupGet(x => x.Items).Returns(new Dictionary()); + 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(); + clients.SetupGet(x => x.Caller).Returns(Mock.Of()); + clients.Setup(x => x.Group(It.IsAny())).Returns(Mock.Of()); + + var links = new Mock(); + links.Setup(x => x.GetLinkByIdAsync(LinkId)) + .ReturnsAsync(new DepartmentLink { DepartmentId = DepartmentId, LinkedDepartmentId = 8, LinkEnabled = true }); + + var calls = new Mock(); + calls.Setup(x => x.GetCallByIdAsync(CallId, It.IsAny())).ReturnsAsync(new Call { DepartmentId = DepartmentId }); + + var channels = new Mock(); + channels.Setup(x => x.GetChannelByIdAsync(ChannelId)) + .ReturnsAsync(new ChatChannel { ChatChannelId = ChannelId, DepartmentId = DepartmentId }); + + var permissions = new Mock(); + permissions.Setup(x => x.IsActiveDepartmentUserAsync(DepartmentId, UserId)).ReturnsAsync(true); + permissions.Setup(x => x.CanAccessChannelAsync(It.IsAny(), 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(), + Mock.Of()), context.Object, clients.Object); + } + + public Mock Groups { get; } = new Mock(); + public EventingHubs.EventingHub Eventing { get; } + public ServicesHubs.EventingHub Services { get; } + public EventingHubs.ChatHub Chat { get; } + + public void CloseConnection() => _connectionAborted.Cancel(); + + private T Attach(T hub, HubCallerContext context, IHubCallerClients clients) where T : Hub + { + hub.Context = context; + hub.Groups = Groups.Object; + hub.Clients = clients; + return hub; + } + } + } +} diff --git a/Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs b/Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs index f20f3b9ae..11e5ab599 100644 --- a/Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs +++ b/Tests/Resgrid.Tests/Web/Eventing/GeolocationVisibilityTests.cs @@ -142,8 +142,8 @@ public async Task Periodic_sync_applies_matrix_changes_to_tracked_connections() var hubContext = new Mock>(); 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" }); @@ -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(), 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(), It.IsAny(), It.IsAny()), Times.Never); + caller.Verify(x => x.SendCoreAsync(It.IsAny(), It.IsAny(), It.IsAny()), Times.Never); + } + [Test] public void Hub_location_publish_methods_are_reserved_for_the_internal_publisher() { @@ -216,7 +230,7 @@ private static (Mock> Clients, List Groups, Mock Caller, GeolocationConnectionTracker Tracker) CreateHub( - IUnitsService unitsService, string userId) + IUnitsService unitsService, string userId, CancellationToken connectionAborted = default) { var tracker = new GeolocationConnectionTracker(); var groups = new Mock(); @@ -226,7 +240,7 @@ private static (Mock> Clients, List(); 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), diff --git a/Web/Resgrid.Web.Eventing/Hubs/ChatHub.cs b/Web/Resgrid.Web.Eventing/Hubs/ChatHub.cs index d1c558dd2..a1e69b879 100644 --- a/Web/Resgrid.Web.Eventing/Hubs/ChatHub.cs +++ b/Web/Resgrid.Web.Eventing/Hubs/ChatHub.cs @@ -75,6 +75,13 @@ private string GetUserId() return Context.User?.FindFirst(ClaimTypes.PrimarySid)?.Value ?? String.Empty; } + /// + /// 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. + /// + private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested; + public override async Task OnConnectedAsync() { var departmentId = GetDepartmentId(); @@ -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}"); @@ -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; @@ -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); diff --git a/Web/Resgrid.Web.Eventing/Hubs/EventingHub.cs b/Web/Resgrid.Web.Eventing/Hubs/EventingHub.cs index 756c6fbc0..77c696d24 100644 --- a/Web/Resgrid.Web.Eventing/Hubs/EventingHub.cs +++ b/Web/Resgrid.Web.Eventing/Hubs/EventingHub.cs @@ -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); } @@ -55,6 +58,9 @@ 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()); } @@ -62,7 +68,7 @@ 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()); } @@ -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)); + + /// + /// 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. + /// + private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested; /// /// Single source of truth for the per-call group name, so a publisher added later cannot drift from diff --git a/Web/Resgrid.Web.Eventing/Hubs/GeolocationHub.cs b/Web/Resgrid.Web.Eventing/Hubs/GeolocationHub.cs index 38f72c8b9..67f6cd8cb 100644 --- a/Web/Resgrid.Web.Eventing/Hubs/GeolocationHub.cs +++ b/Web/Resgrid.Web.Eventing/Hubs/GeolocationHub.cs @@ -65,9 +65,10 @@ private string GetUserId() return Context.User?.FindFirst(ClaimTypes.PrimarySid)?.Value; } + /// Null when the connection closed while this call was running. private GeolocationConnection TrackConnection(int departmentId) { - return _connectionTracker.GetOrAdd(Context.ConnectionId, departmentId, GetUserId()); + return _connectionTracker.Track(Context.ConnectionId, departmentId, GetUserId(), Context.ConnectionAborted); } /// @@ -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); @@ -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); @@ -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); diff --git a/Web/Resgrid.Web.Eventing/Services/GeolocationConnectionTracker.cs b/Web/Resgrid.Web.Eventing/Services/GeolocationConnectionTracker.cs index 6aa58964b..0f7904f45 100644 --- a/Web/Resgrid.Web.Eventing/Services/GeolocationConnectionTracker.cs +++ b/Web/Resgrid.Web.Eventing/Services/GeolocationConnectionTracker.cs @@ -16,9 +16,24 @@ public sealed class GeolocationConnectionTracker private readonly ConcurrentDictionary _connections = new ConcurrentDictionary(StringComparer.Ordinal); - public GeolocationConnection GetOrAdd(string connectionId, int departmentId, string userId) + /// + /// 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. + /// + 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); diff --git a/Web/Resgrid.Web.Services/Hubs/EventingHub.cs b/Web/Resgrid.Web.Services/Hubs/EventingHub.cs index 683773629..73c174e82 100644 --- a/Web/Resgrid.Web.Services/Hubs/EventingHub.cs +++ b/Web/Resgrid.Web.Services/Hubs/EventingHub.cs @@ -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); } @@ -42,6 +45,9 @@ 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()); } @@ -49,10 +55,17 @@ 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()); } + /// + /// 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. + /// + private bool ConnectionClosed => Context.ConnectionAborted.IsCancellationRequested; + public Task PersonnelStatusUpdated(int departmentId, int id) => PublishAsync("personnelStatusUpdated", departmentId, id); diff --git a/Web/Resgrid.Web.Services/Resgrid.Web.Services.xml b/Web/Resgrid.Web.Services/Resgrid.Web.Services.xml index 67722272e..ac6b6eb0b 100644 --- a/Web/Resgrid.Web.Services/Resgrid.Web.Services.xml +++ b/Web/Resgrid.Web.Services/Resgrid.Web.Services.xml @@ -7104,6 +7104,13 @@ the reader happened to see a "Z". + + + 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. + + Gets or sets the on authentication failed.