-
Notifications
You must be signed in to change notification settings - Fork 47
Add ability to observe the number of connections halibut opens for ea… #717
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
a0fcda4
d9b34ea
2aa8583
4d436f2
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,61 +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 IActiveTcpConnectionsLimiter | ||
| { | ||
| IDisposable LeaseActiveTcpConnection(Uri subscriptionId); | ||
|
|
||
| IDisposable CreateUnlimitedLease(); | ||
| } | ||
|
|
||
| public class ActiveTcpConnectionsLimiter : IActiveTcpConnectionsLimiter | ||
| { | ||
| readonly HalibutTimeoutsAndLimits timeoutsAndLimits; | ||
| readonly IConnectionsObserver connectionsObserver; | ||
|
|
||
| Dictionary<Uri, StrongBox<int>> activeConnectionCountPerSubscriptionId = new(); | ||
|
|
||
| public ActiveTcpConnectionsLimiter(HalibutTimeoutsAndLimits timeoutsAndLimits) | ||
| public ActiveTcpConnectionsLimiter(HalibutTimeoutsAndLimits timeoutsAndLimits, IConnectionsObserver connectionsObserver) | ||
| { | ||
| this.timeoutsAndLimits = timeoutsAndLimits; | ||
| this.connectionsObserver = connectionsObserver; | ||
| } | ||
|
|
||
| public IDisposable 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 (the observer is told about every | ||
| //connection either way), we just never reject it | ||
| if (!timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.HasValue) | ||
| { | ||
| return CreateUnlimitedLease(); | ||
| return CreateUnlimitedLease(subscriptionId); | ||
| } | ||
|
|
||
| return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value); | ||
| } | ||
|
|
||
| public IDisposable CreateUnlimitedLease() | ||
| { | ||
| return new UnlimitedAuthorizedTcpConnectionLease(); | ||
| return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, timeoutsAndLimits.MaximumActiveTcpConnectionsPerPollingSubscription.Value, connectionsObserver); | ||
| } | ||
|
|
||
| class UnlimitedAuthorizedTcpConnectionLease : IDisposable | ||
| IDisposable CreateUnlimitedLease(Uri subscriptionId) | ||
| { | ||
| public void Dispose() | ||
| { | ||
| } | ||
| return new LimitingAuthorizedTcpConnectionLease(subscriptionId, activeConnectionCountPerSubscriptionId, int.MaxValue, connectionsObserver); | ||
| } | ||
|
|
||
| class LimitingAuthorizedTcpConnectionLease : IDisposable | ||
| { | ||
| readonly Uri subscriptionId; | ||
| readonly Dictionary<Uri, StrongBox<int>> activeConnectionCountPerSubscriptionId; | ||
| readonly IConnectionsObserver connectionsObserver; | ||
|
|
||
| public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary<Uri, StrongBox<int>> activeConnectionCountPerSubscriptionId, int maximumAcceptedTcpConnectionsPerThumbprint) | ||
| public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary<Uri, StrongBox<int>> activeConnectionCountPerSubscriptionId, int maximumAcceptedTcpConnectionsPerThumbprint, IConnectionsObserver connectionsObserver) | ||
| { | ||
| this.subscriptionId = subscriptionId; | ||
| this.activeConnectionCountPerSubscriptionId = activeConnectionCountPerSubscriptionId; | ||
| this.connectionsObserver = connectionsObserver; | ||
|
|
||
| lock (this.activeConnectionCountPerSubscriptionId) | ||
| { | ||
|
|
@@ -65,17 +62,18 @@ public LimitingAuthorizedTcpConnectionLease(Uri subscriptionId, Dictionary<Uri, | |
| this.activeConnectionCountPerSubscriptionId.Add(subscriptionId, count); | ||
| } | ||
|
|
||
| count.Value++; | ||
| var previousCount = count.Value; | ||
|
|
||
| //validate the new count. If this throws an exception, it'll kill the connection | ||
| if (count.Value > 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++; | ||
|
|
||
| connectionsObserver.ConnectionsCountChangedFor(subscriptionId, previousCount, count.Value); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I guess this means we need to be super careful about what |
||
| } | ||
| } | ||
|
|
||
|
|
@@ -86,16 +84,19 @@ public void Dispose() | |
| if (activeConnectionCountPerSubscriptionId.TryGetValue(subscriptionId, out var count)) | ||
| { | ||
| //decrement the count of authorized connections | ||
| var previousCount = count.Value; | ||
| count.Value--; | ||
|
|
||
| // Remove the key from the dictionary if the value is 0 | ||
| if (count.Value == 0) | ||
| { | ||
| activeConnectionCountPerSubscriptionId.Remove(subscriptionId); | ||
| } | ||
|
|
||
| connectionsObserver.ConnectionsCountChangedFor(subscriptionId, previousCount, count.Value); | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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,20 @@ public interface IConnectionsObserver | |
| /// <summary> | ||
| /// 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. | ||
| /// </summary> | ||
| public void ConnectionClosed(bool authorized); | ||
|
|
||
| /// <summary> | ||
| /// A polling subscriber's connections' count has changed | ||
| /// </summary> | ||
| /// <param name="subscriptionId">The polling subscriber's subscription id.</param> | ||
| /// <param name="previousCount"> | ||
| /// The number of active TCP connections for this subscriptionId immediately before the change | ||
| /// </param> | ||
| /// <param name="currentCount"> | ||
| /// The number of active TCP connections for this subscriptionId immediately after the change | ||
| /// </param> | ||
| public void ConnectionsCountChangedFor(Uri subscriptionId, int previousCount, int currentCount); | ||
|
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I considrerd a pair of
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Make it simpler for the caller is probably the best approach. |
||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -106,32 +106,23 @@ public async Task ExchangeAsServerAsync(Func<RequestMessage, Task<ResponseMessag | |
| { | ||
| var identity = await GetRemoteIdentityAsync(cancellationToken); | ||
|
|
||
| //We might need to limit the connection, so by default, we create an unlimited connection lease | ||
| var limitedConnectionLease = activeTcpConnectionsLimiter.CreateUnlimitedLease(); | ||
|
|
||
| //if the remote identity is a subscriber, we might need to limit their active TCP connections | ||
| if (identity.IdentityType == RemoteIdentityType.Subscriber) | ||
| switch (identity.IdentityType) | ||
| { | ||
| limitedConnectionLease = activeTcpConnectionsLimiter.LeaseActiveTcpConnection(identity.SubscriptionId); | ||
| } | ||
|
|
||
| using (limitedConnectionLease) | ||
| { | ||
| await IdentifyAsServerAsync(identity, cancellationToken); | ||
|
|
||
| switch (identity.IdentityType) | ||
| { | ||
| case RemoteIdentityType.Client: | ||
| await ProcessClientRequestsAsync(incomingRequestProcessor, cancellationToken); | ||
| break; | ||
| case RemoteIdentityType.Subscriber: | ||
| case RemoteIdentityType.Client: | ||
| await IdentifyAsServerAsync(identity, cancellationToken); | ||
| await ProcessClientRequestsAsync(incomingRequestProcessor, cancellationToken); | ||
| break; | ||
| case RemoteIdentityType.Subscriber: | ||
|
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I've moved the connection limit enforcement closer to the code that it applies to.
|
||
| using (activeTcpConnectionsLimiter.LeaseActiveTcpConnection(identity.SubscriptionId)) | ||
| { | ||
| await IdentifyAsServerAsync(identity, cancellationToken); | ||
| var pendingRequestQueue = pendingRequests(identity); | ||
| await ProcessSubscriberAsync(pendingRequestQueue, cancellationToken); | ||
| break; | ||
| default: | ||
| log.Write(EventType.ErrorInIdentify, $"Remote with identify {identity.SubscriptionId} identified itself with an unknown identity type {identity.IdentityType}"); | ||
| throw new ProtocolException("Unexpected remote identity: " + identity.IdentityType); | ||
| } | ||
| } | ||
| default: | ||
| log.Write(EventType.ErrorInIdentify, $"Remote with identify {identity.SubscriptionId} identified itself with an unknown identity type {identity.IdentityType}"); | ||
| throw new ProtocolException("Unexpected remote identity: " + identity.IdentityType); | ||
| } | ||
| } | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
This looks like the information you want.
Maybe have a method that will give you back a copy of this dictionary upon request.