// 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(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(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");