From a0fcda4b8375fa1c904a02433b12ac22d9cec238 Mon Sep 17 00:00:00 2001 From: Pawel Pabich Date: Mon, 17 Aug 2026 16:38:01 +1000 Subject: [PATCH 1/4] Add ability to observe the number of connections halibut opens for each polling subscription --- .../Support/TestConnectionsObserver.cs | 16 ++++++- .../Transport/Protocol/ProtocolFixture.cs | 2 +- .../Transport/SecureClientFixture.cs | 4 +- source/Halibut/HalibutRuntime.cs | 2 +- .../Transport/ActiveTcpConnectionsLimiter.cs | 3 +- .../Observability/IConnectionsObserver.cs | 21 ++++++++- .../Observability/NoOpConnectionsObserver.cs | 10 ++++ .../Protocol/MessageExchangeProtocol.cs | 47 ++++++++++--------- 8 files changed, 73 insertions(+), 32 deletions(-) diff --git a/source/Halibut.Tests/Support/TestConnectionsObserver.cs b/source/Halibut.Tests/Support/TestConnectionsObserver.cs index 9714c3177..29ca38ccb 100644 --- a/source/Halibut.Tests/Support/TestConnectionsObserver.cs +++ b/source/Halibut.Tests/Support/TestConnectionsObserver.cs @@ -10,12 +10,16 @@ public class TestConnectionsObserver : IConnectionsObserver { readonly ConcurrentBag connectionAcceptedAuthorized = new(); readonly ConcurrentBag connectionClosedAuthorized = new(); - + readonly ConcurrentBag connectionAcceptedForSubscriptions = new(); + readonly ConcurrentBag connectionClosedForSubscriptions = new(); + public long ConnectionAcceptedCount => connectionAcceptedAuthorized.Count; public long ConnectionClosedCount => connectionClosedAuthorized.Count; public IReadOnlyList ConnectionAcceptedAuthorized => connectionAcceptedAuthorized.ToList(); public IReadOnlyList ConnectionClosedAuthorized => connectionClosedAuthorized.ToList(); + public IReadOnlyList ConnectionAcceptedForSubscriptions => connectionAcceptedForSubscriptions.ToList(); + public IReadOnlyList ConnectionClosedForSubscriptions => connectionClosedForSubscriptions.ToList(); public void ConnectionAccepted(bool authorized) { @@ -26,5 +30,15 @@ public void ConnectionClosed(bool authorized) { connectionClosedAuthorized.Add(authorized); } + + public void ConnectionAcceptedFor(Uri subscriptionId) + { + connectionAcceptedForSubscriptions.Add(subscriptionId); + } + + public void ConnectionClosedFor(Uri subscriptionId) + { + connectionClosedForSubscriptions.Add(subscriptionId); + } } } \ No newline at end of file diff --git a/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs b/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs index 70b11fdd7..a2e5c7af3 100644 --- a/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs +++ b/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs @@ -27,7 +27,7 @@ public void SetUp() stream.SetRemoteIdentity(new RemoteIdentity(RemoteIdentityType.Server)); var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); var activeConnectionsLimiter = new ActiveTcpConnectionsLimiter(limits); - protocol = new MessageExchangeProtocol(stream, new HalibutTimeoutsAndLimitsForTestsBuilder().Build(), activeConnectionsLimiter, Substitute.For()); + protocol = new MessageExchangeProtocol(stream, new HalibutTimeoutsAndLimitsForTestsBuilder().Build(), activeConnectionsLimiter, NoOpConnectionsObserver.Instance, Substitute.For()); } // TODO - ASYNC ME UP! ExchangeAsClientAsync cancellation diff --git a/source/Halibut.Tests/Transport/SecureClientFixture.cs b/source/Halibut.Tests/Transport/SecureClientFixture.cs index 0df6e4398..aaae5537b 100644 --- a/source/Halibut.Tests/Transport/SecureClientFixture.cs +++ b/source/Halibut.Tests/Transport/SecureClientFixture.cs @@ -75,7 +75,7 @@ public async Task SecureClientClearsPoolWhenAllConnectionsCorrupt() var connection = Substitute.For(); var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits); - connection.Protocol.Returns(new MessageExchangeProtocol(stream, limits, activeConnectionLimiter, log)); + connection.Protocol.Returns(new MessageExchangeProtocol(stream, limits, activeConnectionLimiter, NoOpConnectionsObserver.Instance, log)); await connectionManager.ReleaseConnectionAsync(endpoint, connection, CancellationToken.None); } @@ -109,7 +109,7 @@ static MessageExchangeProtocol GetProtocol(Stream stream, ILog logger) { var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits); - return new MessageExchangeProtocol(new MessageExchangeStream(stream, new MessageSerializerBuilder(new LogFactory()).Build(), new NoOpControlMessageObserver(), limits, logger), limits, activeConnectionLimiter, logger); + return new MessageExchangeProtocol(new MessageExchangeStream(stream, new MessageSerializerBuilder(new LogFactory()).Build(), new NoOpControlMessageObserver(), limits, logger), limits, activeConnectionLimiter, NoOpConnectionsObserver.Instance, logger); } } } \ No newline at end of file diff --git a/source/Halibut/HalibutRuntime.cs b/source/Halibut/HalibutRuntime.cs index 1588239e0..abd11cefa 100644 --- a/source/Halibut/HalibutRuntime.cs +++ b/source/Halibut/HalibutRuntime.cs @@ -119,7 +119,7 @@ public int Listen(int port) ExchangeProtocolBuilder ExchangeProtocolBuilder() { - return (stream, log) => new MessageExchangeProtocol(new MessageExchangeStream(stream, messageSerializer, controlMessageObserver, TimeoutsAndLimits, log), TimeoutsAndLimits, activeTcpConnectionsLimiter, log); + return (stream, log) => new MessageExchangeProtocol(new MessageExchangeStream(stream, messageSerializer, controlMessageObserver, TimeoutsAndLimits, log), TimeoutsAndLimits, activeTcpConnectionsLimiter, connectionsObserver, log); } public int Listen(IPEndPoint endpoint) diff --git a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs index 1f253f763..ea0080e89 100644 --- a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs +++ b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs @@ -10,7 +10,6 @@ public interface IActiveTcpConnectionsLimiter { IDisposable LeaseActiveTcpConnection(Uri subscriptionId); - IDisposable CreateUnlimitedLease(); } public class ActiveTcpConnectionsLimiter : IActiveTcpConnectionsLimiter @@ -35,7 +34,7 @@ public IDisposable LeaseActiveTcpConnection(Uri subscriptionId) return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value); } - public IDisposable CreateUnlimitedLease() + IDisposable CreateUnlimitedLease() { return new UnlimitedAuthorizedTcpConnectionLease(); } diff --git a/source/Halibut/Transport/Observability/IConnectionsObserver.cs b/source/Halibut/Transport/Observability/IConnectionsObserver.cs index b1bf8c694..cddb59183 100644 --- a/source/Halibut/Transport/Observability/IConnectionsObserver.cs +++ b/source/Halibut/Transport/Observability/IConnectionsObserver.cs @@ -1,3 +1,5 @@ +using System; + namespace Halibut.Transport.Observability { public interface IConnectionsObserver @@ -6,7 +8,7 @@ public interface IConnectionsObserver /// The connection has been accepted and no bytes have been read from the wire. /// /// In this context server is anything that listens on a port. - /// + /// /// This is called when any of the following occurs: /// - When a "server" accepts a connection from a polling service (either websocket or regular) /// - When a "server" accepts a connection from a listening client (so in this case the server is the service) @@ -16,8 +18,23 @@ public interface IConnectionsObserver /// /// A previously accepted connection has been closed. /// - /// For every call to ConnectionClosed() their can be at most one call to this method. + /// For every call to ConnectionClosed() their can be at most one call to this method. /// public void ConnectionClosed(bool authorized); + + /// + /// Called once the connection is known to be for a + /// polling subscriber (i.e. after the subscription id has been read off the wire), and only + /// for connections that were not rejected for exceeding the active connection limit. + /// + /// For every call to this method there will be at most one matching call to ConnectionClosedFor() + /// with the same subscriptionId. + /// + public void ConnectionAcceptedFor(Uri subscriptionId); + + /// + /// A previously accepted polling subscriber connection has been closed. + /// + public void ConnectionClosedFor(Uri subscriptionId); } } \ No newline at end of file diff --git a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs index 802e3ed55..e8975f02a 100644 --- a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs +++ b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs @@ -1,3 +1,5 @@ +using System; + namespace Halibut.Transport.Observability { public class NoOpConnectionsObserver : IConnectionsObserver @@ -13,5 +15,13 @@ public void ConnectionAccepted(bool authorized) public void ConnectionClosed(bool authorized) { } + + public void ConnectionAcceptedFor(Uri subscriptionId) + { + } + + public void ConnectionClosedFor(Uri subscriptionId) + { + } } } \ No newline at end of file diff --git a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs index 777b7e268..9d2643f75 100644 --- a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs +++ b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs @@ -5,6 +5,7 @@ using Halibut.Diagnostics; using Halibut.Exceptions; using Halibut.ServiceModel; +using Halibut.Transport.Observability; namespace Halibut.Transport.Protocol { @@ -20,15 +21,17 @@ public class MessageExchangeProtocol readonly IMessageExchangeStream stream; readonly HalibutTimeoutsAndLimits halibutTimeoutsAndLimits; readonly IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter; + readonly IConnectionsObserver connectionsObserver; readonly ILog log; bool identified; volatile bool acceptClientRequests = true; - public MessageExchangeProtocol(IMessageExchangeStream stream, HalibutTimeoutsAndLimits halibutTimeoutsAndLimits, IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter, ILog log) + public MessageExchangeProtocol(IMessageExchangeStream stream, HalibutTimeoutsAndLimits halibutTimeoutsAndLimits, IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter, IConnectionsObserver connectionsObserver, ILog log) { this.stream = stream; this.halibutTimeoutsAndLimits = halibutTimeoutsAndLimits; this.activeTcpConnectionsLimiter = activeTcpConnectionsLimiter; + this.connectionsObserver = connectionsObserver; this.log = log; } @@ -106,32 +109,30 @@ public async Task ExchangeAsServerAsync(Func Date: Mon, 17 Aug 2026 16:53:58 +1000 Subject: [PATCH 2/4] Fix the order of calls. --- source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs index 9d2643f75..7fd64b007 100644 --- a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs +++ b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs @@ -127,8 +127,8 @@ public async Task ExchangeAsServerAsync(Func Date: Tue, 18 Aug 2026 17:04:29 +1000 Subject: [PATCH 3/4] Make bucketing easier on the server side without having to replicate the same state --- .../Support/TestConnectionsObserver.cs | 16 +++---- .../Transport/ActiveTcpConnectionsLimiter.cs | 46 ++++++++++--------- .../Observability/IConnectionsObserver.cs | 15 +++--- .../Observability/NoOpConnectionsObserver.cs | 4 +- .../Protocol/MessageExchangeProtocol.cs | 4 +- 5 files changed, 43 insertions(+), 42 deletions(-) diff --git a/source/Halibut.Tests/Support/TestConnectionsObserver.cs b/source/Halibut.Tests/Support/TestConnectionsObserver.cs index 29ca38ccb..782b801ca 100644 --- a/source/Halibut.Tests/Support/TestConnectionsObserver.cs +++ b/source/Halibut.Tests/Support/TestConnectionsObserver.cs @@ -10,16 +10,16 @@ public class TestConnectionsObserver : IConnectionsObserver { readonly ConcurrentBag connectionAcceptedAuthorized = new(); readonly ConcurrentBag connectionClosedAuthorized = new(); - readonly ConcurrentBag connectionAcceptedForSubscriptions = new(); - readonly ConcurrentBag connectionClosedForSubscriptions = new(); + readonly ConcurrentBag<(Uri SubscriptionId, int CurrentCount)> connectionAcceptedForSubscriptions = new(); + readonly ConcurrentBag<(Uri SubscriptionId, int CurrentCount)> connectionClosedForSubscriptions = new(); public long ConnectionAcceptedCount => connectionAcceptedAuthorized.Count; public long ConnectionClosedCount => connectionClosedAuthorized.Count; public IReadOnlyList ConnectionAcceptedAuthorized => connectionAcceptedAuthorized.ToList(); public IReadOnlyList ConnectionClosedAuthorized => connectionClosedAuthorized.ToList(); - public IReadOnlyList ConnectionAcceptedForSubscriptions => connectionAcceptedForSubscriptions.ToList(); - public IReadOnlyList ConnectionClosedForSubscriptions => connectionClosedForSubscriptions.ToList(); + public IReadOnlyList<(Uri SubscriptionId, int CurrentCount)> ConnectionAcceptedForSubscriptions => connectionAcceptedForSubscriptions.ToList(); + public IReadOnlyList<(Uri SubscriptionId, int CurrentCount)> ConnectionClosedForSubscriptions => connectionClosedForSubscriptions.ToList(); public void ConnectionAccepted(bool authorized) { @@ -31,14 +31,14 @@ public void ConnectionClosed(bool authorized) connectionClosedAuthorized.Add(authorized); } - public void ConnectionAcceptedFor(Uri subscriptionId) + public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount) { - connectionAcceptedForSubscriptions.Add(subscriptionId); + connectionAcceptedForSubscriptions.Add((subscriptionId, currentCount)); } - public void ConnectionClosedFor(Uri subscriptionId) + public void ConnectionClosedFor(Uri subscriptionId, int currentCount) { - connectionClosedForSubscriptions.Add(subscriptionId); + connectionClosedForSubscriptions.Add((subscriptionId, currentCount)); } } } \ No newline at end of file diff --git a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs index ea0080e89..856617da6 100644 --- a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs +++ b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs @@ -6,10 +6,19 @@ namespace Halibut.Transport { - public interface IActiveTcpConnectionsLimiter + public interface IActiveTcpConnectionLease : IDisposable { - IDisposable LeaseActiveTcpConnection(Uri subscriptionId); + /// + /// The number of active TCP connections for the leased subscription, as at the point the lease was + /// created. After Dispose() is called, this reflects the count immediately after this connection + /// was released. + /// + int CurrentCount { get; } + } + public interface IActiveTcpConnectionsLimiter + { + IActiveTcpConnectionLease LeaseActiveTcpConnection(Uri subscriptionId); } public class ActiveTcpConnectionsLimiter : IActiveTcpConnectionsLimiter @@ -23,34 +32,30 @@ public ActiveTcpConnectionsLimiter(HalibutTimeoutsAndLimits timeoutsAndLimits) this.timeoutsAndLimits = timeoutsAndLimits; } - public IDisposable LeaseActiveTcpConnection(Uri subscriptionId) + public IActiveTcpConnectionLease LeaseActiveTcpConnection(Uri subscriptionId) { - //if there is no limit, then we return a NoOp lease (which doesn't limit anything) + //if there is no limit, then we still count the connection (callers rely on the resulting count), + //we just never reject it if (!timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.HasValue) { - return CreateUnlimitedLease(); + return CreateUnlimitedLease(subscriptionId); } return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value); } - IDisposable CreateUnlimitedLease() + IActiveTcpConnectionLease CreateUnlimitedLease(Uri subscriptionId) { - return new UnlimitedAuthorizedTcpConnectionLease(); + return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, int.MaxValue); } - class UnlimitedAuthorizedTcpConnectionLease : IDisposable - { - public void Dispose() - { - } - } - - class LimitingAuthorizedTcpConnectionLease : IDisposable + class LimitingAuthorizedTcpConnectionLease : IActiveTcpConnectionLease { readonly Uri subscriptionId; readonly Dictionary> activeConnectionCountPerSubscriptionId; + public int CurrentCount { get; private set; } + public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary> activeConnectionCountPerSubscriptionId, int maximumAcceptedTcpConnectionsPerThumbprint) { this.subscriptionId = subscriptionId; @@ -64,17 +69,15 @@ public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary maximumAcceptedTcpConnectionsPerThumbprint) + if (count.Value + 1 > maximumAcceptedTcpConnectionsPerThumbprint) { - //decrement as this connection has been rejected - count.Value--; - //throw an exception, bailing on the connection throw new ActiveTcpConnectionsExceededException(this.subscriptionId, $"Exceeded the maximum number ({maximumAcceptedTcpConnectionsPerThumbprint}) of active TCP connections for subscription {subscriptionId}"); } + + count.Value++; + CurrentCount = count.Value; } } @@ -86,6 +89,7 @@ public void Dispose() { //decrement the count of authorized connections count.Value--; + CurrentCount = count.Value; // Remove the key from the dictionary if the value is 0 if (count.Value == 0) diff --git a/source/Halibut/Transport/Observability/IConnectionsObserver.cs b/source/Halibut/Transport/Observability/IConnectionsObserver.cs index cddb59183..48645e51e 100644 --- a/source/Halibut/Transport/Observability/IConnectionsObserver.cs +++ b/source/Halibut/Transport/Observability/IConnectionsObserver.cs @@ -23,18 +23,15 @@ public interface IConnectionsObserver public void ConnectionClosed(bool authorized); /// - /// Called once the connection is known to be for a - /// polling subscriber (i.e. after the subscription id has been read off the wire), and only - /// for connections that were not rejected for exceeding the active connection limit. - /// - /// For every call to this method there will be at most one matching call to ConnectionClosedFor() - /// with the same subscriptionId. + /// The number of active TCP connections for this subscriptionId immediately after this connection + /// was accepted (i.e. including this one). /// - public void ConnectionAcceptedFor(Uri subscriptionId); + public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount); /// - /// A previously accepted polling subscriber connection has been closed. + /// The number of active TCP connections for this subscriptionId immediately after this connection + /// was closed (i.e. excluding this one). /// - public void ConnectionClosedFor(Uri subscriptionId); + public void ConnectionClosedFor(Uri subscriptionId, int currentCount); } } \ No newline at end of file diff --git a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs index e8975f02a..0b6e76ec0 100644 --- a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs +++ b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs @@ -16,11 +16,11 @@ public void ConnectionClosed(bool authorized) { } - public void ConnectionAcceptedFor(Uri subscriptionId) + public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount) { } - public void ConnectionClosedFor(Uri subscriptionId) + public void ConnectionClosedFor(Uri subscriptionId, int currentCount) { } } diff --git a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs index 7fd64b007..8b24be894 100644 --- a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs +++ b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs @@ -119,7 +119,7 @@ public async Task ExchangeAsServerAsync(Func Date: Wed, 19 Aug 2026 09:29:20 +1000 Subject: [PATCH 4/4] We should be able to bucket polling connections correctly now. --- .../Support/TestConnectionsObserver.cs | 14 ++---- .../ActiveTcpConnectionsLimiterFixture.cs | 11 ++-- .../Transport/Protocol/ProtocolFixture.cs | 4 +- .../Transport/SecureClientFixture.cs | 8 +-- source/Halibut/HalibutRuntime.cs | 4 +- .../Transport/ActiveTcpConnectionsLimiter.cs | 50 +++++++++---------- .../Observability/IConnectionsObserver.cs | 18 +++---- .../Observability/NoOpConnectionsObserver.cs | 6 +-- .../Protocol/MessageExchangeProtocol.cs | 14 +----- 9 files changed, 53 insertions(+), 76 deletions(-) diff --git a/source/Halibut.Tests/Support/TestConnectionsObserver.cs b/source/Halibut.Tests/Support/TestConnectionsObserver.cs index 782b801ca..d5e6ec2d9 100644 --- a/source/Halibut.Tests/Support/TestConnectionsObserver.cs +++ b/source/Halibut.Tests/Support/TestConnectionsObserver.cs @@ -10,16 +10,13 @@ public class TestConnectionsObserver : IConnectionsObserver { readonly ConcurrentBag connectionAcceptedAuthorized = new(); readonly ConcurrentBag connectionClosedAuthorized = new(); - readonly ConcurrentBag<(Uri SubscriptionId, int CurrentCount)> connectionAcceptedForSubscriptions = new(); - readonly ConcurrentBag<(Uri SubscriptionId, int CurrentCount)> connectionClosedForSubscriptions = new(); + readonly ConcurrentBag<(Uri SubscriptionId, int PreviousCount, int CurrentCount)> connectionsCountChangedForSubscription = new(); public long ConnectionAcceptedCount => connectionAcceptedAuthorized.Count; public long ConnectionClosedCount => connectionClosedAuthorized.Count; public IReadOnlyList ConnectionAcceptedAuthorized => connectionAcceptedAuthorized.ToList(); public IReadOnlyList ConnectionClosedAuthorized => connectionClosedAuthorized.ToList(); - public IReadOnlyList<(Uri SubscriptionId, int CurrentCount)> ConnectionAcceptedForSubscriptions => connectionAcceptedForSubscriptions.ToList(); - public IReadOnlyList<(Uri SubscriptionId, int CurrentCount)> ConnectionClosedForSubscriptions => connectionClosedForSubscriptions.ToList(); public void ConnectionAccepted(bool authorized) { @@ -31,14 +28,9 @@ public void ConnectionClosed(bool authorized) connectionClosedAuthorized.Add(authorized); } - public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount) + public void ConnectionsCountChangedFor(Uri subscriptionId, int previousCount, int currentCount) { - connectionAcceptedForSubscriptions.Add((subscriptionId, currentCount)); - } - - public void ConnectionClosedFor(Uri subscriptionId, int currentCount) - { - connectionClosedForSubscriptions.Add((subscriptionId, currentCount)); + connectionsCountChangedForSubscription.Add((subscriptionId, previousCount, currentCount)); } } } \ No newline at end of file diff --git a/source/Halibut.Tests/Transport/ActiveTcpConnectionsLimiterFixture.cs b/source/Halibut.Tests/Transport/ActiveTcpConnectionsLimiterFixture.cs index 18230285b..a38032cfa 100644 --- a/source/Halibut.Tests/Transport/ActiveTcpConnectionsLimiterFixture.cs +++ b/source/Halibut.Tests/Transport/ActiveTcpConnectionsLimiterFixture.cs @@ -6,6 +6,7 @@ using Halibut.Diagnostics; using Halibut.Exceptions; using Halibut.Transport; +using Halibut.Transport.Observability; using NUnit.Framework; namespace Halibut.Tests.Transport @@ -22,7 +23,7 @@ public void LimitsConcurrentConnectionsForSingleSubscription() var limiter = new ActiveTcpConnectionsLimiter(new HalibutTimeoutsAndLimits { MaximumActiveTcpConnectionsPerPollingSubscription = limit - }); + }, NoOpConnectionsObserver.Instance); // Act //we create a new URI each time to make sure we aren't doing object reference checks @@ -46,7 +47,7 @@ public void CompletedLeasesAreRemovedFromTheCount() var limiter = new ActiveTcpConnectionsLimiter(new HalibutTimeoutsAndLimits { MaximumActiveTcpConnectionsPerPollingSubscription = limit - }); + }, NoOpConnectionsObserver.Instance); // Act limiter.LeaseActiveTcpConnection(subscription); @@ -76,7 +77,7 @@ public void DoesNotLimitConcurrentConnectionsForDifferentSubscriptions() var limiter = new ActiveTcpConnectionsLimiter(new HalibutTimeoutsAndLimits { MaximumActiveTcpConnectionsPerPollingSubscription = limit - }); + }, NoOpConnectionsObserver.Instance); // Act limiter.LeaseActiveTcpConnection(subscription1); @@ -99,7 +100,7 @@ public async Task ShouldHandleMultiThreading() var limiter = new ActiveTcpConnectionsLimiter(new HalibutTimeoutsAndLimits { MaximumActiveTcpConnectionsPerPollingSubscription = limit - }); + }, NoOpConnectionsObserver.Instance); // Capture how many claims fail with the exception var failures = 0; @@ -140,7 +141,7 @@ public async Task ShouldHandleMultiThreadingWithFakeWorkDuringLease() var limiter = new ActiveTcpConnectionsLimiter(new HalibutTimeoutsAndLimits { MaximumActiveTcpConnectionsPerPollingSubscription = limit - }); + }, NoOpConnectionsObserver.Instance); // Capture how many claims fail with the exception var failures = 0; diff --git a/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs b/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs index a2e5c7af3..be4e395ec 100644 --- a/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs +++ b/source/Halibut.Tests/Transport/Protocol/ProtocolFixture.cs @@ -26,8 +26,8 @@ public void SetUp() stream = new DumpStream(); stream.SetRemoteIdentity(new RemoteIdentity(RemoteIdentityType.Server)); var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); - var activeConnectionsLimiter = new ActiveTcpConnectionsLimiter(limits); - protocol = new MessageExchangeProtocol(stream, new HalibutTimeoutsAndLimitsForTestsBuilder().Build(), activeConnectionsLimiter, NoOpConnectionsObserver.Instance, Substitute.For()); + var activeConnectionsLimiter = new ActiveTcpConnectionsLimiter(limits, NoOpConnectionsObserver.Instance); + protocol = new MessageExchangeProtocol(stream, new HalibutTimeoutsAndLimitsForTestsBuilder().Build(), activeConnectionsLimiter, Substitute.For()); } // TODO - ASYNC ME UP! ExchangeAsClientAsync cancellation diff --git a/source/Halibut.Tests/Transport/SecureClientFixture.cs b/source/Halibut.Tests/Transport/SecureClientFixture.cs index aaae5537b..d8652e4b7 100644 --- a/source/Halibut.Tests/Transport/SecureClientFixture.cs +++ b/source/Halibut.Tests/Transport/SecureClientFixture.cs @@ -74,8 +74,8 @@ public async Task SecureClientClearsPoolWhenAllConnectionsCorrupt() { var connection = Substitute.For(); var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); - var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits); - connection.Protocol.Returns(new MessageExchangeProtocol(stream, limits, activeConnectionLimiter, NoOpConnectionsObserver.Instance, log)); + var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits, NoOpConnectionsObserver.Instance); + connection.Protocol.Returns(new MessageExchangeProtocol(stream, limits, activeConnectionLimiter, log)); await connectionManager.ReleaseConnectionAsync(endpoint, connection, CancellationToken.None); } @@ -108,8 +108,8 @@ public async Task SecureClientClearsPoolWhenAllConnectionsCorrupt() static MessageExchangeProtocol GetProtocol(Stream stream, ILog logger) { var limits = new HalibutTimeoutsAndLimitsForTestsBuilder().Build(); - var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits); - return new MessageExchangeProtocol(new MessageExchangeStream(stream, new MessageSerializerBuilder(new LogFactory()).Build(), new NoOpControlMessageObserver(), limits, logger), limits, activeConnectionLimiter, NoOpConnectionsObserver.Instance, logger); + var activeConnectionLimiter = new ActiveTcpConnectionsLimiter(limits, NoOpConnectionsObserver.Instance); + return new MessageExchangeProtocol(new MessageExchangeStream(stream, new MessageSerializerBuilder(new LogFactory()).Build(), new NoOpControlMessageObserver(), limits, logger), limits, activeConnectionLimiter, logger); } } } \ No newline at end of file diff --git a/source/Halibut/HalibutRuntime.cs b/source/Halibut/HalibutRuntime.cs index abd11cefa..58c062888 100644 --- a/source/Halibut/HalibutRuntime.cs +++ b/source/Halibut/HalibutRuntime.cs @@ -82,7 +82,7 @@ ISecureConnectionObserver secureConnectionObserver connectionManager = new ConnectionManagerAsync(); tcpConnectionFactory = new TcpConnectionFactory(serverCertificate, TimeoutsAndLimits, streamFactory, secureConnectionObserver); - activeTcpConnectionsLimiter = new ActiveTcpConnectionsLimiter(TimeoutsAndLimits); + activeTcpConnectionsLimiter = new ActiveTcpConnectionsLimiter(TimeoutsAndLimits, connectionsObserver); } public ILogFactory Logs => logs; @@ -119,7 +119,7 @@ public int Listen(int port) ExchangeProtocolBuilder ExchangeProtocolBuilder() { - return (stream, log) => new MessageExchangeProtocol(new MessageExchangeStream(stream, messageSerializer, controlMessageObserver, TimeoutsAndLimits, log), TimeoutsAndLimits, activeTcpConnectionsLimiter, connectionsObserver, log); + return (stream, log) => new MessageExchangeProtocol(new MessageExchangeStream(stream, messageSerializer, controlMessageObserver, TimeoutsAndLimits, log), TimeoutsAndLimits, activeTcpConnectionsLimiter, log); } public int Listen(IPEndPoint endpoint) diff --git a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs index 856617da6..87a9e862f 100644 --- a/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs +++ b/source/Halibut/Transport/ActiveTcpConnectionsLimiter.cs @@ -1,65 +1,58 @@ -using System; +using System; using System.Collections.Generic; using System.Runtime.CompilerServices; using Halibut.Diagnostics; using Halibut.Exceptions; +using Halibut.Transport.Observability; namespace Halibut.Transport { - public interface IActiveTcpConnectionLease : IDisposable - { - /// - /// The number of active TCP connections for the leased subscription, as at the point the lease was - /// created. After Dispose() is called, this reflects the count immediately after this connection - /// was released. - /// - int CurrentCount { get; } - } - public interface IActiveTcpConnectionsLimiter { - IActiveTcpConnectionLease LeaseActiveTcpConnection(Uri subscriptionId); + IDisposable LeaseActiveTcpConnection(Uri subscriptionId); } public class ActiveTcpConnectionsLimiter : IActiveTcpConnectionsLimiter { readonly HalibutTimeoutsAndLimits timeoutsAndLimits; + readonly IConnectionsObserver connectionsObserver; Dictionary> activeConnectionCountPerSubscriptionId = new(); - public ActiveTcpConnectionsLimiter(HalibutTimeoutsAndLimits timeoutsAndLimits) + public ActiveTcpConnectionsLimiter(HalibutTimeoutsAndLimits timeoutsAndLimits, IConnectionsObserver connectionsObserver) { this.timeoutsAndLimits = timeoutsAndLimits; + this.connectionsObserver = connectionsObserver; } - public IActiveTcpConnectionLease LeaseActiveTcpConnection(Uri subscriptionId) + public IDisposable LeaseActiveTcpConnection(Uri subscriptionId) { - //if there is no limit, then we still count the connection (callers rely on the resulting count), - //we just never reject it + //if there is no limit, then we still count the connection (the observer is told about every + //connection either way), we just never reject it if (!timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.HasValue) { return CreateUnlimitedLease(subscriptionId); } - return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value); + return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value, connectionsObserver); } - IActiveTcpConnectionLease CreateUnlimitedLease(Uri subscriptionId) + IDisposable CreateUnlimitedLease(Uri subscriptionId) { - return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, int.MaxValue); + return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, int.MaxValue, connectionsObserver); } - class LimitingAuthorizedTcpConnectionLease : IActiveTcpConnectionLease + class LimitingAuthorizedTcpConnectionLease : IDisposable { readonly Uri subscriptionId; readonly Dictionary> activeConnectionCountPerSubscriptionId; + readonly IConnectionsObserver connectionsObserver; - public int CurrentCount { get; private set; } - - public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary> activeConnectionCountPerSubscriptionId, int maximumAcceptedTcpConnectionsPerThumbprint) + public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary> activeConnectionCountPerSubscriptionId, int maximumAcceptedTcpConnectionsPerThumbprint, IConnectionsObserver connectionsObserver) { this.subscriptionId = subscriptionId; this.activeConnectionCountPerSubscriptionId = activeConnectionCountPerSubscriptionId; + this.connectionsObserver = connectionsObserver; lock (this.activeConnectionCountPerSubscriptionId) { @@ -69,6 +62,8 @@ public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary maximumAcceptedTcpConnectionsPerThumbprint) { @@ -77,7 +72,8 @@ public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary - /// The number of active TCP connections for this subscriptionId immediately after this connection - /// was accepted (i.e. including this one). + /// A polling subscriber's connections' count has changed /// - public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount); - - /// - /// The number of active TCP connections for this subscriptionId immediately after this connection - /// was closed (i.e. excluding this one). - /// - public void ConnectionClosedFor(Uri subscriptionId, int currentCount); + /// The polling subscriber's subscription id. + /// + /// The number of active TCP connections for this subscriptionId immediately before the change + /// + /// + /// The number of active TCP connections for this subscriptionId immediately after the change + /// + public void ConnectionsCountChangedFor(Uri subscriptionId, int previousCount, int currentCount); } } \ No newline at end of file diff --git a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs index 0b6e76ec0..2cd6b52fb 100644 --- a/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs +++ b/source/Halibut/Transport/Observability/NoOpConnectionsObserver.cs @@ -16,11 +16,7 @@ public void ConnectionClosed(bool authorized) { } - public void ConnectionAcceptedFor(Uri subscriptionId, int currentCount) - { - } - - public void ConnectionClosedFor(Uri subscriptionId, int currentCount) + public void ConnectionsCountChangedFor(Uri subscriptionId, int previousCount, int currentCount) { } } diff --git a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs index 8b24be894..4f1c7bb9f 100644 --- a/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs +++ b/source/Halibut/Transport/Protocol/MessageExchangeProtocol.cs @@ -5,7 +5,6 @@ using Halibut.Diagnostics; using Halibut.Exceptions; using Halibut.ServiceModel; -using Halibut.Transport.Observability; namespace Halibut.Transport.Protocol { @@ -21,17 +20,15 @@ public class MessageExchangeProtocol readonly IMessageExchangeStream stream; readonly HalibutTimeoutsAndLimits halibutTimeoutsAndLimits; readonly IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter; - readonly IConnectionsObserver connectionsObserver; readonly ILog log; bool identified; volatile bool acceptClientRequests = true; - public MessageExchangeProtocol(IMessageExchangeStream stream, HalibutTimeoutsAndLimits halibutTimeoutsAndLimits, IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter, IConnectionsObserver connectionsObserver, ILog log) + public MessageExchangeProtocol(IMessageExchangeStream stream, HalibutTimeoutsAndLimits halibutTimeoutsAndLimits, IActiveTcpConnectionsLimiter activeTcpConnectionsLimiter, ILog log) { this.stream = stream; this.halibutTimeoutsAndLimits = halibutTimeoutsAndLimits; this.activeTcpConnectionsLimiter = activeTcpConnectionsLimiter; - this.connectionsObserver = connectionsObserver; this.log = log; } @@ -116,20 +113,13 @@ public async Task ExchangeAsServerAsync(Func