Sitelet https://github.com/StackExchange/StackExchange.Redis/pull/3211/files
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 41 additions & 10 deletions docs/Execute.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,33 +50,64 @@ For a single scalar (blob) reply, reading it via `ExecuteResp`/`ScriptEvaluateRe
Leasing the argument buffer
---

The examples above allocate a fresh `RedisKeyOrValue[]` per call, which rather defeats the point of an API whose main selling point is low allocation. On a hot path, rent the array from `ArrayPool<RedisKeyOrValue>.Shared` instead - but the buffer can only be recycled once you know the server has fully received the write. For a synchronous call, that's once `ExecuteResp` has returned successfully *or* thrown `RedisServerException` (the server still definitely got the command - it just responded with an error). Any other exception (a timeout, a dropped connection) can mean a retry is still using the same buffer, so don't recycle it in that case:
The examples above allocate a fresh `RedisKeyOrValue[]` per call, which rather defeats the point of an API whose main selling point is low allocation. On a hot path, rent the array from `ArrayPool<RedisKeyOrValue>.Shared` instead - and you can return it as soon as the call returns:

```csharp
var args = ArrayPool<RedisKeyOrValue>.Shared.Rent(1); // usually larger!
args[0] = (RedisKey)"mykey";
var canReturn = true;
try
{
args[0] = (RedisKey)"mykey";
using RespResult result = db.ExecuteResp("GET", args.AsMemory(0, 1));
// use result...
}
catch (RedisServerException)
finally
{
throw; // the server responded - still safe to recycle below
ArrayPool<RedisKeyOrValue>.Shared.Return(args, clearArray: true);
}
catch
```

No conditions, and nothing to work out from the exception that came back: the arguments are rendered into the request before the call returns, so the library is not looking at your array afterwards. That holds for a call that succeeds, one that throws for any reason at all, and a fire-and-forget call that never waits for a reply.

The async form is the same, and does *not* need the `await` to happen first - `ExecuteRespAsync` renders the arguments before it hands back the task, so the buffer is yours again the moment the method returns:

```csharp
var args = ArrayPool<RedisKeyOrValue>.Shared.Rent(1);
Task<RespResult> pending;
try
{
canReturn = false; // e.g. a timeout/connection failure - a retry may still need this buffer
throw;
args[0] = (RedisKey)"mykey";
pending = db.ExecuteRespAsync("GET", args.AsMemory(0, 1));
}
finally
{
if (canReturn) ArrayPool<RedisKeyOrValue>.Shared.Return(args, clearArray: true);
ArrayPool<RedisKeyOrValue>.Shared.Return(args, clearArray: true); // before awaiting, deliberately
}

using RespResult result = await pending;
```

That also means one buffer can be refilled and reused across a whole run of queued calls - a batch, say - without waiting for any of them:

```csharp
var batch = db.CreateBatch();
var args = ArrayPool<RedisKeyOrValue>.Shared.Rent(3);
var pending = new List<Task<RespResult>>();
for (int i = 0; i < count; i++)
{
args[0] = key;
args[1] = (RedisValue)("field" + i);
args[2] = (RedisValue)("value" + i);
pending.Add(batch.ExecuteRespAsync("HSET", args.AsMemory(0, 3))); // rendered here, not at Execute()
}

