Skip to content

Commit f5b7b32

Browse files
authored
Allow accessing send queue count from WebSocketConnection (#1156)
1 parent 2d2ab4c commit f5b7b32

File tree

7 files changed

+58
-4
lines changed

7 files changed

+58
-4
lines changed

src/Transports.AspNetCore/WebSockets/AsyncMessagePump.cs

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,21 @@ internal class AsyncMessagePump<T>
2222
private readonly Func<T, Task> _callback;
2323
private readonly Queue<ValueTask<T>> _queue = new();
2424

25+
/// <summary>
26+
/// Returns the number of messages in the queue.
27+
/// This count includes any message currently being processed.
28+
/// </summary>
29+
public int Count
30+
{
31+
get
32+
{
33+
lock (_queue)
34+
{
35+
return _queue.Count;
36+
}
37+
}
38+
}
39+
2540
/// <summary>
2641
/// Initializes a new instance with the specified asynchronous callback delegate.
2742
/// </summary>

src/Transports.AspNetCore/WebSockets/WebSocketConnection.cs

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,13 @@ public class WebSocketConnection : IWebSocketConnection
4343
/// <inheritdoc/>
4444
public HttpContext HttpContext { get; }
4545

46+
/// <summary>
47+
/// Returns the number of packets waiting in the send queue, including
48+
/// messages, keep-alive packets, and the close message.
49+
/// This count includes any packet currently being processed.
50+
/// </summary>
51+
protected int SendQueueCount => _pump.Count;
52+
4653
/// <summary>
4754
/// Initializes an instance with the specified parameters.
4855
/// </summary>
@@ -218,7 +225,7 @@ public Task CloseAsync(int eventId, string? description)
218225
/// <remarks>
219226
/// The message is posted to a queue and execution returns immediately.
220227
/// </remarks>
221-
public Task SendMessageAsync(OperationMessage message)
228+
public virtual Task SendMessageAsync(OperationMessage message)
222229
{
223230
_pump.Post(new Message { OperationMessage = message });
224231
return Task.CompletedTask;

tests/ApiApprovalTests/net50+net60+net80/GraphQL.Server.Transports.AspNetCore.approved.txt

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -340,6 +340,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
340340
public Microsoft.AspNetCore.Http.HttpContext HttpContext { get; }
341341
public System.DateTime LastMessageSentAt { get; }
342342
public System.Threading.CancellationToken RequestAborted { get; }
343+
protected int SendQueueCount { get; }
343344
public System.Threading.Tasks.Task CloseAsync() { }
344345
public System.Threading.Tasks.Task CloseAsync(int eventId, string? description) { }
345346
public virtual void Dispose() { }
@@ -348,7 +349,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
348349
protected virtual System.Threading.Tasks.Task OnDispatchMessageAsync(GraphQL.Server.Transports.AspNetCore.WebSockets.IOperationMessageProcessor operationMessageProcessor, GraphQL.Transport.OperationMessage message) { }
349350
protected virtual System.Threading.Tasks.Task OnNonGracefulShutdownAsync(bool receivedCloseMessage, bool sentCloseMessage) { }
350351
protected virtual System.Threading.Tasks.Task OnSendMessageAsync(GraphQL.Transport.OperationMessage message) { }
351-
public System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
352+
public virtual System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
352353
}
353354
}
354355
namespace GraphQL.Server.Transports.AspNetCore.WebSockets.GraphQLWs

tests/ApiApprovalTests/netcoreapp21+netstandard20/GraphQL.Server.Transports.AspNetCore.approved.txt

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -358,6 +358,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
358358
public Microsoft.AspNetCore.Http.HttpContext HttpContext { get; }
359359
public System.DateTime LastMessageSentAt { get; }
360360
public System.Threading.CancellationToken RequestAborted { get; }
361+
protected int SendQueueCount { get; }
361362
public System.Threading.Tasks.Task CloseAsync() { }
362363
public System.Threading.Tasks.Task CloseAsync(int eventId, string? description) { }
363364
public virtual void Dispose() { }
@@ -366,7 +367,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
366367
protected virtual System.Threading.Tasks.Task OnDispatchMessageAsync(GraphQL.Server.Transports.AspNetCore.WebSockets.IOperationMessageProcessor operationMessageProcessor, GraphQL.Transport.OperationMessage message) { }
367368
protected virtual System.Threading.Tasks.Task OnNonGracefulShutdownAsync(bool receivedCloseMessage, bool sentCloseMessage) { }
368369
protected virtual System.Threading.Tasks.Task OnSendMessageAsync(GraphQL.Transport.OperationMessage message) { }
369-
public System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
370+
public virtual System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
370371
}
371372
}
372373
namespace GraphQL.Server.Transports.AspNetCore.WebSockets.GraphQLWs

