231 lines
7.4 KiB
C#
231 lines
7.4 KiB
C#
using System.Text;
|
|
|
|
namespace ZB.MOM.NatsNet.Server;
|
|
|
|
public sealed partial class NatsServer
|
|
{
|
|
internal void SendDomainLeaderElectAdvisory()
|
|
{
|
|
var (_, cluster) = GetJetStreamCluster();
|
|
var meta = cluster?.Meta;
|
|
if (meta == null)
|
|
return;
|
|
|
|
Noticef(
|
|
"JetStream domain leader elected advisory for leader {0} in cluster {1}",
|
|
meta.GroupLeader(),
|
|
CachedClusterName());
|
|
}
|
|
|
|
internal void SendConsumerLostQuorumAdvisory(NatsConsumer? consumer)
|
|
{
|
|
if (consumer == null || !consumer.ShouldSendLostQuorum())
|
|
return;
|
|
|
|
Noticef("JetStream consumer lost quorum advisory for consumer {0} on stream {1}", consumer.Name, consumer.Stream);
|
|
}
|
|
|
|
internal void SendConsumerLeaderElectAdvisory(NatsConsumer? consumer)
|
|
{
|
|
if (consumer == null)
|
|
return;
|
|
|
|
Noticef("JetStream consumer leader elected advisory for consumer {0} on stream {1}", consumer.Name, consumer.Stream);
|
|
}
|
|
|
|
internal void JsClusteredStreamRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage,
|
|
StreamConfigRequest configRequest)
|
|
{
|
|
var (js, cluster) = GetJetStreamCluster();
|
|
if (js == null || cluster == null)
|
|
return;
|
|
|
|
var cfg = configRequest.Config;
|
|
var engine = new JetStreamEngine(js);
|
|
var limitsError = engine.JsClusteredStreamLimitsCheck(account, cfg);
|
|
if (limitsError != null)
|
|
{
|
|
var response = new ApiResponse
|
|
{
|
|
Type = JsApiSubjects.JsApiStreamCreateResponseType,
|
|
Error = limitsError,
|
|
};
|
|
SendAPIErrResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
return;
|
|
}
|
|
|
|
var (group, createError) = engine.CreateGroupForStream(clientInfo, cfg);
|
|
if (group == null || createError != null)
|
|
{
|
|
var response = new ApiResponse
|
|
{
|
|
Type = JsApiSubjects.JsApiStreamCreateResponseType,
|
|
Error = JsApiErrors.NewJSClusterNoPeersError(createError ?? new SelectPeerError { Misc = true }),
|
|
};
|
|
SendAPIErrResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
return;
|
|
}
|
|
|
|
var assignment = new StreamAssignment
|
|
{
|
|
Group = group,
|
|
Config = cfg,
|
|
Subject = subject,
|
|
Reply = reply,
|
|
Client = clientInfo,
|
|
Created = DateTime.UtcNow,
|
|
};
|
|
|
|
if (cluster.Meta != null)
|
|
{
|
|
cluster.Meta.Propose(Encoding.UTF8.GetBytes($"create-stream:{account.Name}:{cfg.Name}"));
|
|
cluster.TrackInflightStreamProposal(account.Name, assignment, deleted: false);
|
|
}
|
|
}
|
|
|
|
internal void JsClusteredStreamUpdateRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage,
|
|
StreamConfig config)
|
|
{
|
|
_ = rawMessage;
|
|
JsClusteredStreamRequest(clientInfo, account, subject, reply, rawMessage, new StreamConfigRequest { Config = config });
|
|
}
|
|
|
|
internal void JsClusteredStreamDeleteRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
string stream,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage)
|
|
{
|
|
_ = rawMessage;
|
|
var (js, cluster) = GetJetStreamCluster();
|
|
if (js == null || cluster?.Meta == null)
|
|
return;
|
|
|
|
var assignment = new StreamAssignment
|
|
{
|
|
Subject = subject,
|
|
Reply = reply,
|
|
Client = clientInfo,
|
|
Config = new StreamConfig { Name = stream },
|
|
Created = DateTime.UtcNow,
|
|
};
|
|
|
|
cluster.Meta.Propose(Encoding.UTF8.GetBytes($"delete-stream:{account.Name}:{stream}"));
|
|
cluster.TrackInflightStreamProposal(account.Name, assignment, deleted: true);
|
|
}
|
|
|
|
internal void JsClusteredStreamPurgeRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
NatsStream? stream,
|
|
string streamName,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage,
|
|
StreamPurgeRequest request)
|
|
{
|
|
_ = stream;
|
|
_ = streamName;
|
|
_ = rawMessage;
|
|
_ = request;
|
|
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiStreamPurgeResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
|
|
internal void JsClusteredStreamRestoreRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
object request,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage)
|
|
{
|
|
_ = request;
|
|
_ = rawMessage;
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiStreamRestoreResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
|
|
internal bool AllPeersOffline(RaftGroup? group)
|
|
{
|
|
if (group == null || group.Peers.Length == 0)
|
|
return false;
|
|
|
|
foreach (var peer in group.Peers)
|
|
{
|
|
if (GetNodeInfo(peer) is { Offline: false })
|
|
return false;
|
|
}
|
|
|
|
return true;
|
|
}
|
|
|
|
internal void JsClusteredStreamListRequest(Account account, ClientInfo clientInfo, string filter, int offset, string subject, string reply, byte[] rawMessage)
|
|
{
|
|
_ = filter;
|
|
_ = offset;
|
|
_ = rawMessage;
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiStreamListResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
|
|
internal void JsClusteredConsumerListRequest(Account account, ClientInfo clientInfo, int offset, string stream, string subject, string reply, byte[] rawMessage)
|
|
{
|
|
_ = offset;
|
|
_ = stream;
|
|
_ = rawMessage;
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiConsumerListResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
|
|
internal void JsClusteredConsumerDeleteRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
string stream,
|
|
string consumer,
|
|
string subject,
|
|
string reply,
|
|
byte[] rawMessage)
|
|
{
|
|
_ = rawMessage;
|
|
var (js, cluster) = GetJetStreamCluster();
|
|
if (js == null || cluster?.Meta == null)
|
|
return;
|
|
|
|
cluster.Meta.Propose(Encoding.UTF8.GetBytes($"delete-consumer:{account.Name}:{stream}:{consumer}"));
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiConsumerDeleteResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
|
|
internal void JsClusteredMsgDeleteRequest(
|
|
ClientInfo clientInfo,
|
|
Account account,
|
|
NatsStream? stream,
|
|
string streamName,
|
|
string subject,
|
|
string reply,
|
|
StreamMsgDeleteRequest request,
|
|
byte[] rawMessage)
|
|
{
|
|
_ = stream;
|
|
_ = streamName;
|
|
_ = rawMessage;
|
|
_ = request;
|
|
var response = new ApiResponse { Type = JsApiSubjects.JsApiMsgDeleteResponseType };
|
|
SendAPIResponse(clientInfo, account, subject, reply, string.Empty, JsonResponse(response));
|
|
}
|
|
}
|