247 lines
8.1 KiB
C#
247 lines
8.1 KiB
C#
using System.Text;
|
|
using System.Text.Json;
|
|
using System.Linq;
|
|
using ZB.MOM.NatsNet.Server.Internal;
|
|
using ZB.MOM.NatsNet.Server.Internal.DataStructures;
|
|
|
|
namespace ZB.MOM.NatsNet.Server;
|
|
|
|
internal sealed partial class NatsStream
|
|
{
|
|
internal static (string Ttl, bool Ok) GetMessageScheduleTTL(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return (string.Empty, true);
|
|
|
|
var ttl = NatsMessageHeaders.GetHeader(NatsHeaderConstants.JsScheduleTtl, hdr);
|
|
if (ttl == null || ttl.Length == 0)
|
|
return (string.Empty, true);
|
|
|
|
var ttlValue = Encoding.ASCII.GetString(ttl);
|
|
var (_, err) = ParseMessageTTL(ttlValue);
|
|
return err == null ? (ttlValue, true) : (string.Empty, false);
|
|
}
|
|
|
|
internal static string GetMessageScheduleTarget(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return string.Empty;
|
|
|
|
var value = NatsMessageHeaders.GetHeader(NatsHeaderConstants.JsScheduleTarget, hdr);
|
|
return value == null || value.Length == 0 ? string.Empty : Encoding.ASCII.GetString(value);
|
|
}
|
|
|
|
internal static string GetMessageScheduleSource(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return string.Empty;
|
|
|
|
var value = NatsMessageHeaders.GetHeader(NatsHeaderConstants.JsScheduleSource, hdr);
|
|
return value == null || value.Length == 0 ? string.Empty : Encoding.ASCII.GetString(value);
|
|
}
|
|
|
|
internal static string GetBatchId(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return string.Empty;
|
|
|
|
var value = NatsMessageHeaders.GetHeader(NatsHeaderConstants.JsBatchId, hdr);
|
|
return value == null || value.Length == 0 ? string.Empty : Encoding.ASCII.GetString(value);
|
|
}
|
|
|
|
internal static (ulong Seq, bool Exists) GetBatchSequence(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return (0, false);
|
|
|
|
var value = NatsMessageHeaders.SliceHeader(NatsHeaderConstants.JsBatchSeq, hdr);
|
|
if (value is null || value.Value.Length == 0)
|
|
return (0, false);
|
|
|
|
var parsed = ServerUtilities.ParseInt64(value.Value.Span);
|
|
return parsed < 0 ? (0, false) : ((ulong)parsed, true);
|
|
}
|
|
|
|
public bool IsClustered()
|
|
{
|
|
_mu.EnterReadLock();
|
|
try
|
|
{
|
|
return IsClusteredInternal();
|
|
}
|
|
finally
|
|
{
|
|
_mu.ExitReadLock();
|
|
}
|
|
}
|
|
|
|
internal bool IsClusteredInternal() => _node is IRaftNode;
|
|
|
|
internal void QueueInbound(IpQueue<InMsg> inbound, string subject, string? reply, byte[]? hdr, byte[]? msg, object? sourceInfo, object? trace)
|
|
{
|
|
_ = sourceInfo;
|
|
_ = trace;
|
|
|
|
var inboundMsg = InMsg.Rent();
|
|
inboundMsg.Subject = subject;
|
|
inboundMsg.Reply = reply;
|
|
inboundMsg.Hdr = hdr;
|
|
inboundMsg.Msg = msg;
|
|
|
|
var (_, error) = inbound.Push(inboundMsg);
|
|
if (error != null)
|
|
inboundMsg.ReturnToPool();
|
|
}
|
|
|
|
internal JsApiMsgGetResponse ProcessDirectGetRequest(string reply, byte[]? hdr, byte[]? msg)
|
|
{
|
|
var response = new JsApiMsgGetResponse();
|
|
if (string.IsNullOrWhiteSpace(reply))
|
|
return response;
|
|
|
|
if (JetStreamVersioning.ErrorOnRequiredApiLevel(GetRequiredApiLevelHeader(hdr)))
|
|
{
|
|
response.Error = JsApiErrors.NewJSRequiredApiLevelError();
|
|
return response;
|
|
}
|
|
|
|
if (msg == null || msg.Length == 0)
|
|
{
|
|
response.Error = JsApiErrors.NewJSBadRequestError();
|
|
return response;
|
|
}
|
|
|
|
JsApiMsgGetRequest? request;
|
|
try
|
|
{
|
|
request = JsonSerializer.Deserialize<JsApiMsgGetRequest>(msg);
|
|
}
|
|
catch (JsonException)
|
|
{
|
|
response.Error = JsApiErrors.NewJSBadRequestError();
|
|
return response;
|
|
}
|
|
|
|
if (request == null)
|
|
{
|
|
response.Error = JsApiErrors.NewJSBadRequestError();
|
|
return response;
|
|
}
|
|
|
|
return GetDirectRequest(request, reply);
|
|
}
|
|
|
|
internal JsApiMsgGetResponse ProcessDirectGetLastBySubjectRequest(string subject, string reply, byte[]? hdr, byte[]? msg)
|
|
{
|
|
var response = new JsApiMsgGetResponse();
|
|
if (string.IsNullOrWhiteSpace(reply))
|
|
return response;
|
|
|
|
if (JetStreamVersioning.ErrorOnRequiredApiLevel(GetRequiredApiLevelHeader(hdr)))
|
|
{
|
|
response.Error = JsApiErrors.NewJSRequiredApiLevelError();
|
|
return response;
|
|
}
|
|
|
|
var request = msg == null || msg.Length == 0
|
|
? new JsApiMsgGetRequest()
|
|
: JsonSerializer.Deserialize<JsApiMsgGetRequest>(msg) ?? new JsApiMsgGetRequest();
|
|
|
|
var key = ExtractDirectGetLastBySubjectKey(subject);
|
|
if (string.IsNullOrEmpty(key))
|
|
{
|
|
response.Error = JsApiErrors.NewJSBadRequestError();
|
|
return response;
|
|
}
|
|
|
|
request.LastFor = key;
|
|
return GetDirectRequest(request, reply);
|
|
}
|
|
|
|
internal JsApiMsgGetResponse GetDirectMulti(JsApiMsgGetRequest request, string reply)
|
|
{
|
|
_ = reply;
|
|
|
|
if (Store == null)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
var filters = request.MultiLastFor ?? [];
|
|
if (filters.Length == 0)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
var (seqs, error) = Store.MultiLastSeqs(filters, request.UpToSeq, 1024);
|
|
if (error != null || seqs.Length == 0)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
var firstSeq = seqs[0];
|
|
var loaded = Store.LoadMsg(firstSeq, new StoreMsg());
|
|
if (loaded == null)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
return new JsApiMsgGetResponse { Message = ToStoredMsg(loaded) };
|
|
}
|
|
|
|
internal JsApiMsgGetResponse GetDirectRequest(JsApiMsgGetRequest request, string reply)
|
|
{
|
|
_ = reply;
|
|
if (request.MultiLastFor is { Length: > 0 })
|
|
return GetDirectMulti(request, reply);
|
|
|
|
if (Store == null)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
StoreMsg? loaded = null;
|
|
if (request.Seq > 0)
|
|
{
|
|
loaded = Store.LoadMsg(request.Seq, new StoreMsg());
|
|
}
|
|
else if (!string.IsNullOrWhiteSpace(request.LastFor))
|
|
{
|
|
loaded = Store.LoadLastMsg(request.LastFor!, new StoreMsg());
|
|
}
|
|
else if (!string.IsNullOrWhiteSpace(request.NextFor))
|
|
{
|
|
var (sm, _) = Store.LoadNextMsg(request.NextFor!, SubscriptionIndex.SubjectHasWildcard(request.NextFor!), request.Seq, new StoreMsg());
|
|
loaded = sm;
|
|
}
|
|
else if (request.StartTime.HasValue)
|
|
{
|
|
var seq = Store.GetSeqFromTime(request.StartTime.Value);
|
|
if (seq > 0)
|
|
loaded = Store.LoadMsg(seq, new StoreMsg());
|
|
}
|
|
|
|
if (loaded == null)
|
|
return new JsApiMsgGetResponse { Error = JsApiErrors.NewJSNoMessageFoundError() };
|
|
|
|
return new JsApiMsgGetResponse { Message = ToStoredMsg(loaded) };
|
|
}
|
|
|
|
private static StoredMsg ToStoredMsg(StoreMsg loaded) => new()
|
|
{
|
|
Subject = loaded.Subject,
|
|
Sequence = loaded.Seq,
|
|
Header = loaded.Hdr,
|
|
Data = loaded.Msg,
|
|
Time = DateTimeOffset.FromUnixTimeMilliseconds(loaded.Ts / 1_000_000L).UtcDateTime,
|
|
};
|
|
|
|
private static string? GetRequiredApiLevelHeader(byte[]? hdr)
|
|
{
|
|
if (hdr == null || hdr.Length == 0)
|
|
return null;
|
|
|
|
var value = NatsMessageHeaders.GetHeader(JsApiSubjects.JsRequiredApiLevel, hdr);
|
|
return value == null || value.Length == 0 ? null : Encoding.ASCII.GetString(value);
|
|
}
|
|
|
|
private static string ExtractDirectGetLastBySubjectKey(string subject)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(subject))
|
|
return string.Empty;
|
|
|
|
var parts = subject.Split('.');
|
|
return parts.Length <= 5 ? string.Empty : string.Join('.', parts.Skip(5));
|
|
}
|
|
}
|