Files
natsdotnet/tests/NATS.Server.Tests/JetStreamApiProtocolIntegrationTests.cs

69 lines
2.1 KiB
C#

using Microsoft.Extensions.Logging.Abstractions;
using NATS.Client.Core;
using NATS.Server.Configuration;
namespace NATS.Server.Tests;
public class JetStreamApiProtocolIntegrationTests
{
[Fact]
public async Task Js_api_request_over_pub_reply_returns_response_message()
{
await using var server = await ServerFixture.StartJetStreamEnabledAsync();
var response = await server.RequestAsync("$JS.API.INFO", "{}", timeoutMs: 1000);
response.ShouldContain("\"error\"");
}
}
internal sealed class ServerFixture : IAsyncDisposable
{
private readonly NatsServer _server;
private readonly CancellationTokenSource _cts;
private ServerFixture(NatsServer server, CancellationTokenSource cts)
{
_server = server;
_cts = cts;
}
public static async Task<ServerFixture> StartJetStreamEnabledAsync()
{
var options = new NatsOptions
{
Host = "127.0.0.1",
Port = 0,
JetStream = new JetStreamOptions
{
StoreDir = Path.Combine(Path.GetTempPath(), $"nats-js-proto-{Guid.NewGuid():N}"),
MaxMemoryStore = 1024 * 1024,
MaxFileStore = 10 * 1024 * 1024,
},
};
var server = new NatsServer(options, NullLoggerFactory.Instance);
var cts = new CancellationTokenSource();
_ = server.StartAsync(cts.Token);
await server.WaitForReadyAsync();
return new ServerFixture(server, cts);
}
public async Task<string> RequestAsync(string subject, string payload, int timeoutMs)
{
await using var conn = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{_server.Port}" });
await conn.ConnectAsync();
using var timeout = new CancellationTokenSource(TimeSpan.FromMilliseconds(timeoutMs));
var response = await conn.RequestAsync<string, string>(subject, payload, cancellationToken: timeout.Token);
return response.Data ?? string.Empty;
}
public async ValueTask DisposeAsync()
{
await _cts.CancelAsync();
_server.Dispose();
_cts.Dispose();
}
}