From 09097a8b7852e4cc4b9efc6fb8bd266181523585 Mon Sep 17 00:00:00 2001 From: Aetherance Date: Sun, 9 Aug 2026 20:52:37 +0800 Subject: [PATCH] feat(api): expose cluster membership service --- cmd/kv/main.go | 2 + engine/storage/raft_storage/cluster_api.go | 154 ++++ .../storage/raft_storage/cluster_api_test.go | 112 +++ engine/storage/raft_storage/membership.go | 11 + engine/storage/raft_storage/raft.go | 32 + engine/storage/raft_storage/raft_storage.go | 17 + proto/pkg/clusterpb/clusterpb.pb.go | 819 +++++++++++++++++- proto/pkg/clusterpb/clusterpb_grpc.pb.go | 311 +++++++ proto/proto/clusterpb.proto | 76 ++ 9 files changed, 1519 insertions(+), 15 deletions(-) create mode 100644 engine/storage/raft_storage/cluster_api.go create mode 100644 engine/storage/raft_storage/cluster_api_test.go create mode 100644 proto/pkg/clusterpb/clusterpb_grpc.pb.go diff --git a/cmd/kv/main.go b/cmd/kv/main.go index 89f3c13..be87197 100644 --- a/cmd/kv/main.go +++ b/cmd/kv/main.go @@ -11,6 +11,7 @@ import ( "github.com/Aetherance/kv/engine/config" "github.com/Aetherance/kv/engine/storage/raft_storage" + "github.com/Aetherance/kv/proto/pkg/clusterpb" "github.com/Aetherance/kv/proto/pkg/kvpb" rspb "github.com/Aetherance/kv/proto/pkg/raft_serverpb" "github.com/Aetherance/kv/server" @@ -46,6 +47,7 @@ func main() { grpcServer := grpc.NewServer() kvpb.RegisterKvServer(grpcServer, server.NewServer(rs)) rspb.RegisterRaftServiceServer(grpcServer, rs) + clusterpb.RegisterClusterServer(grpcServer, rs) lis, err := net.Listen("tcp", addr) if err != nil { diff --git a/engine/storage/raft_storage/cluster_api.go b/engine/storage/raft_storage/cluster_api.go new file mode 100644 index 0000000..479918f --- /dev/null +++ b/engine/storage/raft_storage/cluster_api.go @@ -0,0 +1,154 @@ +package raft_storage + +import ( + "context" + "errors" + "strings" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/proto" + + "github.com/Aetherance/kv/proto/pkg/clusterpb" + "github.com/Aetherance/kv/proto/pkg/raftpb" + "github.com/Aetherance/kv/raft" +) + +func (rs *RaftStorage) MemberList(ctx context.Context, _ *clusterpb.MemberListRequest) (*clusterpb.MemberListResponse, error) { + clusterStatus, err := rs.status(ctx) + if err != nil { + return nil, membershipRPCError(err) + } + return memberListResponse(clusterStatus), nil +} + +func (rs *RaftStorage) MemberAdd(ctx context.Context, request *clusterpb.MemberAddRequest) (*clusterpb.MemberAddResponse, error) { + if request == nil || !request.Learner { + return nil, status.Error(codes.InvalidArgument, "members must first be added as learners") + } + member := &clusterpb.Member{Id: request.Id, RaftAddress: strings.TrimSpace(request.RaftAddress)} + if err := rs.proposeMembership(ctx, raftpb.ConfChangeType_AddLearnerNode, member); err != nil { + return nil, membershipRPCError(err) + } + list, err := rs.MemberList(ctx, &clusterpb.MemberListRequest{}) + return &clusterpb.MemberAddResponse{Cluster: list}, err +} + +func (rs *RaftStorage) MemberPromote(ctx context.Context, request *clusterpb.MemberPromoteRequest) (*clusterpb.MemberPromoteResponse, error) { + if request == nil || request.Id == 0 { + return nil, status.Error(codes.InvalidArgument, "member ID must be non-zero") + } + clusterStatus, err := rs.status(ctx) + if err != nil { + return nil, membershipRPCError(err) + } + member, _ := findMember(clusterStatus.metadata, request.Id) + if member == nil { + return nil, membershipRPCError(errMemberNotFound) + } + if err := rs.proposeMembership(ctx, raftpb.ConfChangeType_AddNode, member); err != nil { + return nil, membershipRPCError(err) + } + list, err := rs.MemberList(ctx, &clusterpb.MemberListRequest{}) + return &clusterpb.MemberPromoteResponse{Cluster: list}, err +} + +func (rs *RaftStorage) MemberRemove(ctx context.Context, request *clusterpb.MemberRemoveRequest) (*clusterpb.MemberRemoveResponse, error) { + if request == nil || request.Id == 0 { + return nil, status.Error(codes.InvalidArgument, "member ID must be non-zero") + } + clusterStatus, err := rs.status(ctx) + if err != nil { + return nil, membershipRPCError(err) + } + member, _ := findMember(clusterStatus.metadata, request.Id) + if member == nil { + member = &clusterpb.Member{Id: request.Id} + } + if err := rs.proposeMembership(ctx, raftpb.ConfChangeType_RemoveNode, member); err != nil { + return nil, membershipRPCError(err) + } + list, err := rs.MemberList(ctx, &clusterpb.MemberListRequest{}) + return &clusterpb.MemberRemoveResponse{Cluster: list}, err +} + +func (rs *RaftStorage) MemberUpdate(ctx context.Context, request *clusterpb.MemberUpdateRequest) (*clusterpb.MemberUpdateResponse, error) { + if request == nil { + return nil, status.Error(codes.InvalidArgument, "request is required") + } + member := &clusterpb.Member{Id: request.Id, RaftAddress: strings.TrimSpace(request.RaftAddress)} + if err := rs.proposeMembership(ctx, raftpb.ConfChangeType_UpdateNode, member); err != nil { + return nil, membershipRPCError(err) + } + list, err := rs.MemberList(ctx, &clusterpb.MemberListRequest{}) + return &clusterpb.MemberUpdateResponse{Cluster: list}, err +} + +func (rs *RaftStorage) MemberStatus(ctx context.Context, request *clusterpb.MemberStatusRequest) (*clusterpb.MemberStatusResponse, error) { + if request == nil || request.Id == 0 { + return nil, status.Error(codes.InvalidArgument, "member ID must be non-zero") + } + clusterStatus, err := rs.status(ctx) + if err != nil { + return nil, membershipRPCError(err) + } + member, _ := findMember(clusterStatus.metadata, request.Id) + if member == nil { + return nil, membershipRPCError(errMemberNotFound) + } + return &clusterpb.MemberStatusResponse{ + LeaderId: clusterStatus.leaderID, + CommitIndex: clusterStatus.commitIndex, + Member: memberInfo(clusterStatus, member), + }, nil +} + +func memberListResponse(clusterStatus *clusterStatus) *clusterpb.MemberListResponse { + response := &clusterpb.MemberListResponse{ + ClusterId: clusterStatus.metadata.ClusterId, + LeaderId: clusterStatus.leaderID, + ConfRevision: clusterStatus.metadata.ConfRevision, + Members: make([]*clusterpb.MemberInfo, 0, len(clusterStatus.metadata.Members)), + } + for _, member := range clusterStatus.metadata.Members { + response.Members = append(response.Members, memberInfo(clusterStatus, member)) + } + return response +} + +func memberInfo(clusterStatus *clusterStatus, member *clusterpb.Member) *clusterpb.MemberInfo { + progress, exists := clusterStatus.progress[member.Id] + active := exists && progress.RecentActive + if member.Id == clusterStatus.leaderID { + active = true + } + return &clusterpb.MemberInfo{ + Member: proto.Clone(member).(*clusterpb.Member), + Role: memberRole(clusterStatus.confState, member.Id), + Active: active, + MatchIndex: progress.Match, + } +} + +func membershipRPCError(err error) error { + if err == nil { + return nil + } + var notLeader *NotLeaderError + switch { + case errors.As(err, ¬Leader): + return status.Error(codes.FailedPrecondition, err.Error()) + case errors.Is(err, errMemberNotFound): + return status.Error(codes.NotFound, err.Error()) + case errors.Is(err, errMemberRemoved), errors.Is(err, errMemberAlreadyExists), errors.Is(err, errAddressAlreadyExists): + return status.Error(codes.AlreadyExists, err.Error()) + case errors.Is(err, errLearnerNotReady), errors.Is(err, errLastVoter), errors.Is(err, errTooManyLearners), errors.Is(err, errUnsafeReconfiguration), errors.Is(err, raft.ErrConfChangePending): + return status.Error(codes.FailedPrecondition, err.Error()) + case errors.Is(err, context.Canceled): + return status.Error(codes.Canceled, err.Error()) + case errors.Is(err, context.DeadlineExceeded): + return status.Error(codes.DeadlineExceeded, err.Error()) + default: + return status.Error(codes.InvalidArgument, err.Error()) + } +} diff --git a/engine/storage/raft_storage/cluster_api_test.go b/engine/storage/raft_storage/cluster_api_test.go new file mode 100644 index 0000000..ccf49cb --- /dev/null +++ b/engine/storage/raft_storage/cluster_api_test.go @@ -0,0 +1,112 @@ +package raft_storage + +import ( + "context" + "testing" + + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + + "github.com/Aetherance/kv/proto/pkg/clusterpb" +) + +func TestClusterMembershipAPI(t *testing.T) { + store := startMembershipStore(t) + ctx := context.Background() + + list, err := store.MemberList(ctx, &clusterpb.MemberListRequest{}) + if err != nil { + t.Fatalf("list initial members: %v", err) + } + assertMemberList(t, list, 42, 1, 0, []memberExpectation{{ + id: 1, address: "127.0.0.1:1", role: clusterpb.MemberRole_MemberRoleVoter, active: true, + }}) + + if _, err := store.MemberAdd(ctx, &clusterpb.MemberAddRequest{ + Id: 2, RaftAddress: "127.0.0.1:2", + }); status.Code(err) != codes.InvalidArgument { + t.Fatalf("direct voter add code = %s, want InvalidArgument", status.Code(err)) + } + added, err := store.MemberAdd(ctx, &clusterpb.MemberAddRequest{ + Id: 2, RaftAddress: " 127.0.0.1:2 ", Learner: true, + }) + if err != nil { + t.Fatalf("add learner: %v", err) + } + assertMemberList(t, added.Cluster, 42, 1, 1, []memberExpectation{ + {id: 1, address: "127.0.0.1:1", role: clusterpb.MemberRole_MemberRoleVoter, active: true}, + {id: 2, address: "127.0.0.1:2", role: clusterpb.MemberRole_MemberRoleLearner}, + }) + + memberStatus, err := store.MemberStatus(ctx, &clusterpb.MemberStatusRequest{Id: 2}) + if err != nil { + t.Fatalf("learner status: %v", err) + } + if memberStatus.LeaderId != 1 || memberStatus.Member.Member.Id != 2 || memberStatus.Member.Role != clusterpb.MemberRole_MemberRoleLearner { + t.Fatalf("unexpected learner status: %v", memberStatus) + } + + updated, err := store.MemberUpdate(ctx, &clusterpb.MemberUpdateRequest{Id: 2, RaftAddress: "127.0.0.1:22"}) + if err != nil { + t.Fatalf("update learner: %v", err) + } + assertMemberList(t, updated.Cluster, 42, 1, 2, []memberExpectation{ + {id: 1, address: "127.0.0.1:1", role: clusterpb.MemberRole_MemberRoleVoter, active: true}, + {id: 2, address: "127.0.0.1:22", role: clusterpb.MemberRole_MemberRoleLearner}, + }) + + removed, err := store.MemberRemove(ctx, &clusterpb.MemberRemoveRequest{Id: 2}) + if err != nil { + t.Fatalf("remove learner: %v", err) + } + assertMemberList(t, removed.Cluster, 42, 1, 3, []memberExpectation{{ + id: 1, address: "127.0.0.1:1", role: clusterpb.MemberRole_MemberRoleVoter, active: true, + }}) + removed, err = store.MemberRemove(ctx, &clusterpb.MemberRemoveRequest{Id: 2}) + if err != nil || removed.Cluster.ConfRevision != 3 { + t.Fatalf("idempotent remove = %v, %v", removed, err) + } + if _, err := store.MemberAdd(ctx, &clusterpb.MemberAddRequest{ + Id: 2, RaftAddress: "127.0.0.1:22", Learner: true, + }); status.Code(err) != codes.AlreadyExists { + t.Fatalf("re-add removed member code = %s, want AlreadyExists", status.Code(err)) + } +} + +func TestClusterMembershipAPIErrors(t *testing.T) { + store := startMembershipStore(t) + ctx := context.Background() + + if _, err := store.MemberStatus(ctx, &clusterpb.MemberStatusRequest{Id: 99}); status.Code(err) != codes.NotFound { + t.Fatalf("unknown member status code = %s, want NotFound", status.Code(err)) + } + if _, err := store.MemberRemove(ctx, &clusterpb.MemberRemoveRequest{Id: 99}); status.Code(err) != codes.NotFound { + t.Fatalf("unknown member remove code = %s, want NotFound", status.Code(err)) + } + if _, err := store.MemberRemove(ctx, &clusterpb.MemberRemoveRequest{Id: 1}); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("last voter remove code = %s, want FailedPrecondition", status.Code(err)) + } +} + +type memberExpectation struct { + id uint64 + address string + role clusterpb.MemberRole + active bool +} + +func assertMemberList(t *testing.T, response *clusterpb.MemberListResponse, clusterID, leaderID, revision uint64, members []memberExpectation) { + t.Helper() + if response == nil || response.ClusterId != clusterID || response.LeaderId != leaderID || response.ConfRevision != revision { + t.Fatalf("unexpected cluster response: %v", response) + } + if len(response.Members) != len(members) { + t.Fatalf("members = %v, want %v", response.Members, members) + } + for index, want := range members { + got := response.Members[index] + if got.Member.Id != want.id || got.Member.RaftAddress != want.address || got.Role != want.role || got.Active != want.active { + t.Fatalf("member %d = %v, want %+v", index, got, want) + } + } +} diff --git a/engine/storage/raft_storage/membership.go b/engine/storage/raft_storage/membership.go index c3cbde3..2cf59b2 100644 --- a/engine/storage/raft_storage/membership.go +++ b/engine/storage/raft_storage/membership.go @@ -138,6 +138,17 @@ func roleOf(state *raftpb.ConfState, id uint64) raftMemberRole { return memberRoleUnknown } +func memberRole(state *raftpb.ConfState, id uint64) clusterpb.MemberRole { + switch roleOf(state, id) { + case memberRoleVoter: + return clusterpb.MemberRole_MemberRoleVoter + case memberRoleLearner: + return clusterpb.MemberRole_MemberRoleLearner + default: + return clusterpb.MemberRole_MemberRoleUnknown + } +} + func (rs *RaftStorage) validateConfChange(change *raftpb.ConfChange) error { var context clusterpb.ConfChangeContext if change == nil || proto.Unmarshal(change.GetContext(), &context) != nil || context.Member == nil || context.Member.Id != change.NodeId { diff --git a/engine/storage/raft_storage/raft.go b/engine/storage/raft_storage/raft.go index cbc4808..c1d7907 100644 --- a/engine/storage/raft_storage/raft.go +++ b/engine/storage/raft_storage/raft.go @@ -33,6 +33,7 @@ const ( opStep raftOperation = iota opPropose opProposeConfChange + opStatus ) type raftEvent struct { @@ -40,12 +41,23 @@ type raftEvent struct { message *raftpb.Message data []byte change *raftpb.ConfChange + status *clusterStatus done chan error } +type clusterStatus struct { + metadata *clusterpb.ClusterMetadata + confState *raftpb.ConfState + leaderID uint64 + commitIndex uint64 + progress map[uint64]raft.Progress +} + func (rs *RaftStorage) run(ctx context.Context) { err := rs.runLoop(ctx) + rs.lifecycleMu.Lock() rs.runErr = err + rs.lifecycleMu.Unlock() pendingErr := err if pendingErr == nil { @@ -107,6 +119,16 @@ func (rs *RaftStorage) handleEvent(event raftEvent) (error, bool) { if err := rs.node.ProposeConfChange(event.change); err != nil { return err, false } + case opStatus: + if event.status == nil { + return errors.New("raft storage: nil status target"), false + } + event.status.metadata = cloneClusterMetadata(rs.state.cluster) + event.status.confState = rs.node.ConfState() + event.status.leaderID = rs.node.LeaderID() + event.status.commitIndex = rs.node.CommitIndex() + event.status.progress = rs.node.GetProgress() + return nil, false default: return errors.New("raft storage: unknown operation"), false } @@ -194,6 +216,14 @@ func (rs *RaftStorage) proposeConfChangeData(ctx context.Context, change *raftpb return rs.submit(ctx, raftEvent{op: opProposeConfChange, change: change}) } +func (rs *RaftStorage) status(ctx context.Context) (*clusterStatus, error) { + status := new(clusterStatus) + if err := rs.submit(ctx, raftEvent{op: opStatus, status: status}); err != nil { + return nil, err + } + return status, nil +} + func (rs *RaftStorage) submit(ctx context.Context, event raftEvent) error { if rs.inbox == nil || rs.done == nil { return errStopped @@ -218,6 +248,8 @@ func (rs *RaftStorage) submit(ctx context.Context, event raftEvent) error { } func (rs *RaftStorage) stoppedError() error { + rs.lifecycleMu.Lock() + defer rs.lifecycleMu.Unlock() if rs.runErr != nil { return rs.runErr } diff --git a/engine/storage/raft_storage/raft_storage.go b/engine/storage/raft_storage/raft_storage.go index 6f6c66f..68710cb 100644 --- a/engine/storage/raft_storage/raft_storage.go +++ b/engine/storage/raft_storage/raft_storage.go @@ -2,6 +2,7 @@ package raft_storage import ( "context" + cryptorand "crypto/rand" "encoding/binary" "errors" "fmt" @@ -42,6 +43,7 @@ type requestID struct { // state machines. type RaftStorage struct { rspb.UnimplementedRaftServiceServer + clusterpb.UnimplementedClusterServer config *config.Config clusterID uint64 @@ -166,6 +168,12 @@ func (rs *RaftStorage) Start() error { rs.state = state rs.node = rawNode rs.clusterID = state.cluster.ClusterId + seed, err := proposalSequenceSeed() + if err != nil { + cleanup() + return fmt.Errorf("raft storage: initialize proposal sequence: %w", err) + } + rs.sequence.Store(seed) rs.transport = NewServerTransport(rs.clusterID, clusterAddresses(state.cluster)) rs.inbox = make(chan raftEvent, 256) @@ -377,6 +385,15 @@ func decodeProposal(data []byte) (requestID, []byte, error) { }, data[20:], nil } +func proposalSequenceSeed() (uint64, error) { + var encoded [8]byte + if _, err := cryptorand.Read(encoded[:]); err != nil { + return 0, err + } + // Keep enough headroom that a process cannot realistically wrap the counter. + return binary.BigEndian.Uint64(encoded[:]) & ((uint64(1) << 62) - 1), nil +} + func captureKVSnapshot(txn *badger.Txn, cluster *clusterpb.ClusterMetadata) ([]byte, error) { snapshot := &rspb.RaftSnapshotData{Cluster: cloneClusterMetadata(cluster)} for _, cf := range engine_util.CFs { diff --git a/proto/pkg/clusterpb/clusterpb.pb.go b/proto/pkg/clusterpb/clusterpb.pb.go index e684a13..44748a6 100644 --- a/proto/pkg/clusterpb/clusterpb.pb.go +++ b/proto/pkg/clusterpb/clusterpb.pb.go @@ -21,6 +21,55 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) +type MemberRole int32 + +const ( + MemberRole_MemberRoleUnknown MemberRole = 0 + MemberRole_MemberRoleVoter MemberRole = 1 + MemberRole_MemberRoleLearner MemberRole = 2 +) + +// Enum value maps for MemberRole. +var ( + MemberRole_name = map[int32]string{ + 0: "MemberRoleUnknown", + 1: "MemberRoleVoter", + 2: "MemberRoleLearner", + } + MemberRole_value = map[string]int32{ + "MemberRoleUnknown": 0, + "MemberRoleVoter": 1, + "MemberRoleLearner": 2, + } +) + +func (x MemberRole) Enum() *MemberRole { + p := new(MemberRole) + *p = x + return p +} + +func (x MemberRole) String() string { + return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x)) +} + +func (MemberRole) Descriptor() protoreflect.EnumDescriptor { + return file_clusterpb_proto_enumTypes[0].Descriptor() +} + +func (MemberRole) Type() protoreflect.EnumType { + return &file_clusterpb_proto_enumTypes[0] +} + +func (x MemberRole) Number() protoreflect.EnumNumber { + return protoreflect.EnumNumber(x) +} + +// Deprecated: Use MemberRole.Descriptor instead. +func (MemberRole) EnumDescriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{0} +} + type Member struct { state protoimpl.MessageState `protogen:"open.v1"` Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` @@ -205,6 +254,658 @@ func (x *ConfChangeContext) GetMember() *Member { return nil } +type MemberInfo struct { + state protoimpl.MessageState `protogen:"open.v1"` + Member *Member `protobuf:"bytes,1,opt,name=member,proto3" json:"member,omitempty"` + Role MemberRole `protobuf:"varint,2,opt,name=role,proto3,enum=clusterpb.MemberRole" json:"role,omitempty"` + Active bool `protobuf:"varint,3,opt,name=active,proto3" json:"active,omitempty"` + MatchIndex uint64 `protobuf:"varint,4,opt,name=match_index,json=matchIndex,proto3" json:"match_index,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberInfo) Reset() { + *x = MemberInfo{} + mi := &file_clusterpb_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberInfo) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberInfo) ProtoMessage() {} + +func (x *MemberInfo) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[3] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberInfo.ProtoReflect.Descriptor instead. +func (*MemberInfo) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{3} +} + +func (x *MemberInfo) GetMember() *Member { + if x != nil { + return x.Member + } + return nil +} + +func (x *MemberInfo) GetRole() MemberRole { + if x != nil { + return x.Role + } + return MemberRole_MemberRoleUnknown +} + +func (x *MemberInfo) GetActive() bool { + if x != nil { + return x.Active + } + return false +} + +func (x *MemberInfo) GetMatchIndex() uint64 { + if x != nil { + return x.MatchIndex + } + return 0 +} + +type MemberListRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberListRequest) Reset() { + *x = MemberListRequest{} + mi := &file_clusterpb_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberListRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberListRequest) ProtoMessage() {} + +func (x *MemberListRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[4] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberListRequest.ProtoReflect.Descriptor instead. +func (*MemberListRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{4} +} + +type MemberListResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + ClusterId uint64 `protobuf:"varint,1,opt,name=cluster_id,json=clusterId,proto3" json:"cluster_id,omitempty"` + LeaderId uint64 `protobuf:"varint,2,opt,name=leader_id,json=leaderId,proto3" json:"leader_id,omitempty"` + ConfRevision uint64 `protobuf:"varint,3,opt,name=conf_revision,json=confRevision,proto3" json:"conf_revision,omitempty"` + Members []*MemberInfo `protobuf:"bytes,4,rep,name=members,proto3" json:"members,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberListResponse) Reset() { + *x = MemberListResponse{} + mi := &file_clusterpb_proto_msgTypes[5] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberListResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberListResponse) ProtoMessage() {} + +func (x *MemberListResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[5] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberListResponse.ProtoReflect.Descriptor instead. +func (*MemberListResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{5} +} + +func (x *MemberListResponse) GetClusterId() uint64 { + if x != nil { + return x.ClusterId + } + return 0 +} + +func (x *MemberListResponse) GetLeaderId() uint64 { + if x != nil { + return x.LeaderId + } + return 0 +} + +func (x *MemberListResponse) GetConfRevision() uint64 { + if x != nil { + return x.ConfRevision + } + return 0 +} + +func (x *MemberListResponse) GetMembers() []*MemberInfo { + if x != nil { + return x.Members + } + return nil +} + +type MemberAddRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` + RaftAddress string `protobuf:"bytes,2,opt,name=raft_address,json=raftAddress,proto3" json:"raft_address,omitempty"` + Learner bool `protobuf:"varint,3,opt,name=learner,proto3" json:"learner,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberAddRequest) Reset() { + *x = MemberAddRequest{} + mi := &file_clusterpb_proto_msgTypes[6] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberAddRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberAddRequest) ProtoMessage() {} + +func (x *MemberAddRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[6] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberAddRequest.ProtoReflect.Descriptor instead. +func (*MemberAddRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{6} +} + +func (x *MemberAddRequest) GetId() uint64 { + if x != nil { + return x.Id + } + return 0 +} + +func (x *MemberAddRequest) GetRaftAddress() string { + if x != nil { + return x.RaftAddress + } + return "" +} + +func (x *MemberAddRequest) GetLearner() bool { + if x != nil { + return x.Learner + } + return false +} + +type MemberAddResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Cluster *MemberListResponse `protobuf:"bytes,1,opt,name=cluster,proto3" json:"cluster,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberAddResponse) Reset() { + *x = MemberAddResponse{} + mi := &file_clusterpb_proto_msgTypes[7] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberAddResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberAddResponse) ProtoMessage() {} + +func (x *MemberAddResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[7] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberAddResponse.ProtoReflect.Descriptor instead. +func (*MemberAddResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{7} +} + +func (x *MemberAddResponse) GetCluster() *MemberListResponse { + if x != nil { + return x.Cluster + } + return nil +} + +type MemberPromoteRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberPromoteRequest) Reset() { + *x = MemberPromoteRequest{} + mi := &file_clusterpb_proto_msgTypes[8] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberPromoteRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberPromoteRequest) ProtoMessage() {} + +func (x *MemberPromoteRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[8] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberPromoteRequest.ProtoReflect.Descriptor instead. +func (*MemberPromoteRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{8} +} + +func (x *MemberPromoteRequest) GetId() uint64 { + if x != nil { + return x.Id + } + return 0 +} + +type MemberPromoteResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Cluster *MemberListResponse `protobuf:"bytes,1,opt,name=cluster,proto3" json:"cluster,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberPromoteResponse) Reset() { + *x = MemberPromoteResponse{} + mi := &file_clusterpb_proto_msgTypes[9] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberPromoteResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberPromoteResponse) ProtoMessage() {} + +func (x *MemberPromoteResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[9] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberPromoteResponse.ProtoReflect.Descriptor instead. +func (*MemberPromoteResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{9} +} + +func (x *MemberPromoteResponse) GetCluster() *MemberListResponse { + if x != nil { + return x.Cluster + } + return nil +} + +type MemberRemoveRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberRemoveRequest) Reset() { + *x = MemberRemoveRequest{} + mi := &file_clusterpb_proto_msgTypes[10] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberRemoveRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberRemoveRequest) ProtoMessage() {} + +func (x *MemberRemoveRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[10] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberRemoveRequest.ProtoReflect.Descriptor instead. +func (*MemberRemoveRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{10} +} + +func (x *MemberRemoveRequest) GetId() uint64 { + if x != nil { + return x.Id + } + return 0 +} + +type MemberRemoveResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Cluster *MemberListResponse `protobuf:"bytes,1,opt,name=cluster,proto3" json:"cluster,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberRemoveResponse) Reset() { + *x = MemberRemoveResponse{} + mi := &file_clusterpb_proto_msgTypes[11] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberRemoveResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberRemoveResponse) ProtoMessage() {} + +func (x *MemberRemoveResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[11] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberRemoveResponse.ProtoReflect.Descriptor instead. +func (*MemberRemoveResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{11} +} + +func (x *MemberRemoveResponse) GetCluster() *MemberListResponse { + if x != nil { + return x.Cluster + } + return nil +} + +type MemberUpdateRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` + RaftAddress string `protobuf:"bytes,2,opt,name=raft_address,json=raftAddress,proto3" json:"raft_address,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberUpdateRequest) Reset() { + *x = MemberUpdateRequest{} + mi := &file_clusterpb_proto_msgTypes[12] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberUpdateRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberUpdateRequest) ProtoMessage() {} + +func (x *MemberUpdateRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[12] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberUpdateRequest.ProtoReflect.Descriptor instead. +func (*MemberUpdateRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{12} +} + +func (x *MemberUpdateRequest) GetId() uint64 { + if x != nil { + return x.Id + } + return 0 +} + +func (x *MemberUpdateRequest) GetRaftAddress() string { + if x != nil { + return x.RaftAddress + } + return "" +} + +type MemberUpdateResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Cluster *MemberListResponse `protobuf:"bytes,1,opt,name=cluster,proto3" json:"cluster,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberUpdateResponse) Reset() { + *x = MemberUpdateResponse{} + mi := &file_clusterpb_proto_msgTypes[13] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberUpdateResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberUpdateResponse) ProtoMessage() {} + +func (x *MemberUpdateResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[13] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberUpdateResponse.ProtoReflect.Descriptor instead. +func (*MemberUpdateResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{13} +} + +func (x *MemberUpdateResponse) GetCluster() *MemberListResponse { + if x != nil { + return x.Cluster + } + return nil +} + +type MemberStatusRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id uint64 `protobuf:"varint,1,opt,name=id,proto3" json:"id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberStatusRequest) Reset() { + *x = MemberStatusRequest{} + mi := &file_clusterpb_proto_msgTypes[14] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberStatusRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberStatusRequest) ProtoMessage() {} + +func (x *MemberStatusRequest) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[14] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberStatusRequest.ProtoReflect.Descriptor instead. +func (*MemberStatusRequest) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{14} +} + +func (x *MemberStatusRequest) GetId() uint64 { + if x != nil { + return x.Id + } + return 0 +} + +type MemberStatusResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + LeaderId uint64 `protobuf:"varint,1,opt,name=leader_id,json=leaderId,proto3" json:"leader_id,omitempty"` + CommitIndex uint64 `protobuf:"varint,2,opt,name=commit_index,json=commitIndex,proto3" json:"commit_index,omitempty"` + Member *MemberInfo `protobuf:"bytes,3,opt,name=member,proto3" json:"member,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MemberStatusResponse) Reset() { + *x = MemberStatusResponse{} + mi := &file_clusterpb_proto_msgTypes[15] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MemberStatusResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MemberStatusResponse) ProtoMessage() {} + +func (x *MemberStatusResponse) ProtoReflect() protoreflect.Message { + mi := &file_clusterpb_proto_msgTypes[15] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MemberStatusResponse.ProtoReflect.Descriptor instead. +func (*MemberStatusResponse) Descriptor() ([]byte, []int) { + return file_clusterpb_proto_rawDescGZIP(), []int{15} +} + +func (x *MemberStatusResponse) GetLeaderId() uint64 { + if x != nil { + return x.LeaderId + } + return 0 +} + +func (x *MemberStatusResponse) GetCommitIndex() uint64 { + if x != nil { + return x.CommitIndex + } + return 0 +} + +func (x *MemberStatusResponse) GetMember() *MemberInfo { + if x != nil { + return x.Member + } + return nil +} + var File_clusterpb_proto protoreflect.FileDescriptor const file_clusterpb_proto_rawDesc = "" + @@ -223,7 +924,59 @@ const file_clusterpb_proto_rawDesc = "" + "\vproposer_id\x18\x01 \x01(\x04R\n" + "proposerId\x12\x1a\n" + "\bsequence\x18\x02 \x01(\x04R\bsequence\x12)\n" + - "\x06member\x18\x03 \x01(\v2\x11.clusterpb.MemberR\x06memberB8Z6github.com/Aetherance/kv/proto/pkg/clusterpb;clusterpbb\x06proto3" + "\x06member\x18\x03 \x01(\v2\x11.clusterpb.MemberR\x06member\"\x9b\x01\n" + + "\n" + + "MemberInfo\x12)\n" + + "\x06member\x18\x01 \x01(\v2\x11.clusterpb.MemberR\x06member\x12)\n" + + "\x04role\x18\x02 \x01(\x0e2\x15.clusterpb.MemberRoleR\x04role\x12\x16\n" + + "\x06active\x18\x03 \x01(\bR\x06active\x12\x1f\n" + + "\vmatch_index\x18\x04 \x01(\x04R\n" + + "matchIndex\"\x13\n" + + "\x11MemberListRequest\"\xa6\x01\n" + + "\x12MemberListResponse\x12\x1d\n" + + "\n" + + "cluster_id\x18\x01 \x01(\x04R\tclusterId\x12\x1b\n" + + "\tleader_id\x18\x02 \x01(\x04R\bleaderId\x12#\n" + + "\rconf_revision\x18\x03 \x01(\x04R\fconfRevision\x12/\n" + + "\amembers\x18\x04 \x03(\v2\x15.clusterpb.MemberInfoR\amembers\"_\n" + + "\x10MemberAddRequest\x12\x0e\n" + + "\x02id\x18\x01 \x01(\x04R\x02id\x12!\n" + + "\fraft_address\x18\x02 \x01(\tR\vraftAddress\x12\x18\n" + + "\alearner\x18\x03 \x01(\bR\alearner\"L\n" + + "\x11MemberAddResponse\x127\n" + + "\acluster\x18\x01 \x01(\v2\x1d.clusterpb.MemberListResponseR\acluster\"&\n" + + "\x14MemberPromoteRequest\x12\x0e\n" + + "\x02id\x18\x01 \x01(\x04R\x02id\"P\n" + + "\x15MemberPromoteResponse\x127\n" + + "\acluster\x18\x01 \x01(\v2\x1d.clusterpb.MemberListResponseR\acluster\"%\n" + + "\x13MemberRemoveRequest\x12\x0e\n" + + "\x02id\x18\x01 \x01(\x04R\x02id\"O\n" + + "\x14MemberRemoveResponse\x127\n" + + "\acluster\x18\x01 \x01(\v2\x1d.clusterpb.MemberListResponseR\acluster\"H\n" + + "\x13MemberUpdateRequest\x12\x0e\n" + + "\x02id\x18\x01 \x01(\x04R\x02id\x12!\n" + + "\fraft_address\x18\x02 \x01(\tR\vraftAddress\"O\n" + + "\x14MemberUpdateResponse\x127\n" + + "\acluster\x18\x01 \x01(\v2\x1d.clusterpb.MemberListResponseR\acluster\"%\n" + + "\x13MemberStatusRequest\x12\x0e\n" + + "\x02id\x18\x01 \x01(\x04R\x02id\"\x85\x01\n" + + "\x14MemberStatusResponse\x12\x1b\n" + + "\tleader_id\x18\x01 \x01(\x04R\bleaderId\x12!\n" + + "\fcommit_index\x18\x02 \x01(\x04R\vcommitIndex\x12-\n" + + "\x06member\x18\x03 \x01(\v2\x15.clusterpb.MemberInfoR\x06member*O\n" + + "\n" + + "MemberRole\x12\x15\n" + + "\x11MemberRoleUnknown\x10\x00\x12\x13\n" + + "\x0fMemberRoleVoter\x10\x01\x12\x15\n" + + "\x11MemberRoleLearner\x10\x022\xef\x03\n" + + "\aCluster\x12K\n" + + "\n" + + "MemberList\x12\x1c.clusterpb.MemberListRequest\x1a\x1d.clusterpb.MemberListResponse\"\x00\x12H\n" + + "\tMemberAdd\x12\x1b.clusterpb.MemberAddRequest\x1a\x1c.clusterpb.MemberAddResponse\"\x00\x12T\n" + + "\rMemberPromote\x12\x1f.clusterpb.MemberPromoteRequest\x1a .clusterpb.MemberPromoteResponse\"\x00\x12Q\n" + + "\fMemberRemove\x12\x1e.clusterpb.MemberRemoveRequest\x1a\x1f.clusterpb.MemberRemoveResponse\"\x00\x12Q\n" + + "\fMemberUpdate\x12\x1e.clusterpb.MemberUpdateRequest\x1a\x1f.clusterpb.MemberUpdateResponse\"\x00\x12Q\n" + + "\fMemberStatus\x12\x1e.clusterpb.MemberStatusRequest\x1a\x1f.clusterpb.MemberStatusResponse\"\x00B8Z6github.com/Aetherance/kv/proto/pkg/clusterpb;clusterpbb\x06proto3" var ( file_clusterpb_proto_rawDescOnce sync.Once @@ -237,20 +990,55 @@ func file_clusterpb_proto_rawDescGZIP() []byte { return file_clusterpb_proto_rawDescData } -var file_clusterpb_proto_msgTypes = make([]protoimpl.MessageInfo, 3) +var file_clusterpb_proto_enumTypes = make([]protoimpl.EnumInfo, 1) +var file_clusterpb_proto_msgTypes = make([]protoimpl.MessageInfo, 16) var file_clusterpb_proto_goTypes = []any{ - (*Member)(nil), // 0: clusterpb.Member - (*ClusterMetadata)(nil), // 1: clusterpb.ClusterMetadata - (*ConfChangeContext)(nil), // 2: clusterpb.ConfChangeContext + (MemberRole)(0), // 0: clusterpb.MemberRole + (*Member)(nil), // 1: clusterpb.Member + (*ClusterMetadata)(nil), // 2: clusterpb.ClusterMetadata + (*ConfChangeContext)(nil), // 3: clusterpb.ConfChangeContext + (*MemberInfo)(nil), // 4: clusterpb.MemberInfo + (*MemberListRequest)(nil), // 5: clusterpb.MemberListRequest + (*MemberListResponse)(nil), // 6: clusterpb.MemberListResponse + (*MemberAddRequest)(nil), // 7: clusterpb.MemberAddRequest + (*MemberAddResponse)(nil), // 8: clusterpb.MemberAddResponse + (*MemberPromoteRequest)(nil), // 9: clusterpb.MemberPromoteRequest + (*MemberPromoteResponse)(nil), // 10: clusterpb.MemberPromoteResponse + (*MemberRemoveRequest)(nil), // 11: clusterpb.MemberRemoveRequest + (*MemberRemoveResponse)(nil), // 12: clusterpb.MemberRemoveResponse + (*MemberUpdateRequest)(nil), // 13: clusterpb.MemberUpdateRequest + (*MemberUpdateResponse)(nil), // 14: clusterpb.MemberUpdateResponse + (*MemberStatusRequest)(nil), // 15: clusterpb.MemberStatusRequest + (*MemberStatusResponse)(nil), // 16: clusterpb.MemberStatusResponse } var file_clusterpb_proto_depIdxs = []int32{ - 0, // 0: clusterpb.ClusterMetadata.members:type_name -> clusterpb.Member - 0, // 1: clusterpb.ConfChangeContext.member:type_name -> clusterpb.Member - 2, // [2:2] is the sub-list for method output_type - 2, // [2:2] is the sub-list for method input_type - 2, // [2:2] is the sub-list for extension type_name - 2, // [2:2] is the sub-list for extension extendee - 0, // [0:2] is the sub-list for field type_name + 1, // 0: clusterpb.ClusterMetadata.members:type_name -> clusterpb.Member + 1, // 1: clusterpb.ConfChangeContext.member:type_name -> clusterpb.Member + 1, // 2: clusterpb.MemberInfo.member:type_name -> clusterpb.Member + 0, // 3: clusterpb.MemberInfo.role:type_name -> clusterpb.MemberRole + 4, // 4: clusterpb.MemberListResponse.members:type_name -> clusterpb.MemberInfo + 6, // 5: clusterpb.MemberAddResponse.cluster:type_name -> clusterpb.MemberListResponse + 6, // 6: clusterpb.MemberPromoteResponse.cluster:type_name -> clusterpb.MemberListResponse + 6, // 7: clusterpb.MemberRemoveResponse.cluster:type_name -> clusterpb.MemberListResponse + 6, // 8: clusterpb.MemberUpdateResponse.cluster:type_name -> clusterpb.MemberListResponse + 4, // 9: clusterpb.MemberStatusResponse.member:type_name -> clusterpb.MemberInfo + 5, // 10: clusterpb.Cluster.MemberList:input_type -> clusterpb.MemberListRequest + 7, // 11: clusterpb.Cluster.MemberAdd:input_type -> clusterpb.MemberAddRequest + 9, // 12: clusterpb.Cluster.MemberPromote:input_type -> clusterpb.MemberPromoteRequest + 11, // 13: clusterpb.Cluster.MemberRemove:input_type -> clusterpb.MemberRemoveRequest + 13, // 14: clusterpb.Cluster.MemberUpdate:input_type -> clusterpb.MemberUpdateRequest + 15, // 15: clusterpb.Cluster.MemberStatus:input_type -> clusterpb.MemberStatusRequest + 6, // 16: clusterpb.Cluster.MemberList:output_type -> clusterpb.MemberListResponse + 8, // 17: clusterpb.Cluster.MemberAdd:output_type -> clusterpb.MemberAddResponse + 10, // 18: clusterpb.Cluster.MemberPromote:output_type -> clusterpb.MemberPromoteResponse + 12, // 19: clusterpb.Cluster.MemberRemove:output_type -> clusterpb.MemberRemoveResponse + 14, // 20: clusterpb.Cluster.MemberUpdate:output_type -> clusterpb.MemberUpdateResponse + 16, // 21: clusterpb.Cluster.MemberStatus:output_type -> clusterpb.MemberStatusResponse + 16, // [16:22] is the sub-list for method output_type + 10, // [10:16] is the sub-list for method input_type + 10, // [10:10] is the sub-list for extension type_name + 10, // [10:10] is the sub-list for extension extendee + 0, // [0:10] is the sub-list for field type_name } func init() { file_clusterpb_proto_init() } @@ -263,13 +1051,14 @@ func file_clusterpb_proto_init() { File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_clusterpb_proto_rawDesc), len(file_clusterpb_proto_rawDesc)), - NumEnums: 0, - NumMessages: 3, + NumEnums: 1, + NumMessages: 16, NumExtensions: 0, - NumServices: 0, + NumServices: 1, }, GoTypes: file_clusterpb_proto_goTypes, DependencyIndexes: file_clusterpb_proto_depIdxs, + EnumInfos: file_clusterpb_proto_enumTypes, MessageInfos: file_clusterpb_proto_msgTypes, }.Build() File_clusterpb_proto = out.File diff --git a/proto/pkg/clusterpb/clusterpb_grpc.pb.go b/proto/pkg/clusterpb/clusterpb_grpc.pb.go new file mode 100644 index 0000000..d12507e --- /dev/null +++ b/proto/pkg/clusterpb/clusterpb_grpc.pb.go @@ -0,0 +1,311 @@ +// Code generated by protoc-gen-go-grpc. DO NOT EDIT. +// versions: +// - protoc-gen-go-grpc v1.6.1 +// - protoc v7.35.0 +// source: clusterpb.proto + +package clusterpb + +import ( + context "context" + grpc "google.golang.org/grpc" + codes "google.golang.org/grpc/codes" + status "google.golang.org/grpc/status" +) + +// This is a compile-time assertion to ensure that this generated file +// is compatible with the grpc package it is being compiled against. +// Requires gRPC-Go v1.64.0 or later. +const _ = grpc.SupportPackageIsVersion9 + +const ( + Cluster_MemberList_FullMethodName = "/clusterpb.Cluster/MemberList" + Cluster_MemberAdd_FullMethodName = "/clusterpb.Cluster/MemberAdd" + Cluster_MemberPromote_FullMethodName = "/clusterpb.Cluster/MemberPromote" + Cluster_MemberRemove_FullMethodName = "/clusterpb.Cluster/MemberRemove" + Cluster_MemberUpdate_FullMethodName = "/clusterpb.Cluster/MemberUpdate" + Cluster_MemberStatus_FullMethodName = "/clusterpb.Cluster/MemberStatus" +) + +// ClusterClient is the client API for Cluster service. +// +// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. +type ClusterClient interface { + MemberList(ctx context.Context, in *MemberListRequest, opts ...grpc.CallOption) (*MemberListResponse, error) + MemberAdd(ctx context.Context, in *MemberAddRequest, opts ...grpc.CallOption) (*MemberAddResponse, error) + MemberPromote(ctx context.Context, in *MemberPromoteRequest, opts ...grpc.CallOption) (*MemberPromoteResponse, error) + MemberRemove(ctx context.Context, in *MemberRemoveRequest, opts ...grpc.CallOption) (*MemberRemoveResponse, error) + MemberUpdate(ctx context.Context, in *MemberUpdateRequest, opts ...grpc.CallOption) (*MemberUpdateResponse, error) + MemberStatus(ctx context.Context, in *MemberStatusRequest, opts ...grpc.CallOption) (*MemberStatusResponse, error) +} + +type clusterClient struct { + cc grpc.ClientConnInterface +} + +func NewClusterClient(cc grpc.ClientConnInterface) ClusterClient { + return &clusterClient{cc} +} + +func (c *clusterClient) MemberList(ctx context.Context, in *MemberListRequest, opts ...grpc.CallOption) (*MemberListResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberListResponse) + err := c.cc.Invoke(ctx, Cluster_MemberList_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *clusterClient) MemberAdd(ctx context.Context, in *MemberAddRequest, opts ...grpc.CallOption) (*MemberAddResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberAddResponse) + err := c.cc.Invoke(ctx, Cluster_MemberAdd_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *clusterClient) MemberPromote(ctx context.Context, in *MemberPromoteRequest, opts ...grpc.CallOption) (*MemberPromoteResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberPromoteResponse) + err := c.cc.Invoke(ctx, Cluster_MemberPromote_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *clusterClient) MemberRemove(ctx context.Context, in *MemberRemoveRequest, opts ...grpc.CallOption) (*MemberRemoveResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberRemoveResponse) + err := c.cc.Invoke(ctx, Cluster_MemberRemove_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *clusterClient) MemberUpdate(ctx context.Context, in *MemberUpdateRequest, opts ...grpc.CallOption) (*MemberUpdateResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberUpdateResponse) + err := c.cc.Invoke(ctx, Cluster_MemberUpdate_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *clusterClient) MemberStatus(ctx context.Context, in *MemberStatusRequest, opts ...grpc.CallOption) (*MemberStatusResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(MemberStatusResponse) + err := c.cc.Invoke(ctx, Cluster_MemberStatus_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +// ClusterServer is the server API for Cluster service. +// All implementations must embed UnimplementedClusterServer +// for forward compatibility. +type ClusterServer interface { + MemberList(context.Context, *MemberListRequest) (*MemberListResponse, error) + MemberAdd(context.Context, *MemberAddRequest) (*MemberAddResponse, error) + MemberPromote(context.Context, *MemberPromoteRequest) (*MemberPromoteResponse, error) + MemberRemove(context.Context, *MemberRemoveRequest) (*MemberRemoveResponse, error) + MemberUpdate(context.Context, *MemberUpdateRequest) (*MemberUpdateResponse, error) + MemberStatus(context.Context, *MemberStatusRequest) (*MemberStatusResponse, error) + mustEmbedUnimplementedClusterServer() +} + +// UnimplementedClusterServer must be embedded to have +// forward compatible implementations. +// +// NOTE: this should be embedded by value instead of pointer to avoid a nil +// pointer dereference when methods are called. +type UnimplementedClusterServer struct{} + +func (UnimplementedClusterServer) MemberList(context.Context, *MemberListRequest) (*MemberListResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberList not implemented") +} +func (UnimplementedClusterServer) MemberAdd(context.Context, *MemberAddRequest) (*MemberAddResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberAdd not implemented") +} +func (UnimplementedClusterServer) MemberPromote(context.Context, *MemberPromoteRequest) (*MemberPromoteResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberPromote not implemented") +} +func (UnimplementedClusterServer) MemberRemove(context.Context, *MemberRemoveRequest) (*MemberRemoveResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberRemove not implemented") +} +func (UnimplementedClusterServer) MemberUpdate(context.Context, *MemberUpdateRequest) (*MemberUpdateResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberUpdate not implemented") +} +func (UnimplementedClusterServer) MemberStatus(context.Context, *MemberStatusRequest) (*MemberStatusResponse, error) { + return nil, status.Error(codes.Unimplemented, "method MemberStatus not implemented") +} +func (UnimplementedClusterServer) mustEmbedUnimplementedClusterServer() {} +func (UnimplementedClusterServer) testEmbeddedByValue() {} + +// UnsafeClusterServer may be embedded to opt out of forward compatibility for this service. +// Use of this interface is not recommended, as added methods to ClusterServer will +// result in compilation errors. +type UnsafeClusterServer interface { + mustEmbedUnimplementedClusterServer() +} + +func RegisterClusterServer(s grpc.ServiceRegistrar, srv ClusterServer) { + // If the following call panics, it indicates UnimplementedClusterServer was + // embedded by pointer and is nil. This will cause panics if an + // unimplemented method is ever invoked, so we test this at initialization + // time to prevent it from happening at runtime later due to I/O. + if t, ok := srv.(interface{ testEmbeddedByValue() }); ok { + t.testEmbeddedByValue() + } + s.RegisterService(&Cluster_ServiceDesc, srv) +} + +func _Cluster_MemberList_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberListRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberList(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberList_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberList(ctx, req.(*MemberListRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Cluster_MemberAdd_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberAddRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberAdd(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberAdd_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberAdd(ctx, req.(*MemberAddRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Cluster_MemberPromote_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberPromoteRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberPromote(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberPromote_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberPromote(ctx, req.(*MemberPromoteRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Cluster_MemberRemove_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberRemoveRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberRemove(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberRemove_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberRemove(ctx, req.(*MemberRemoveRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Cluster_MemberUpdate_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberUpdateRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberUpdate(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberUpdate_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberUpdate(ctx, req.(*MemberUpdateRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Cluster_MemberStatus_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(MemberStatusRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(ClusterServer).MemberStatus(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Cluster_MemberStatus_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(ClusterServer).MemberStatus(ctx, req.(*MemberStatusRequest)) + } + return interceptor(ctx, in, info, handler) +} + +// Cluster_ServiceDesc is the grpc.ServiceDesc for Cluster service. +// It's only intended for direct use with grpc.RegisterService, +// and not to be introspected or modified (even as a copy) +var Cluster_ServiceDesc = grpc.ServiceDesc{ + ServiceName: "clusterpb.Cluster", + HandlerType: (*ClusterServer)(nil), + Methods: []grpc.MethodDesc{ + { + MethodName: "MemberList", + Handler: _Cluster_MemberList_Handler, + }, + { + MethodName: "MemberAdd", + Handler: _Cluster_MemberAdd_Handler, + }, + { + MethodName: "MemberPromote", + Handler: _Cluster_MemberPromote_Handler, + }, + { + MethodName: "MemberRemove", + Handler: _Cluster_MemberRemove_Handler, + }, + { + MethodName: "MemberUpdate", + Handler: _Cluster_MemberUpdate_Handler, + }, + { + MethodName: "MemberStatus", + Handler: _Cluster_MemberStatus_Handler, + }, + }, + Streams: []grpc.StreamDesc{}, + Metadata: "clusterpb.proto", +} diff --git a/proto/proto/clusterpb.proto b/proto/proto/clusterpb.proto index 5c9c696..22fd1b5 100644 --- a/proto/proto/clusterpb.proto +++ b/proto/proto/clusterpb.proto @@ -25,3 +25,79 @@ message ConfChangeContext { uint64 sequence = 2; Member member = 3; } + +enum MemberRole { + MemberRoleUnknown = 0; + MemberRoleVoter = 1; + MemberRoleLearner = 2; +} + +message MemberInfo { + Member member = 1; + MemberRole role = 2; + bool active = 3; + uint64 match_index = 4; +} + +message MemberListRequest {} + +message MemberListResponse { + uint64 cluster_id = 1; + uint64 leader_id = 2; + uint64 conf_revision = 3; + repeated MemberInfo members = 4; +} + +message MemberAddRequest { + uint64 id = 1; + string raft_address = 2; + bool learner = 3; +} + +message MemberAddResponse { + MemberListResponse cluster = 1; +} + +message MemberPromoteRequest { + uint64 id = 1; +} + +message MemberPromoteResponse { + MemberListResponse cluster = 1; +} + +message MemberRemoveRequest { + uint64 id = 1; +} + +message MemberRemoveResponse { + MemberListResponse cluster = 1; +} + +message MemberUpdateRequest { + uint64 id = 1; + string raft_address = 2; +} + +message MemberUpdateResponse { + MemberListResponse cluster = 1; +} + +message MemberStatusRequest { + uint64 id = 1; +} + +message MemberStatusResponse { + uint64 leader_id = 1; + uint64 commit_index = 2; + MemberInfo member = 3; +} + +service Cluster { + rpc MemberList(MemberListRequest) returns (MemberListResponse) {} + rpc MemberAdd(MemberAddRequest) returns (MemberAddResponse) {} + rpc MemberPromote(MemberPromoteRequest) returns (MemberPromoteResponse) {} + rpc MemberRemove(MemberRemoveRequest) returns (MemberRemoveResponse) {} + rpc MemberUpdate(MemberUpdateRequest) returns (MemberUpdateResponse) {} + rpc MemberStatus(MemberStatusRequest) returns (MemberStatusResponse) {} +}