fail quick connect closed on an unreachable valkey and restore the idempotent exchange
This commit is contained in:
@@ -4,6 +4,7 @@ using System.Globalization;
|
||||
using System.IO;
|
||||
using System.Net;
|
||||
using System.Net.Sockets;
|
||||
using System.Text;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using StackExchange.Redis;
|
||||
@@ -25,6 +26,7 @@ public sealed class RedisFaultProxy : IAsyncDisposable
|
||||
private readonly int _port;
|
||||
|
||||
private volatile bool _cut;
|
||||
private volatile byte[]? _cutAfterMarker;
|
||||
|
||||
private RedisFaultProxy(TcpListener listener, int port, string targetHost, int targetPort)
|
||||
{
|
||||
@@ -74,11 +76,23 @@ public sealed class RedisFaultProxy : IAsyncDisposable
|
||||
DropLiveConnections();
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Arms a cut for the moment after a command containing <paramref name="marker"/> has been forwarded
|
||||
/// and answered, so a test can take Redis away between two round trips of one operation rather than
|
||||
/// only before or after all of them.
|
||||
/// </summary>
|
||||
/// <param name="marker">Text that identifies the command to cut after.</param>
|
||||
public void CutAfterForwarding(string marker) => _cutAfterMarker = Encoding.UTF8.GetBytes(marker);
|
||||
|
||||
/// <summary>
|
||||
/// Lets connections through again. Clients reconnect on their own schedule, so callers have to wait
|
||||
/// for the connection to come back rather than assume it already has.
|
||||
/// </summary>
|
||||
public void Restore() => _cut = false;
|
||||
public void Restore()
|
||||
{
|
||||
_cutAfterMarker = null;
|
||||
_cut = false;
|
||||
}
|
||||
|
||||
/// <inheritdoc/>
|
||||
public async ValueTask DisposeAsync()
|
||||
@@ -136,10 +150,17 @@ public sealed class RedisFaultProxy : IAsyncDisposable
|
||||
_live[client] = 0;
|
||||
_live[upstream] = 0;
|
||||
|
||||
// Registered first, then rechecked: a cut concurrent with this connect would otherwise drop
|
||||
// the live connections before this pair joined them and leave it running through the outage.
|
||||
if (_cut)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
var clientStream = client.GetStream();
|
||||
var upstreamStream = upstream.GetStream();
|
||||
await Task.WhenAny(
|
||||
CopyAsync(clientStream, upstreamStream),
|
||||
CopyFromClientAsync(clientStream, upstreamStream),
|
||||
CopyAsync(upstreamStream, clientStream)).ConfigureAwait(false);
|
||||
}
|
||||
catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException)
|
||||
@@ -157,6 +178,38 @@ public sealed class RedisFaultProxy : IAsyncDisposable
|
||||
}
|
||||
}
|
||||
|
||||
private async Task CopyFromClientAsync(NetworkStream from, NetworkStream to)
|
||||
{
|
||||
var buffer = new byte[16 * 1024];
|
||||
try
|
||||
{
|
||||
while (true)
|
||||
{
|
||||
var read = await from.ReadAsync(buffer, _cts.Token).ConfigureAwait(false);
|
||||
if (read == 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
await to.WriteAsync(buffer.AsMemory(0, read), _cts.Token).ConfigureAwait(false);
|
||||
|
||||
var marker = _cutAfterMarker;
|
||||
if (marker is not null && buffer.AsSpan(0, read).IndexOf(marker) >= 0)
|
||||
{
|
||||
_cutAfterMarker = null;
|
||||
|
||||
// Long enough for the server to have applied the command that was just forwarded.
|
||||
await Task.Delay(TimeSpan.FromMilliseconds(250), _cts.Token).ConfigureAwait(false);
|
||||
Cut();
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception exception) when (exception is IOException or SocketException or OperationCanceledException or ObjectDisposedException)
|
||||
{
|
||||
}
|
||||
}
|
||||
|
||||
private async Task CopyAsync(NetworkStream from, NetworkStream to)
|
||||
{
|
||||
try
|
||||
|
||||
Reference in New Issue
Block a user