62 lines
1.8 KiB
C#
62 lines
1.8 KiB
C#
using System.Text;
|
|
using System.Text.Json;
|
|
|
|
namespace NATS.Server.JetStream.Api.Handlers;
|
|
|
|
public static class DirectApiHandlers
|
|
{
|
|
private const string Prefix = JetStreamApiSubjects.DirectGet;
|
|
|
|
public static JetStreamApiResponse HandleGet(string subject, ReadOnlySpan<byte> payload, StreamManager streamManager)
|
|
{
|
|
var streamName = ExtractTrailingToken(subject, Prefix);
|
|
if (streamName == null)
|
|
return JetStreamApiResponse.NotFound(subject);
|
|
|
|
var sequence = ParseSequence(payload);
|
|
if (sequence == 0)
|
|
return JetStreamApiResponse.ErrorResponse(400, "sequence required");
|
|
|
|
var message = streamManager.GetMessage(streamName, sequence);
|
|
if (message == null)
|
|
return JetStreamApiResponse.NotFound(subject);
|
|
|
|
return new JetStreamApiResponse
|
|
{
|
|
DirectMessage = new JetStreamDirectMessage
|
|
{
|
|
Sequence = message.Sequence,
|
|
Subject = message.Subject,
|
|
Payload = Encoding.UTF8.GetString(message.Payload.Span),
|
|
},
|
|
};
|
|
}
|
|
|
|
private static string? ExtractTrailingToken(string subject, string prefix)
|
|
{
|
|
if (!subject.StartsWith(prefix, StringComparison.Ordinal))
|
|
return null;
|
|
|
|
var token = subject[prefix.Length..].Trim();
|
|
return token.Length == 0 ? null : token;
|
|
}
|
|
|
|
private static ulong ParseSequence(ReadOnlySpan<byte> payload)
|
|
{
|
|
if (payload.IsEmpty)
|
|
return 0;
|
|
|
|
try
|
|
{
|
|
using var doc = JsonDocument.Parse(payload.ToArray());
|
|
if (doc.RootElement.TryGetProperty("seq", out var seqEl) && seqEl.TryGetUInt64(out var sequence))
|
|
return sequence;
|
|
}
|
|
catch (JsonException)
|
|
{
|
|
}
|
|
|
|
return 0;
|
|
}
|
|
}
|