Add C# NATS-over-WebSocket example (core + JetStream replay)

This commit is contained in:
2026-08-25 14:43:55 +00:00
parent 1f42637fcf
commit 69033acd74
3 changed files with 90 additions and 0 deletions
+4
View File
@@ -1 +1,5 @@
deploy.json
# .NET build output (for examples/csharp)
bin/
obj/
+12
View File
@@ -0,0 +1,12 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net8.0</TargetFramework>
<Nullable>enable</Nullable>
<ImplicitUsings>enable</ImplicitUsings>
<RootNamespace>NatsExample</RootNamespace>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="NATS.Net" Version="2.5.*" />
</ItemGroup>
</Project>
+74
View File
@@ -0,0 +1,74 @@
// Small NATS-over-WebSocket example for the deployed server.
//
// dotnet run --project examples/csharp
//
// Env vars:
// NATS_URL default: wss://nats.tes.gd
// NATS_USER default: admin
// NATS_PASSWORD required
using System.Text;
using NATS.Client.Core;
using NATS.Client.JetStream;
using NATS.Client.JetStream.Models;
var url = Environment.GetEnvironmentVariable("NATS_URL") ?? "wss://nats.tes.gd";
var user = Environment.GetEnvironmentVariable("NATS_USER") ?? "admin";
var pass = Environment.GetEnvironmentVariable("NATS_PASSWORD")
?? throw new InvalidOperationException("NATS_PASSWORD is required");
// 1. Connect. wss:// is understood natively; credentials go via AuthOpts.
var opts = new NatsOpts
{
Url = url,
AuthOpts = new NatsAuthOpts { Username = user, Password = pass },
};
await using var nats = new NatsConnection(opts);
await nats.ConnectAsync();
Console.WriteLine($"connected server={nats.ServerInfo?.Name} version={nats.ServerInfo?.Version} jetstream={nats.ServerInfo?.JetStreamAvailable}");
// 2. Core pub/sub — fire-and-forget events, no persistence.
var subject = $"demo.{Guid.NewGuid():N}";
using var subCts = new CancellationTokenSource(TimeSpan.FromSeconds(5));
var subTask = Task.Run(async () =>
{
await foreach (var msg in nats.SubscribeAsync<string>(subject, cancellationToken: subCts.Token))
{
Console.WriteLine($"core recv subject={msg.Subject} data=\"{msg.Data}\"");
subCts.Cancel();
}
});
await Task.Delay(200); // give the subscription a tick to register
await nats.PublishAsync(subject, "hello-core");
try { await subTask; } catch (OperationCanceledException) { }
// 3. JetStream — persistent events. Create a stream, publish a few, then read them back.
var js = new NatsJSContext(nats);
const string streamName = "DEMO_EVENTS";
await js.CreateOrUpdateStreamAsync(new StreamConfig(streamName, new[] { "events.>" }));
for (var i = 1; i <= 5; i++)
{
var ack = await js.PublishAsync($"events.item.{i}", Encoding.UTF8.GetBytes($"payload-{i}"));
Console.WriteLine($"js publish subject=events.item.{i} seq={ack.Seq}");
}
// 4. Replay from the beginning with a fresh ephemeral consumer.
// DeliverPolicy.All ignores prior consumers — you get every message in the stream.
var consumer = await js.CreateOrUpdateConsumerAsync(streamName, new ConsumerConfig
{
Name = $"replay-{Guid.NewGuid():N}",
DeliverPolicy = ConsumerConfigDeliverPolicy.All,
AckPolicy = ConsumerConfigAckPolicy.Explicit,
});
var fetchOpts = new NatsJSFetchOpts { MaxMsgs = 10, Expires = TimeSpan.FromSeconds(2) };
await foreach (var msg in consumer.FetchAsync<byte[]>(opts: fetchOpts))
{
Console.WriteLine($"js replay seq={msg.Metadata?.Sequence.Stream} data=\"{Encoding.UTF8.GetString(msg.Data!)}\"");
await msg.AckAsync();
}
// 5. Clean up the demo stream so re-runs stay idempotent.
await js.DeleteStreamAsync(streamName);
Console.WriteLine("done");