Skip to content
Draft
2 changes: 1 addition & 1 deletion Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
<!-- primary library -->
<PackageVersion Include="NetTopologySuite" Version="2.6.0" />
<PackageVersion Include="System.Text.Json" Version="10.0.12" />
<PackageVersion Include="StackExchange.Redis" Version="3.3.1" />
<PackageVersion Include="StackExchange.Redis" Version="4.0.99-alpha" />
<!-- tests, etc -->
<PackageVersion Include="BouncyCastle.Cryptography" Version="2.7.0" />
<PackageVersion Include="coverlet.collector" Version="10.1.0" />
Expand Down
1 change: 1 addition & 0 deletions src/NRedisStack/CountMinSketch/CmsCommandBuilder.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

namespace NRedisStack;

[Obsolete("The command builders are superseded by the command groups (db.CountMinSketch); they will be removed in a future version.")]
public static class CmsCommandBuilder
{
public static SerializedCommand IncrBy(RedisKey key, RedisValue item, long increment)
Expand Down
63 changes: 42 additions & 21 deletions src/NRedisStack/CountMinSketch/CmsCommands.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,55 +3,76 @@

namespace NRedisStack;

/// <summary>
/// The original synchronous Count-Min Sketch surface, served from the <see cref="RespCountMinSketch"/> group.
/// </summary>
/// <remarks>
/// The group is async-only, as StackExchange.Redis's own groups are. The synchronous methods here send through a
/// <b>blocking</b> context (<c>db.Context.Blocking()</c>): every send then completes on the calling thread and the
/// <see cref="ValueTask"/> comes back already completed, so <c>GetAwaiter().GetResult()</c> is a read, not a wait.
/// This is how the synchronous <see cref="IDatabase"/> methods are implemented too.
/// </remarks>
public class CmsCommands : CmsCommandsAsync, ICmsCommands
{
private readonly IDatabase _db;
private readonly RespCountMinSketch _blocking;

public CmsCommands(IDatabase db) : base(db)
{
_db = db;
// once, not per call: a context is immutable and the database's does not change, and Blocking()
// allocates a new one each time it is asked - the same reason RedisDatabase keeps its own
_blocking = db.Context.Blocking().CountMinSketch;
}

/// <inheritdoc/>
public long IncrBy(RedisKey key, RedisValue item, long increment)
private RespCountMinSketch Group
{
return _db.Execute(CmsCommandBuilder.IncrBy(key, item, increment)).ToLong();
get
{
_db.SetLibraryInfoOnce();
return _blocking;
}
}

/// <inheritdoc/>
public long IncrBy(RedisKey key, RedisValue item, long increment)
=> Group.IncrByAsync(key, item, increment).GetAwaiter().GetResult();

/// <inheritdoc/>
public long[] IncrBy(RedisKey key, Tuple<RedisValue, long>[] itemIncrements)
{
return _db.Execute(CmsCommandBuilder.IncrBy(key, itemIncrements)).ToLongArray();
}
=> ToArray(Group.IncrByAsync(key, ToPairs(itemIncrements)));

/// <inheritdoc/>
public CmsInformation Info(RedisKey key)
{
var info = _db.Execute(CmsCommandBuilder.Info(key));
return info.ToCmsInfo();
}
=> Group.InfoAsync(key).GetAwaiter().GetResult();

/// <inheritdoc/>
public bool InitByDim(RedisKey key, long width, long depth)
{
return _db.Execute(CmsCommandBuilder.InitByDim(key, width, depth)).OKtoBoolean();
}
=> True(Group.InitByDimAsync(key, width, depth));

/// <inheritdoc/>
public bool InitByProb(RedisKey key, double error, double probability)
{
return _db.Execute(CmsCommandBuilder.InitByProb(key, error, probability)).OKtoBoolean();
}
=> True(Group.InitByProbAsync(key, error, probability));

/// <inheritdoc/>
public bool Merge(RedisValue destination, long numKeys, RedisValue[] source, long[]? weight = null)
{
return _db.Execute(CmsCommandBuilder.Merge(destination, numKeys, source, weight)).OKtoBoolean();
}
=> True(Group.MergeAsync(ToKey(destination), ToKeys(source, numKeys), weight ?? default(ReadOnlySpan<long>)));

/// <inheritdoc/>
public long[] Query(RedisKey key, params RedisValue[] items)
=> ToArray(Group.QueryAsync(key, items));

// the old surface reported "OK" as true, and threw otherwise; the new one just completes or throws
private static bool True(ValueTask completed)
{
completed.GetAwaiter().GetResult();
return true;
}

// the old surface promised an array the caller owns; the group hands back a pooled lease
private static long[] ToArray(ValueTask<ReadOnlyLease<long>> completed)
{
return _db.Execute(CmsCommandBuilder.Query(key, items)).ToLongArray();
using var lease = completed.GetAwaiter().GetResult();
return lease.ToArray();
}
}
}
98 changes: 74 additions & 24 deletions src/NRedisStack/CountMinSketch/CmsCommandsAsync.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,55 +3,105 @@

namespace NRedisStack;

