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)

Reply via email to