Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
34 changes: 28 additions & 6 deletions src/NATS.Client.KeyValueStore/NatsKVContext.cs
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,30 @@ public enum NatsKVStorageType
/// </summary>
public class NatsKVContext : INatsKVContext
{
private const string KvStreamNamePrefix = "KV_";
internal const string KvStreamNamePrefix = "KV_";
private static readonly int KvStreamNamePrefixLen = KvStreamNamePrefix.Length;
private static readonly Regex ValidBucketRegex = new(pattern: @"\A[a-zA-Z0-9_-]+\z", RegexOptions.Compiled);
private readonly NatsKVOpts _opts;

/// <summary>
/// Create a new Key Value Store context
/// </summary>
/// <param name="context">JetStream context</param>
public NatsKVContext(INatsJSContext context) => JetStreamContext = context;
/// <param name="opts">Context options</param>
public NatsKVContext(INatsJSContext context, NatsKVOpts opts)
{
JetStreamContext = context;
_opts = opts;
}

/// <summary>
/// Create a new Key Value Store context
/// </summary>
/// <param name="context">JetStream context</param>
public NatsKVContext(INatsJSContext context)
: this(context, NatsKVOpts.Default)
{
}

/// <inheritdoc />
public INatsJSContext JetStreamContext { get; }
Expand All @@ -45,7 +60,7 @@ public async ValueTask<INatsKVStore> CreateStoreAsync(NatsKVConfig config, Cance

var stream = await JetStreamContext.CreateStreamAsync(streamConfig, cancellationToken);

return new NatsKVStore(config.Bucket, JetStreamContext, stream);
return new NatsKVStore(config.Bucket, JetStreamContext, stream, _opts);
}

/// <inheritdoc />
Expand All @@ -61,7 +76,7 @@ public async ValueTask<INatsKVStore> GetStoreAsync(string bucket, CancellationTo
}

// TODO: KV mirror
return new NatsKVStore(bucket, JetStreamContext, stream);
return new NatsKVStore(bucket, JetStreamContext, stream, _opts);
}

/// <inheritdoc />
Expand All @@ -73,7 +88,7 @@ public async ValueTask<INatsKVStore> UpdateStoreAsync(NatsKVConfig config, Cance

var stream = await JetStreamContext.UpdateStreamAsync(streamConfig, cancellationToken);

return new NatsKVStore(config.Bucket, JetStreamContext, stream);
return new NatsKVStore(config.Bucket, JetStreamContext, stream, _opts);
}

/// <inheritdoc />
Expand All @@ -85,7 +100,7 @@ public async ValueTask<INatsKVStore> CreateOrUpdateStoreAsync(NatsKVConfig confi

var stream = await JetStreamContext.CreateOrUpdateStreamAsync(streamConfig, cancellationToken);

return new NatsKVStore(config.Bucket, JetStreamContext, stream);
return new NatsKVStore(config.Bucket, JetStreamContext, stream, _opts);
}

/// <inheritdoc />
Expand Down Expand Up @@ -251,3 +266,10 @@ private static StreamConfig CreateStreamConfig(NatsKVConfig config)
return streamConfig;
}
}

public class NatsKVOpts
{
public static readonly NatsKVOpts Default = new();

public bool UseDirectGetApiWithKeysInSubject { get; init; }
}
20 changes: 18 additions & 2 deletions src/NATS.Client.KeyValueStore/NatsKVStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -56,14 +56,18 @@ public class NatsKVStore : INatsKVStore
private static readonly NatsKVException KeyCannotStartOrEndWithPeriodException = new("Key cannot start or end with a period");
private static readonly NatsKVException KeyContainsInvalidCharactersException = new("Key contains invalid characters");
private readonly INatsJSStream _stream;
private readonly NatsKVOpts _opts;
private readonly string _kvBucket;
private readonly string _streamName;

internal NatsKVStore(string bucket, INatsJSContext context, INatsJSStream stream)
internal NatsKVStore(string bucket, INatsJSContext context, INatsJSStream stream, NatsKVOpts opts)
{
Bucket = bucket;
JetStreamContext = context;
_stream = stream;
_opts = opts;
_kvBucket = $"$KV.{Bucket}.";
_streamName = NatsKVContext.KvStreamNamePrefix + Bucket;
}

/// <inheritdoc />
Expand Down Expand Up @@ -333,7 +337,19 @@ public async ValueTask<NatsResult<NatsKVEntry<T>>> TryGetEntryAsync<T>(string ke

if (_stream.Info.Config.AllowDirect)
{
var direct = await _stream.GetDirectAsync<T>(request, serializer, cancellationToken);
NatsMsg<T> direct;
if (_opts.UseDirectGetApiWithKeysInSubject)
{
direct = await JetStreamContext.Connection.RequestAsync<object, T>(
subject: $"{JetStreamContext.Opts.Prefix}.DIRECT.GET.{_streamName}.{keySubject}",
data: null,
replySerializer: serializer,
cancellationToken: cancellationToken);
}
else
{
direct = await _stream.GetDirectAsync<T>(request, serializer, cancellationToken);
}

if (direct is { Headers: { } headers } msg)
{
Expand Down
63 changes: 63 additions & 0 deletions tests/NATS.Client.KeyValueStore.Tests/DirectGetTest.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
using NATS.Client.Core.Tests;

namespace NATS.Client.KeyValueStore.Tests;

public class DirectGetTest(ITestOutputHelper output)
{
[Fact]
public async Task API_subject_test()
{
await using var server = await NatsServer.StartJSAsync();
var (nats1, proxy) = server.CreateProxiedClientConnection();
await using var nats = nats1;

var js = new NatsJSContext(nats);

var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
var cancellationToken = cts.Token;

// default
{
var kv = new NatsKVContext(js);

var store = await kv.CreateStoreAsync(new NatsKVConfig("b1"), cancellationToken);
await store.PutAsync("x", 1, cancellationToken: cancellationToken);
await store.PutAsync("x", 2, cancellationToken: cancellationToken);

await proxy.FlushFramesAsync(nats);

var entry = await store.GetEntryAsync<int>("x", cancellationToken: cancellationToken);
Assert.Equal(2, entry.Value);

var proto = proxy.ClientFrames[0].Message;
Assert.StartsWith("PUB $JS.API.DIRECT.GET.KV_b1 _INBOX.", proto);
Assert.EndsWith("""␍␊{"last_by_subj":"$KV.b1.x"}""", proto);
foreach (var proxyFrame in proxy.ClientFrames)
{
output.WriteLine(proxyFrame.Message);
}
}

// key in api subject
{
var kv = new NatsKVContext(js, new NatsKVOpts { UseDirectGetApiWithKeysInSubject = true });

var store = await kv.CreateStoreAsync(new NatsKVConfig("b1"), cancellationToken);
await store.PutAsync("x", 1, cancellationToken: cancellationToken);
await store.PutAsync("x", 2, cancellationToken: cancellationToken);

await proxy.FlushFramesAsync(nats);

var entry = await store.GetEntryAsync<int>("x", cancellationToken: cancellationToken);
Assert.Equal(2, entry.Value);

var proto = proxy.ClientFrames[0].Message;
Assert.StartsWith("PUB $JS.API.DIRECT.GET.KV_b1.$KV.b1.x _INBOX.", proto);
Assert.EndsWith(""" 0␍␊""", proto);
foreach (var proxyFrame in proxy.ClientFrames)
{
output.WriteLine(proxyFrame.Message);
}
}
}
}