feat: route publishes to jetstream with puback
This commit is contained in:
30
src/NATS.Server/JetStream/Publish/JetStreamPublisher.cs
Normal file
30
src/NATS.Server/JetStream/Publish/JetStreamPublisher.cs
Normal file
@@ -0,0 +1,30 @@
|
||||
namespace NATS.Server.JetStream.Publish;
|
||||
|
||||
public sealed class JetStreamPublisher
|
||||
{
|
||||
private readonly StreamManager _streamManager;
|
||||
|
||||
public JetStreamPublisher(StreamManager streamManager)
|
||||
{
|
||||
_streamManager = streamManager;
|
||||
}
|
||||
|
||||
public bool TryCapture(string subject, ReadOnlyMemory<byte> payload, out PubAck ack)
|
||||
{
|
||||
var stream = _streamManager.FindBySubject(subject);
|
||||
if (stream == null)
|
||||
{
|
||||
ack = new PubAck();
|
||||
return false;
|
||||
}
|
||||
|
||||
var seq = stream.Store.AppendAsync(subject, payload, default).GetAwaiter().GetResult();
|
||||
ack = new PubAck
|
||||
{
|
||||
Stream = stream.Config.Name,
|
||||
Seq = seq,
|
||||
};
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user