diff --git a/cluster/cluster.go b/cluster/cluster.go index cd03868..8601992 100644 --- a/cluster/cluster.go +++ b/cluster/cluster.go @@ -62,6 +62,7 @@ type Cluster struct { locker sync.RWMutex //结点与服务关系保护锁 mapRpc map[string]*NodeRpcInfo //nodeId mapServiceNode map[string]map[string]struct{} //map[serviceName]map[NodeId] + mapTemplateServiceNode map[string]map[string]struct{} //map[templateServiceName]map[serviceName]nodeId callSet rpc.CallSet rpcNats rpc.RpcNats @@ -137,6 +138,20 @@ func (cls *Cluster) delServiceNode(serviceName string, nodeId string) { return } + //处理模板服务 + splitServiceName := strings.Split(serviceName,":") + if len(splitServiceName) == 2 { + serviceName = splitServiceName[0] + templateServiceName := splitServiceName[1] + + mapService := cls.mapTemplateServiceNode[templateServiceName] + delete(mapService,serviceName) + + if len(cls.mapTemplateServiceNode[templateServiceName]) == 0 { + delete(cls.mapTemplateServiceNode,templateServiceName) + } + } + mapNode := cls.mapServiceNode[serviceName] delete(mapNode, nodeId) if len(mapNode) == 0 { @@ -171,7 +186,20 @@ func (cls *Cluster) serviceDiscoverySetNodeInfo(nodeInfo *NodeInfo) { continue } mapDuplicate[serviceName] = nil - if _, ok := cls.mapServiceNode[serviceName]; ok == false { + + //如果是模板服务,则记录模板关系 + splitServiceName := strings.Split(serviceName,":") + if len(splitServiceName) == 2 { + serviceName = splitServiceName[0] + templateServiceName := splitServiceName[1] + //记录模板 + if _, ok = cls.mapTemplateServiceNode[templateServiceName]; ok == false { + cls.mapTemplateServiceNode[templateServiceName]=map[string]struct{}{} + } + cls.mapTemplateServiceNode[templateServiceName][serviceName] = struct{}{} + } + + if _, ok = cls.mapServiceNode[serviceName]; ok == false { cls.mapServiceNode[serviceName] = make(map[string]struct{}, 1) } cls.mapServiceNode[serviceName][nodeInfo.NodeId] = struct{}{} @@ -259,25 +287,29 @@ func (cls *Cluster) GetRpcClient(nodeId string) (*rpc.Client,bool) { return cls.getRpcClient(nodeId) } -func GetRpcClient(nodeId string, serviceMethod string,filterRetire bool, clientList []*rpc.Client) (error, int) { +func GetNodeIdByTemplateService(templateServiceName string, rpcClientList []*rpc.Client, filterRetire bool) (error, []*rpc.Client) { + return GetCluster().GetNodeIdByTemplateService(templateServiceName, rpcClientList, filterRetire) +} + +func GetRpcClient(nodeId string, serviceMethod string,filterRetire bool, clientList []*rpc.Client) (error, []*rpc.Client) { if nodeId != rpc.NodeIdNull { pClient,retire := GetCluster().GetRpcClient(nodeId) if pClient == nil { - return fmt.Errorf("cannot find nodeid %d!", nodeId), 0 + return fmt.Errorf("cannot find nodeid %d!", nodeId), nil } //如果需要筛选掉退休结点 if filterRetire == true && retire == true { - return fmt.Errorf("cannot find nodeid %d!", nodeId), 0 + return fmt.Errorf("cannot find nodeid %d!", nodeId), nil } - clientList[0] = pClient - return nil, 1 + clientList = append(clientList,pClient) + return nil, clientList } findIndex := strings.Index(serviceMethod, ".") if findIndex == -1 { - return fmt.Errorf("servicemethod param %s is error!", serviceMethod), 0 + return fmt.Errorf("servicemethod param %s is error!", serviceMethod), nil } serviceName := serviceMethod[:findIndex] @@ -376,6 +408,27 @@ func GetNodeByServiceName(serviceName string) map[string]struct{} { return mapNodeId } +// GetNodeByTemplateServiceName 通过模板服务名获取服务名,返回 map[serviceName真实服务名]NodeId +func GetNodeByTemplateServiceName(templateServiceName string) map[string]string { + cluster.locker.RLock() + defer cluster.locker.RUnlock() + + mapServiceName := cluster.mapTemplateServiceNode[templateServiceName] + mapNodeId := make(map[string]string,9) + for serviceName := range mapServiceName { + mapNode, ok := cluster.mapServiceNode[serviceName] + if ok == false { + return nil + } + + for nodeId,_ := range mapNode { + mapNodeId[serviceName] = nodeId + } + } + + return mapNodeId +} + func (cls *Cluster) GetGlobalCfg() interface{} { return cls.globalCfg } diff --git a/cluster/parsecfg.go b/cluster/parsecfg.go index 9c5d4fd..3628b5b 100644 --- a/cluster/parsecfg.go +++ b/cluster/parsecfg.go @@ -270,7 +270,7 @@ func (cls *Cluster) readLocalClusterConfig(nodeId string) (DiscoveryInfo, []Node } if nodeId != rpc.NodeIdNull && (len(nodeInfoList) != 1) { - return discoveryInfo, nil,rpcMode, fmt.Errorf("%d configurations were found for the configuration with node ID %d!", len(nodeInfoList), nodeId) + return discoveryInfo, nil,rpcMode, fmt.Errorf("nodeid %s configuration error in NodeList", nodeId) } for i, _ := range nodeInfoList { @@ -382,12 +382,24 @@ func (cls *Cluster) parseLocalCfg() { cls.mapRpc[cls.localNodeInfo.NodeId] = &rpcInfo - for _, sName := range cls.localNodeInfo.ServiceList { - if _, ok := cls.mapServiceNode[sName]; ok == false { - cls.mapServiceNode[sName] = make(map[string]struct{}) + for _, serviceName := range cls.localNodeInfo.ServiceList { + splitServiceName := strings.Split(serviceName,":") + if len(splitServiceName) == 2 { + serviceName = splitServiceName[0] + templateServiceName := splitServiceName[1] + //记录模板 + if _, ok := cls.mapTemplateServiceNode[templateServiceName]; ok == false { + cls.mapTemplateServiceNode[templateServiceName]=map[string]struct{}{} + } + cls.mapTemplateServiceNode[templateServiceName][serviceName] = struct{}{} } - cls.mapServiceNode[sName][cls.localNodeInfo.NodeId] = struct{}{} + + if _, ok := cls.mapServiceNode[serviceName]; ok == false { + cls.mapServiceNode[serviceName] = make(map[string]struct{}) + } + + cls.mapServiceNode[serviceName][cls.localNodeInfo.NodeId] = struct{}{} } } @@ -403,6 +415,7 @@ func (cls *Cluster) InitCfg(localNodeId string) error { cls.localServiceCfg = map[string]interface{}{} cls.mapRpc = map[string]*NodeRpcInfo{} cls.mapServiceNode = map[string]map[string]struct{}{} + cls.mapTemplateServiceNode = map[string]map[string]struct{}{} //加载本地结点的NodeList配置 discoveryInfo, nodeInfoList,rpcMode, err := cls.readLocalClusterConfig(localNodeId) @@ -436,12 +449,37 @@ func (cls *Cluster) IsConfigService(serviceName string) bool { return ok } +func (cls *Cluster) GetNodeIdByTemplateService(templateServiceName string, rpcClientList []*rpc.Client, filterRetire bool) (error, []*rpc.Client) { + cls.locker.RLock() + defer cls.locker.RUnlock() -func (cls *Cluster) GetNodeIdByService(serviceName string, rpcClientList []*rpc.Client, filterRetire bool) (error, int) { + mapServiceName := cls.mapTemplateServiceNode[templateServiceName] + for serviceName := range mapServiceName { + mapNodeId, ok := cls.mapServiceNode[serviceName] + if ok == true { + for nodeId, _ := range mapNodeId { + pClient,retire := GetCluster().getRpcClient(nodeId) + if pClient == nil || pClient.IsConnected() == false { + continue + } + + //如果需要筛选掉退休的,对retire状态的结点略过 + if filterRetire == true && retire == true { + continue + } + + rpcClientList = append(rpcClientList,pClient) + } + } + } + + return nil, rpcClientList +} + +func (cls *Cluster) GetNodeIdByService(serviceName string, rpcClientList []*rpc.Client, filterRetire bool) (error, []*rpc.Client) { cls.locker.RLock() defer cls.locker.RUnlock() mapNodeId, ok := cls.mapServiceNode[serviceName] - count := 0 if ok == true { for nodeId, _ := range mapNodeId { pClient,retire := GetCluster().getRpcClient(nodeId) @@ -454,15 +492,11 @@ func (cls *Cluster) GetNodeIdByService(serviceName string, rpcClientList []*rpc. continue } - rpcClientList[count] = pClient - count++ - if count >= cap(rpcClientList) { - break - } + rpcClientList = append(rpcClientList,pClient) } } - return nil, count + return nil, rpcClientList } func (cls *Cluster) GetServiceCfg(serviceName string) interface{} { diff --git a/node/node.go b/node/node.go index 70e710f..d28e9ca 100644 --- a/node/node.go +++ b/node/node.go @@ -25,6 +25,7 @@ import ( var sig chan os.Signal var nodeId string var preSetupService []service.IService //预安装 +var preSetupTemplateService []func()service.IService var profilerInterval time.Duration var bValid bool var configDir = "./config/" @@ -169,6 +170,31 @@ func initNode(id string) { serviceOrder := cluster.GetCluster().GetLocalNodeInfo().ServiceList for _,serviceName:= range serviceOrder{ bSetup := false + + //判断是否有配置模板服务 + splitServiceName := strings.Split(serviceName,":") + if len(splitServiceName) == 2 { + serviceName = splitServiceName[0] + templateServiceName := splitServiceName[1] + for _,newSer := range preSetupTemplateService { + ser := newSer() + ser.OnSetup(ser) + if ser.GetName() == templateServiceName { + ser.SetName(serviceName) + ser.Init(ser,cluster.GetRpcClient,cluster.GetRpcServer,cluster.GetCluster().GetServiceCfg(ser.GetName())) + service.Setup(ser) + + bSetup = true + break + } + } + + if bSetup == false{ + log.Error("Template service not found",log.String("service name",serviceName),log.String("template service name",templateServiceName)) + os.Exit(1) + } + } + for _, s := range preSetupService { if s.GetName() != serviceName { continue @@ -367,6 +393,12 @@ func Setup(s ...service.IService) { } } +func SetupTemplate(fs ...func()service.IService){ + for _, f := range fs { + preSetupTemplateService = append(preSetupTemplateService, f) + } +} + func GetService(serviceName string) service.IService { return service.GetService(serviceName) } diff --git a/rpc/messagequeue.pb.go b/rpc/messagequeue.pb.go index a181b13..dfe5202 100644 --- a/rpc/messagequeue.pb.go +++ b/rpc/messagequeue.pb.go @@ -1,8 +1,8 @@ // Code generated by protoc-gen-go. DO NOT EDIT. // versions: // protoc-gen-go v1.31.0 -// protoc v3.11.4 -// source: test/rpc/messagequeue.proto +// protoc v4.24.0 +// source: rpcproto/messagequeue.proto package rpc @@ -50,11 +50,11 @@ func (x SubscribeType) String() string { } func (SubscribeType) Descriptor() protoreflect.EnumDescriptor { - return file_test_rpc_messagequeue_proto_enumTypes[0].Descriptor() + return file_rpcproto_messagequeue_proto_enumTypes[0].Descriptor() } func (SubscribeType) Type() protoreflect.EnumType { - return &file_test_rpc_messagequeue_proto_enumTypes[0] + return &file_rpcproto_messagequeue_proto_enumTypes[0] } func (x SubscribeType) Number() protoreflect.EnumNumber { @@ -63,7 +63,7 @@ func (x SubscribeType) Number() protoreflect.EnumNumber { // Deprecated: Use SubscribeType.Descriptor instead. func (SubscribeType) EnumDescriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{0} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{0} } type SubscribeMethod int32 @@ -96,11 +96,11 @@ func (x SubscribeMethod) String() string { } func (SubscribeMethod) Descriptor() protoreflect.EnumDescriptor { - return file_test_rpc_messagequeue_proto_enumTypes[1].Descriptor() + return file_rpcproto_messagequeue_proto_enumTypes[1].Descriptor() } func (SubscribeMethod) Type() protoreflect.EnumType { - return &file_test_rpc_messagequeue_proto_enumTypes[1] + return &file_rpcproto_messagequeue_proto_enumTypes[1] } func (x SubscribeMethod) Number() protoreflect.EnumNumber { @@ -109,7 +109,7 @@ func (x SubscribeMethod) Number() protoreflect.EnumNumber { // Deprecated: Use SubscribeMethod.Descriptor instead. func (SubscribeMethod) EnumDescriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{1} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{1} } type DBQueuePopReq struct { @@ -127,7 +127,7 @@ type DBQueuePopReq struct { func (x *DBQueuePopReq) Reset() { *x = DBQueuePopReq{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[0] + mi := &file_rpcproto_messagequeue_proto_msgTypes[0] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -140,7 +140,7 @@ func (x *DBQueuePopReq) String() string { func (*DBQueuePopReq) ProtoMessage() {} func (x *DBQueuePopReq) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[0] + mi := &file_rpcproto_messagequeue_proto_msgTypes[0] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -153,7 +153,7 @@ func (x *DBQueuePopReq) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueuePopReq.ProtoReflect.Descriptor instead. func (*DBQueuePopReq) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{0} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{0} } func (x *DBQueuePopReq) GetCustomerId() string { @@ -203,7 +203,7 @@ type DBQueuePopRes struct { func (x *DBQueuePopRes) Reset() { *x = DBQueuePopRes{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[1] + mi := &file_rpcproto_messagequeue_proto_msgTypes[1] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -216,7 +216,7 @@ func (x *DBQueuePopRes) String() string { func (*DBQueuePopRes) ProtoMessage() {} func (x *DBQueuePopRes) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[1] + mi := &file_rpcproto_messagequeue_proto_msgTypes[1] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -229,7 +229,7 @@ func (x *DBQueuePopRes) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueuePopRes.ProtoReflect.Descriptor instead. func (*DBQueuePopRes) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{1} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{1} } func (x *DBQueuePopRes) GetQueueName() string { @@ -255,7 +255,7 @@ type DBQueueSubscribeReq struct { SubType SubscribeType `protobuf:"varint,1,opt,name=SubType,proto3,enum=SubscribeType" json:"SubType,omitempty"` //订阅类型 Method SubscribeMethod `protobuf:"varint,2,opt,name=Method,proto3,enum=SubscribeMethod" json:"Method,omitempty"` //订阅方法 CustomerId string `protobuf:"bytes,3,opt,name=CustomerId,proto3" json:"CustomerId,omitempty"` //消费者Id - FromNodeId int32 `protobuf:"varint,4,opt,name=FromNodeId,proto3" json:"FromNodeId,omitempty"` + FromNodeId string `protobuf:"bytes,4,opt,name=FromNodeId,proto3" json:"FromNodeId,omitempty"` RpcMethod string `protobuf:"bytes,5,opt,name=RpcMethod,proto3" json:"RpcMethod,omitempty"` TopicName string `protobuf:"bytes,6,opt,name=TopicName,proto3" json:"TopicName,omitempty"` //主题名称 StartIndex uint64 `protobuf:"varint,7,opt,name=StartIndex,proto3" json:"StartIndex,omitempty"` //开始位置 ,格式前4位是时间戳秒,后面是序号。如果填0时,服务自动修改成:(4bit 当前时间秒)| (0000 4bit) @@ -265,7 +265,7 @@ type DBQueueSubscribeReq struct { func (x *DBQueueSubscribeReq) Reset() { *x = DBQueueSubscribeReq{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[2] + mi := &file_rpcproto_messagequeue_proto_msgTypes[2] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -278,7 +278,7 @@ func (x *DBQueueSubscribeReq) String() string { func (*DBQueueSubscribeReq) ProtoMessage() {} func (x *DBQueueSubscribeReq) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[2] + mi := &file_rpcproto_messagequeue_proto_msgTypes[2] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -291,7 +291,7 @@ func (x *DBQueueSubscribeReq) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueueSubscribeReq.ProtoReflect.Descriptor instead. func (*DBQueueSubscribeReq) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{2} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{2} } func (x *DBQueueSubscribeReq) GetSubType() SubscribeType { @@ -315,11 +315,11 @@ func (x *DBQueueSubscribeReq) GetCustomerId() string { return "" } -func (x *DBQueueSubscribeReq) GetFromNodeId() int32 { +func (x *DBQueueSubscribeReq) GetFromNodeId() string { if x != nil { return x.FromNodeId } - return 0 + return "" } func (x *DBQueueSubscribeReq) GetRpcMethod() string { @@ -359,7 +359,7 @@ type DBQueueSubscribeRes struct { func (x *DBQueueSubscribeRes) Reset() { *x = DBQueueSubscribeRes{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[3] + mi := &file_rpcproto_messagequeue_proto_msgTypes[3] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -372,7 +372,7 @@ func (x *DBQueueSubscribeRes) String() string { func (*DBQueueSubscribeRes) ProtoMessage() {} func (x *DBQueueSubscribeRes) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[3] + mi := &file_rpcproto_messagequeue_proto_msgTypes[3] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -385,7 +385,7 @@ func (x *DBQueueSubscribeRes) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueueSubscribeRes.ProtoReflect.Descriptor instead. func (*DBQueueSubscribeRes) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{3} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{3} } type DBQueuePublishReq struct { @@ -400,7 +400,7 @@ type DBQueuePublishReq struct { func (x *DBQueuePublishReq) Reset() { *x = DBQueuePublishReq{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[4] + mi := &file_rpcproto_messagequeue_proto_msgTypes[4] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -413,7 +413,7 @@ func (x *DBQueuePublishReq) String() string { func (*DBQueuePublishReq) ProtoMessage() {} func (x *DBQueuePublishReq) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[4] + mi := &file_rpcproto_messagequeue_proto_msgTypes[4] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -426,7 +426,7 @@ func (x *DBQueuePublishReq) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueuePublishReq.ProtoReflect.Descriptor instead. func (*DBQueuePublishReq) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{4} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{4} } func (x *DBQueuePublishReq) GetTopicName() string { @@ -452,7 +452,7 @@ type DBQueuePublishRes struct { func (x *DBQueuePublishRes) Reset() { *x = DBQueuePublishRes{} if protoimpl.UnsafeEnabled { - mi := &file_test_rpc_messagequeue_proto_msgTypes[5] + mi := &file_rpcproto_messagequeue_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -465,7 +465,7 @@ func (x *DBQueuePublishRes) String() string { func (*DBQueuePublishRes) ProtoMessage() {} func (x *DBQueuePublishRes) ProtoReflect() protoreflect.Message { - mi := &file_test_rpc_messagequeue_proto_msgTypes[5] + mi := &file_rpcproto_messagequeue_proto_msgTypes[5] if protoimpl.UnsafeEnabled && x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -478,13 +478,13 @@ func (x *DBQueuePublishRes) ProtoReflect() protoreflect.Message { // Deprecated: Use DBQueuePublishRes.ProtoReflect.Descriptor instead. func (*DBQueuePublishRes) Descriptor() ([]byte, []int) { - return file_test_rpc_messagequeue_proto_rawDescGZIP(), []int{5} + return file_rpcproto_messagequeue_proto_rawDescGZIP(), []int{5} } -var File_test_rpc_messagequeue_proto protoreflect.FileDescriptor +var File_rpcproto_messagequeue_proto protoreflect.FileDescriptor -var file_test_rpc_messagequeue_proto_rawDesc = []byte{ - 0x0a, 0x1b, 0x74, 0x65, 0x73, 0x74, 0x2f, 0x72, 0x70, 0x63, 0x2f, 0x6d, 0x65, 0x73, 0x73, 0x61, +var file_rpcproto_messagequeue_proto_rawDesc = []byte{ + 0x0a, 0x1b, 0x72, 0x70, 0x63, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x2f, 0x6d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x71, 0x75, 0x65, 0x75, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, 0xa3, 0x01, 0x0a, 0x0d, 0x44, 0x42, 0x51, 0x75, 0x65, 0x75, 0x65, 0x50, 0x6f, 0x70, 0x52, 0x65, 0x71, 0x12, 0x1e, 0x0a, 0x0a, 0x43, 0x75, 0x73, 0x74, 0x6f, 0x6d, 0x65, 0x72, 0x49, 0x64, 0x18, 0x01, 0x20, @@ -510,7 +510,7 @@ var file_test_rpc_messagequeue_proto_rawDesc = []byte{ 0x6f, 0x64, 0x52, 0x06, 0x4d, 0x65, 0x74, 0x68, 0x6f, 0x64, 0x12, 0x1e, 0x0a, 0x0a, 0x43, 0x75, 0x73, 0x74, 0x6f, 0x6d, 0x65, 0x72, 0x49, 0x64, 0x18, 0x03, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0a, 0x43, 0x75, 0x73, 0x74, 0x6f, 0x6d, 0x65, 0x72, 0x49, 0x64, 0x12, 0x1e, 0x0a, 0x0a, 0x46, 0x72, - 0x6f, 0x6d, 0x4e, 0x6f, 0x64, 0x65, 0x49, 0x64, 0x18, 0x04, 0x20, 0x01, 0x28, 0x05, 0x52, 0x0a, + 0x6f, 0x6d, 0x4e, 0x6f, 0x64, 0x65, 0x49, 0x64, 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, 0x0a, 0x46, 0x72, 0x6f, 0x6d, 0x4e, 0x6f, 0x64, 0x65, 0x49, 0x64, 0x12, 0x1c, 0x0a, 0x09, 0x52, 0x70, 0x63, 0x4d, 0x65, 0x74, 0x68, 0x6f, 0x64, 0x18, 0x05, 0x20, 0x01, 0x28, 0x09, 0x52, 0x09, 0x52, 0x70, 0x63, 0x4d, 0x65, 0x74, 0x68, 0x6f, 0x64, 0x12, 0x1c, 0x0a, 0x09, 0x54, 0x6f, 0x70, 0x69, @@ -539,20 +539,20 @@ var file_test_rpc_messagequeue_proto_rawDesc = []byte{ } var ( - file_test_rpc_messagequeue_proto_rawDescOnce sync.Once - file_test_rpc_messagequeue_proto_rawDescData = file_test_rpc_messagequeue_proto_rawDesc + file_rpcproto_messagequeue_proto_rawDescOnce sync.Once + file_rpcproto_messagequeue_proto_rawDescData = file_rpcproto_messagequeue_proto_rawDesc ) -func file_test_rpc_messagequeue_proto_rawDescGZIP() []byte { - file_test_rpc_messagequeue_proto_rawDescOnce.Do(func() { - file_test_rpc_messagequeue_proto_rawDescData = protoimpl.X.CompressGZIP(file_test_rpc_messagequeue_proto_rawDescData) +func file_rpcproto_messagequeue_proto_rawDescGZIP() []byte { + file_rpcproto_messagequeue_proto_rawDescOnce.Do(func() { + file_rpcproto_messagequeue_proto_rawDescData = protoimpl.X.CompressGZIP(file_rpcproto_messagequeue_proto_rawDescData) }) - return file_test_rpc_messagequeue_proto_rawDescData + return file_rpcproto_messagequeue_proto_rawDescData } -var file_test_rpc_messagequeue_proto_enumTypes = make([]protoimpl.EnumInfo, 2) -var file_test_rpc_messagequeue_proto_msgTypes = make([]protoimpl.MessageInfo, 6) -var file_test_rpc_messagequeue_proto_goTypes = []interface{}{ +var file_rpcproto_messagequeue_proto_enumTypes = make([]protoimpl.EnumInfo, 2) +var file_rpcproto_messagequeue_proto_msgTypes = make([]protoimpl.MessageInfo, 6) +var file_rpcproto_messagequeue_proto_goTypes = []interface{}{ (SubscribeType)(0), // 0: SubscribeType (SubscribeMethod)(0), // 1: SubscribeMethod (*DBQueuePopReq)(nil), // 2: DBQueuePopReq @@ -562,7 +562,7 @@ var file_test_rpc_messagequeue_proto_goTypes = []interface{}{ (*DBQueuePublishReq)(nil), // 6: DBQueuePublishReq (*DBQueuePublishRes)(nil), // 7: DBQueuePublishRes } -var file_test_rpc_messagequeue_proto_depIdxs = []int32{ +var file_rpcproto_messagequeue_proto_depIdxs = []int32{ 0, // 0: DBQueueSubscribeReq.SubType:type_name -> SubscribeType 1, // 1: DBQueueSubscribeReq.Method:type_name -> SubscribeMethod 2, // [2:2] is the sub-list for method output_type @@ -572,13 +572,13 @@ var file_test_rpc_messagequeue_proto_depIdxs = []int32{ 0, // [0:2] is the sub-list for field type_name } -func init() { file_test_rpc_messagequeue_proto_init() } -func file_test_rpc_messagequeue_proto_init() { - if File_test_rpc_messagequeue_proto != nil { +func init() { file_rpcproto_messagequeue_proto_init() } +func file_rpcproto_messagequeue_proto_init() { + if File_rpcproto_messagequeue_proto != nil { return } if !protoimpl.UnsafeEnabled { - file_test_rpc_messagequeue_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueuePopReq); i { case 0: return &v.state @@ -590,7 +590,7 @@ func file_test_rpc_messagequeue_proto_init() { return nil } } - file_test_rpc_messagequeue_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[1].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueuePopRes); i { case 0: return &v.state @@ -602,7 +602,7 @@ func file_test_rpc_messagequeue_proto_init() { return nil } } - file_test_rpc_messagequeue_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[2].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueueSubscribeReq); i { case 0: return &v.state @@ -614,7 +614,7 @@ func file_test_rpc_messagequeue_proto_init() { return nil } } - file_test_rpc_messagequeue_proto_msgTypes[3].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[3].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueueSubscribeRes); i { case 0: return &v.state @@ -626,7 +626,7 @@ func file_test_rpc_messagequeue_proto_init() { return nil } } - file_test_rpc_messagequeue_proto_msgTypes[4].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[4].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueuePublishReq); i { case 0: return &v.state @@ -638,7 +638,7 @@ func file_test_rpc_messagequeue_proto_init() { return nil } } - file_test_rpc_messagequeue_proto_msgTypes[5].Exporter = func(v interface{}, i int) interface{} { + file_rpcproto_messagequeue_proto_msgTypes[5].Exporter = func(v interface{}, i int) interface{} { switch v := v.(*DBQueuePublishRes); i { case 0: return &v.state @@ -655,19 +655,19 @@ func file_test_rpc_messagequeue_proto_init() { out := protoimpl.TypeBuilder{ File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), - RawDescriptor: file_test_rpc_messagequeue_proto_rawDesc, + RawDescriptor: file_rpcproto_messagequeue_proto_rawDesc, NumEnums: 2, NumMessages: 6, NumExtensions: 0, NumServices: 0, }, - GoTypes: file_test_rpc_messagequeue_proto_goTypes, - DependencyIndexes: file_test_rpc_messagequeue_proto_depIdxs, - EnumInfos: file_test_rpc_messagequeue_proto_enumTypes, - MessageInfos: file_test_rpc_messagequeue_proto_msgTypes, + GoTypes: file_rpcproto_messagequeue_proto_goTypes, + DependencyIndexes: file_rpcproto_messagequeue_proto_depIdxs, + EnumInfos: file_rpcproto_messagequeue_proto_enumTypes, + MessageInfos: file_rpcproto_messagequeue_proto_msgTypes, }.Build() - File_test_rpc_messagequeue_proto = out.File - file_test_rpc_messagequeue_proto_rawDesc = nil - file_test_rpc_messagequeue_proto_goTypes = nil - file_test_rpc_messagequeue_proto_depIdxs = nil + File_rpcproto_messagequeue_proto = out.File + file_rpcproto_messagequeue_proto_rawDesc = nil + file_rpcproto_messagequeue_proto_goTypes = nil + file_rpcproto_messagequeue_proto_depIdxs = nil } diff --git a/rpc/messagequeue.proto b/rpc/messagequeue.proto index 0f0628b..3f8b298 100644 --- a/rpc/messagequeue.proto +++ b/rpc/messagequeue.proto @@ -31,7 +31,7 @@ message DBQueueSubscribeReq { SubscribeType SubType = 1; //订阅类型 SubscribeMethod Method = 2; //订阅方法 string CustomerId = 3; //消费者Id - int32 FromNodeId = 4; + string FromNodeId = 4; string RpcMethod = 5; string TopicName = 6; //主题名称 uint64 StartIndex = 7; //开始位置 ,格式前4位是时间戳秒,后面是序号。如果填0时,服务自动修改成:(4bit 当前时间秒)| (0000 4bit) diff --git a/rpc/rpchandler.go b/rpc/rpchandler.go index 4b72151..1d7deab 100644 --- a/rpc/rpchandler.go +++ b/rpc/rpchandler.go @@ -13,9 +13,9 @@ import ( "time" ) -const maxClusterNode int = 128 +const maxClusterNode int = 32 -type FuncRpcClient func(nodeId string, serviceMethod string,filterRetire bool, client []*Client) (error, int) +type FuncRpcClient func(nodeId string, serviceMethod string,filterRetire bool, client []*Client) (error, []*Client) type FuncRpcServer func() IServer const NodeIdNull = "" @@ -63,7 +63,7 @@ type RpcHandler struct { funcRpcClient FuncRpcClient funcRpcServer FuncRpcServer - pClientList []*Client + //pClientList []*Client } //type TriggerRpcConnEvent func(bConnect bool, clientSeq uint32, nodeId string) @@ -135,7 +135,6 @@ func (handler *RpcHandler) InitRpcHandler(rpcHandler IRpcHandler, getClientFun F handler.mapFunctions = map[string]RpcMethodInfo{} handler.funcRpcClient = getClientFun handler.funcRpcServer = getServerFun - handler.pClientList = make([]*Client, maxClusterNode) handler.RegisterRpc(rpcHandler) } @@ -274,7 +273,7 @@ func (handler *RpcHandler) HandlerRpcRequest(request *RpcRequest) { //普通的rpc请求 v, ok := handler.mapFunctions[request.RpcRequestData.GetServiceMethod()] if ok == false { - err := "RpcHandler " + handler.rpcHandler.GetName() + "cannot find " + request.RpcRequestData.GetServiceMethod() + err := "RpcHandler " + handler.rpcHandler.GetName() + " cannot find " + request.RpcRequestData.GetServiceMethod() log.Error("HandlerRpcRequest cannot find serviceMethod",log.String("RpcHandlerName",handler.rpcHandler.GetName()),log.String("serviceMethod",request.RpcRequestData.GetServiceMethod())) if request.requestHandle != nil { request.requestHandle(nil, RpcError(err)) @@ -435,9 +434,9 @@ func (handler *RpcHandler) CallMethod(client *Client,ServiceMethod string, param } func (handler *RpcHandler) goRpc(processor IRpcProcessor, bCast bool, nodeId string, serviceMethod string, args interface{}) error { - var pClientList [maxClusterNode]*Client - err, count := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList[:]) - if count == 0 { + pClientList :=make([]*Client,0,maxClusterNode) + err, pClientList := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList) + if len(pClientList) == 0 { if err != nil { log.Error("call serviceMethod is failed",log.String("serviceMethod",serviceMethod),log.ErrorAttr("error",err)) } else { @@ -446,13 +445,13 @@ func (handler *RpcHandler) goRpc(processor IRpcProcessor, bCast bool, nodeId str return err } - if count > 1 && bCast == false { + if len(pClientList) > 1 && bCast == false { log.Error("cannot call serviceMethod more then 1 node",log.String("serviceMethod",serviceMethod)) return errors.New("cannot call more then 1 node") } //2.rpcClient调用 - for i := 0; i < count; i++ { + for i := 0; i < len(pClientList); i++ { pCall := pClientList[i].Go(pClientList[i].GetTargetNodeId(),DefaultRpcTimeout,handler.rpcHandler,true, serviceMethod, args, nil) if pCall.Err != nil { err = pCall.Err @@ -465,16 +464,16 @@ func (handler *RpcHandler) goRpc(processor IRpcProcessor, bCast bool, nodeId str } func (handler *RpcHandler) callRpc(timeout time.Duration,nodeId string, serviceMethod string, args interface{}, reply interface{}) error { - var pClientList [maxClusterNode]*Client - err, count := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList[:]) + pClientList :=make([]*Client,0,maxClusterNode) + err, pClientList := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList) if err != nil { log.Error("Call serviceMethod is failed",log.ErrorAttr("error",err)) return err - } else if count <= 0 { + } else if len(pClientList) <= 0 { err = errors.New("Call serviceMethod is error:cannot find " + serviceMethod) log.Error("cannot find serviceMethod",log.String("serviceMethod",serviceMethod)) return err - } else if count > 1 { + } else if len(pClientList) > 1 { log.Error("Cannot call more then 1 node!",log.String("serviceMethod",serviceMethod)) return errors.New("cannot call more then 1 node") } @@ -509,9 +508,9 @@ func (handler *RpcHandler) asyncCallRpc(timeout time.Duration,nodeId string, ser } reply := reflect.New(fVal.Type().In(0).Elem()).Interface() - var pClientList [2]*Client - err, count := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList[:]) - if count == 0 || err != nil { + pClientList :=make([]*Client,0,1) + err, pClientList := handler.funcRpcClient(nodeId, serviceMethod,false, pClientList[:]) + if len(pClientList) == 0 || err != nil { if err == nil { if nodeId != NodeIdNull { err = fmt.Errorf("cannot find %s from nodeId %d",serviceMethod,nodeId) @@ -524,7 +523,7 @@ func (handler *RpcHandler) asyncCallRpc(timeout time.Duration,nodeId string, ser return emptyCancelRpc,nil } - if count > 1 { + if len(pClientList) > 1 { err := errors.New("cannot call more then 1 node") fVal.Call([]reflect.Value{reflect.ValueOf(reply), reflect.ValueOf(err)}) log.Error("cannot call more then 1 node",log.String("serviceMethod",serviceMethod)) @@ -593,12 +592,13 @@ func (handler *RpcHandler) CastGo(serviceMethod string, args interface{}) error func (handler *RpcHandler) RawGoNode(rpcProcessorType RpcProcessorType, nodeId string, rpcMethodId uint32, serviceName string, rawArgs []byte) error { processor := GetProcessor(uint8(rpcProcessorType)) - err, count := handler.funcRpcClient(nodeId, serviceName,false, handler.pClientList) - if count == 0 || err != nil { + pClientList := make([]*Client,0,1) + err, pClientList := handler.funcRpcClient(nodeId, serviceName,false, pClientList) + if len(pClientList) == 0 || err != nil { log.Error("call serviceMethod is failed",log.ErrorAttr("error",err)) return err } - if count > 1 { + if len(pClientList) > 1 { err := errors.New("cannot call more then 1 node") log.Error("cannot call more then 1 node",log.String("serviceName",serviceName)) return err @@ -606,14 +606,14 @@ func (handler *RpcHandler) RawGoNode(rpcProcessorType RpcProcessorType, nodeId s //2.rpcClient调用 //如果调用本结点服务 - for i := 0; i < count; i++ { + for i := 0; i < len(pClientList); i++ { //跨node调用 - pCall := handler.pClientList[i].RawGo(handler.pClientList[i].GetTargetNodeId(),DefaultRpcTimeout,handler.rpcHandler,processor, true, rpcMethodId, serviceName, rawArgs, nil) + pCall := pClientList[i].RawGo(pClientList[i].GetTargetNodeId(),DefaultRpcTimeout,handler.rpcHandler,processor, true, rpcMethodId, serviceName, rawArgs, nil) if pCall.Err != nil { err = pCall.Err } - handler.pClientList[i].RemovePending(pCall.Seq) + pClientList[i].RemovePending(pCall.Seq) ReleaseCall(pCall) } diff --git a/service/service.go b/service/service.go index 68f6dea..1042c7c 100644 --- a/service/service.go +++ b/service/service.go @@ -19,6 +19,8 @@ import ( var timerDispatcherLen = 100000 var maxServiceEventChannelNum = 2000000 + + type IService interface { concurrent.IConcurrent Init(iService IService,getClientFun rpc.FuncRpcClient,getServerFun rpc.FuncRpcServer,serviceCfg interface{}) diff --git a/sysservice/messagequeueservice/CustomerSubscriber.go b/sysservice/messagequeueservice/CustomerSubscriber.go index 94578b4..313b0d5 100644 --- a/sysservice/messagequeueservice/CustomerSubscriber.go +++ b/sysservice/messagequeueservice/CustomerSubscriber.go @@ -16,7 +16,7 @@ type CustomerSubscriber struct { rpc.IRpcHandler topic string subscriber *Subscriber - fromNodeId int + fromNodeId string callBackRpcMethod string serviceName string StartIndex uint64 @@ -37,7 +37,7 @@ const ( MethodLast SubscribeMethod = 1 //Last模式,以该消费者上次记录的位置开始订阅 ) -func (cs *CustomerSubscriber) trySetSubscriberBaseInfo(rpcHandler rpc.IRpcHandler, ss *Subscriber, topic string, subscribeMethod SubscribeMethod, customerId string, fromNodeId int, callBackRpcMethod string, startIndex uint64, oneBatchQuantity int32) error { +func (cs *CustomerSubscriber) trySetSubscriberBaseInfo(rpcHandler rpc.IRpcHandler, ss *Subscriber, topic string, subscribeMethod SubscribeMethod, customerId string, fromNodeId string, callBackRpcMethod string, startIndex uint64, oneBatchQuantity int32) error { cs.subscriber = ss cs.fromNodeId = fromNodeId cs.callBackRpcMethod = callBackRpcMethod @@ -85,7 +85,7 @@ func (cs *CustomerSubscriber) trySetSubscriberBaseInfo(rpcHandler rpc.IRpcHandle } // 开始订阅 -func (cs *CustomerSubscriber) Subscribe(rpcHandler rpc.IRpcHandler, ss *Subscriber, topic string, subscribeMethod SubscribeMethod, customerId string, fromNodeId int, callBackRpcMethod string, startIndex uint64, oneBatchQuantity int32) error { +func (cs *CustomerSubscriber) Subscribe(rpcHandler rpc.IRpcHandler, ss *Subscriber, topic string, subscribeMethod SubscribeMethod, customerId string, fromNodeId string, callBackRpcMethod string, startIndex uint64, oneBatchQuantity int32) error { err := cs.trySetSubscriberBaseInfo(rpcHandler, ss, topic, subscribeMethod, customerId, fromNodeId, callBackRpcMethod, startIndex, oneBatchQuantity) if err != nil { return err diff --git a/sysservice/messagequeueservice/MessageQueueService.go b/sysservice/messagequeueservice/MessageQueueService.go index 90a4d38..80ff8cc 100644 --- a/sysservice/messagequeueservice/MessageQueueService.go +++ b/sysservice/messagequeueservice/MessageQueueService.go @@ -122,5 +122,5 @@ func (ms *MessageQueueService) RPC_Publish(inParam *rpc.DBQueuePublishReq, outPa func (ms *MessageQueueService) RPC_Subscribe(req *rpc.DBQueueSubscribeReq, res *rpc.DBQueueSubscribeRes) error { topicRoom := ms.GetTopicRoom(req.TopicName) - return topicRoom.TopicSubscribe(ms.GetRpcHandler(), req.SubType, int32(req.Method), int(req.FromNodeId), req.RpcMethod, req.TopicName, req.CustomerId, req.StartIndex, req.OneBatchQuantity) + return topicRoom.TopicSubscribe(ms.GetRpcHandler(), req.SubType, int32(req.Method), req.FromNodeId, req.RpcMethod, req.TopicName, req.CustomerId, req.StartIndex, req.OneBatchQuantity) } diff --git a/sysservice/messagequeueservice/Subscriber.go b/sysservice/messagequeueservice/Subscriber.go index 42ad824..4deff8c 100644 --- a/sysservice/messagequeueservice/Subscriber.go +++ b/sysservice/messagequeueservice/Subscriber.go @@ -31,7 +31,7 @@ func (ss *Subscriber) PersistTopicData(topic string, topics []TopicData, retryCo return ss.dataPersist.PersistTopicData(topic, topics, retryCount) } -func (ss *Subscriber) TopicSubscribe(rpcHandler rpc.IRpcHandler, subScribeType rpc.SubscribeType, subscribeMethod SubscribeMethod, fromNodeId int, callBackRpcMethod string, topic string, customerId string, StartIndex uint64, oneBatchQuantity int32) error { +func (ss *Subscriber) TopicSubscribe(rpcHandler rpc.IRpcHandler, subScribeType rpc.SubscribeType, subscribeMethod SubscribeMethod, fromNodeId string, callBackRpcMethod string, topic string, customerId string, StartIndex uint64, oneBatchQuantity int32) error { //取消订阅时 if subScribeType == rpc.SubscribeType_Unsubscribe { ss.UnSubscribe(customerId)