FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

improve gateway stability against new discord error conditions by akiraveliara · Pull Request #2455 · DSharpPlus/DSharpPlus · GitHub

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

Filter by extension

Filter by extension .cs  (4) All 1 file type selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
9 changes: 9 additions & 0 deletions DSharpPlus/Clients/MultiShardOrchestrator.cs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Linq;
using System.Threading.RateLimiting;
using System.Threading.Tasks;

using DSharpPlus.Entities;
Expand All @@ -27,6 +28,7 @@ public sealed class MultiShardOrchestrator : IShardOrchestrator
private readonly IServiceProvider serviceProvider;
private readonly IPayloadDecompressor decompressor;
private readonly ILogger<IShardOrchestrator> logger;
private readonly ConcurrencyLimiter reconnectConcurrencyLimiter;

private uint shardCount;
private uint stride;
Expand Down Expand Up @@ -67,6 +69,12 @@ ILogger<IShardOrchestrator> logger
this.serviceProvider = serviceProvider;
this.decompressor = decompressor;
this.logger = logger;

this.reconnectConcurrencyLimiter = new(new()
{
PermitLimit = 1,
QueueLimit = int.MaxValue
});
}

/// <inheritdoc/>
Expand Down Expand Up @@ -251,6 +259,7 @@ private async Task RunGatewayAsync(int shardNumber, string uri, DiscordActivity?
break;
}

using RateLimitLease lease = await this.reconnectConcurrencyLimiter.AcquireAsync();
gatewayTask = this.shards[shardNumber].ReconnectAsync();
}
}
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
using System;

namespace DSharpPlus.Exceptions;

/// <summary>
/// Indicates that the gateway has entered an invalid state and must reconnect.
/// </summary>
public sealed class InvalidGatewayStateException(string message) : Exception(message);
35 changes: 31 additions & 4 deletions DSharpPlus/Net/Gateway/GatewayClient.cs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using System.Threading.Tasks;

using DSharpPlus.Entities;
using DSharpPlus.Exceptions;
using DSharpPlus.Net.Abstractions;
using DSharpPlus.Net.Gateway.Compression;
using DSharpPlus.Net.Serialization;
Expand Down Expand Up @@ -328,6 +329,8 @@ await WriteAsync
/// </summary>
private async Task HandleEventsAsync(CancellationToken ct)
{
int consecutiveInvalidEvents = 0;

try
{
while (!ct.IsCancellationRequested)
Expand All @@ -344,7 +347,30 @@ private async Task HandleEventsAsync(CancellationToken ct)
continue;
}

GatewayPayload? payload = await ProcessAndDeserializeTransportFrameAsync(frame);
GatewayPayload? payload = null;

try
{
payload = await ProcessAndDeserializeTransportFrameAsync(frame);
consecutiveInvalidEvents = 0;
}
catch (JsonException)
{
if (frame.TryGetMessage(out byte[]? message))
{
this.logger.LogError("Received undeserializable gateway payload: {payload}", Encoding.UTF8.GetString(message));
consecutiveInvalidEvents++;
}
else
{
throw;
}

if (consecutiveInvalidEvents > 5)
{
throw new InvalidGatewayStateException("Received over five consecutive invalid events, the connection likely entered an invalid state and should reconnect.");
}
}

if (payload is null)
{
Expand All @@ -363,7 +389,8 @@ private async Task HandleEventsAsync(CancellationToken ct)
}
}

private async Task HandleEventCoreAsync(GatewayPayload payload)
[AsyncMethodBuilder(typeof(PoolingAsyncValueTaskMethodBuilder))]
private async ValueTask HandleEventCoreAsync(GatewayPayload payload)
{
switch (payload.OpCode)
{
Expand Down Expand Up @@ -451,7 +478,7 @@ private async Task HandleEventCoreAsync(GatewayPayload payload)
/// </summary>
private async Task<bool> TryResumeAsync()
{
if (this.resumeUrl is null || this.sessionId is null)
if (this.resumeUrl is null || this.sessionId is null || !this.IsConnected)
{
return false;
}
Expand Down Expand Up @@ -581,7 +608,7 @@ private async Task HandleErrorAndAttemptToResumeAsync(TransportFrame frame)
bool success = errorCode switch
{
< 4000 => await HandleSystemErrorAsync(errorCode),
(>= 4000 and <= 4002) or 4005 or 4008 => await TryResumeAsync(),
(>= 4000 and <= 4002) or 4005 or 4008 or 5000 => await TryResumeAsync(),
_ => false
};

Expand Down
6 changes: 6 additions & 0 deletions DSharpPlus/Net/Gateway/TransportService.cs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ await this.socket.CloseAsync
}
catch (WebSocketException) { }
catch (OperationCanceledException) { }
catch (InvalidOperationException) { }

break;
}
Expand Down Expand Up @@ -167,6 +168,11 @@ public async ValueTask<TransportFrame> ReadAsync()

this.metrics.RecordGatewayEventReceived(this.writer.WrittenCount);

if (this.socket.CloseStatus is not null || this.writer.WrittenCount == 0)
{
return new(((int?)this.socket.CloseStatus) ?? 5000, WebSocketMessageType.Close);
}

if (!this.decompressor.TryDecompress(this.writer.WrittenSpan, this.decompressedWriter))
{
throw new InvalidDataException("Failed to decompress a gateway payload.");
Expand Down
Loading

Back | FazBrowse Home | New Git URL