/// <summary>
/// The original Count-Min Sketch surface, served from the <see cref="RespCountMinSketch"/> group.
/// </summary>
/// <remarks>
/// Every method here is a thin proxy: <c>db.CountMinSketch.XxxAsync(...)</c> plus whatever adaptation the
/// old signature needs. <see cref="IDatabaseAsync"/> carries the context directly, so this works for a
/// database, a batch, a transaction, and the async-only retry wrapper alike. The <see cref="Task"/> shape is
/// kept because this is legacy API; new code should use the group's <see cref="ValueTask"/> methods, which
/// can complete synchronously (client-side caching). The conversion is <c>db.AsTask(...)</c>, which returns
/// the task <see cref="IDatabaseAsync"/>'s own methods would: it carries the database's async state, and a
/// fault is marked observed as it happens, so a discarded task (<c>_ = tran.Cms.InitByDimAsync(...)</c>, then a
/// transaction that aborts) never raises <see cref="TaskScheduler.UnobservedTaskException"/>.
/// </remarks>
public class CmsCommandsAsync : ICmsCommandsAsync
{
private readonly IDatabaseAsync _db;
private readonly RespCountMinSketch _group;

public CmsCommandsAsync(IDatabaseAsync db)
{
_db = db;
_group = db.CountMinSketch;
}

/// <inheritdoc/>
public async Task<long> IncrByAsync(RedisKey key, RedisValue item, long increment)
private RespCountMinSketch Group
{
return (await _db.ExecuteAsync(CmsCommandBuilder.IncrBy(key, item, increment))).ToLong();
get
{
_db.SetLibraryInfoOnce();
return _group;
}
}

/// <inheritdoc/>
public async Task<long[]> IncrByAsync(RedisKey key, Tuple<RedisValue, long>[] itemIncrements)
{
return (await _db.ExecuteAsync(CmsCommandBuilder.IncrBy(key, itemIncrements))).ToLongArray();
}
public Task<long> IncrByAsync(RedisKey key, RedisValue item, long increment)
=> _db.AsTask(Group.IncrByAsync(key, item, increment));

/// <inheritdoc/>
public async Task<CmsInformation> InfoAsync(RedisKey key)
{
var info = await _db.ExecuteAsync(CmsCommandBuilder.Info(key));
return info.ToCmsInfo();
}
public Task<long[]> IncrByAsync(RedisKey key, Tuple<RedisValue, long>[] itemIncrements)
=> _db.AsTask(ToArray(Group.IncrByAsync(key, ToPairs(itemIncrements))));

/// <inheritdoc/>
public Task<CmsInformation> InfoAsync(RedisKey key)
=> _db.AsTask(Group.InfoAsync(key));

/// <inheritdoc/>
public Task<bool> InitByDimAsync(RedisKey key, long width, long depth)
=> _db.AsTask(True(Group.InitByDimAsync(key, width, depth)));

/// <inheritdoc/>
public async Task<bool> InitByDimAsync(RedisKey key, long width, long depth)
public Task<bool> InitByProbAsync(RedisKey key, double error, double probability)
=> _db.AsTask(True(Group.InitByProbAsync(key, error, probability)));

/// <inheritdoc/>
public Task<bool> MergeAsync(RedisValue destination, long numKeys, RedisValue[] source, long[]? weight = null)
=> _db.AsTask(True(Group.MergeAsync(ToKey(destination), ToKeys(source, numKeys), weight ?? default(ReadOnlySpan<long>))));

/// <inheritdoc/>
public Task<long[]> QueryAsync(RedisKey key, params RedisValue[] items)
=> _db.AsTask(ToArray(Group.QueryAsync(key, items)));

// the old surface reported "OK" as true, and threw otherwise; the new one just completes or throws
private static async ValueTask<bool> True(ValueTask pending)
{
return (await _db.ExecuteAsync(CmsCommandBuilder.InitByDim(key, width, depth))).OKtoBoolean();
await pending.ConfigureAwait(false);
return true;
}

/// <inheritdoc/>
public async Task<bool> InitByProbAsync(RedisKey key, double error, double probability)
// the old surface promised an array the caller owns; the group hands back a pooled lease
private static async ValueTask<long[]> ToArray(ValueTask<ReadOnlyLease<long>> pending)
{
return (await _db.ExecuteAsync(CmsCommandBuilder.InitByProb(key, error, probability))).OKtoBoolean();
using var lease = await pending.ConfigureAwait(false);
return lease.ToArray();
}

/// <inheritdoc/>
public async Task<bool> MergeAsync(RedisValue destination, long numKeys, RedisValue[] source, long[]? weight = null)
internal static (RedisValue Item, long Increment)[] ToPairs(Tuple<RedisValue, long>[] itemIncrements)
{
return (await _db.ExecuteAsync(CmsCommandBuilder.Merge(destination, numKeys, source, weight))).OKtoBoolean();
if (itemIncrements is null) throw new ArgumentNullException(nameof(itemIncrements));
var pairs = new (RedisValue, long)[itemIncrements.Length];
for (int i = 0; i < pairs.Length; i++)
{
pairs[i] = (itemIncrements[i].Item1, itemIncrements[i].Item2);
}
return pairs;
}

/// <inheritdoc/>
public async Task<long[]> QueryAsync(RedisKey key, params RedisValue[] items)
// CMS.MERGE's old signature typed its keys as values, so they never took part in cluster routing or
// key prefixing; the group types them correctly, and these adapt the old spelling
internal static RedisKey ToKey(RedisValue value) => (byte[]?)value;

internal static RedisKey[] ToKeys(RedisValue[] source, long numKeys)
{
return (await _db.ExecuteAsync(CmsCommandBuilder.Query(key, items))).ToLongArray();
if (source is null) throw new ArgumentNullException(nameof(source));
if (numKeys != source.Length) throw new ArgumentOutOfRangeException(nameof(numKeys), "numKeys must match the number of source sketches.");
var keys = new RedisKey[source.Length];
for (int i = 0; i < keys.Length; i++)
{
keys[i] = ToKey(source[i]);
}
return keys;
}
}
}
Loading
Loading