tests/ApiApprovalTests/netcoreapp31/GraphQL.Server.Transports.AspNetCore.approved.txt

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -340,6 +340,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
340340
public Microsoft.AspNetCore.Http.HttpContext HttpContext { get; }
341341
public System.DateTime LastMessageSentAt { get; }
342342
public System.Threading.CancellationToken RequestAborted { get; }
343+
protected int SendQueueCount { get; }
343344
public System.Threading.Tasks.Task CloseAsync() { }
344345
public System.Threading.Tasks.Task CloseAsync(int eventId, string? description) { }
345346
public virtual void Dispose() { }
@@ -348,7 +349,7 @@ namespace GraphQL.Server.Transports.AspNetCore.WebSockets
348349
protected virtual System.Threading.Tasks.Task OnDispatchMessageAsync(GraphQL.Server.Transports.AspNetCore.WebSockets.IOperationMessageProcessor operationMessageProcessor, GraphQL.Transport.OperationMessage message) { }
349350
protected virtual System.Threading.Tasks.Task OnNonGracefulShutdownAsync(bool receivedCloseMessage, bool sentCloseMessage) { }
350351
protected virtual System.Threading.Tasks.Task OnSendMessageAsync(GraphQL.Transport.OperationMessage message) { }
351-
public System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
352+
public virtual System.Threading.Tasks.Task SendMessageAsync(GraphQL.Transport.OperationMessage message) { }
352353
}
353354
}
354355
namespace GraphQL.Server.Transports.AspNetCore.WebSockets.GraphQLWs

tests/Transports.AspNetCore.Tests/WebSockets/TestWebSocketConnection.cs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,4 +25,6 @@ public Task Do_OnCloseOutputAsync(WebSocketCloseStatus closeStatus, string? clos
2525

2626
public TimeSpan Get_DefaultDisconnectionTimeout
2727
=> DefaultDisconnectionTimeout;
28+
29+
public int Get_SendQueueCount => base.SendQueueCount;
2830
}

tests/Transports.AspNetCore.Tests/WebSockets/WebSocketConnectionTests.cs

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -597,18 +597,44 @@ public async Task CloseConnectionAsync_Specific()
597597
public async Task SendMessageAsync()
598598
{
599599
var message = new OperationMessage();
600+
_mockConnection.Setup(x => x.SendMessageAsync(It.IsAny<OperationMessage>())).CallBase().Verifiable();
600601
_mockConnection.Protected().Setup<Task>("OnSendMessageAsync", message)
601602
.Returns(Task.CompletedTask).Verifiable();
602603
await _connection.SendMessageAsync(message);
603604
_mockConnection.Verify();
604605
}
605606

607+
[Fact]
608+
public async Task MessageCountAsync()
609+
{
610+
var tc = new TaskCompletionSource<bool>();
611+
var message = new OperationMessage();
612+
_mockConnection.Setup(x => x.SendMessageAsync(It.IsAny<OperationMessage>())).CallBase().Verifiable();
613+
_mockConnection.Protected().Setup<Task>("OnSendMessageAsync", message)
614+
.Returns(tc.Task).Verifiable();
615+
await _connection.SendMessageAsync(message);
616+
_connection.Get_SendQueueCount.ShouldBe(1);
617+
await _connection.SendMessageAsync(message);
618+
_connection.Get_SendQueueCount.ShouldBe(2);
619+
tc.SetResult(true);
620+
for (int i = 0; i < 100; i++)
621+
{
622+
if (_connection.Get_SendQueueCount != 0)
623+
await Task.Delay(100);
624+
else
625+
break;
626+
}
627+
_connection.Get_SendQueueCount.ShouldBe(0);
628+
_mockConnection.Verify();
629+
}
630+
606631
[Fact]
607632
public async Task LastMessageSentAt()
608633
{
609634
var oldTime = _connection.LastMessageSentAt;
610635
await Task.Delay(100);
611636
var message = new OperationMessage();
637+
_mockConnection.Setup(x => x.SendMessageAsync(It.IsAny<OperationMessage>())).CallBase().Verifiable();
612638
_mockConnection.Protected().Setup<Task>("OnSendMessageAsync", message)
613639
.Returns(Task.CompletedTask).Verifiable();
614640
await _connection.SendMessageAsync(message);
@@ -623,6 +649,7 @@ public async Task DoNotSendMessagesAfterOutputIsClosed()
623649
{
624650
// send a message
625651
var message = new OperationMessage();
652+
_mockConnection.Setup(x => x.SendMessageAsync(It.IsAny<OperationMessage>())).CallBase().Verifiable();
626653
_mockConnection.Protected().SetupGet<TimeSpan>("DefaultDisconnectionTimeout").CallBase().Verifiable();
627654
_mockConnection.Protected().Setup<Task>("OnSendMessageAsync", message)
628655
.Returns(Task.CompletedTask).Verifiable();

0 commit comments

Comments
 (0)