diff --git a/internal/cli/cli.go b/internal/cli/cli.go index 0374a48e..a2819014 100644 --- a/internal/cli/cli.go +++ b/internal/cli/cli.go @@ -231,8 +231,8 @@ func (cli *CLI) AddMachine(ctx context.Context, remoteMachine RemoteMachine, clu } otherMachines := make([]*pb.MachineInfo, 0, len(listResp.Machines)-1) for _, m := range listResp.Machines { - if m.Id != addResp.Machine.Id { - otherMachines = append(otherMachines, m) + if m.Machine.Id != addResp.Machine.Id { + otherMachines = append(otherMachines, m.Machine) } } @@ -341,11 +341,12 @@ func (cli *CLI) ListMachines(ctx context.Context, clusterName string) error { // Print the list of machines in a table format. tw := tabwriter.NewWriter(os.Stdout, 0, 0, 3, ' ', 0) // Print header. - if _, err = fmt.Fprintln(tw, "NAME\tADDRESS\tPUBLIC KEY\tENDPOINTS"); err != nil { + if _, err = fmt.Fprintln(tw, "NAME\tSTATE\tADDRESS\tPUBLIC KEY\tENDPOINTS"); err != nil { return fmt.Errorf("write header: %w", err) } // Print rows. - for _, m := range listResp.Machines { + for _, member := range listResp.Machines { + m := member.Machine subnet, _ := m.Network.Subnet.ToPrefix() subnet = netip.PrefixFrom(network.MachineIP(subnet), subnet.Bits()) endpoints := make([]string, len(m.Network.Endpoints)) @@ -355,10 +356,18 @@ func (cli *CLI) ListMachines(ctx context.Context, clusterName string) error { } publicKey := secret.Secret(m.Network.PublicKey) if _, err = fmt.Fprintf( - tw, "%s\t%s\t%s\t%s\n", m.Name, subnet, publicKey, strings.Join(endpoints, ", "), + tw, "%s\t%s\t%s\t%s\t%s\n", m.Name, capitalise(member.State.String()), subnet, publicKey, strings.Join(endpoints, ", "), ); err != nil { return fmt.Errorf("write row: %w", err) } } return tw.Flush() } + +// capitalise returns a string where the first character is upper case, and the rest is lower case. +func capitalise(s string) string { + if s == "" { + return "" + } + return strings.ToUpper(s[:1]) + strings.ToLower(s[1:]) +} diff --git a/internal/cli/service.go b/internal/cli/service.go index 6d87d613..4dfb6b00 100644 --- a/internal/cli/service.go +++ b/internal/cli/service.go @@ -25,17 +25,17 @@ func (cli *CLI) RunService(ctx context.Context, clusterName string, opts *Servic _ = c.Close() }() - machResp, err := c.ListMachines(ctx, &emptypb.Empty{}) + listResp, err := c.ListMachines(ctx, &emptypb.Empty{}) if err != nil { return fmt.Errorf("list machines: %w", err) } // TODO: update ListMachine endpoint to return machine status based on the Corrosion member list. - machineIP, _ := machResp.Machines[0].Network.ManagementIp.ToAddr() + machineIP, _ := listResp.Machines[0].Machine.Network.ManagementIp.ToAddr() if opts.Machine != "" { - for _, m := range machResp.Machines { - if m.Name == opts.Machine || m.Id == opts.Machine { - machineIP, _ = m.Network.ManagementIp.ToAddr() + for _, m := range listResp.Machines { + if m.Machine.Name == opts.Machine || m.Machine.Id == opts.Machine { + machineIP, _ = m.Machine.Network.ManagementIp.ToAddr() break } } diff --git a/internal/machine/api/pb/cluster.pb.go b/internal/machine/api/pb/cluster.pb.go index 4fd1352e..94f4feb6 100644 --- a/internal/machine/api/pb/cluster.pb.go +++ b/internal/machine/api/pb/cluster.pb.go @@ -21,6 +21,63 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) +type MachineMember_MembershipState int32 + +const ( + MachineMember_UNKNOWN MachineMember_MembershipState = 0 + // The member is active. + MachineMember_UP MachineMember_MembershipState = 1 + // The member is active, but at least one cluster member suspects its down. For all purposes, + // a SUSPECT member is treated as if it were UP until either it refutes the suspicion (becoming UP) + // or fails to do so (being declared DOWN). + MachineMember_SUSPECT MachineMember_MembershipState = 2 + // The member is confirmed DOWN. + MachineMember_DOWN MachineMember_MembershipState = 3 +) + +// Enum value maps for MachineMember_MembershipState. +var ( + MachineMember_MembershipState_name = map[int32]string{ + 0: "UNKNOWN", + 1: "UP", + 2: "SUSPECT", + 3: "DOWN", + } + MachineMember_MembershipState_value = map[string]int32{ + "UNKNOWN": 0, + "UP": 1, + "SUSPECT": 2, + "DOWN": 3, + } +) + +func (x MachineMember_MembershipState) Enum() *MachineMember_MembershipState { + p := new(MachineMember_MembershipState) + *p = x + return p +} + +func (x MachineMember_MembershipState) String() string { + return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x)) +} + +func (MachineMember_MembershipState) Descriptor() protoreflect.EnumDescriptor { + return file_internal_machine_api_pb_cluster_proto_enumTypes[0].Descriptor() +} + +func (MachineMember_MembershipState) Type() protoreflect.EnumType { + return &file_internal_machine_api_pb_cluster_proto_enumTypes[0] +} + +func (x MachineMember_MembershipState) Number() protoreflect.EnumNumber { + return protoreflect.EnumNumber(x) +} + +// Deprecated: Use MachineMember_MembershipState.Descriptor instead. +func (MachineMember_MembershipState) EnumDescriptor() ([]byte, []int) { + return file_internal_machine_api_pb_cluster_proto_rawDescGZIP(), []int{2, 0} +} + type AddMachineRequest struct { state protoimpl.MessageState sizeCache protoimpl.SizeCache @@ -123,18 +180,73 @@ func (x *AddMachineResponse) GetMachine() *MachineInfo { return nil } +type MachineMember struct { + state protoimpl.MessageState + sizeCache protoimpl.SizeCache + unknownFields protoimpl.UnknownFields + + Machine *MachineInfo `protobuf:"bytes,1,opt,name=machine,proto3" json:"machine,omitempty"` + State MachineMember_MembershipState `protobuf:"varint,2,opt,name=state,proto3,enum=api.MachineMember_MembershipState" json:"state,omitempty"` +} + +func (x *MachineMember) Reset() { + *x = MachineMember{} + if protoimpl.UnsafeEnabled { + mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) + } +} + +func (x *MachineMember) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MachineMember) ProtoMessage() {} + +func (x *MachineMember) ProtoReflect() protoreflect.Message { + mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[2] + if protoimpl.UnsafeEnabled && x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use MachineMember.ProtoReflect.Descriptor instead. +func (*MachineMember) Descriptor() ([]byte, []int) { + return file_internal_machine_api_pb_cluster_proto_rawDescGZIP(), []int{2} +} + +func (x *MachineMember) GetMachine() *MachineInfo { + if x != nil { + return x.Machine + } + return nil +} + +func (x *MachineMember) GetState() MachineMember_MembershipState { + if x != nil { + return x.State + } + return MachineMember_UNKNOWN +} + type ListMachinesResponse struct { state protoimpl.MessageState sizeCache protoimpl.SizeCache unknownFields protoimpl.UnknownFields - Machines []*MachineInfo `protobuf:"bytes,1,rep,name=machines,proto3" json:"machines,omitempty"` + Machines []*MachineMember `protobuf:"bytes,1,rep,name=machines,proto3" json:"machines,omitempty"` } func (x *ListMachinesResponse) Reset() { *x = ListMachinesResponse{} if protoimpl.UnsafeEnabled { - mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[2] + mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -147,7 +259,7 @@ func (x *ListMachinesResponse) String() string { func (*ListMachinesResponse) ProtoMessage() {} func (x *ListMachinesResponse) ProtoReflect() protoreflect.Message { - mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[2] + mi := &file_internal_machine_api_pb_cluster_proto_msgTypes[3] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -160,10 +272,10 @@ func (x *ListMachinesResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use ListMachinesResponse.ProtoReflect.Descriptor instead. func (*ListMachinesResponse) Descriptor() ([]byte, []int) { - return file_internal_machine_api_pb_cluster_proto_rawDescGZIP(), []int{2} + return file_internal_machine_api_pb_cluster_proto_rawDescGZIP(), []int{3} } -func (x *ListMachinesResponse) GetMachines() []*MachineInfo { +func (x *ListMachinesResponse) GetMachines() []*MachineMember { if x != nil { return x.Machines } @@ -177,36 +289,50 @@ var file_internal_machine_api_pb_cluster_proto_rawDesc = []byte{ 0x6e, 0x65, 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x70, 0x62, 0x2f, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x12, 0x03, 0x61, 0x70, 0x69, 0x1a, 0x1b, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2f, 0x65, 0x6d, - 0x70, 0x74, 0x79, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, 0x25, 0x69, 0x6e, 0x74, 0x65, 0x72, + 0x70, 0x74, 0x79, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, 0x24, 0x69, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x2f, 0x61, 0x70, 0x69, 0x2f, - 0x70, 0x62, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, - 0x22, 0x55, 0x0a, 0x11, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, - 0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x12, 0x0a, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x01, 0x20, - 0x01, 0x28, 0x09, 0x52, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x12, 0x2c, 0x0a, 0x07, 0x6e, 0x65, 0x74, - 0x77, 0x6f, 0x72, 0x6b, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x12, 0x2e, 0x61, 0x70, 0x69, - 0x2e, 0x4e, 0x65, 0x74, 0x77, 0x6f, 0x72, 0x6b, 0x43, 0x6f, 0x6e, 0x66, 0x69, 0x67, 0x52, 0x07, - 0x6e, 0x65, 0x74, 0x77, 0x6f, 0x72, 0x6b, 0x22, 0x40, 0x0a, 0x12, 0x41, 0x64, 0x64, 0x4d, 0x61, - 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x2a, 0x0a, - 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x10, - 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x49, 0x6e, 0x66, 0x6f, - 0x52, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x22, 0x44, 0x0a, 0x14, 0x4c, 0x69, 0x73, - 0x74, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, - 0x65, 0x12, 0x2c, 0x0a, 0x08, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x18, 0x01, 0x20, - 0x03, 0x28, 0x0b, 0x32, 0x10, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, - 0x65, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x08, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x32, - 0x8b, 0x01, 0x0a, 0x07, 0x43, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x12, 0x3d, 0x0a, 0x0a, 0x41, - 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x12, 0x16, 0x2e, 0x61, 0x70, 0x69, 0x2e, - 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, - 0x74, 0x1a, 0x17, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, - 0x6e, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x41, 0x0a, 0x0c, 0x4c, 0x69, - 0x73, 0x74, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x12, 0x16, 0x2e, 0x67, 0x6f, 0x6f, - 0x67, 0x6c, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, - 0x74, 0x79, 0x1a, 0x19, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4c, 0x69, 0x73, 0x74, 0x4d, 0x61, 0x63, - 0x68, 0x69, 0x6e, 0x65, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x42, 0x37, 0x5a, - 0x35, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x70, 0x73, 0x76, 0x69, - 0x64, 0x65, 0x72, 0x73, 0x6b, 0x69, 0x2f, 0x75, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x69, - 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x2f, - 0x61, 0x70, 0x69, 0x2f, 0x70, 0x62, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, + 0x70, 0x62, 0x2f, 0x63, 0x6f, 0x6d, 0x6d, 0x6f, 0x6e, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x1a, + 0x25, 0x69, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, + 0x65, 0x2f, 0x61, 0x70, 0x69, 0x2f, 0x70, 0x62, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, + 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, 0x55, 0x0a, 0x11, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, + 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x12, 0x0a, 0x04, 0x6e, + 0x61, 0x6d, 0x65, 0x18, 0x01, 0x20, 0x01, 0x28, 0x09, 0x52, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x12, + 0x2c, 0x0a, 0x07, 0x6e, 0x65, 0x74, 0x77, 0x6f, 0x72, 0x6b, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0b, + 0x32, 0x12, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4e, 0x65, 0x74, 0x77, 0x6f, 0x72, 0x6b, 0x43, 0x6f, + 0x6e, 0x66, 0x69, 0x67, 0x52, 0x07, 0x6e, 0x65, 0x74, 0x77, 0x6f, 0x72, 0x6b, 0x22, 0x40, 0x0a, + 0x12, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, + 0x6e, 0x73, 0x65, 0x12, 0x2a, 0x0a, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x18, 0x01, + 0x20, 0x01, 0x28, 0x0b, 0x32, 0x10, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, + 0x6e, 0x65, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x22, + 0xb4, 0x01, 0x0a, 0x0d, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x4d, 0x65, 0x6d, 0x62, 0x65, + 0x72, 0x12, 0x2a, 0x0a, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x18, 0x01, 0x20, 0x01, + 0x28, 0x0b, 0x32, 0x10, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, + 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x12, 0x38, 0x0a, + 0x05, 0x73, 0x74, 0x61, 0x74, 0x65, 0x18, 0x02, 0x20, 0x01, 0x28, 0x0e, 0x32, 0x22, 0x2e, 0x61, + 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x4d, 0x65, 0x6d, 0x62, 0x65, 0x72, + 0x2e, 0x4d, 0x65, 0x6d, 0x62, 0x65, 0x72, 0x73, 0x68, 0x69, 0x70, 0x53, 0x74, 0x61, 0x74, 0x65, + 0x52, 0x05, 0x73, 0x74, 0x61, 0x74, 0x65, 0x22, 0x3d, 0x0a, 0x0f, 0x4d, 0x65, 0x6d, 0x62, 0x65, + 0x72, 0x73, 0x68, 0x69, 0x70, 0x53, 0x74, 0x61, 0x74, 0x65, 0x12, 0x0b, 0x0a, 0x07, 0x55, 0x4e, + 0x4b, 0x4e, 0x4f, 0x57, 0x4e, 0x10, 0x00, 0x12, 0x06, 0x0a, 0x02, 0x55, 0x50, 0x10, 0x01, 0x12, + 0x0b, 0x0a, 0x07, 0x53, 0x55, 0x53, 0x50, 0x45, 0x43, 0x54, 0x10, 0x02, 0x12, 0x08, 0x0a, 0x04, + 0x44, 0x4f, 0x57, 0x4e, 0x10, 0x03, 0x22, 0x46, 0x0a, 0x14, 0x4c, 0x69, 0x73, 0x74, 0x4d, 0x61, + 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x2e, + 0x0a, 0x08, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x18, 0x01, 0x20, 0x03, 0x28, 0x0b, + 0x32, 0x12, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x4d, 0x65, + 0x6d, 0x62, 0x65, 0x72, 0x52, 0x08, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x32, 0x8b, + 0x01, 0x0a, 0x07, 0x43, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x12, 0x3d, 0x0a, 0x0a, 0x41, 0x64, + 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x12, 0x16, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x41, + 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, + 0x1a, 0x17, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, + 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x41, 0x0a, 0x0c, 0x4c, 0x69, 0x73, + 0x74, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x73, 0x12, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, + 0x6c, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, + 0x79, 0x1a, 0x19, 0x2e, 0x61, 0x70, 0x69, 0x2e, 0x4c, 0x69, 0x73, 0x74, 0x4d, 0x61, 0x63, 0x68, + 0x69, 0x6e, 0x65, 0x73, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x42, 0x37, 0x5a, 0x35, + 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x70, 0x73, 0x76, 0x69, 0x64, + 0x65, 0x72, 0x73, 0x6b, 0x69, 0x2f, 0x75, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x69, 0x6e, + 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x2f, 0x61, + 0x70, 0x69, 0x2f, 0x70, 0x62, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, } var ( @@ -221,28 +347,33 @@ func file_internal_machine_api_pb_cluster_proto_rawDescGZIP() []byte { return file_internal_machine_api_pb_cluster_proto_rawDescData } -var file_internal_machine_api_pb_cluster_proto_msgTypes = make([]protoimpl.MessageInfo, 3) +var file_internal_machine_api_pb_cluster_proto_enumTypes = make([]protoimpl.EnumInfo, 1) +var file_internal_machine_api_pb_cluster_proto_msgTypes = make([]protoimpl.MessageInfo, 4) var file_internal_machine_api_pb_cluster_proto_goTypes = []any{ - (*AddMachineRequest)(nil), // 0: api.AddMachineRequest - (*AddMachineResponse)(nil), // 1: api.AddMachineResponse - (*ListMachinesResponse)(nil), // 2: api.ListMachinesResponse - (*NetworkConfig)(nil), // 3: api.NetworkConfig - (*MachineInfo)(nil), // 4: api.MachineInfo - (*emptypb.Empty)(nil), // 5: google.protobuf.Empty + (MachineMember_MembershipState)(0), // 0: api.MachineMember.MembershipState + (*AddMachineRequest)(nil), // 1: api.AddMachineRequest + (*AddMachineResponse)(nil), // 2: api.AddMachineResponse + (*MachineMember)(nil), // 3: api.MachineMember + (*ListMachinesResponse)(nil), // 4: api.ListMachinesResponse + (*NetworkConfig)(nil), // 5: api.NetworkConfig + (*MachineInfo)(nil), // 6: api.MachineInfo + (*emptypb.Empty)(nil), // 7: google.protobuf.Empty } var file_internal_machine_api_pb_cluster_proto_depIdxs = []int32{ - 3, // 0: api.AddMachineRequest.network:type_name -> api.NetworkConfig - 4, // 1: api.AddMachineResponse.machine:type_name -> api.MachineInfo - 4, // 2: api.ListMachinesResponse.machines:type_name -> api.MachineInfo - 0, // 3: api.Cluster.AddMachine:input_type -> api.AddMachineRequest - 5, // 4: api.Cluster.ListMachines:input_type -> google.protobuf.Empty - 1, // 5: api.Cluster.AddMachine:output_type -> api.AddMachineResponse - 2, // 6: api.Cluster.ListMachines:output_type -> api.ListMachinesResponse - 5, // [5:7] is the sub-list for method output_type - 3, // [3:5] is the sub-list for method input_type - 3, // [3:3] is the sub-list for extension type_name - 3, // [3:3] is the sub-list for extension extendee - 0, // [0:3] is the sub-list for field type_name + 5, // 0: api.AddMachineRequest.network:type_name -> api.NetworkConfig + 6, // 1: api.AddMachineResponse.machine:type_name -> api.MachineInfo + 6, // 2: api.MachineMember.machine:type_name -> api.MachineInfo + 0, // 3: api.MachineMember.state:type_name -> api.MachineMember.MembershipState + 3, // 4: api.ListMachinesResponse.machines:type_name -> api.MachineMember + 1, // 5: api.Cluster.AddMachine:input_type -> api.AddMachineRequest + 7, // 6: api.Cluster.ListMachines:input_type -> google.protobuf.Empty + 2, // 7: api.Cluster.AddMachine:output_type -> api.AddMachineResponse + 4, // 8: api.Cluster.ListMachines:output_type -> api.ListMachinesResponse + 7, // [7:9] is the sub-list for method output_type + 5, // [5:7] is the sub-list for method input_type + 5, // [5:5] is the sub-list for extension type_name + 5, // [5:5] is the sub-list for extension extendee + 0, // [0:5] is the sub-list for field type_name } func init() { file_internal_machine_api_pb_cluster_proto_init() } @@ -250,6 +381,7 @@ func file_internal_machine_api_pb_cluster_proto_init() { if File_internal_machine_api_pb_cluster_proto != nil { return } + file_internal_machine_api_pb_common_proto_init() file_internal_machine_api_pb_machine_proto_init() if !protoimpl.UnsafeEnabled { file_internal_machine_api_pb_cluster_proto_msgTypes[0].Exporter = func(v any, i int) any { @@ -277,6 +409,18 @@ func file_internal_machine_api_pb_cluster_proto_init() { } } file_internal_machine_api_pb_cluster_proto_msgTypes[2].Exporter = func(v any, i int) any { + switch v := v.(*MachineMember); i { + case 0: + return &v.state + case 1: + return &v.sizeCache + case 2: + return &v.unknownFields + default: + return nil + } + } + file_internal_machine_api_pb_cluster_proto_msgTypes[3].Exporter = func(v any, i int) any { switch v := v.(*ListMachinesResponse); i { case 0: return &v.state @@ -294,13 +438,14 @@ func file_internal_machine_api_pb_cluster_proto_init() { File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: file_internal_machine_api_pb_cluster_proto_rawDesc, - NumEnums: 0, - NumMessages: 3, + NumEnums: 1, + NumMessages: 4, NumExtensions: 0, NumServices: 1, }, GoTypes: file_internal_machine_api_pb_cluster_proto_goTypes, DependencyIndexes: file_internal_machine_api_pb_cluster_proto_depIdxs, + EnumInfos: file_internal_machine_api_pb_cluster_proto_enumTypes, MessageInfos: file_internal_machine_api_pb_cluster_proto_msgTypes, }.Build() File_internal_machine_api_pb_cluster_proto = out.File diff --git a/internal/machine/api/pb/cluster.proto b/internal/machine/api/pb/cluster.proto index 811cba4d..ff84d7fc 100644 --- a/internal/machine/api/pb/cluster.proto +++ b/internal/machine/api/pb/cluster.proto @@ -21,6 +21,23 @@ message AddMachineResponse { MachineInfo machine = 1; } +message MachineMember { + MachineInfo machine = 1; + + enum MembershipState { + UNKNOWN = 0; + // The member is active. + UP = 1; + // The member is active, but at least one cluster member suspects its down. For all purposes, + // a SUSPECT member is treated as if it were UP until either it refutes the suspicion (becoming UP) + // or fails to do so (being declared DOWN). + SUSPECT = 2; + // The member is confirmed DOWN. + DOWN = 3; + } + MembershipState state = 2; +} + message ListMachinesResponse { - repeated MachineInfo machines = 1; + repeated MachineMember machines = 1; } diff --git a/internal/machine/cluster/cluster.go b/internal/machine/cluster/cluster.go index ab3c5be9..7ad4d035 100644 --- a/internal/machine/cluster/cluster.go +++ b/internal/machine/cluster/cluster.go @@ -11,6 +11,7 @@ import ( "log/slog" "net/netip" "time" + "uncloud/internal/corrosion" "uncloud/internal/machine/api/pb" "uncloud/internal/machine/network" "uncloud/internal/machine/store" @@ -20,11 +21,22 @@ import ( type Cluster struct { pb.UnimplementedClusterServer - store *store.Store + store *store.Store + corroAdmin *corrosion.AdminClient + // machineID is the ID of the current machine that is running the cluster service. + machineID string } -func NewCluster(store *store.Store) *Cluster { - return &Cluster{store: store} +func NewCluster(store *store.Store, corroAdmin *corrosion.AdminClient) *Cluster { + return &Cluster{ + store: store, + corroAdmin: corroAdmin, + } +} + +// UpdateMachineID updates the current machine ID that is running the cluster service. +func (c *Cluster) UpdateMachineID(mid string) { + c.machineID = mid } func (c *Cluster) Init(ctx context.Context, network netip.Prefix) error { @@ -175,6 +187,7 @@ func (c *Cluster) AddMachine(ctx context.Context, req *pb.AddMachineRequest) (*p return resp, nil } +// ListMachines lists all machines in the cluster including their membership states. func (c *Cluster) ListMachines(ctx context.Context, _ *emptypb.Empty) (*pb.ListMachinesResponse, error) { if err := c.checkInitialised(ctx); err != nil { return nil, err @@ -184,5 +197,55 @@ func (c *Cluster) ListMachines(ctx context.Context, _ *emptypb.Empty) (*pb.ListM if err != nil { return nil, status.Error(codes.Internal, err.Error()) } - return &pb.ListMachinesResponse{Machines: machines}, nil + + states, err := c.corroAdmin.ClusterMembershipStates(true) + if err != nil { + return nil, status.Errorf(codes.Internal, "get cluster membership states: %v", err) + } + + members := make([]*pb.MachineMember, len(machines)) + for i, m := range machines { + // If the machine is not in the cluster membership states or its state is not ALIVE or SUSPECT, it is DOWN. + // The exception is the current machine which is always UP as it is serving this request. + state := pb.MachineMember_DOWN + addr, _ := m.Network.ManagementIp.ToAddr() + for _, s := range states { + if s.Addr.Addr().Compare(addr) == 0 && + (s.State == corrosion.MembershipStateAlive || s.State == corrosion.MembershipStateSuspect) { + state = pb.MachineMember_UP + break + } + } + // If the machine is the current machine, it is UP. + if m.Id == c.machineID { + state = pb.MachineMember_UP + } + members[i] = &pb.MachineMember{ + Machine: m, + State: state, + } + } + + return &pb.ListMachinesResponse{Machines: members}, nil } + +//func (c *Cluster) ListServices(ctx context.Context, _ *emptypb.Empty) (*pb.ListServicesResponse, error) { +// if err := c.checkInitialised(ctx); err != nil { +// return nil, err +// } +// +// return &pb.ListServicesResponse{ +// Messages: []*pb.Services{ +// { +// Services: []*pb.Service{ +// { +// Name: "service1", +// }, +// { +// Name: "service2", +// }, +// }, +// }, +// }, +// }, nil +//} diff --git a/internal/machine/machine.go b/internal/machine/machine.go index 9644a309..8f4a7fd4 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -256,7 +256,7 @@ func (m *Machine) Run(ctx context.Context) error { // Signal that the machine is ready. close(m.started) - // Control loop for managing the network controller. + // Control loop for managing components that depend on the machine being initialised as a cluster member. errGroup.Go( func() error { if !m.Initialised() { @@ -276,6 +276,9 @@ func (m *Machine) Run(ctx context.Context) error { // It can be reset when leaving the cluster and then re-initialised again with a new configuration. case <-m.initialised: var err error + + m.cluster.UpdateMachineID(m.state.ID) + // Ensure the corrosion config is up to date, including a new gossip address if the machine // has just joined a cluster. if err = m.configureCorrosion(); err != nil {