diff --git a/seaweed-volume/proto/master.proto b/seaweed-volume/proto/master.proto index f4125452c..0fbfe198a 100644 --- a/seaweed-volume/proto/master.proto +++ b/seaweed-volume/proto/master.proto @@ -45,6 +45,8 @@ service Seaweed { } rpc ReleaseAdminToken (ReleaseAdminTokenRequest) returns (ReleaseAdminTokenResponse) { } + rpc GetAdminLockStatus (GetAdminLockStatusRequest) returns (GetAdminLockStatusResponse) { + } rpc Ping (PingRequest) returns (PingResponse) { } rpc RaftListClusterServers (RaftListClusterServersRequest) returns (RaftListClusterServersResponse) { @@ -464,6 +466,15 @@ message ReleaseAdminTokenRequest { message ReleaseAdminTokenResponse { } +message GetAdminLockStatusRequest { + string lock_name = 1; +} +message GetAdminLockStatusResponse { + bool is_locked = 1; + string client_name = 2; + string message = 3; +} + message PingRequest { string target = 1; // default to ping itself string target_type = 2; diff --git a/weed/pb/master.proto b/weed/pb/master.proto index f4125452c..0fbfe198a 100644 --- a/weed/pb/master.proto +++ b/weed/pb/master.proto @@ -45,6 +45,8 @@ service Seaweed { } rpc ReleaseAdminToken (ReleaseAdminTokenRequest) returns (ReleaseAdminTokenResponse) { } + rpc GetAdminLockStatus (GetAdminLockStatusRequest) returns (GetAdminLockStatusResponse) { + } rpc Ping (PingRequest) returns (PingResponse) { } rpc RaftListClusterServers (RaftListClusterServersRequest) returns (RaftListClusterServersResponse) { @@ -464,6 +466,15 @@ message ReleaseAdminTokenRequest { message ReleaseAdminTokenResponse { } +message GetAdminLockStatusRequest { + string lock_name = 1; +} +message GetAdminLockStatusResponse { + bool is_locked = 1; + string client_name = 2; + string message = 3; +} + message PingRequest { string target = 1; // default to ping itself string target_type = 2; diff --git a/weed/pb/master_pb/master.pb.go b/weed/pb/master_pb/master.pb.go index 641595d33..955d2397f 100644 --- a/weed/pb/master_pb/master.pb.go +++ b/weed/pb/master_pb/master.pb.go @@ -3597,6 +3597,110 @@ func (*ReleaseAdminTokenResponse) Descriptor() ([]byte, []int) { return file_master_proto_rawDescGZIP(), []int{51} } +type GetAdminLockStatusRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + LockName string `protobuf:"bytes,1,opt,name=lock_name,json=lockName,proto3" json:"lock_name,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetAdminLockStatusRequest) Reset() { + *x = GetAdminLockStatusRequest{} + mi := &file_master_proto_msgTypes[52] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetAdminLockStatusRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetAdminLockStatusRequest) ProtoMessage() {} + +func (x *GetAdminLockStatusRequest) ProtoReflect() protoreflect.Message { + mi := &file_master_proto_msgTypes[52] + 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 GetAdminLockStatusRequest.ProtoReflect.Descriptor instead. +func (*GetAdminLockStatusRequest) Descriptor() ([]byte, []int) { + return file_master_proto_rawDescGZIP(), []int{52} +} + +func (x *GetAdminLockStatusRequest) GetLockName() string { + if x != nil { + return x.LockName + } + return "" +} + +type GetAdminLockStatusResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + IsLocked bool `protobuf:"varint,1,opt,name=is_locked,json=isLocked,proto3" json:"is_locked,omitempty"` + ClientName string `protobuf:"bytes,2,opt,name=client_name,json=clientName,proto3" json:"client_name,omitempty"` + Message string `protobuf:"bytes,3,opt,name=message,proto3" json:"message,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetAdminLockStatusResponse) Reset() { + *x = GetAdminLockStatusResponse{} + mi := &file_master_proto_msgTypes[53] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetAdminLockStatusResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetAdminLockStatusResponse) ProtoMessage() {} + +func (x *GetAdminLockStatusResponse) ProtoReflect() protoreflect.Message { + mi := &file_master_proto_msgTypes[53] + 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 GetAdminLockStatusResponse.ProtoReflect.Descriptor instead. +func (*GetAdminLockStatusResponse) Descriptor() ([]byte, []int) { + return file_master_proto_rawDescGZIP(), []int{53} +} + +func (x *GetAdminLockStatusResponse) GetIsLocked() bool { + if x != nil { + return x.IsLocked + } + return false +} + +func (x *GetAdminLockStatusResponse) GetClientName() string { + if x != nil { + return x.ClientName + } + return "" +} + +func (x *GetAdminLockStatusResponse) GetMessage() string { + if x != nil { + return x.Message + } + return "" +} + type PingRequest struct { state protoimpl.MessageState `protogen:"open.v1"` Target string `protobuf:"bytes,1,opt,name=target,proto3" json:"target,omitempty"` // default to ping itself @@ -3607,7 +3711,7 @@ type PingRequest struct { func (x *PingRequest) Reset() { *x = PingRequest{} - mi := &file_master_proto_msgTypes[52] + mi := &file_master_proto_msgTypes[54] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3619,7 +3723,7 @@ func (x *PingRequest) String() string { func (*PingRequest) ProtoMessage() {} func (x *PingRequest) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[52] + mi := &file_master_proto_msgTypes[54] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3632,7 +3736,7 @@ func (x *PingRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use PingRequest.ProtoReflect.Descriptor instead. func (*PingRequest) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{52} + return file_master_proto_rawDescGZIP(), []int{54} } func (x *PingRequest) GetTarget() string { @@ -3660,7 +3764,7 @@ type PingResponse struct { func (x *PingResponse) Reset() { *x = PingResponse{} - mi := &file_master_proto_msgTypes[53] + mi := &file_master_proto_msgTypes[55] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3672,7 +3776,7 @@ func (x *PingResponse) String() string { func (*PingResponse) ProtoMessage() {} func (x *PingResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[53] + mi := &file_master_proto_msgTypes[55] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3685,7 +3789,7 @@ func (x *PingResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use PingResponse.ProtoReflect.Descriptor instead. func (*PingResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{53} + return file_master_proto_rawDescGZIP(), []int{55} } func (x *PingResponse) GetStartTimeNs() int64 { @@ -3720,7 +3824,7 @@ type RaftAddServerRequest struct { func (x *RaftAddServerRequest) Reset() { *x = RaftAddServerRequest{} - mi := &file_master_proto_msgTypes[54] + mi := &file_master_proto_msgTypes[56] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3732,7 +3836,7 @@ func (x *RaftAddServerRequest) String() string { func (*RaftAddServerRequest) ProtoMessage() {} func (x *RaftAddServerRequest) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[54] + mi := &file_master_proto_msgTypes[56] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3745,7 +3849,7 @@ func (x *RaftAddServerRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftAddServerRequest.ProtoReflect.Descriptor instead. func (*RaftAddServerRequest) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{54} + return file_master_proto_rawDescGZIP(), []int{56} } func (x *RaftAddServerRequest) GetId() string { @@ -3777,7 +3881,7 @@ type RaftAddServerResponse struct { func (x *RaftAddServerResponse) Reset() { *x = RaftAddServerResponse{} - mi := &file_master_proto_msgTypes[55] + mi := &file_master_proto_msgTypes[57] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3789,7 +3893,7 @@ func (x *RaftAddServerResponse) String() string { func (*RaftAddServerResponse) ProtoMessage() {} func (x *RaftAddServerResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[55] + mi := &file_master_proto_msgTypes[57] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3802,7 +3906,7 @@ func (x *RaftAddServerResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftAddServerResponse.ProtoReflect.Descriptor instead. func (*RaftAddServerResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{55} + return file_master_proto_rawDescGZIP(), []int{57} } type RaftRemoveServerRequest struct { @@ -3815,7 +3919,7 @@ type RaftRemoveServerRequest struct { func (x *RaftRemoveServerRequest) Reset() { *x = RaftRemoveServerRequest{} - mi := &file_master_proto_msgTypes[56] + mi := &file_master_proto_msgTypes[58] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3827,7 +3931,7 @@ func (x *RaftRemoveServerRequest) String() string { func (*RaftRemoveServerRequest) ProtoMessage() {} func (x *RaftRemoveServerRequest) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[56] + mi := &file_master_proto_msgTypes[58] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3840,7 +3944,7 @@ func (x *RaftRemoveServerRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftRemoveServerRequest.ProtoReflect.Descriptor instead. func (*RaftRemoveServerRequest) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{56} + return file_master_proto_rawDescGZIP(), []int{58} } func (x *RaftRemoveServerRequest) GetId() string { @@ -3865,7 +3969,7 @@ type RaftRemoveServerResponse struct { func (x *RaftRemoveServerResponse) Reset() { *x = RaftRemoveServerResponse{} - mi := &file_master_proto_msgTypes[57] + mi := &file_master_proto_msgTypes[59] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3877,7 +3981,7 @@ func (x *RaftRemoveServerResponse) String() string { func (*RaftRemoveServerResponse) ProtoMessage() {} func (x *RaftRemoveServerResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[57] + mi := &file_master_proto_msgTypes[59] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3890,7 +3994,7 @@ func (x *RaftRemoveServerResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftRemoveServerResponse.ProtoReflect.Descriptor instead. func (*RaftRemoveServerResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{57} + return file_master_proto_rawDescGZIP(), []int{59} } type RaftListClusterServersRequest struct { @@ -3901,7 +4005,7 @@ type RaftListClusterServersRequest struct { func (x *RaftListClusterServersRequest) Reset() { *x = RaftListClusterServersRequest{} - mi := &file_master_proto_msgTypes[58] + mi := &file_master_proto_msgTypes[60] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3913,7 +4017,7 @@ func (x *RaftListClusterServersRequest) String() string { func (*RaftListClusterServersRequest) ProtoMessage() {} func (x *RaftListClusterServersRequest) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[58] + mi := &file_master_proto_msgTypes[60] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3926,7 +4030,7 @@ func (x *RaftListClusterServersRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftListClusterServersRequest.ProtoReflect.Descriptor instead. func (*RaftListClusterServersRequest) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{58} + return file_master_proto_rawDescGZIP(), []int{60} } type RaftListClusterServersResponse struct { @@ -3938,7 +4042,7 @@ type RaftListClusterServersResponse struct { func (x *RaftListClusterServersResponse) Reset() { *x = RaftListClusterServersResponse{} - mi := &file_master_proto_msgTypes[59] + mi := &file_master_proto_msgTypes[61] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3950,7 +4054,7 @@ func (x *RaftListClusterServersResponse) String() string { func (*RaftListClusterServersResponse) ProtoMessage() {} func (x *RaftListClusterServersResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[59] + mi := &file_master_proto_msgTypes[61] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3963,7 +4067,7 @@ func (x *RaftListClusterServersResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftListClusterServersResponse.ProtoReflect.Descriptor instead. func (*RaftListClusterServersResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{59} + return file_master_proto_rawDescGZIP(), []int{61} } func (x *RaftListClusterServersResponse) GetClusterServers() []*RaftListClusterServersResponse_ClusterServers { @@ -3983,7 +4087,7 @@ type RaftLeadershipTransferRequest struct { func (x *RaftLeadershipTransferRequest) Reset() { *x = RaftLeadershipTransferRequest{} - mi := &file_master_proto_msgTypes[60] + mi := &file_master_proto_msgTypes[62] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3995,7 +4099,7 @@ func (x *RaftLeadershipTransferRequest) String() string { func (*RaftLeadershipTransferRequest) ProtoMessage() {} func (x *RaftLeadershipTransferRequest) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[60] + mi := &file_master_proto_msgTypes[62] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4008,7 +4112,7 @@ func (x *RaftLeadershipTransferRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftLeadershipTransferRequest.ProtoReflect.Descriptor instead. func (*RaftLeadershipTransferRequest) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{60} + return file_master_proto_rawDescGZIP(), []int{62} } func (x *RaftLeadershipTransferRequest) GetTargetId() string { @@ -4035,7 +4139,7 @@ type RaftLeadershipTransferResponse struct { func (x *RaftLeadershipTransferResponse) Reset() { *x = RaftLeadershipTransferResponse{} - mi := &file_master_proto_msgTypes[61] + mi := &file_master_proto_msgTypes[63] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4047,7 +4151,7 @@ func (x *RaftLeadershipTransferResponse) String() string { func (*RaftLeadershipTransferResponse) ProtoMessage() {} func (x *RaftLeadershipTransferResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[61] + mi := &file_master_proto_msgTypes[63] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4060,7 +4164,7 @@ func (x *RaftLeadershipTransferResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use RaftLeadershipTransferResponse.ProtoReflect.Descriptor instead. func (*RaftLeadershipTransferResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{61} + return file_master_proto_rawDescGZIP(), []int{63} } func (x *RaftLeadershipTransferResponse) GetPreviousLeader() string { @@ -4085,7 +4189,7 @@ type VolumeGrowResponse struct { func (x *VolumeGrowResponse) Reset() { *x = VolumeGrowResponse{} - mi := &file_master_proto_msgTypes[62] + mi := &file_master_proto_msgTypes[64] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4097,7 +4201,7 @@ func (x *VolumeGrowResponse) String() string { func (*VolumeGrowResponse) ProtoMessage() {} func (x *VolumeGrowResponse) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[62] + mi := &file_master_proto_msgTypes[64] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4110,7 +4214,7 @@ func (x *VolumeGrowResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use VolumeGrowResponse.ProtoReflect.Descriptor instead. func (*VolumeGrowResponse) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{62} + return file_master_proto_rawDescGZIP(), []int{64} } type SuperBlockExtra_ErasureCoding struct { @@ -4124,7 +4228,7 @@ type SuperBlockExtra_ErasureCoding struct { func (x *SuperBlockExtra_ErasureCoding) Reset() { *x = SuperBlockExtra_ErasureCoding{} - mi := &file_master_proto_msgTypes[67] + mi := &file_master_proto_msgTypes[69] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4136,7 +4240,7 @@ func (x *SuperBlockExtra_ErasureCoding) String() string { func (*SuperBlockExtra_ErasureCoding) ProtoMessage() {} func (x *SuperBlockExtra_ErasureCoding) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[67] + mi := &file_master_proto_msgTypes[69] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4185,7 +4289,7 @@ type LookupVolumeResponse_VolumeIdLocation struct { func (x *LookupVolumeResponse_VolumeIdLocation) Reset() { *x = LookupVolumeResponse_VolumeIdLocation{} - mi := &file_master_proto_msgTypes[68] + mi := &file_master_proto_msgTypes[70] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4197,7 +4301,7 @@ func (x *LookupVolumeResponse_VolumeIdLocation) String() string { func (*LookupVolumeResponse_VolumeIdLocation) ProtoMessage() {} func (x *LookupVolumeResponse_VolumeIdLocation) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[68] + mi := &file_master_proto_msgTypes[70] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4251,7 +4355,7 @@ type LookupEcVolumeResponse_EcShardIdLocation struct { func (x *LookupEcVolumeResponse_EcShardIdLocation) Reset() { *x = LookupEcVolumeResponse_EcShardIdLocation{} - mi := &file_master_proto_msgTypes[74] + mi := &file_master_proto_msgTypes[76] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4263,7 +4367,7 @@ func (x *LookupEcVolumeResponse_EcShardIdLocation) String() string { func (*LookupEcVolumeResponse_EcShardIdLocation) ProtoMessage() {} func (x *LookupEcVolumeResponse_EcShardIdLocation) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[74] + mi := &file_master_proto_msgTypes[76] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4306,7 +4410,7 @@ type ListClusterNodesResponse_ClusterNode struct { func (x *ListClusterNodesResponse_ClusterNode) Reset() { *x = ListClusterNodesResponse_ClusterNode{} - mi := &file_master_proto_msgTypes[75] + mi := &file_master_proto_msgTypes[77] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4318,7 +4422,7 @@ func (x *ListClusterNodesResponse_ClusterNode) String() string { func (*ListClusterNodesResponse_ClusterNode) ProtoMessage() {} func (x *ListClusterNodesResponse_ClusterNode) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[75] + mi := &file_master_proto_msgTypes[77] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4381,7 +4485,7 @@ type RaftListClusterServersResponse_ClusterServers struct { func (x *RaftListClusterServersResponse_ClusterServers) Reset() { *x = RaftListClusterServersResponse_ClusterServers{} - mi := &file_master_proto_msgTypes[76] + mi := &file_master_proto_msgTypes[78] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -4393,7 +4497,7 @@ func (x *RaftListClusterServersResponse_ClusterServers) String() string { func (*RaftListClusterServersResponse_ClusterServers) ProtoMessage() {} func (x *RaftListClusterServersResponse_ClusterServers) ProtoReflect() protoreflect.Message { - mi := &file_master_proto_msgTypes[76] + mi := &file_master_proto_msgTypes[78] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -4406,7 +4510,7 @@ func (x *RaftListClusterServersResponse_ClusterServers) ProtoReflect() protorefl // Deprecated: Use RaftListClusterServersResponse_ClusterServers.ProtoReflect.Descriptor instead. func (*RaftListClusterServersResponse_ClusterServers) Descriptor() ([]byte, []int) { - return file_master_proto_rawDescGZIP(), []int{59, 0} + return file_master_proto_rawDescGZIP(), []int{61, 0} } func (x *RaftListClusterServersResponse_ClusterServers) GetId() string { @@ -4806,7 +4910,14 @@ const file_master_proto_rawDesc = "" + "\x0eprevious_token\x18\x01 \x01(\x03R\rpreviousToken\x12,\n" + "\x12previous_lock_time\x18\x02 \x01(\x03R\x10previousLockTime\x12\x1b\n" + "\tlock_name\x18\x03 \x01(\tR\blockName\"\x1b\n" + - "\x19ReleaseAdminTokenResponse\"F\n" + + "\x19ReleaseAdminTokenResponse\"8\n" + + "\x19GetAdminLockStatusRequest\x12\x1b\n" + + "\tlock_name\x18\x01 \x01(\tR\blockName\"t\n" + + "\x1aGetAdminLockStatusResponse\x12\x1b\n" + + "\tis_locked\x18\x01 \x01(\bR\bisLocked\x12\x1f\n" + + "\vclient_name\x18\x02 \x01(\tR\n" + + "clientName\x12\x18\n" + + "\amessage\x18\x03 \x01(\tR\amessage\"F\n" + "\vPingRequest\x12\x16\n" + "\x06target\x18\x01 \x01(\tR\x06target\x12\x1f\n" + "\vtarget_type\x18\x02 \x01(\tR\n" + @@ -4840,7 +4951,7 @@ const file_master_proto_rawDesc = "" + "\x0fprevious_leader\x18\x01 \x01(\tR\x0epreviousLeader\x12\x1d\n" + "\n" + "new_leader\x18\x02 \x01(\tR\tnewLeader\"\x14\n" + - "\x12VolumeGrowResponse2\xc6\x10\n" + + "\x12VolumeGrowResponse2\xab\x11\n" + "\aSeaweed\x12I\n" + "\rSendHeartbeat\x12\x14.master_pb.Heartbeat\x1a\x1c.master_pb.HeartbeatResponse\"\x00(\x010\x01\x12X\n" + "\rKeepConnected\x12\x1f.master_pb.KeepConnectedRequest\x1a .master_pb.KeepConnectedResponse\"\x00(\x010\x01\x12Q\n" + @@ -4861,7 +4972,8 @@ const file_master_proto_rawDesc = "" + "\x16GetMasterConfiguration\x12(.master_pb.GetMasterConfigurationRequest\x1a).master_pb.GetMasterConfigurationResponse\"\x00\x12]\n" + "\x10ListClusterNodes\x12\".master_pb.ListClusterNodesRequest\x1a#.master_pb.ListClusterNodesResponse\"\x00\x12Z\n" + "\x0fLeaseAdminToken\x12!.master_pb.LeaseAdminTokenRequest\x1a\".master_pb.LeaseAdminTokenResponse\"\x00\x12`\n" + - "\x11ReleaseAdminToken\x12#.master_pb.ReleaseAdminTokenRequest\x1a$.master_pb.ReleaseAdminTokenResponse\"\x00\x129\n" + + "\x11ReleaseAdminToken\x12#.master_pb.ReleaseAdminTokenRequest\x1a$.master_pb.ReleaseAdminTokenResponse\"\x00\x12c\n" + + "\x12GetAdminLockStatus\x12$.master_pb.GetAdminLockStatusRequest\x1a%.master_pb.GetAdminLockStatusResponse\"\x00\x129\n" + "\x04Ping\x12\x16.master_pb.PingRequest\x1a\x17.master_pb.PingResponse\"\x00\x12o\n" + "\x16RaftListClusterServers\x12(.master_pb.RaftListClusterServersRequest\x1a).master_pb.RaftListClusterServersResponse\"\x00\x12T\n" + "\rRaftAddServer\x12\x1f.master_pb.RaftAddServerRequest\x1a .master_pb.RaftAddServerResponse\"\x00\x12]\n" + @@ -4882,7 +4994,7 @@ func file_master_proto_rawDescGZIP() []byte { return file_master_proto_rawDescData } -var file_master_proto_msgTypes = make([]protoimpl.MessageInfo, 77) +var file_master_proto_msgTypes = make([]protoimpl.MessageInfo, 79) var file_master_proto_goTypes = []any{ (*DiskTag)(nil), // 0: master_pb.DiskTag (*Heartbeat)(nil), // 1: master_pb.Heartbeat @@ -4936,32 +5048,34 @@ var file_master_proto_goTypes = []any{ (*LeaseAdminTokenResponse)(nil), // 49: master_pb.LeaseAdminTokenResponse (*ReleaseAdminTokenRequest)(nil), // 50: master_pb.ReleaseAdminTokenRequest (*ReleaseAdminTokenResponse)(nil), // 51: master_pb.ReleaseAdminTokenResponse - (*PingRequest)(nil), // 52: master_pb.PingRequest - (*PingResponse)(nil), // 53: master_pb.PingResponse - (*RaftAddServerRequest)(nil), // 54: master_pb.RaftAddServerRequest - (*RaftAddServerResponse)(nil), // 55: master_pb.RaftAddServerResponse - (*RaftRemoveServerRequest)(nil), // 56: master_pb.RaftRemoveServerRequest - (*RaftRemoveServerResponse)(nil), // 57: master_pb.RaftRemoveServerResponse - (*RaftListClusterServersRequest)(nil), // 58: master_pb.RaftListClusterServersRequest - (*RaftListClusterServersResponse)(nil), // 59: master_pb.RaftListClusterServersResponse - (*RaftLeadershipTransferRequest)(nil), // 60: master_pb.RaftLeadershipTransferRequest - (*RaftLeadershipTransferResponse)(nil), // 61: master_pb.RaftLeadershipTransferResponse - (*VolumeGrowResponse)(nil), // 62: master_pb.VolumeGrowResponse - nil, // 63: master_pb.Heartbeat.MaxVolumeCountsEntry - nil, // 64: master_pb.Heartbeat.DiskTotalBytesEntry - nil, // 65: master_pb.Heartbeat.DiskFreeBytesEntry - nil, // 66: master_pb.StorageBackend.PropertiesEntry - (*SuperBlockExtra_ErasureCoding)(nil), // 67: master_pb.SuperBlockExtra.ErasureCoding - (*LookupVolumeResponse_VolumeIdLocation)(nil), // 68: master_pb.LookupVolumeResponse.VolumeIdLocation - nil, // 69: master_pb.DiskInfo.MaxVolumeCountByDiskEntry - nil, // 70: master_pb.DataNodeInfo.DiskInfosEntry - nil, // 71: master_pb.RackInfo.DiskInfosEntry - nil, // 72: master_pb.DataCenterInfo.DiskInfosEntry - nil, // 73: master_pb.TopologyInfo.DiskInfosEntry - (*LookupEcVolumeResponse_EcShardIdLocation)(nil), // 74: master_pb.LookupEcVolumeResponse.EcShardIdLocation - (*ListClusterNodesResponse_ClusterNode)(nil), // 75: master_pb.ListClusterNodesResponse.ClusterNode - (*RaftListClusterServersResponse_ClusterServers)(nil), // 76: master_pb.RaftListClusterServersResponse.ClusterServers - (*volume_server_pb.VolumeServerState)(nil), // 77: volume_server_pb.VolumeServerState + (*GetAdminLockStatusRequest)(nil), // 52: master_pb.GetAdminLockStatusRequest + (*GetAdminLockStatusResponse)(nil), // 53: master_pb.GetAdminLockStatusResponse + (*PingRequest)(nil), // 54: master_pb.PingRequest + (*PingResponse)(nil), // 55: master_pb.PingResponse + (*RaftAddServerRequest)(nil), // 56: master_pb.RaftAddServerRequest + (*RaftAddServerResponse)(nil), // 57: master_pb.RaftAddServerResponse + (*RaftRemoveServerRequest)(nil), // 58: master_pb.RaftRemoveServerRequest + (*RaftRemoveServerResponse)(nil), // 59: master_pb.RaftRemoveServerResponse + (*RaftListClusterServersRequest)(nil), // 60: master_pb.RaftListClusterServersRequest + (*RaftListClusterServersResponse)(nil), // 61: master_pb.RaftListClusterServersResponse + (*RaftLeadershipTransferRequest)(nil), // 62: master_pb.RaftLeadershipTransferRequest + (*RaftLeadershipTransferResponse)(nil), // 63: master_pb.RaftLeadershipTransferResponse + (*VolumeGrowResponse)(nil), // 64: master_pb.VolumeGrowResponse + nil, // 65: master_pb.Heartbeat.MaxVolumeCountsEntry + nil, // 66: master_pb.Heartbeat.DiskTotalBytesEntry + nil, // 67: master_pb.Heartbeat.DiskFreeBytesEntry + nil, // 68: master_pb.StorageBackend.PropertiesEntry + (*SuperBlockExtra_ErasureCoding)(nil), // 69: master_pb.SuperBlockExtra.ErasureCoding + (*LookupVolumeResponse_VolumeIdLocation)(nil), // 70: master_pb.LookupVolumeResponse.VolumeIdLocation + nil, // 71: master_pb.DiskInfo.MaxVolumeCountByDiskEntry + nil, // 72: master_pb.DataNodeInfo.DiskInfosEntry + nil, // 73: master_pb.RackInfo.DiskInfosEntry + nil, // 74: master_pb.DataCenterInfo.DiskInfosEntry + nil, // 75: master_pb.TopologyInfo.DiskInfosEntry + (*LookupEcVolumeResponse_EcShardIdLocation)(nil), // 76: master_pb.LookupEcVolumeResponse.EcShardIdLocation + (*ListClusterNodesResponse_ClusterNode)(nil), // 77: master_pb.ListClusterNodesResponse.ClusterNode + (*RaftListClusterServersResponse_ClusterServers)(nil), // 78: master_pb.RaftListClusterServersResponse.ClusterServers + (*volume_server_pb.VolumeServerState)(nil), // 79: volume_server_pb.VolumeServerState } var file_master_proto_depIdxs = []int32{ 3, // 0: master_pb.Heartbeat.volumes:type_name -> master_pb.VolumeInformationMessage @@ -4970,36 +5084,36 @@ var file_master_proto_depIdxs = []int32{ 5, // 3: master_pb.Heartbeat.ec_shards:type_name -> master_pb.VolumeEcShardInformationMessage 5, // 4: master_pb.Heartbeat.new_ec_shards:type_name -> master_pb.VolumeEcShardInformationMessage 5, // 5: master_pb.Heartbeat.deleted_ec_shards:type_name -> master_pb.VolumeEcShardInformationMessage - 63, // 6: master_pb.Heartbeat.max_volume_counts:type_name -> master_pb.Heartbeat.MaxVolumeCountsEntry - 77, // 7: master_pb.Heartbeat.state:type_name -> volume_server_pb.VolumeServerState + 65, // 6: master_pb.Heartbeat.max_volume_counts:type_name -> master_pb.Heartbeat.MaxVolumeCountsEntry + 79, // 7: master_pb.Heartbeat.state:type_name -> volume_server_pb.VolumeServerState 0, // 8: master_pb.Heartbeat.disk_tags:type_name -> master_pb.DiskTag - 64, // 9: master_pb.Heartbeat.disk_total_bytes:type_name -> master_pb.Heartbeat.DiskTotalBytesEntry - 65, // 10: master_pb.Heartbeat.disk_free_bytes:type_name -> master_pb.Heartbeat.DiskFreeBytesEntry + 66, // 9: master_pb.Heartbeat.disk_total_bytes:type_name -> master_pb.Heartbeat.DiskTotalBytesEntry + 67, // 10: master_pb.Heartbeat.disk_free_bytes:type_name -> master_pb.Heartbeat.DiskFreeBytesEntry 6, // 11: master_pb.HeartbeatResponse.storage_backends:type_name -> master_pb.StorageBackend - 66, // 12: master_pb.StorageBackend.properties:type_name -> master_pb.StorageBackend.PropertiesEntry - 67, // 13: master_pb.SuperBlockExtra.erasure_coding:type_name -> master_pb.SuperBlockExtra.ErasureCoding + 68, // 12: master_pb.StorageBackend.properties:type_name -> master_pb.StorageBackend.PropertiesEntry + 69, // 13: master_pb.SuperBlockExtra.erasure_coding:type_name -> master_pb.SuperBlockExtra.ErasureCoding 10, // 14: master_pb.KeepConnectedResponse.volume_location:type_name -> master_pb.VolumeLocation 11, // 15: master_pb.KeepConnectedResponse.cluster_node_update:type_name -> master_pb.ClusterNodeUpdate 13, // 16: master_pb.KeepConnectedResponse.lock_ring_update:type_name -> master_pb.LockRingUpdate - 68, // 17: master_pb.LookupVolumeResponse.volume_id_locations:type_name -> master_pb.LookupVolumeResponse.VolumeIdLocation + 70, // 17: master_pb.LookupVolumeResponse.volume_id_locations:type_name -> master_pb.LookupVolumeResponse.VolumeIdLocation 16, // 18: master_pb.AssignResponse.replicas:type_name -> master_pb.Location 16, // 19: master_pb.AssignResponse.location:type_name -> master_pb.Location 22, // 20: master_pb.CollectionListResponse.collections:type_name -> master_pb.Collection 3, // 21: master_pb.DiskInfo.volume_infos:type_name -> master_pb.VolumeInformationMessage 5, // 22: master_pb.DiskInfo.ec_shard_infos:type_name -> master_pb.VolumeEcShardInformationMessage - 69, // 23: master_pb.DiskInfo.max_volume_count_by_disk:type_name -> master_pb.DiskInfo.MaxVolumeCountByDiskEntry - 70, // 24: master_pb.DataNodeInfo.diskInfos:type_name -> master_pb.DataNodeInfo.DiskInfosEntry + 71, // 23: master_pb.DiskInfo.max_volume_count_by_disk:type_name -> master_pb.DiskInfo.MaxVolumeCountByDiskEntry + 72, // 24: master_pb.DataNodeInfo.diskInfos:type_name -> master_pb.DataNodeInfo.DiskInfosEntry 28, // 25: master_pb.RackInfo.data_node_infos:type_name -> master_pb.DataNodeInfo - 71, // 26: master_pb.RackInfo.diskInfos:type_name -> master_pb.RackInfo.DiskInfosEntry + 73, // 26: master_pb.RackInfo.diskInfos:type_name -> master_pb.RackInfo.DiskInfosEntry 29, // 27: master_pb.DataCenterInfo.rack_infos:type_name -> master_pb.RackInfo - 72, // 28: master_pb.DataCenterInfo.diskInfos:type_name -> master_pb.DataCenterInfo.DiskInfosEntry + 74, // 28: master_pb.DataCenterInfo.diskInfos:type_name -> master_pb.DataCenterInfo.DiskInfosEntry 30, // 29: master_pb.TopologyInfo.data_center_infos:type_name -> master_pb.DataCenterInfo - 73, // 30: master_pb.TopologyInfo.diskInfos:type_name -> master_pb.TopologyInfo.DiskInfosEntry + 75, // 30: master_pb.TopologyInfo.diskInfos:type_name -> master_pb.TopologyInfo.DiskInfosEntry 31, // 31: master_pb.VolumeListResponse.topology_info:type_name -> master_pb.TopologyInfo - 74, // 32: master_pb.LookupEcVolumeResponse.shard_id_locations:type_name -> master_pb.LookupEcVolumeResponse.EcShardIdLocation + 76, // 32: master_pb.LookupEcVolumeResponse.shard_id_locations:type_name -> master_pb.LookupEcVolumeResponse.EcShardIdLocation 6, // 33: master_pb.GetMasterConfigurationResponse.storage_backends:type_name -> master_pb.StorageBackend - 75, // 34: master_pb.ListClusterNodesResponse.cluster_nodes:type_name -> master_pb.ListClusterNodesResponse.ClusterNode - 76, // 35: master_pb.RaftListClusterServersResponse.cluster_servers:type_name -> master_pb.RaftListClusterServersResponse.ClusterServers + 77, // 34: master_pb.ListClusterNodesResponse.cluster_nodes:type_name -> master_pb.ListClusterNodesResponse.ClusterNode + 78, // 35: master_pb.RaftListClusterServersResponse.cluster_servers:type_name -> master_pb.RaftListClusterServersResponse.ClusterServers 16, // 36: master_pb.LookupVolumeResponse.VolumeIdLocation.locations:type_name -> master_pb.Location 27, // 37: master_pb.DataNodeInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo 27, // 38: master_pb.RackInfo.DiskInfosEntry.value:type_name -> master_pb.DiskInfo @@ -5024,38 +5138,40 @@ var file_master_proto_depIdxs = []int32{ 46, // 57: master_pb.Seaweed.ListClusterNodes:input_type -> master_pb.ListClusterNodesRequest 48, // 58: master_pb.Seaweed.LeaseAdminToken:input_type -> master_pb.LeaseAdminTokenRequest 50, // 59: master_pb.Seaweed.ReleaseAdminToken:input_type -> master_pb.ReleaseAdminTokenRequest - 52, // 60: master_pb.Seaweed.Ping:input_type -> master_pb.PingRequest - 58, // 61: master_pb.Seaweed.RaftListClusterServers:input_type -> master_pb.RaftListClusterServersRequest - 54, // 62: master_pb.Seaweed.RaftAddServer:input_type -> master_pb.RaftAddServerRequest - 56, // 63: master_pb.Seaweed.RaftRemoveServer:input_type -> master_pb.RaftRemoveServerRequest - 60, // 64: master_pb.Seaweed.RaftLeadershipTransfer:input_type -> master_pb.RaftLeadershipTransferRequest - 18, // 65: master_pb.Seaweed.VolumeGrow:input_type -> master_pb.VolumeGrowRequest - 2, // 66: master_pb.Seaweed.SendHeartbeat:output_type -> master_pb.HeartbeatResponse - 12, // 67: master_pb.Seaweed.KeepConnected:output_type -> master_pb.KeepConnectedResponse - 15, // 68: master_pb.Seaweed.LookupVolume:output_type -> master_pb.LookupVolumeResponse - 19, // 69: master_pb.Seaweed.Assign:output_type -> master_pb.AssignResponse - 19, // 70: master_pb.Seaweed.StreamAssign:output_type -> master_pb.AssignResponse - 21, // 71: master_pb.Seaweed.Statistics:output_type -> master_pb.StatisticsResponse - 24, // 72: master_pb.Seaweed.CollectionList:output_type -> master_pb.CollectionListResponse - 26, // 73: master_pb.Seaweed.CollectionDelete:output_type -> master_pb.CollectionDeleteResponse - 33, // 74: master_pb.Seaweed.VolumeList:output_type -> master_pb.VolumeListResponse - 35, // 75: master_pb.Seaweed.LookupEcVolume:output_type -> master_pb.LookupEcVolumeResponse - 37, // 76: master_pb.Seaweed.VacuumVolume:output_type -> master_pb.VacuumVolumeResponse - 39, // 77: master_pb.Seaweed.DisableVacuum:output_type -> master_pb.DisableVacuumResponse - 41, // 78: master_pb.Seaweed.EnableVacuum:output_type -> master_pb.EnableVacuumResponse - 43, // 79: master_pb.Seaweed.VolumeMarkReadonly:output_type -> master_pb.VolumeMarkReadonlyResponse - 45, // 80: master_pb.Seaweed.GetMasterConfiguration:output_type -> master_pb.GetMasterConfigurationResponse - 47, // 81: master_pb.Seaweed.ListClusterNodes:output_type -> master_pb.ListClusterNodesResponse - 49, // 82: master_pb.Seaweed.LeaseAdminToken:output_type -> master_pb.LeaseAdminTokenResponse - 51, // 83: master_pb.Seaweed.ReleaseAdminToken:output_type -> master_pb.ReleaseAdminTokenResponse - 53, // 84: master_pb.Seaweed.Ping:output_type -> master_pb.PingResponse - 59, // 85: master_pb.Seaweed.RaftListClusterServers:output_type -> master_pb.RaftListClusterServersResponse - 55, // 86: master_pb.Seaweed.RaftAddServer:output_type -> master_pb.RaftAddServerResponse - 57, // 87: master_pb.Seaweed.RaftRemoveServer:output_type -> master_pb.RaftRemoveServerResponse - 61, // 88: master_pb.Seaweed.RaftLeadershipTransfer:output_type -> master_pb.RaftLeadershipTransferResponse - 62, // 89: master_pb.Seaweed.VolumeGrow:output_type -> master_pb.VolumeGrowResponse - 66, // [66:90] is the sub-list for method output_type - 42, // [42:66] is the sub-list for method input_type + 52, // 60: master_pb.Seaweed.GetAdminLockStatus:input_type -> master_pb.GetAdminLockStatusRequest + 54, // 61: master_pb.Seaweed.Ping:input_type -> master_pb.PingRequest + 60, // 62: master_pb.Seaweed.RaftListClusterServers:input_type -> master_pb.RaftListClusterServersRequest + 56, // 63: master_pb.Seaweed.RaftAddServer:input_type -> master_pb.RaftAddServerRequest + 58, // 64: master_pb.Seaweed.RaftRemoveServer:input_type -> master_pb.RaftRemoveServerRequest + 62, // 65: master_pb.Seaweed.RaftLeadershipTransfer:input_type -> master_pb.RaftLeadershipTransferRequest + 18, // 66: master_pb.Seaweed.VolumeGrow:input_type -> master_pb.VolumeGrowRequest + 2, // 67: master_pb.Seaweed.SendHeartbeat:output_type -> master_pb.HeartbeatResponse + 12, // 68: master_pb.Seaweed.KeepConnected:output_type -> master_pb.KeepConnectedResponse + 15, // 69: master_pb.Seaweed.LookupVolume:output_type -> master_pb.LookupVolumeResponse + 19, // 70: master_pb.Seaweed.Assign:output_type -> master_pb.AssignResponse + 19, // 71: master_pb.Seaweed.StreamAssign:output_type -> master_pb.AssignResponse + 21, // 72: master_pb.Seaweed.Statistics:output_type -> master_pb.StatisticsResponse + 24, // 73: master_pb.Seaweed.CollectionList:output_type -> master_pb.CollectionListResponse + 26, // 74: master_pb.Seaweed.CollectionDelete:output_type -> master_pb.CollectionDeleteResponse + 33, // 75: master_pb.Seaweed.VolumeList:output_type -> master_pb.VolumeListResponse + 35, // 76: master_pb.Seaweed.LookupEcVolume:output_type -> master_pb.LookupEcVolumeResponse + 37, // 77: master_pb.Seaweed.VacuumVolume:output_type -> master_pb.VacuumVolumeResponse + 39, // 78: master_pb.Seaweed.DisableVacuum:output_type -> master_pb.DisableVacuumResponse + 41, // 79: master_pb.Seaweed.EnableVacuum:output_type -> master_pb.EnableVacuumResponse + 43, // 80: master_pb.Seaweed.VolumeMarkReadonly:output_type -> master_pb.VolumeMarkReadonlyResponse + 45, // 81: master_pb.Seaweed.GetMasterConfiguration:output_type -> master_pb.GetMasterConfigurationResponse + 47, // 82: master_pb.Seaweed.ListClusterNodes:output_type -> master_pb.ListClusterNodesResponse + 49, // 83: master_pb.Seaweed.LeaseAdminToken:output_type -> master_pb.LeaseAdminTokenResponse + 51, // 84: master_pb.Seaweed.ReleaseAdminToken:output_type -> master_pb.ReleaseAdminTokenResponse + 53, // 85: master_pb.Seaweed.GetAdminLockStatus:output_type -> master_pb.GetAdminLockStatusResponse + 55, // 86: master_pb.Seaweed.Ping:output_type -> master_pb.PingResponse + 61, // 87: master_pb.Seaweed.RaftListClusterServers:output_type -> master_pb.RaftListClusterServersResponse + 57, // 88: master_pb.Seaweed.RaftAddServer:output_type -> master_pb.RaftAddServerResponse + 59, // 89: master_pb.Seaweed.RaftRemoveServer:output_type -> master_pb.RaftRemoveServerResponse + 63, // 90: master_pb.Seaweed.RaftLeadershipTransfer:output_type -> master_pb.RaftLeadershipTransferResponse + 64, // 91: master_pb.Seaweed.VolumeGrow:output_type -> master_pb.VolumeGrowResponse + 67, // [67:92] is the sub-list for method output_type + 42, // [42:67] is the sub-list for method input_type 42, // [42:42] is the sub-list for extension type_name 42, // [42:42] is the sub-list for extension extendee 0, // [0:42] is the sub-list for field type_name @@ -5072,7 +5188,7 @@ func file_master_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_master_proto_rawDesc), len(file_master_proto_rawDesc)), NumEnums: 0, - NumMessages: 77, + NumMessages: 79, NumExtensions: 0, NumServices: 1, }, diff --git a/weed/pb/master_pb/master_grpc.pb.go b/weed/pb/master_pb/master_grpc.pb.go index f13d53bca..81dfd58b1 100644 --- a/weed/pb/master_pb/master_grpc.pb.go +++ b/weed/pb/master_pb/master_grpc.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go-grpc. DO NOT EDIT. // versions: -// - protoc-gen-go-grpc v1.5.1 -// - protoc v6.33.4 +// - protoc-gen-go-grpc v1.6.2 +// - protoc v7.35.0 // source: master.proto package master_pb @@ -37,6 +37,7 @@ const ( Seaweed_ListClusterNodes_FullMethodName = "/master_pb.Seaweed/ListClusterNodes" Seaweed_LeaseAdminToken_FullMethodName = "/master_pb.Seaweed/LeaseAdminToken" Seaweed_ReleaseAdminToken_FullMethodName = "/master_pb.Seaweed/ReleaseAdminToken" + Seaweed_GetAdminLockStatus_FullMethodName = "/master_pb.Seaweed/GetAdminLockStatus" Seaweed_Ping_FullMethodName = "/master_pb.Seaweed/Ping" Seaweed_RaftListClusterServers_FullMethodName = "/master_pb.Seaweed/RaftListClusterServers" Seaweed_RaftAddServer_FullMethodName = "/master_pb.Seaweed/RaftAddServer" @@ -67,6 +68,7 @@ type SeaweedClient interface { ListClusterNodes(ctx context.Context, in *ListClusterNodesRequest, opts ...grpc.CallOption) (*ListClusterNodesResponse, error) LeaseAdminToken(ctx context.Context, in *LeaseAdminTokenRequest, opts ...grpc.CallOption) (*LeaseAdminTokenResponse, error) ReleaseAdminToken(ctx context.Context, in *ReleaseAdminTokenRequest, opts ...grpc.CallOption) (*ReleaseAdminTokenResponse, error) + GetAdminLockStatus(ctx context.Context, in *GetAdminLockStatusRequest, opts ...grpc.CallOption) (*GetAdminLockStatusResponse, error) Ping(ctx context.Context, in *PingRequest, opts ...grpc.CallOption) (*PingResponse, error) RaftListClusterServers(ctx context.Context, in *RaftListClusterServersRequest, opts ...grpc.CallOption) (*RaftListClusterServersResponse, error) RaftAddServer(ctx context.Context, in *RaftAddServerRequest, opts ...grpc.CallOption) (*RaftAddServerResponse, error) @@ -272,6 +274,16 @@ func (c *seaweedClient) ReleaseAdminToken(ctx context.Context, in *ReleaseAdminT return out, nil } +func (c *seaweedClient) GetAdminLockStatus(ctx context.Context, in *GetAdminLockStatusRequest, opts ...grpc.CallOption) (*GetAdminLockStatusResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(GetAdminLockStatusResponse) + err := c.cc.Invoke(ctx, Seaweed_GetAdminLockStatus_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + func (c *seaweedClient) Ping(ctx context.Context, in *PingRequest, opts ...grpc.CallOption) (*PingResponse, error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) out := new(PingResponse) @@ -354,6 +366,7 @@ type SeaweedServer interface { ListClusterNodes(context.Context, *ListClusterNodesRequest) (*ListClusterNodesResponse, error) LeaseAdminToken(context.Context, *LeaseAdminTokenRequest) (*LeaseAdminTokenResponse, error) ReleaseAdminToken(context.Context, *ReleaseAdminTokenRequest) (*ReleaseAdminTokenResponse, error) + GetAdminLockStatus(context.Context, *GetAdminLockStatusRequest) (*GetAdminLockStatusResponse, error) Ping(context.Context, *PingRequest) (*PingResponse, error) RaftListClusterServers(context.Context, *RaftListClusterServersRequest) (*RaftListClusterServersResponse, error) RaftAddServer(context.Context, *RaftAddServerRequest) (*RaftAddServerResponse, error) @@ -371,76 +384,79 @@ type SeaweedServer interface { type UnimplementedSeaweedServer struct{} func (UnimplementedSeaweedServer) SendHeartbeat(grpc.BidiStreamingServer[Heartbeat, HeartbeatResponse]) error { - return status.Errorf(codes.Unimplemented, "method SendHeartbeat not implemented") + return status.Error(codes.Unimplemented, "method SendHeartbeat not implemented") } func (UnimplementedSeaweedServer) KeepConnected(grpc.BidiStreamingServer[KeepConnectedRequest, KeepConnectedResponse]) error { - return status.Errorf(codes.Unimplemented, "method KeepConnected not implemented") + return status.Error(codes.Unimplemented, "method KeepConnected not implemented") } func (UnimplementedSeaweedServer) LookupVolume(context.Context, *LookupVolumeRequest) (*LookupVolumeResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method LookupVolume not implemented") + return nil, status.Error(codes.Unimplemented, "method LookupVolume not implemented") } func (UnimplementedSeaweedServer) Assign(context.Context, *AssignRequest) (*AssignResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method Assign not implemented") + return nil, status.Error(codes.Unimplemented, "method Assign not implemented") } func (UnimplementedSeaweedServer) StreamAssign(grpc.BidiStreamingServer[AssignRequest, AssignResponse]) error { - return status.Errorf(codes.Unimplemented, "method StreamAssign not implemented") + return status.Error(codes.Unimplemented, "method StreamAssign not implemented") } func (UnimplementedSeaweedServer) Statistics(context.Context, *StatisticsRequest) (*StatisticsResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method Statistics not implemented") + return nil, status.Error(codes.Unimplemented, "method Statistics not implemented") } func (UnimplementedSeaweedServer) CollectionList(context.Context, *CollectionListRequest) (*CollectionListResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method CollectionList not implemented") + return nil, status.Error(codes.Unimplemented, "method CollectionList not implemented") } func (UnimplementedSeaweedServer) CollectionDelete(context.Context, *CollectionDeleteRequest) (*CollectionDeleteResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method CollectionDelete not implemented") + return nil, status.Error(codes.Unimplemented, "method CollectionDelete not implemented") } func (UnimplementedSeaweedServer) VolumeList(context.Context, *VolumeListRequest) (*VolumeListResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method VolumeList not implemented") + return nil, status.Error(codes.Unimplemented, "method VolumeList not implemented") } func (UnimplementedSeaweedServer) LookupEcVolume(context.Context, *LookupEcVolumeRequest) (*LookupEcVolumeResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method LookupEcVolume not implemented") + return nil, status.Error(codes.Unimplemented, "method LookupEcVolume not implemented") } func (UnimplementedSeaweedServer) VacuumVolume(context.Context, *VacuumVolumeRequest) (*VacuumVolumeResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method VacuumVolume not implemented") + return nil, status.Error(codes.Unimplemented, "method VacuumVolume not implemented") } func (UnimplementedSeaweedServer) DisableVacuum(context.Context, *DisableVacuumRequest) (*DisableVacuumResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method DisableVacuum not implemented") + return nil, status.Error(codes.Unimplemented, "method DisableVacuum not implemented") } func (UnimplementedSeaweedServer) EnableVacuum(context.Context, *EnableVacuumRequest) (*EnableVacuumResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method EnableVacuum not implemented") + return nil, status.Error(codes.Unimplemented, "method EnableVacuum not implemented") } func (UnimplementedSeaweedServer) VolumeMarkReadonly(context.Context, *VolumeMarkReadonlyRequest) (*VolumeMarkReadonlyResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method VolumeMarkReadonly not implemented") + return nil, status.Error(codes.Unimplemented, "method VolumeMarkReadonly not implemented") } func (UnimplementedSeaweedServer) GetMasterConfiguration(context.Context, *GetMasterConfigurationRequest) (*GetMasterConfigurationResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method GetMasterConfiguration not implemented") + return nil, status.Error(codes.Unimplemented, "method GetMasterConfiguration not implemented") } func (UnimplementedSeaweedServer) ListClusterNodes(context.Context, *ListClusterNodesRequest) (*ListClusterNodesResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method ListClusterNodes not implemented") + return nil, status.Error(codes.Unimplemented, "method ListClusterNodes not implemented") } func (UnimplementedSeaweedServer) LeaseAdminToken(context.Context, *LeaseAdminTokenRequest) (*LeaseAdminTokenResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method LeaseAdminToken not implemented") + return nil, status.Error(codes.Unimplemented, "method LeaseAdminToken not implemented") } func (UnimplementedSeaweedServer) ReleaseAdminToken(context.Context, *ReleaseAdminTokenRequest) (*ReleaseAdminTokenResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method ReleaseAdminToken not implemented") + return nil, status.Error(codes.Unimplemented, "method ReleaseAdminToken not implemented") +} +func (UnimplementedSeaweedServer) GetAdminLockStatus(context.Context, *GetAdminLockStatusRequest) (*GetAdminLockStatusResponse, error) { + return nil, status.Error(codes.Unimplemented, "method GetAdminLockStatus not implemented") } func (UnimplementedSeaweedServer) Ping(context.Context, *PingRequest) (*PingResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method Ping not implemented") + return nil, status.Error(codes.Unimplemented, "method Ping not implemented") } func (UnimplementedSeaweedServer) RaftListClusterServers(context.Context, *RaftListClusterServersRequest) (*RaftListClusterServersResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method RaftListClusterServers not implemented") + return nil, status.Error(codes.Unimplemented, "method RaftListClusterServers not implemented") } func (UnimplementedSeaweedServer) RaftAddServer(context.Context, *RaftAddServerRequest) (*RaftAddServerResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method RaftAddServer not implemented") + return nil, status.Error(codes.Unimplemented, "method RaftAddServer not implemented") } func (UnimplementedSeaweedServer) RaftRemoveServer(context.Context, *RaftRemoveServerRequest) (*RaftRemoveServerResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method RaftRemoveServer not implemented") + return nil, status.Error(codes.Unimplemented, "method RaftRemoveServer not implemented") } func (UnimplementedSeaweedServer) RaftLeadershipTransfer(context.Context, *RaftLeadershipTransferRequest) (*RaftLeadershipTransferResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method RaftLeadershipTransfer not implemented") + return nil, status.Error(codes.Unimplemented, "method RaftLeadershipTransfer not implemented") } func (UnimplementedSeaweedServer) VolumeGrow(context.Context, *VolumeGrowRequest) (*VolumeGrowResponse, error) { - return nil, status.Errorf(codes.Unimplemented, "method VolumeGrow not implemented") + return nil, status.Error(codes.Unimplemented, "method VolumeGrow not implemented") } func (UnimplementedSeaweedServer) mustEmbedUnimplementedSeaweedServer() {} func (UnimplementedSeaweedServer) testEmbeddedByValue() {} @@ -453,7 +469,7 @@ type UnsafeSeaweedServer interface { } func RegisterSeaweedServer(s grpc.ServiceRegistrar, srv SeaweedServer) { - // If the following call pancis, it indicates UnimplementedSeaweedServer was + // If the following call panics, it indicates UnimplementedSeaweedServer 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. @@ -754,6 +770,24 @@ func _Seaweed_ReleaseAdminToken_Handler(srv interface{}, ctx context.Context, de return interceptor(ctx, in, info, handler) } +func _Seaweed_GetAdminLockStatus_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(GetAdminLockStatusRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(SeaweedServer).GetAdminLockStatus(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Seaweed_GetAdminLockStatus_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(SeaweedServer).GetAdminLockStatus(ctx, req.(*GetAdminLockStatusRequest)) + } + return interceptor(ctx, in, info, handler) +} + func _Seaweed_Ping_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(PingRequest) if err := dec(in); err != nil { @@ -929,6 +963,10 @@ var Seaweed_ServiceDesc = grpc.ServiceDesc{ MethodName: "ReleaseAdminToken", Handler: _Seaweed_ReleaseAdminToken_Handler, }, + { + MethodName: "GetAdminLockStatus", + Handler: _Seaweed_GetAdminLockStatus_Handler, + }, { MethodName: "Ping", Handler: _Seaweed_Ping_Handler, diff --git a/weed/server/master_grpc_server_admin.go b/weed/server/master_grpc_server_admin.go index 11e8e0f83..53dce2048 100644 --- a/weed/server/master_grpc_server_admin.go +++ b/weed/server/master_grpc_server_admin.go @@ -154,12 +154,31 @@ func (ms *MasterServer) LeaseAdminToken(ctx context.Context, req *master_pb.Leas func (ms *MasterServer) ReleaseAdminToken(ctx context.Context, req *master_pb.ReleaseAdminTokenRequest) (*master_pb.ReleaseAdminTokenResponse, error) { resp := &master_pb.ReleaseAdminTokenResponse{} + // a follower holds no lock state; answering success here would fake a + // release while the leader keeps the lock until it expires + if !ms.Topo.IsLeader() { + return resp, raft.NotLeaderError + } if ms.adminLocks.isValidToken(req.LockName, time.Unix(0, req.PreviousLockTime), req.PreviousToken) { ms.adminLocks.deleteLock(req.LockName) } return resp, nil } +func (ms *MasterServer) GetAdminLockStatus(ctx context.Context, req *master_pb.GetAdminLockStatusRequest) (*master_pb.GetAdminLockStatusResponse, error) { + resp := &master_pb.GetAdminLockStatusResponse{} + if !ms.Topo.IsLeader() { + return resp, raft.NotLeaderError + } + // isLocked reports the last holder even after the lease expired + if clientName, message, isLocked := ms.adminLocks.isLocked(req.LockName); isLocked { + resp.IsLocked = true + resp.ClientName = clientName + resp.Message = message + } + return resp, nil +} + // isKnownPingTarget reports whether target is a peer that the master has // learned about as part of cluster membership. Restricting Ping to known // peers avoids turning the RPC into a generic outbound dialer. The lookups diff --git a/weed/shell/command_cluster_status.go b/weed/shell/command_cluster_status.go index 68a8cd3aa..b3e7a29d6 100644 --- a/weed/shell/command_cluster_status.go +++ b/weed/shell/command_cluster_status.go @@ -53,6 +53,9 @@ type ClusterStatusPrinter struct { maxParallelization int locked bool + lockHeld bool + lockHolder string + lockMessage string collections []string topology *master_pb.TopologyInfo volumeSizeLimitMb uint64 @@ -91,12 +94,17 @@ func (c *commandClusterStatus) Do(args []string, commandEnv *CommandEnv, writer return err } + lockHolder, lockMessage, lockHeld := commandEnv.shellLockHolder() + sp := &ClusterStatusPrinter{ writer: writer, humanize: *humanize, maxParallelization: *maxParallelization, locked: commandEnv.isLocked(), + lockHeld: lockHeld, + lockHolder: lockHolder, + lockMessage: lockMessage, collections: collections, topology: topology, volumeSizeLimitMb: volumeSizeLimitMb, @@ -339,8 +347,14 @@ func (sp *ClusterStatusPrinter) printClusterInfo() { } status := "unlocked" - if sp.locked { + if sp.lockHeld || sp.locked { status = "LOCKED" + if sp.lockHolder != "" { + status += " by " + sp.lockHolder + } + if sp.lockMessage != "" { + status += " (" + sp.lockMessage + ")" + } } sp.write("cluster:") diff --git a/weed/shell/command_cluster_status_test.go b/weed/shell/command_cluster_status_test.go index 93756075e..045ac4752 100644 --- a/weed/shell/command_cluster_status_test.go +++ b/weed/shell/command_cluster_status_test.go @@ -9,12 +9,15 @@ import ( func TestPrintClusterInfo(t *testing.T) { testCases := []struct { - topology *master_pb.TopologyInfo - humanize bool - want string + topology *master_pb.TopologyInfo + humanize bool + lockHeld bool + lockHolder string + lockMessage string + want string }{ { - testTopology1, true, + testTopology1, true, false, "", "", `cluster: id: test_topo_1 status: unlocked @@ -24,13 +27,33 @@ func TestPrintClusterInfo(t *testing.T) { `, }, { - testTopology1, false, + testTopology1, false, false, "", "", `cluster: id: test_topo_1 status: unlocked nodes: 5 topology: 5 DC(s), 5 disk(s) on 6 rack(s) +`, + }, + { + testTopology1, true, true, "192.168.1.5", "", + `cluster: + id: test_topo_1 + status: LOCKED by 192.168.1.5 + nodes: 5 + topology: 5 DCs, 5 disks on 6 racks + +`, + }, + { + testTopology1, true, true, "192.168.1.5", "ec.encode", + `cluster: + id: test_topo_1 + status: LOCKED by 192.168.1.5 (ec.encode) + nodes: 5 + topology: 5 DCs, 5 disks on 6 racks + `, }, } @@ -38,9 +61,12 @@ func TestPrintClusterInfo(t *testing.T) { for _, tc := range testCases { var buf bytes.Buffer sp := &ClusterStatusPrinter{ - writer: &buf, - humanize: tc.humanize, - topology: tc.topology, + writer: &buf, + humanize: tc.humanize, + topology: tc.topology, + lockHeld: tc.lockHeld, + lockHolder: tc.lockHolder, + lockMessage: tc.lockMessage, } sp.printClusterInfo() got := buf.String() diff --git a/weed/shell/command_lock_unlock.go b/weed/shell/command_lock_unlock.go index 982bb9120..68f47fa07 100644 --- a/weed/shell/command_lock_unlock.go +++ b/weed/shell/command_lock_unlock.go @@ -1,6 +1,7 @@ package shell import ( + "fmt" "io" "github.com/seaweedfs/seaweedfs/weed/util" @@ -32,8 +33,27 @@ func (c *commandLock) HasTag(CommandTag) bool { func (c *commandLock) Do(args []string, commandEnv *CommandEnv, writer io.Writer) (err error) { + waited := false + if !commandEnv.locker.IsLocked() { + if holder, message, held := commandEnv.shellLockHolder(); held { + waited = true + if holder == "" { + holder = "another client" + } + if message != "" { + fmt.Fprintf(writer, "waiting for lock held by %s: %s\n", holder, message) + } else { + fmt.Fprintf(writer, "waiting for lock held by %s\n", holder) + } + } + } + commandEnv.locker.RequestLock(util.DetectedHostAddress()) + if waited { + fmt.Fprintln(writer, "lock acquired") + } + return nil } diff --git a/weed/shell/commands.go b/weed/shell/commands.go index 5cc7e704f..896a2d81b 100644 --- a/weed/shell/commands.go +++ b/weed/shell/commands.go @@ -5,8 +5,11 @@ import ( "fmt" "io" "strings" + "time" + "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/operation" + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" "github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb" "github.com/seaweedfs/seaweedfs/weed/storage/needle_map" @@ -48,7 +51,7 @@ func NewCommandEnv(options *ShellOptions) *CommandEnv { option: options, noLock: false, } - ce.locker = exclusive_locks.NewExclusiveLocker(ce.MasterClient, "shell") + ce.locker = exclusive_locks.NewExclusiveLocker(ce.MasterClient, cluster.AdminShellLockName) return ce } @@ -107,6 +110,24 @@ func (ce *CommandEnv) isLocked() bool { return ce.locker.IsLocked() } +// shellLockHolder asks the master who currently holds the cluster-wide shell +// lock. Best effort: masters without GetAdminLockStatus report no holder, and +// each attempt is bounded so an unresponsive master cannot hang the shell. +func (ce *CommandEnv) shellLockHolder() (clientName string, message string, held bool) { + ce.MasterClient.WithClient(false, func(client master_pb.SeaweedClient) error { + attemptCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + resp, err := client.GetAdminLockStatus(attemptCtx, &master_pb.GetAdminLockStatusRequest{ + LockName: cluster.AdminShellLockName, + }) + if err == nil && resp.IsLocked { + clientName, message, held = resp.ClientName, resp.Message, true + } + return err + }) + return +} + func (ce *CommandEnv) checkDirectory(path string) error { dir, name := util.FullPath(path).DirAndName() diff --git a/weed/wdclient/exclusive_locks/exclusive_locker.go b/weed/wdclient/exclusive_locks/exclusive_locker.go index e61774ae8..f7c56adb2 100644 --- a/weed/wdclient/exclusive_locks/exclusive_locker.go +++ b/weed/wdclient/exclusive_locks/exclusive_locker.go @@ -2,6 +2,7 @@ package exclusive_locks import ( "context" + "sync" "sync/atomic" "time" @@ -14,6 +15,9 @@ const ( RenewInterval = 4 * time.Second SafeRenewInterval = 3 * time.Second InitLockInterval = 1 * time.Second + // bounds each lease and renew RPC attempt so an unresponsive master + // cannot hold l.mu (and block ReleaseLock) indefinitely + rpcTimeout = 3 * time.Second ) type ExclusiveLocker struct { @@ -24,6 +28,9 @@ type ExclusiveLocker struct { lockName string message string clientName string + // serializes renew and release RPCs: a renewal in flight during a release + // would re-create the lock on the master and leave it held until expiry + mu sync.Mutex // Each lock has and only has one goroutine renewGoroutineRunning atomic.Bool } @@ -52,13 +59,12 @@ func (l *ExclusiveLocker) RequestLock(clientName string) { return } - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - // retry to get the lease for { if err := l.masterClient.WithClient(false, func(client master_pb.SeaweedClient) error { - resp, err := client.LeaseAdminToken(ctx, &master_pb.LeaseAdminTokenRequest{ + attemptCtx, cancel := context.WithTimeout(context.Background(), rpcTimeout) + defer cancel() + resp, err := client.LeaseAdminToken(attemptCtx, &master_pb.LeaseAdminTokenRequest{ PreviousToken: atomic.LoadInt64(&l.token), PreviousLockTime: atomic.LoadInt64(&l.lockTsNs), LockName: l.lockName, @@ -77,66 +83,85 @@ func (l *ExclusiveLocker) RequestLock(clientName string) { } } - l.isLocked.Store(true) + l.mu.Lock() l.clientName = clientName + l.isLocked.Store(true) + l.mu.Unlock() glog.V(1).Infof("Acquired lock %s", l.lockName) // Each lock has and only has one goroutine if l.renewGoroutineRunning.CompareAndSwap(false, true) { // start a goroutine to renew the lease go func() { + defer l.renewGoroutineRunning.Store(false) ctx2, cancel2 := context.WithCancel(context.Background()) defer cancel2() for { - if l.isLocked.Load() { - if err := l.masterClient.WithClient(false, func(client master_pb.SeaweedClient) error { - resp, err := client.LeaseAdminToken(ctx2, &master_pb.LeaseAdminTokenRequest{ - PreviousToken: atomic.LoadInt64(&l.token), - PreviousLockTime: atomic.LoadInt64(&l.lockTsNs), - LockName: l.lockName, - ClientName: l.clientName, - Message: l.message, - }) - if err == nil { - atomic.StoreInt64(&l.token, resp.Token) - atomic.StoreInt64(&l.lockTsNs, resp.LockTsNs) - glog.V(2).Infof("Renewed lock %s: ts %d token %d", l.lockName, l.lockTsNs, l.token) - } - return err - }); err != nil { - glog.Warningf("Failed to renew lock %s: %v", l.lockName, err) - l.isLocked.Store(false) - return - } else { - time.Sleep(RenewInterval) - } - } else { - time.Sleep(RenewInterval) + if err := l.renewLease(ctx2); err != nil { + glog.Warningf("Failed to renew lock %s: %v", l.lockName, err) + l.isLocked.Store(false) + return } + time.Sleep(RenewInterval) } }() } } +func (l *ExclusiveLocker) renewLease(ctx context.Context) error { + l.mu.Lock() + defer l.mu.Unlock() + if !l.isLocked.Load() { + return nil + } + return l.masterClient.WithClient(false, func(client master_pb.SeaweedClient) error { + attemptCtx, cancel := context.WithTimeout(ctx, rpcTimeout) + defer cancel() + resp, err := client.LeaseAdminToken(attemptCtx, &master_pb.LeaseAdminTokenRequest{ + PreviousToken: atomic.LoadInt64(&l.token), + PreviousLockTime: atomic.LoadInt64(&l.lockTsNs), + LockName: l.lockName, + ClientName: l.clientName, + Message: l.message, + }) + if err == nil { + atomic.StoreInt64(&l.token, resp.Token) + atomic.StoreInt64(&l.lockTsNs, resp.LockTsNs) + glog.V(2).Infof("Renewed lock %s: ts %d token %d", l.lockName, l.lockTsNs, l.token) + } + return err + }) +} + func (l *ExclusiveLocker) ReleaseLock() { + l.mu.Lock() + defer l.mu.Unlock() + l.isLocked.Store(false) l.clientName = "" + prevToken := atomic.LoadInt64(&l.token) + prevLockTsNs := atomic.LoadInt64(&l.lockTsNs) + ctx, cancel := context.WithCancel(context.Background()) defer cancel() + // single unbounded attempt: a release cut short by a deadline leaves the + // lock held until it expires, turning a slow unlock into a ghost lock l.masterClient.WithClient(false, func(client master_pb.SeaweedClient) error { client.ReleaseAdminToken(ctx, &master_pb.ReleaseAdminTokenRequest{ - PreviousToken: atomic.LoadInt64(&l.token), - PreviousLockTime: atomic.LoadInt64(&l.lockTsNs), + PreviousToken: prevToken, + PreviousLockTime: prevLockTsNs, LockName: l.lockName, }) return nil }) - atomic.StoreInt64(&l.token, 0) - atomic.StoreInt64(&l.lockTsNs, 0) + // compare on clear: a RequestLock racing a slow release must not have its + // fresh token zeroed + atomic.CompareAndSwapInt64(&l.token, prevToken, 0) + atomic.CompareAndSwapInt64(&l.lockTsNs, prevLockTsNs, 0) } func (l *ExclusiveLocker) SetMessage(message string) {