This is an automated email from the ASF dual-hosted git repository. liujun pushed a commit to branch 3.1 in repository https://gitbox.apache.org/repos/asf/dubbo-go.git
commit 01ac8a5ca0c642dcfd21d5a8356a11d96d084e6a Merge: 59bef4dd1 830fd4aba Author: chickenlj <[email protected]> AuthorDate: Tue Mar 28 11:01:56 2023 +0800 Merge branch '3.0' into 3.1 # Conflicts: # metadata/report/nacos/report.go # registry/servicediscovery/service_discovery_registry.go common/constant/default.go | 10 + common/constant/key.go | 1 + common/url.go | 2 + config/instance/metadata_report_test.go | 7 +- config/metric_config.go | 1 + config/protocol_config.go | 16 + config/protocol_config_test.go | 3 + config/root_config.go | 13 + config/service_config.go | 9 +- config/testdata/config/protocol/application.yaml | 7 +- config/tls_config.go | 52 ++ .../tls_config_test.go | 31 +- config_center/nacos/impl.go | 16 + config_center/zookeeper/impl.go | 10 + go.mod | 26 +- go.sum | 697 +++++++++++++++++++-- metadata/mapping/memory/service_name_mapping.go | 7 +- metadata/mapping/metadata/service_name_mapping.go | 11 +- metadata/mapping/mock_service_name_mapping.go | 7 +- metadata/mapping/service_name_mapping.go | 4 +- metadata/report/etcd/report.go | 7 +- metadata/report/nacos/report.go | 56 +- metadata/report/report.go | 6 +- metadata/report/zookeeper/report.go | 7 +- metadata/service/remote/service_test.go | 7 +- protocol/dubbo3/dubbo3_invoker.go | 17 +- protocol/dubbo3/dubbo3_protocol.go | 18 +- protocol/grpc/client.go | 17 +- protocol/grpc/client_test.go | 7 + protocol/grpc/grpc_protocol.go | 4 - protocol/grpc/server.go | 15 +- registry/event.go | 27 + registry/nacos/service_discovery_test.go | 4 +- registry/polaris/registry.go | 2 +- registry/polaris/service_discovery.go | 11 +- registry/polaris/utils.go | 4 - .../service_mapping_changed_listener.go | 23 +- .../servicediscovery/service_discovery_registry.go | 77 ++- .../service_instances_changed_listener_impl.go | 2 +- .../service_mapping_change_listener_impl.go | 108 ++++ remoting/zookeeper/listener.go | 5 + tools/dubbogo-cli/cmd/show.go | 15 +- tools/dubbogo-cli/go.mod | 35 +- tools/dubbogo-cli/go.sum | 3 +- 44 files changed, 1201 insertions(+), 206 deletions(-) diff --cc metadata/report/nacos/report.go index e55793ea1,c4487f878..2b753f770 --- a/metadata/report/nacos/report.go +++ b/metadata/report/nacos/report.go @@@ -209,9 -210,37 +210,37 @@@ func (n *nacosMetadataReport) getConfig return cfg, nil } + func (n *nacosMetadataReport) addListener(key string, group string, notify registry.MappingListener) error { + return n.client.Client().ListenConfig(vo.ConfigParam{ + DataId: key, + Group: group, + OnChange: func(namespace, group, dataId, data string) { + go callback(notify, dataId, data) + }, + }) + } + + func callback(notify registry.MappingListener, dataId, data string) { + appNames := strings.Split(data, constant.CommaSeparator) + set := gxset.NewSet() + for _, app := range appNames { + set.Add(app) + } + if err := notify.OnEvent(registry.NewServiceMappingChangedEvent(dataId, set)); err != nil { + logger.Errorf("serviceMapping callback err: %s", err.Error()) + } + } + + func (n *nacosMetadataReport) removeServiceMappingListener(key string, group string) error { + return n.client.Client().CancelListenConfig(vo.ConfigParam{ + DataId: key, + Group: group, + }) + } + // RegisterServiceAppMapping map the specified Dubbo service interface to current Dubbo app name func (n *nacosMetadataReport) RegisterServiceAppMapping(key string, group string, value string) error { - oldVal, err := n.getConfig(vo.ConfigParam{ + oldVal, _ := n.getConfig(vo.ConfigParam{ DataId: key, Group: group, }) diff --cc registry/servicediscovery/service_discovery_registry.go index a7b8e5592,6854439fa..5266f4c26 --- a/registry/servicediscovery/service_discovery_registry.go +++ b/registry/servicediscovery/service_discovery_registry.go @@@ -197,21 -205,25 +205,31 @@@ func (s *ServiceDiscoveryRegistry) Subs if err != nil { return perrors.WithMessage(err, "subscribe url error: "+url.String()) } - services := s.getServices(url) + + mappingListener := NewMappingListener(s.url, url, s.subscribedServices, notify) + services := s.getServices(url, mappingListener) if services.Empty() { + return nil + return perrors.Errorf("Should has at least one way to know which services this interface belongs to,"+ + " either specify 'provided-by' for reference or enable metadata-report center subscription url:%s", url.String()) } + // first notify + mappingListener.OnEvent(registry.NewServiceMappingChangedEvent(url.ServiceKey(), services)) + return nil + } + + func (s *ServiceDiscoveryRegistry) SubscribeURL(url *common.URL, notify registry.NotifyListener, services *gxset.HashSet) { // FIXME ServiceNames.String() is not good + var err error serviceNamesKey := services.String() - protocolServiceKey := url.ServiceKey() + ":" + url.Protocol + protocol := "tri" // consume "tri" protocol by default, other protocols need to be specified on reference/consumer explicitly + if url.Protocol != "" { + protocol = url.Protocol + } + protocolServiceKey := url.ServiceKey() + ":" + protocol listener := s.serviceListeners[serviceNamesKey] if listener == nil { - listener = event.NewServiceInstancesChangedListener(services) + listener = NewServiceInstancesChangedListener(services) for _, serviceNameTmp := range services.Values() { serviceName := serviceNameTmp.(string) instances := s.serviceDiscovery.GetInstances(serviceName)
