feat(batch42): implement foundation helpers — msgtrace, monitor helpers, scheduler
Group A: Create MsgTrace.cs with TraceCompressionType, MsgTraceState (factory methods, pipeline event helpers, sendEvent/sendEventFromJetStream, trace header injection) and MsgTraceHelper (sample, genHeaderMapIfTraceHeadersPresent, initAndSendIngressErrEvent, isMsgTraceEnabled, msgTraceSupport). Adds Trace field to PublishArgument and Trace accessor to ParseContext. Group C: Create MonitorHelpers.cs with GatewayzOptions, Gatewayz, RemoteGatewayz, AccountGatewayz, ExtImport, ExtServiceLatency types; plus 25 standalone helper functions (newSubsDetailList, newSubsList, createProxyInfo, makePeerCerts, decodeBool, decodeUint64, decodeInt, decodeState, decodeSubs, newSubDetail, newClientSubDetail, myUptime, tlsCertNotAfter, urlsToStrings, getPinnedCertsAsSlice, getMonitorGWOptions, createOutboundRemoteGatewayz, createOutboundAccountsGatewayz, createAccountOutboundGatewayz, createInboundAccountsGatewayz, createInboundAccountGatewayz, ResponseHandler, handleResponse, newExtServiceLatency, newExtImport). Group D: Implement GetScheduledMessages in MsgScheduling; add Seq field to InMsg for out-of-band scheduling sort. Group B (GatewayInterestMode.String) already complete.
This commit is contained in:
@@ -14,6 +14,7 @@
|
||||
// Adapted from server/scheduler.go in the NATS server Go source.
|
||||
|
||||
using System.Buffers.Binary;
|
||||
using ZB.MOM.NatsNet.Server;
|
||||
using ZB.MOM.NatsNet.Server.Internal.DataStructures;
|
||||
|
||||
namespace ZB.MOM.NatsNet.Server.Internal;
|
||||
@@ -198,7 +199,124 @@ public sealed class MsgScheduling
|
||||
_timer = new Timer(_ => _run(), null, fireIn, Timeout.InfiniteTimeSpan);
|
||||
}
|
||||
|
||||
// getScheduledMessages is deferred to session 08/19 — requires JetStream inMsg, StoreMsg types.
|
||||
/// <summary>
|
||||
/// Processes expired schedule entries and returns the set of messages to be delivered.
|
||||
/// Each message is retrieved from storage, headers are cleaned and augmented, and the
|
||||
/// subject is replaced with the schedule target. Messages are returned sorted by
|
||||
/// sequence number.
|
||||
/// Mirrors Go <c>MsgScheduling.getScheduledMessages</c> in server/scheduler.go.
|
||||
/// </summary>
|
||||
/// <param name="loadMsg">
|
||||
/// Callback that loads a stored message by sequence number.
|
||||
/// The <c>StoreMsg</c> reuse buffer may be passed; returns <c>null</c> if not found.
|
||||
/// </param>
|
||||
/// <param name="loadLast">
|
||||
/// Callback that loads the last stored message for a given subject.
|
||||
/// Returns <c>null</c> if not found.
|
||||
/// </param>
|
||||
public List<InMsg> GetScheduledMessages(
|
||||
Func<ulong, ZB.MOM.NatsNet.Server.StoreMsg, ZB.MOM.NatsNet.Server.StoreMsg?> loadMsg,
|
||||
Func<string, ZB.MOM.NatsNet.Server.StoreMsg, ZB.MOM.NatsNet.Server.StoreMsg?> loadLast)
|
||||
{
|
||||
var smv = new ZB.MOM.NatsNet.Server.StoreMsg();
|
||||
List<InMsg>? msgs = null;
|
||||
|
||||
_ttls.ExpireTasks((seq, ts) =>
|
||||
{
|
||||
var sm = loadMsg(seq, smv);
|
||||
if (sm != null)
|
||||
{
|
||||
// Already in-flight for this subject — skip.
|
||||
var subj = sm.Subject;
|
||||
if (IsInflight(subj))
|
||||
return false;
|
||||
|
||||
// Validate the schedule pattern header.
|
||||
var patternBytes = NatsMessageHeaders.GetHeader(
|
||||
NatsHeaderConstants.JsSchedulePattern, sm.Hdr);
|
||||
if (patternBytes == null || patternBytes.Length == 0)
|
||||
{
|
||||
Remove(seq);
|
||||
return true;
|
||||
}
|
||||
var pattern = System.Text.Encoding.ASCII.GetString(patternBytes);
|
||||
|
||||
var (next, repeat, ok) = ParseMsgSchedule(pattern, ts);
|
||||
if (!ok)
|
||||
{
|
||||
Remove(seq);
|
||||
return true;
|
||||
}
|
||||
|
||||
var (ttl, ttlOk) = ZB.MOM.NatsNet.Server.NatsStream.GetMessageScheduleTTL(sm.Hdr);
|
||||
if (!ttlOk)
|
||||
{
|
||||
Remove(seq);
|
||||
return true;
|
||||
}
|
||||
|
||||
var target = ZB.MOM.NatsNet.Server.NatsStream.GetMessageScheduleTarget(sm.Hdr);
|
||||
if (string.IsNullOrEmpty(target))
|
||||
{
|
||||
Remove(seq);
|
||||
return true;
|
||||
}
|
||||
|
||||
var source = ZB.MOM.NatsNet.Server.NatsStream.GetMessageScheduleSource(sm.Hdr);
|
||||
if (!string.IsNullOrEmpty(source))
|
||||
{
|
||||
sm = loadLast(source, smv);
|
||||
if (sm == null)
|
||||
{
|
||||
Remove(seq);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
// Copy headers and body — message lives beyond this callback.
|
||||
var hdr = sm.Hdr.Length > 0 ? (byte[])sm.Hdr.Clone() : [];
|
||||
var msg = sm.Msg.Length > 0 ? (byte[])sm.Msg.Clone() : [];
|
||||
|
||||
// Strip schedule-specific headers.
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPresent(hdr, NatsHeaderConstants.JsSchedulePattern) ?? [];
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPrefixPresent(hdr, "Nats-Schedule-") ?? [];
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPrefixPresent(hdr, "Nats-Expected-") ?? [];
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPresent(hdr, NatsHeaderConstants.JsMsgId) ?? [];
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPresent(hdr, NatsHeaderConstants.JsMessageTtl) ?? [];
|
||||
hdr = NatsMessageHeaders.RemoveHeaderIfPresent(hdr, NatsHeaderConstants.JsMsgRollup) ?? [];
|
||||
|
||||
// Add scheduler-specific headers.
|
||||
hdr = NatsMessageHeaders.GenHeader(hdr, NatsHeaderConstants.JsScheduler, subj);
|
||||
if (!repeat)
|
||||
{
|
||||
hdr = NatsMessageHeaders.GenHeader(hdr, NatsHeaderConstants.JsScheduleNext,
|
||||
NatsHeaderConstants.JsScheduleNextPurge);
|
||||
}
|
||||
else
|
||||
{
|
||||
hdr = NatsMessageHeaders.GenHeader(hdr, NatsHeaderConstants.JsScheduleNext,
|
||||
next.ToString("yyyy-MM-ddTHH:mm:ssK"));
|
||||
}
|
||||
if (!string.IsNullOrEmpty(ttl))
|
||||
hdr = NatsMessageHeaders.GenHeader(hdr, NatsHeaderConstants.JsMessageTtl, ttl);
|
||||
|
||||
msgs ??= [];
|
||||
msgs.Add(new InMsg { Seq = seq, Subject = target, Hdr = hdr, Msg = msg });
|
||||
MarkInflight(subj);
|
||||
return false;
|
||||
}
|
||||
|
||||
Remove(seq);
|
||||
return true;
|
||||
});
|
||||
|
||||
if (msgs == null)
|
||||
return [];
|
||||
|
||||
// THW is unordered — sort by sequence before returning.
|
||||
msgs.Sort((a, b) => a.Seq.CompareTo(b.Seq));
|
||||
return msgs;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// Encodes the current schedule state to a binary snapshot.
|
||||
|
||||
Reference in New Issue
Block a user