75 lines
2.9 KiB
C#
75 lines
2.9 KiB
C#
// 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");
|