ArrayPool<RedisKeyOrValue>.Shared.Return(args, clearArray: true);
batch.Execute();
await Task.WhenAll(pending);
```

The same rule applies to `ExecuteRespAsync` - just `await` the call before deciding whether it's safe to return the lease.
One caveat, for values built over memory you own: a `RedisValue` can wrap a `ReadOnlyMemory<byte>`, and rendering copies the *value*, not the bytes behind it. Returning the `RedisKeyOrValue[]` is safe; overwriting the byte buffer a `RedisValue` points at, before the request has been written, is not. Values built from `string`, `byte[]` you don't then mutate, or numbers are unaffected.

This last gap is known, and is expected to close in a future update: the same work that moves serialization onto the calling thread consumes the payload bytes before the call returns too, at which point the rule becomes simply "everything you passed is yours again when the call returns", with no exception for memory-backed values.

The original `Execute`/`ExecuteAsync` overload
---
Expand Down
25 changes: 25 additions & 0 deletions docs/Scripting.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,31 @@ RedisValue[]? values = parent.ReadPastArray(static (ref r) => r.ReadRedisValue()

This is equivalent to the manual loop above, just without needing to write it out yourself, and capturing the results as an array.

Leasing the keys and values
---

`ScriptEvaluateResp` takes its keys and values as `ReadOnlyMemory<>`, so on a hot path you can rent those arrays rather than allocating a pair per call - and return them as soon as the call returns, with no conditions attached:

```csharp
var keys = ArrayPool<RedisKey>.Shared.Rent(1);
var values = ArrayPool<RedisValue>.Shared.Rent(2);
try
{
keys[0] = "mykey";
values[0] = 123;
values[1] = 456;
using RespResult result = db.ScriptEvaluateResp(script, keys.AsMemory(0, 1), values.AsMemory(0, 2));
// use result...
}
finally
{
ArrayPool<RedisKey>.Shared.Return(keys, clearArray: true);
ArrayPool<RedisValue>.Shared.Return(values, clearArray: true);
}
```

The arguments are rendered into the request before the call returns - before the task is handed back, for the async form - so the library is not reading your arrays afterwards, whether the call succeeded, threw, or was fire-and-forget. See [Leasing the argument buffer](Execute#leasing-the-argument-buffer) for the async and batched shapes, and for the one caveat: a `RedisValue` wrapping a `ReadOnlyMemory<byte>` is rendered by value, so the array is yours again but the bytes behind such a value are not - a known gap, expected to close in a future update.

Ad-hoc commands
---

Expand Down
52 changes: 49 additions & 3 deletions eng/StackExchange.Redis.Build/AutoDatabaseGenerator.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
using ParamInfo = (string Name, string Type, Microsoft.CodeAnalysis.RefKind RefKind, bool IsParams, bool IsOptional, bool HasDefault, string? Default);
using MethodInfo = (string Name, string ReturnType, StackExchange.Redis.Build.BasicArray<(string Name, string Type, Microsoft.CodeAnalysis.RefKind RefKind, bool IsParams, bool IsOptional, bool HasDefault, string? Default)> Parameters, StackExchange.Redis.Build.BasicArray<string> TypeArgs);
using InterfaceInfo = (string Name, string Namespace, StackExchange.Redis.Build.AutoDatabaseGenerator.KnownInterfaces KnownType, StackExchange.Redis.Build.BasicArray<(string Name, string ReturnType, StackExchange.Redis.Build.BasicArray<(string Name, string Type, Microsoft.CodeAnalysis.RefKind RefKind, bool IsParams, bool IsOptional, bool HasDefault, string? Default)> Parameters, StackExchange.Redis.Build.BasicArray<string> TypeArgs)> Methods);
using ClassInfo = (string Name, string Namespace, StackExchange.Redis.Build.AutoDatabaseGenerator.KnownInterfaces Interfaces, bool IsMutator);
using ClassInfo = (string Name, string Namespace, StackExchange.Redis.Build.AutoDatabaseGenerator.KnownInterfaces Interfaces, bool IsMutator, bool Replays);

namespace StackExchange.Redis.Build;

Expand Down Expand Up @@ -219,8 +219,20 @@ static bool HasAutoDatabaseAttrib(INamedTypeSymbol symbol)
}
}

// [AutoDatabase(Replays = true)] - the owning database can invoke a captured operation more than
// once, so captured Memory<T> arguments have to be copies rather than the caller's own buffer
bool replays = false;
foreach (var attrib in cls.GetAttributes())
{
if (attrib.AttributeClass?.Name is not ("AutoDatabaseAttribute" or "AutoDatabase")) continue;
foreach (var named in attrib.NamedArguments)
{
if (named.Key == "Replays" && named.Value.Value is bool b) replays = b;
}
}

var ns = cls.ContainingNamespace.ToDisplayString(SymbolDisplayFormat.CSharpErrorMessageFormat);
return (cls.Name, ns, known, isMutator);
return (cls.Name, ns, known, isMutator, replays);
}

[Flags]
Expand All @@ -243,6 +255,17 @@ internal enum KnownInterfaces
// not a replayable server command; wrappers want their own behaviour here (e.g. a connection-group
// resolving against the currently-active member and returning null when there is none)
// - streaming returns (IEnumerable<T> / IAsyncEnumerable<T>) whose execution is deferred
// Memory<T>/ReadOnlyMemory<T> arguments still belong to the caller once the call returns; a database
// that replays has to capture a copy. Matched on the open generic prefix so a nullable or an array *of*
// them - neither of which is the caller's own buffer in the same way - does not match.
private static bool IsMemoryArg(string type)
=> type.StartsWith("System.ReadOnlyMemory<", StringComparison.Ordinal)
|| type.StartsWith("System.Memory<", StringComparison.Ordinal);

// the element type, for the ArrayPool we rent the copy from
private static string MemoryElement(string type)
=> type.Substring(type.IndexOf('<') + 1, type.Length - type.IndexOf('<') - 2);

private static bool SkipMethod(MethodInfo method)
=> !method.TypeArgs.IsEmpty
|| method.Name.Contains("Wait")
Expand Down Expand Up @@ -501,12 +524,22 @@ void AppendTupleTypes()
if (flagsArg >= 0)
{
writer.Append(firstBase ? " : " : ", ").Append("global::StackExchange.Redis.IFlaggedRedisArgs");
firstBase = false;
}

// a replaying funnel constrains its TState to IDisposable so it can give cloned
// arguments back; every one of its captures implements it, most as a no-op
if (cls.Replays)
{
writer.Append(firstBase ? " : " : ", ").Append("global::System.IDisposable");
}
writer.NewLine().Append("{").Indent();
for (int p = 0; p < parameters.Length; p++)
{
var clone = cls.Replays && IsMemoryArg(parameters[p].Type);
writer.NewLine().Append("public ").Append(cls.IsMutator && NeedsMap(parameters[p].Type) ? "" : "readonly ").Append(parameters[p].Type)
.Append(" Arg").Append(p).Append(" = arg").Append(p).Append(";");
.Append(" Arg").Append(p).Append(" = ")
.Append(clone ? "global::StackExchange.Redis.CapturedArgs.Clone(arg" : "arg").Append(p).Append(clone ? ");" : ";");
}

if (flagsArg >= 0)
Expand All @@ -518,6 +551,19 @@ void AppendTupleTypes()
.Append(flagsArg).Append(";");
}

if (cls.Replays)
{
// emitted even when there is nothing to give back: the funnel's TState constraint
// needs every capture to satisfy it, and an empty one inlines away
writer.NewLine().Append(cls.IsMutator ? "readonly " : "").Append("public void Dispose()").NewLine().Append("{").Indent();
for (int p = 0; p < parameters.Length; p++)
{
if (!IsMemoryArg(parameters[p].Type)) continue;
writer.NewLine().Append("global::StackExchange.Redis.CapturedArgs.Release(in Arg").Append(p).Append(");");
}
writer.Outdent().NewLine().Append("}");
}

if (!cls.IsMutator)
{
// nothing rewrites keys here, so there is no Map/UnMapper to emit at all
Expand Down
22 changes: 14 additions & 8 deletions src/StackExchange.Redis/APITypes/HashImport.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System;
using System.Buffers;
using System.Collections.Generic;
using System.Diagnostics.CodeAnalysis;
using System.Runtime.CompilerServices;
Expand Down Expand Up @@ -196,31 +197,36 @@ private async Task SafeDiscardAsync(ServerEndPoint server, int db)

// HIMPORT SET <key> <field-set> <value...>: the user-facing per-row import. Carries a reference to its field-set so
// the write path can inject a PREPARE the first time this field-set is seen on a connection (see PhysicalBridge).
internal sealed class HashImportSetMessage : Message.CommandKeyBase
internal sealed class HashImportSetMessage : Message.CommandKeyBase, IRenderedArgsOwner
{
private readonly HashImport _fieldSet;
private readonly ReadOnlyMemory<RedisValue> _values;

public HashImportSetMessage(int db, CommandFlags flags, HashImport fieldSet, in RedisKey key, ReadOnlyMemory<RedisValue> values)
// rendered at construction rather than aliased: this is the per-row bulk-import API, so reusing one
// values buffer per row is the intended usage - and a batch defers every write to Execute(), by which
// point that buffer holds only the last row. Not readonly; see RenderedArgs.
private RenderedArgs _values;

public HashImportSetMessage(int db, CommandFlags flags, HashImport fieldSet, in RedisKey key, ReadOnlyMemory<RedisValue> values, MemoryPool<byte>? pool)
: base(db, flags, RedisCommand.HIMPORT, key)
{
_fieldSet = fieldSet;
_values = values;
_values = RenderedArgs.Create(default, values.Span, pool);
}

void IRenderedArgsOwner.ReleaseRenderedArgs() => RenderedArgs.Recycle(ref _values);

internal HashImport FieldSet => _fieldSet;

protected override void WriteImpl(in MessageWriter writer)
{
var values = _values.Span;
writer.WriteHeader(RedisCommand.HIMPORT, 3 + values.Length);
writer.WriteHeader(RedisCommand.HIMPORT, ArgCount);
writer.WriteBulkString(RedisLiterals.SET);
writer.Write(Key);
_fieldSet.WriteName(writer);
for (int i = 0; i < values.Length; i++) writer.WriteBulkString(values[i]);
_values.WriteTo(writer);
}

public override int ArgCount => 3 + _values.Length;
public override int ArgCount => 3 + _values.Count;
}

// HIMPORT PREPARE <field-set> <field...>: injected fire-and-forget ahead of the first SET for a field-set on a
Expand Down
12 changes: 12 additions & 0 deletions src/StackExchange.Redis/AutoDatabase.cs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,18 @@ namespace StackExchange.Redis;
[AttributeUsage(AttributeTargets.Class | AttributeTargets.Interface, AllowMultiple = false)]
internal sealed class AutoDatabaseAttribute : Attribute
{
/// <summary>
/// Whether the owning database can invoke a captured operation more than once, i.e. it replays.
/// </summary>
/// <remarks>
/// A replaying database outlives the call it captured: it re-invokes with the same captured arguments
/// after the caller has long since regained control. Arguments passed as <c>Memory</c>/<c>ReadOnlyMemory</c>
/// still belong to the caller, who is entitled to reuse that buffer once the call returns - so a replay
/// would read whatever is in it by then. When this is set, the generated capture takes a pooled shallow
/// copy of those arguments and the funnel disposes it; when it is not, the arguments are captured as-is,
/// because a single forward cannot outlive the call.
/// </remarks>
public bool Replays { get; set; }
}

// Implemented by a generated captured-arguments struct only when the owning auto-database implements
Expand Down
87 changes: 54 additions & 33 deletions src/StackExchange.Redis/Availability/RetryDatabase.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

namespace StackExchange.Redis.Availability;

[AutoDatabase]
[AutoDatabase(Replays = true)]
internal partial class RetryDatabase : IDatabaseAsync, IInternalDatabaseAsync
// IRedisArgsMutator <==== if we ever want to support key-mapping
{
Expand Down Expand Up @@ -82,57 +82,78 @@ ITransactionAsync IDatabaseAsync.CreateTransaction(object? asyncState)
// note: the state cannot be taken by `in` here (async methods forbid by-ref parameters);
// only the projection takes it by readonly-ref, avoiding a copy per attempt
private async Task<TResult> ExecuteAsync<TState, TResult>(TState state, AutoDatabaseAsyncOperation<TState, TResult> operation)
where TState : struct
where TState : struct, IDisposable
{
/* key mapping, not used currently
* state.MapInPlace(this); */

int attempt = 0;
TResult result;
// note we need to capture this *before* the attempt - otherwise the failover could happen
// between the failed attempt and fetching this, and we'd miss it
CancellationToken failover = GetNextFailover();
while (true)
// the capture owns a pooled copy of any Memory<T> arguments it holds: the caller's own buffer is
// theirs again once the call returns, and we may still be replaying long after that. Disposed here
// rather than by the capture's creator, because this is the only point that knows the last attempt
// is done. The IDisposable constraint makes that call constrained rather than boxed.
try
{
try
int attempt = 0;
TResult result;
// note we need to capture this *before* the attempt - otherwise the failover could happen
// between the failed attempt and fetching this, and we'd miss it
CancellationToken failover = GetNextFailover();
while (true)
{
result = await operation(in state, _inner).ConfigureAwait(false);
break;
}
catch (Exception ex) when (_controller.CanRetry(++attempt, ex, ref failover, out var delay))
{
await _controller.FailoverOrDelayAsync(delay).ConfigureAwait(false);
try
{
result = await operation(in state, _inner).ConfigureAwait(false);
break;
}
catch (Exception ex) when (_controller.CanRetry(++attempt, ex, ref failover, out var delay))
{
await _controller.FailoverOrDelayAsync(delay).ConfigureAwait(false);
}
}

/* key mapping, not used currently
// post-process results outside the loop
return this.UnMap(state, result);*/
return result;
}
finally
{
state.Dispose();
}
/* key mapping, not used currently
// post-process results outside the loop
return this.UnMap(state, result);*/
return result;
}

private async Task ExecuteAsync<TState>(TState state, AutoDatabaseAsyncOperation<TState> operation)
where TState : struct
where TState : struct, IDisposable
{
/* key mapping, not used currently
* state.MapInPlace(this); */

int attempt = 0;
// note we need to capture this *before* the attempt - otherwise the failover could happen
// between the failed attempt and fetching this, and we'd miss it
CancellationToken failover = GetNextFailover();
while (true)
// see the TResult overload
try
{
try
int attempt = 0;
// note we need to capture this *before* the attempt - otherwise the failover could happen
// between the failed attempt and fetching this, and we'd miss it
CancellationToken failover = GetNextFailover();
while (true)
{
await operation(in state, _inner).ConfigureAwait(false);
break;
}
catch (Exception ex) when (_controller.CanRetry(++attempt, ex, ref failover, out var delay))
{
await _controller.FailoverOrDelayAsync(delay).ConfigureAwait(false);
try
{
await operation(in state, _inner).ConfigureAwait(false);
break;
}
catch (Exception ex) when (_controller.CanRetry(++attempt, ex, ref failover, out var delay))
{
await _controller.FailoverOrDelayAsync(delay).ConfigureAwait(false);
}
}

// (nothing to post-process)
}
finally
{
state.Dispose();
}
// (nothing to post-process)
}

// forwarding is not using: these decorators must implement the interface in full, and the
Expand Down
Loading
Loading