This is an automated email from the ASF dual-hosted git repository.

Alanxtl pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/dubbo-go.git


The following commit(s) were added to refs/heads/develop by this push:
     new 788f55894 Feat: enhance mapping observability (#3605)
788f55894 is described below

commit 788f55894264c2d6885bc473806ed88e828ca715
Author: Modo <[email protected]>
AuthorDate: Thu Aug 20 13:36:04 2026 +0800

    Feat: enhance mapping observability (#3605)
    
    * add metrics and logs for mapping operation
    
    * fix ci
    
    * add test covering MetadataMetricCollector.start()
    
    * fix: address review comments
    
    * refactor: move mapping metrics to ServiceNameMapping
    
    refactor: move mapping metrics to ServiceNameMapping
    
    * fix ci
    
    * fix empty event in early-return and add test
    
    * fix: hide sensitive information
---
 metadata/mapping/metadata/service_name_mapping.go  |  99 ++++++--
 .../mapping/metadata/service_name_mapping_test.go  | 267 +++++++++++++++++++++
 metadata/report_instance_test.go                   |  15 +-
 metrics/metadata/collector.go                      |  55 +++++
 metrics/metadata/collector_test.go                 | 239 ++++++++++++++++++
 metrics/metadata/metric_set.go                     |  46 +++-
 6 files changed, 696 insertions(+), 25 deletions(-)

diff --git a/metadata/mapping/metadata/service_name_mapping.go 
b/metadata/mapping/metadata/service_name_mapping.go
index 51c1be308..bd3225efb 100644
--- a/metadata/mapping/metadata/service_name_mapping.go
+++ b/metadata/mapping/metadata/service_name_mapping.go
@@ -25,6 +25,7 @@ import (
 
 import (
        gxset "github.com/dubbogo/gost/container/set"
+       "github.com/dubbogo/gost/log/logger"
 
        perrors "github.com/pkg/errors"
 )
@@ -36,6 +37,8 @@ import (
        "dubbo.apache.org/dubbo-go/v3/metadata"
        "dubbo.apache.org/dubbo-go/v3/metadata/mapping"
        "dubbo.apache.org/dubbo-go/v3/metadata/report"
+       "dubbo.apache.org/dubbo-go/v3/metrics"
+       metadataMetrics "dubbo.apache.org/dubbo-go/v3/metrics/metadata"
 )
 
 const DefaultGroup = "mapping"
@@ -71,20 +74,36 @@ type ServiceNameMapping struct {
 }
 
 // Map will map the service to this application-level service
-func (d *ServiceNameMapping) Map(url *common.URL) error {
+func (d *ServiceNameMapping) Map(url *common.URL) (err error) {
        serviceInterface := url.GetParam(constant.InterfaceKey, "")
        appName := url.GetParam(constant.ApplicationKey, "")
+
+       event := 
metadataMetrics.NewMetadataMetricTimeEvent(metadataMetrics.MetadataMappingRegister)
+       event.Attachment[constant.InterfaceKey] = serviceInterface
+       event.Attachment[constant.GroupKey] = DefaultGroup
+       event.Attachment[constant.ApplicationKey] = appName
+       defer func() {
+               event.Succ = err == nil
+               event.End = time.Now()
+               metrics.Publish(event)
+       }()
+
        // url is the service url,not the registry url,this url has no registry 
id info,can not get where to write mapping,so write all
        // if the mapping can hold a report instance, it can write once
        metadataReports := metadata.GetMetadataReports()
        if len(metadataReports) == 0 {
-               return perrors.New("can not registering mapping to remote cause 
no metadata report instance found")
+               err = perrors.New("can not registering mapping to remote cause 
no metadata report instance found")
+               logger.Errorf("[Metadata][Mapping] register failed interface=%s 
application=%s group=%s reports=0 err=%v", serviceInterface, appName, 
DefaultGroup, err)
+               return err
        }
+
        for _, metadataReport := range metadataReports {
                if err := registerWithRetry(metadataReport, serviceInterface, 
DefaultGroup, appName); err != nil {
+                       logger.Errorf("[Metadata][Mapping] register failed 
interface=%s application=%s group=%s reports=%d err=%v", serviceInterface, 
appName, DefaultGroup, len(metadataReports), err)
                        return err
                }
        }
+       logger.Debugf("[Metadata][Mapping] register succeeded interface=%s 
application=%s group=%s reports=%d", serviceInterface, appName, DefaultGroup, 
len(metadataReports))
        return nil
 }
 
@@ -117,27 +136,52 @@ func backoff(attempt int) time.Duration {
 }
 
 // Get will return the application-level services. If not found, the empty set 
will be returned.
-func (d *ServiceNameMapping) Get(url *common.URL, listener 
mapping.MappingListener) (*gxset.HashSet, error) {
+func (d *ServiceNameMapping) Get(url *common.URL, listener 
mapping.MappingListener) (result *gxset.HashSet, err error) {
        serviceInterface := url.GetParam(constant.InterfaceKey, "")
+
+       operation := "get"
+       eventName := metadataMetrics.MetadataMappingGet
+       if listener != nil {
+               operation = "listen"
+               eventName = metadataMetrics.MetadataMappingListen
+       }
+
+       event := metadataMetrics.NewMetadataMetricTimeEvent(eventName)
+       event.Attachment[constant.InterfaceKey] = serviceInterface
+       event.Attachment[constant.GroupKey] = DefaultGroup
+       var errs []error
+       defer func() {
+               event.Succ = err == nil
+               event.Partial = err == nil && len(errs) > 0
+               event.End = time.Now()
+               metrics.Publish(event)
+       }()
+
        metadataReports := metadata.GetMetadataReports()
        if len(metadataReports) == 0 {
-               return nil, perrors.New("can not get mapping in remote cause no 
metadata report instance found")
+               err = perrors.New("can not get mapping in remote cause no 
metadata report instance found")
+               logger.Warnf("[Metadata][Mapping] get failed interface=%s 
group=%s reports=0 err=%v", serviceInterface, DefaultGroup, err)
+               return nil, err
        }
+
        // Attach the listener to the stable primary report only 
(GetMetadataReport uses
        // a deterministic selection: prefer "default", otherwise lexicographic 
first).
        // GetMetadataReports() iterates a map so its order is 
non-deterministic; using
        // i==0 as the anchor would bind the listener to a random backend each 
run.
        primaryReport := metadata.GetMetadataReport()
-       var result *gxset.HashSet
-       var errs []error
-       for _, metadataReport := range metadataReports {
+       for i, metadataReport := range metadataReports {
                var reportListener mapping.MappingListener
                if metadataReport == primaryReport {
                        reportListener = listener
                }
-               set, err := 
metadataReport.GetServiceAppMapping(serviceInterface, DefaultGroup, 
reportListener)
-               if err != nil {
-                       errs = append(errs, err)
+               set, getErr := 
metadataReport.GetServiceAppMapping(serviceInterface, DefaultGroup, 
reportListener)
+               if getErr != nil {
+                       errs = append(errs, getErr)
+                       reportURL := ""
+                       if u := metadataReport.URL(); u != nil {
+                               reportURL = u.Protocol + "://" + u.Address()
+                       }
+                       logger.Warnf("[Metadata][Mapping] %s report %d/%d 
failed interface=%s group=%s url=%s err=%v", operation, i+1, 
len(metadataReports), serviceInterface, DefaultGroup, reportURL, getErr)
                        continue
                }
                if result == nil {
@@ -147,8 +191,13 @@ func (d *ServiceNameMapping) Get(url *common.URL, listener 
mapping.MappingListen
                }
        }
        if result == nil {
-               return nil, errors.Join(errs...)
+               if err = errors.Join(errs...); err != nil {
+                       logger.Warnf("[Metadata][Mapping] %s failed 
interface=%s group=%s reports=%d err=%v", operation, serviceInterface, 
DefaultGroup, len(metadataReports), err)
+                       return nil, err
+               }
+               return nil, nil
        }
+       logger.Debugf("[Metadata][Mapping] %s succeeded interface=%s group=%s 
reports=%d apps=%d", operation, serviceInterface, DefaultGroup, 
len(metadataReports), result.Size())
        return result, nil
 }
 
@@ -158,17 +207,35 @@ func (d *ServiceNameMapping) Get(url *common.URL, 
listener mapping.MappingListen
 // joined together so the caller can see the full failure picture. The
 // intent is to avoid leaving stale entries in any registry due to a transient
 // error in one of the others.
-func (d *ServiceNameMapping) Remove(url *common.URL) error {
+func (d *ServiceNameMapping) Remove(url *common.URL) (err error) {
        serviceInterface := url.GetParam(constant.InterfaceKey, "")
+
+       event := 
metadataMetrics.NewMetadataMetricTimeEvent(metadataMetrics.MetadataMappingRemove)
+       event.Attachment[constant.InterfaceKey] = serviceInterface
+       event.Attachment[constant.GroupKey] = DefaultGroup
+       defer func() {
+               event.Succ = err == nil
+               event.End = time.Now()
+               metrics.Publish(event)
+       }()
+
        metadataReports := metadata.GetMetadataReports()
        if len(metadataReports) == 0 {
-               return perrors.New("can not remove mapping in remote cause no 
metadata report instance found")
+               err = perrors.New("can not remove mapping in remote cause no 
metadata report instance found")
+               logger.Warnf("[Metadata][Mapping] remove failed interface=%s 
group=%s reports=0 err=%v", serviceInterface, DefaultGroup, err)
+               return err
        }
+
        var errs []error
        for _, metadataReport := range metadataReports {
-               if err := 
metadataReport.RemoveServiceAppMappingListener(serviceInterface, DefaultGroup); 
err != nil {
-                       errs = append(errs, err)
+               if removeErr := 
metadataReport.RemoveServiceAppMappingListener(serviceInterface, DefaultGroup); 
removeErr != nil {
+                       errs = append(errs, removeErr)
                }
        }
-       return errors.Join(errs...)
+       if err = errors.Join(errs...); err != nil {
+               logger.Warnf("[Metadata][Mapping] remove failed interface=%s 
group=%s reports=%d err=%v", serviceInterface, DefaultGroup, 
len(metadataReports), err)
+               return err
+       }
+       logger.Debugf("[Metadata][Mapping] remove succeeded interface=%s 
group=%s reports=%d", serviceInterface, DefaultGroup, len(metadataReports))
+       return nil
 }
diff --git a/metadata/mapping/metadata/service_name_mapping_test.go 
b/metadata/mapping/metadata/service_name_mapping_test.go
index fb543c947..7c5b384d7 100644
--- a/metadata/mapping/metadata/service_name_mapping_test.go
+++ b/metadata/mapping/metadata/service_name_mapping_test.go
@@ -19,6 +19,8 @@ package metadata
 
 import (
        "errors"
+       "fmt"
+       "strings"
        "sync"
        "testing"
 )
@@ -26,6 +28,7 @@ import (
 import (
        gxset "github.com/dubbogo/gost/container/set"
        "github.com/dubbogo/gost/gof/observer"
+       gostlogger "github.com/dubbogo/gost/log/logger"
 
        "github.com/stretchr/testify/assert"
        "github.com/stretchr/testify/mock"
@@ -40,6 +43,8 @@ import (
        "dubbo.apache.org/dubbo-go/v3/metadata/info"
        "dubbo.apache.org/dubbo-go/v3/metadata/mapping"
        "dubbo.apache.org/dubbo-go/v3/metadata/report"
+       "dubbo.apache.org/dubbo-go/v3/metrics"
+       metricsMetadata "dubbo.apache.org/dubbo-go/v3/metrics/metadata"
 )
 
 func TestGetNameMappingInstance(t *testing.T) {
@@ -62,6 +67,45 @@ func TestNoReportInstance(t *testing.T) {
        require.Error(t, err, "test Remove with no report instance")
 }
 
+func TestServiceNameMappingNoReportMetersPerBusinessOperation(t *testing.T) {
+       metadata.ClearMetadataReportInstances()
+       t.Cleanup(metadata.ClearMetadataReportInstances)
+       serviceNameMappingOnce = sync.Once{}
+       serviceNameMappingInstance = nil
+
+       ch := make(chan metrics.MetricsEvent, 10)
+       metrics.Subscribe(constant.MetricsMetadata, ch)
+       defer metrics.Unsubscribe(constant.MetricsMetadata)
+
+       ins := GetNameMappingInstance()
+       serviceUrl := common.NewURLWithOptions(
+               common.WithInterface("org.example.NoReportService"),
+               common.WithParamsValue(constant.ApplicationKey, 
"no-report-app"),
+       )
+
+       err := ins.Map(serviceUrl)
+       require.Error(t, err)
+       wantAttachment := mappingAttachment("org.example.NoReportService")
+       wantAttachment[constant.ApplicationKey] = "no-report-app"
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRegister, false, false, wantAttachment)
+       assert.Empty(t, ch)
+
+       _, err = ins.Get(serviceUrl, nil)
+       require.Error(t, err)
+       assertMappingMetricEvent(t, <-ch, metricsMetadata.MetadataMappingGet, 
false, false, mappingAttachment("org.example.NoReportService"))
+       assert.Empty(t, ch)
+
+       _, err = ins.Get(serviceUrl, &listener{})
+       require.Error(t, err)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingListen, false, false, 
mappingAttachment("org.example.NoReportService"))
+       assert.Empty(t, ch)
+
+       err = ins.Remove(serviceUrl)
+       require.Error(t, err)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRemove, false, false, 
mappingAttachment("org.example.NoReportService"))
+       assert.Empty(t, ch)
+}
+
 func TestServiceNameMappingGet(t *testing.T) {
        ins := GetNameMappingInstance()
        lis := &listener{}
@@ -247,6 +291,207 @@ func 
TestServiceNameMappingRemoveContinuesAfterPartialFailure(t *testing.T) {
        r2.AssertExpectations(t)
 }
 
+func TestServiceNameMappingGetMetersPerBusinessOperation(t *testing.T) {
+       metadata.ClearMetadataReportInstances()
+       t.Cleanup(metadata.ClearMetadataReportInstances)
+       serviceNameMappingOnce = sync.Once{}
+       serviceNameMappingInstance = nil
+
+       r1 := initMockWithId(t, "reg-get-a")
+       r2 := initMockWithId(t, "reg-get-b")
+
+       ch := make(chan metrics.MetricsEvent, 10)
+       metrics.Subscribe(constant.MetricsMetadata, ch)
+       defer metrics.Unsubscribe(constant.MetricsMetadata)
+
+       ins := GetNameMappingInstance()
+       serviceUrl := common.NewURLWithOptions(
+               common.WithInterface("org.example.MeteredService"),
+       )
+
+       // a listen fans out to every report but must count as a single listen 
event
+       r1.On("GetServiceAppMapping").Return(gxset.NewSet("app-a"), nil).Once()
+       r2.On("GetServiceAppMapping").Return(gxset.NewSet("app-b"), nil).Once()
+       apps, err := ins.Get(serviceUrl, &listener{})
+       require.NoError(t, err)
+       assert.Equal(t, 2, apps.Size())
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingListen, true, false, 
mappingAttachment("org.example.MeteredService"))
+       assert.Empty(t, ch)
+
+       // a plain get counts as a single get event
+       r1.On("GetServiceAppMapping").Return(gxset.NewSet("app-a"), nil).Once()
+       r2.On("GetServiceAppMapping").Return(gxset.NewSet("app-b"), nil).Once()
+       apps, err = ins.Get(serviceUrl, nil)
+       require.NoError(t, err)
+       assert.Equal(t, 2, apps.Size())
+       assertMappingMetricEvent(t, <-ch, metricsMetadata.MetadataMappingGet, 
true, false, mappingAttachment("org.example.MeteredService"))
+       assert.Empty(t, ch)
+
+       // partial failure is still a successful business operation (err == 
nil),
+       // but the event must carry the partial flag
+       getErr := errors.New("r1 failure")
+       r1.On("GetServiceAppMapping").Return(gxset.NewSet(), getErr).Once()
+       r2.On("GetServiceAppMapping").Return(gxset.NewSet("app-b"), nil).Once()
+       apps, err = ins.Get(serviceUrl, &listener{})
+       require.NoError(t, err)
+       assert.False(t, apps.Empty())
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingListen, true, true, 
mappingAttachment("org.example.MeteredService"))
+       assert.Empty(t, ch)
+
+       // all reports failing produces a single failed event
+       r1.On("GetServiceAppMapping").Return(gxset.NewSet(), getErr).Once()
+       r2.On("GetServiceAppMapping").Return(gxset.NewSet(), errors.New("r2 
failure")).Once()
+       _, err = ins.Get(serviceUrl, &listener{})
+       require.Error(t, err)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingListen, false, false, 
mappingAttachment("org.example.MeteredService"))
+       assert.Empty(t, ch)
+
+       r1.AssertExpectations(t)
+       r2.AssertExpectations(t)
+}
+
+func TestServiceNameMappingGetDoesNotLogMetadataReportSecrets(t *testing.T) {
+       metadata.ClearMetadataReportInstances()
+       t.Cleanup(metadata.ClearMetadataReportInstances)
+       serviceNameMappingOnce = sync.Once{}
+       serviceNameMappingInstance = nil
+
+       const (
+               username    = "review-user"
+               password    = "review-secret"
+               accessToken = "review-access-token"
+       )
+       reportURL, err := common.NewURL("nacos://" + username + ":" + password 
+ "@127.0.0.1:8848?access-token=" + accessToken)
+       require.NoError(t, err)
+       mockReport := initMockWithId(t, "reg-secret")
+       mockReport.reportURL = reportURL
+       mockReport.On("GetServiceAppMapping").Return(gxset.NewSet(), 
errors.New("get failure")).Once()
+
+       previousLogger := gostlogger.GetLogger()
+       capture := &captureMappingWarnLogger{Logger: previousLogger}
+       gostlogger.SetLogger(capture)
+       t.Cleanup(func() {
+               gostlogger.SetLogger(previousLogger)
+       })
+
+       serviceURL := 
common.NewURLWithOptions(common.WithInterface("org.example.SecretService"))
+       _, err = GetNameMappingInstance().Get(serviceURL, nil)
+       require.Error(t, err)
+
+       logOutput := capture.String()
+       assert.Contains(t, logOutput, "url=nacos://127.0.0.1:8848")
+       assert.NotContains(t, logOutput, username)
+       assert.NotContains(t, logOutput, password)
+       assert.NotContains(t, logOutput, accessToken)
+       mockReport.AssertExpectations(t)
+}
+
+func TestServiceNameMappingMapMetersPerBusinessOperation(t *testing.T) {
+       metadata.ClearMetadataReportInstances()
+       t.Cleanup(metadata.ClearMetadataReportInstances)
+       serviceNameMappingOnce = sync.Once{}
+       serviceNameMappingInstance = nil
+
+       r1 := initMockWithId(t, "reg-map-a")
+       r2 := initMockWithId(t, "reg-map-b")
+
+       ch := make(chan metrics.MetricsEvent, 10)
+       metrics.Subscribe(constant.MetricsMetadata, ch)
+       defer metrics.Unsubscribe(constant.MetricsMetadata)
+
+       ins := GetNameMappingInstance()
+       serviceUrl := common.NewURLWithOptions(
+               common.WithInterface("org.example.RegisterService"),
+               common.WithParamsValue(constant.ApplicationKey, "reg-app"),
+       )
+       wantAttachment := mappingAttachment("org.example.RegisterService")
+       wantAttachment[constant.ApplicationKey] = "reg-app"
+
+       // all reports succeed: one register event
+       r1.On("RegisterServiceAppMapping").Return(nil).Once()
+       r2.On("RegisterServiceAppMapping").Return(nil).Once()
+       err := ins.Map(serviceUrl)
+       require.NoError(t, err)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRegister, true, false, wantAttachment)
+       assert.Empty(t, ch)
+
+       // first report fails: Map stops immediately, one failed register event
+       regErr := errors.New("r1 failure")
+       // report iteration order is non-deterministic, so both may be the 
first to
+       // fail; the call-count assertion below pins down that only one is 
invoked.
+       r1.On("RegisterServiceAppMapping").Return(regErr).Maybe()
+       r2.On("RegisterServiceAppMapping").Return(regErr).Maybe()
+       callsBefore := len(r1.Calls) + len(r2.Calls)
+       err = ins.Map(serviceUrl)
+       require.Error(t, err)
+       require.Equal(t, callsBefore+1, len(r1.Calls)+len(r2.Calls), "Map must 
stop at the first failing report")
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRegister, false, false, wantAttachment)
+       assert.Empty(t, ch)
+
+       r1.AssertExpectations(t)
+       r2.AssertExpectations(t)
+}
+
+func TestServiceNameMappingRemoveMetersPerBusinessOperation(t *testing.T) {
+       metadata.ClearMetadataReportInstances()
+       t.Cleanup(metadata.ClearMetadataReportInstances)
+       serviceNameMappingOnce = sync.Once{}
+       serviceNameMappingInstance = nil
+
+       r1 := initMockWithId(t, "reg-remove-a")
+       r2 := initMockWithId(t, "reg-remove-b")
+
+       ch := make(chan metrics.MetricsEvent, 10)
+       metrics.Subscribe(constant.MetricsMetadata, ch)
+       defer metrics.Unsubscribe(constant.MetricsMetadata)
+
+       ins := GetNameMappingInstance()
+       serviceUrl := common.NewURLWithOptions(
+               common.WithInterface("org.example.RemoveService"),
+       )
+       wantAttachment := mappingAttachment("org.example.RemoveService")
+
+       // all reports succeed: one remove event
+       r1.On("RemoveServiceAppMappingListener").Return(nil).Once()
+       r2.On("RemoveServiceAppMappingListener").Return(nil).Once()
+       err := ins.Remove(serviceUrl)
+       require.NoError(t, err)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRemove, true, false, wantAttachment)
+       assert.Empty(t, ch)
+
+       // partial failure: best-effort still calls every report, one failed 
remove event
+       removeErr := errors.New("r1 failure")
+       r1.On("RemoveServiceAppMappingListener").Return(removeErr).Once()
+       r2.On("RemoveServiceAppMappingListener").Return(nil).Once()
+       err = ins.Remove(serviceUrl)
+       require.ErrorIs(t, err, removeErr)
+       assertMappingMetricEvent(t, <-ch, 
metricsMetadata.MetadataMappingRemove, false, false, wantAttachment)
+       assert.Empty(t, ch)
+
+       r1.AssertExpectations(t)
+       r2.AssertExpectations(t)
+}
+
+func mappingAttachment(serviceInterface string) map[string]string {
+       return map[string]string{
+               constant.InterfaceKey: serviceInterface,
+               constant.GroupKey:     DefaultGroup,
+       }
+}
+
+func assertMappingMetricEvent(t *testing.T, metricEvent metrics.MetricsEvent, 
name metricsMetadata.MetricName, succ bool, partial bool, wantAttachment 
map[string]string) {
+       t.Helper()
+       assert.Equal(t, constant.MetricsMetadata, metricEvent.Type())
+       event, ok := metricEvent.(*metricsMetadata.MetadataMetricEvent)
+       assert.True(t, ok)
+       assert.Equal(t, name, event.Name)
+       assert.Equal(t, succ, event.Succ)
+       assert.Equal(t, partial, event.Partial)
+       assert.Equal(t, wantAttachment, event.Attachment)
+       assert.NotNil(t, event.Start)
+       assert.NotNil(t, event.End)
+}
+
 type listener struct {
 }
 
@@ -259,6 +504,7 @@ func (l listener) Stop() {
 
 type mockMetadataReport struct {
        mock.Mock
+       reportURL *common.URL
 }
 
 func (m *mockMetadataReport) CreateMetadataReport(*common.URL) 
report.MetadataReport {
@@ -301,6 +547,27 @@ func (m *mockMetadataReport) ListAppRevisions(string) 
([]report.AppRevision, err
 }
 
 func (m *mockMetadataReport) URL() *common.URL {
+       if m.reportURL != nil {
+               return m.reportURL
+       }
        u, _ := common.NewURL("mock://127.0.0.1:8848")
        return u
 }
+
+type captureMappingWarnLogger struct {
+       gostlogger.Logger
+       mu    sync.Mutex
+       warns []string
+}
+
+func (l *captureMappingWarnLogger) Warnf(template string, args ...any) {
+       l.mu.Lock()
+       defer l.mu.Unlock()
+       l.warns = append(l.warns, fmt.Sprintf(template, args...))
+}
+
+func (l *captureMappingWarnLogger) String() string {
+       l.mu.Lock()
+       defer l.mu.Unlock()
+       return strings.Join(l.warns, "\n")
+}
diff --git a/metadata/report_instance_test.go b/metadata/report_instance_test.go
index cae9c0456..8b2cefab5 100644
--- a/metadata/report_instance_test.go
+++ b/metadata/report_instance_test.go
@@ -124,13 +124,24 @@ func TestDelegateMetadataReportGetServiceAppMapping(t 
*testing.T) {
        mockReport := new(mockMetadataReport)
        defer mockReport.AssertExpectations(t)
        delegate := &DelegateMetadataReport{instance: mockReport}
-       t.Run("normal", func(t *testing.T) {
+       t.Run("get normal", func(t *testing.T) {
+               mockReport.On("GetServiceAppMapping").Return(gxset.NewSet(), 
nil).Once()
+               got, err := delegate.GetServiceAppMapping("dubbo", "dev", nil)
+               require.NoError(t, err)
+               assert.True(t, got.Empty())
+       })
+       t.Run("get error", func(t *testing.T) {
+               mockReport.On("GetServiceAppMapping").Return(gxset.NewSet(), 
errors.New("mock error")).Once()
+               _, err := delegate.GetServiceAppMapping("dubbo", "dev", nil)
+               require.Error(t, err)
+       })
+       t.Run("listen normal", func(t *testing.T) {
                mockReport.On("GetServiceAppMapping").Return(gxset.NewSet(), 
nil).Once()
                got, err := delegate.GetServiceAppMapping("dubbo", "dev", 
&listener{})
                require.NoError(t, err)
                assert.True(t, got.Empty())
        })
-       t.Run("error", func(t *testing.T) {
+       t.Run("listen error", func(t *testing.T) {
                mockReport.On("GetServiceAppMapping").Return(gxset.NewSet(), 
errors.New("mock error")).Once()
                _, err := delegate.GetServiceAppMapping("dubbo", "dev", 
&listener{})
                require.Error(t, err)
diff --git a/metrics/metadata/collector.go b/metrics/metadata/collector.go
index 8a08e0ff4..d0d0a9aa8 100644
--- a/metrics/metadata/collector.go
+++ b/metrics/metadata/collector.go
@@ -58,6 +58,14 @@ func (c *MetadataMetricCollector) start() {
                                        c.handleMetadataSub(event)
                                case SubscribeServiceRt:
                                        c.handleSubscribeService(event)
+                               case MetadataMappingRegister:
+                                       c.handleMetadataMappingRegister(event)
+                               case MetadataMappingGet:
+                                       c.handleMetadataMappingGet(event)
+                               case MetadataMappingListen:
+                                       c.handleMetadataMappingListen(event)
+                               case MetadataMappingRemove:
+                                       c.handleMetadataMappingRemove(event)
                                default:
                                }
                        }
@@ -88,9 +96,56 @@ func (c *MetadataMetricCollector) 
handleSubscribeService(event *MetadataMetricEv
        c.R.Rt(metrics.NewMetricId(subscribeServiceRt, level), 
&metrics.RtOpts{}).Observe(event.CostMs())
 }
 
+func (c *MetadataMetricCollector) handleMetadataMappingRegister(event 
*MetadataMetricEvent) {
+       level := newMetadataMappingMetricLevel(event.Attachment)
+       c.StateCount(metadataMappingRegisterNum, 
metadataMappingRegisterSucceed, metadataMappingRegisterFailed, level, 
event.Succ)
+       c.R.Rt(metrics.NewMetricId(metadataMappingRegisterRt, level), 
&metrics.RtOpts{}).Observe(event.CostMs())
+}
+
+func (c *MetadataMetricCollector) handleMetadataMappingGet(event 
*MetadataMetricEvent) {
+       level := newMetadataMappingMetricLevel(event.Attachment)
+       c.StateCount(metadataMappingGetNum, metadataMappingGetSucceed, 
metadataMappingGetFailed, level, event.Succ && !event.Partial)
+       c.R.Rt(metrics.NewMetricId(metadataMappingGetRt, level), 
&metrics.RtOpts{}).Observe(event.CostMs())
+}
+
+func (c *MetadataMetricCollector) handleMetadataMappingListen(event 
*MetadataMetricEvent) {
+       level := newMetadataMappingMetricLevel(event.Attachment)
+       c.StateCount(metadataMappingListenNum, metadataMappingListenSucceed, 
metadataMappingListenFailed, level, event.Succ && !event.Partial)
+       c.R.Rt(metrics.NewMetricId(metadataMappingListenRt, level), 
&metrics.RtOpts{}).Observe(event.CostMs())
+}
+
+func (c *MetadataMetricCollector) handleMetadataMappingRemove(event 
*MetadataMetricEvent) {
+       level := newMetadataMappingMetricLevel(event.Attachment)
+       c.StateCount(metadataMappingRemoveNum, metadataMappingRemoveSucceed, 
metadataMappingRemoveFailed, level, event.Succ)
+       c.R.Rt(metrics.NewMetricId(metadataMappingRemoveRt, level), 
&metrics.RtOpts{}).Observe(event.CostMs())
+}
+
+type metadataMappingMetricLevel struct {
+       *metrics.ApplicationMetricLevel
+       attachment map[string]string
+}
+
+func newMetadataMappingMetricLevel(attachment map[string]string) 
metadataMappingMetricLevel {
+       return metadataMappingMetricLevel{
+               ApplicationMetricLevel: metrics.GetApplicationLevel(),
+               attachment:             attachment,
+       }
+}
+
+func (m metadataMappingMetricLevel) Tags() map[string]string {
+       tags := m.ApplicationMetricLevel.Tags()
+       tags[constant.TagInterface] = m.attachment[constant.InterfaceKey]
+       tags[constant.TagGroup] = m.attachment[constant.GroupKey]
+       if app := m.attachment[constant.ApplicationKey]; app != "" {
+               tags[constant.TagApplicationName] = app
+       }
+       return tags
+}
+
 type MetadataMetricEvent struct {
        Name       MetricName
        Succ       bool
+       Partial    bool
        Start      time.Time
        End        time.Time
        Attachment map[string]string
diff --git a/metrics/metadata/collector_test.go 
b/metrics/metadata/collector_test.go
index 59c29ba90..4e9c4125f 100644
--- a/metrics/metadata/collector_test.go
+++ b/metrics/metadata/collector_test.go
@@ -18,16 +18,20 @@
 package metadata
 
 import (
+       "fmt"
+       "sync"
        "testing"
        "time"
 )
 
 import (
        "github.com/stretchr/testify/assert"
+       "github.com/stretchr/testify/require"
 )
 
 import (
        "dubbo.apache.org/dubbo-go/v3/common/constant"
+       "dubbo.apache.org/dubbo-go/v3/metrics"
 )
 
 func TestMetadataMetricEventType(t *testing.T) {
@@ -62,3 +66,238 @@ func TestNewMetadataMetricTimeEvent(t *testing.T) {
        assert.NotNil(t, event.Attachment)
        assert.Empty(t, event.Attachment)
 }
+
+func TestMetadataMetricCollectorHandleMapping(t *testing.T) {
+       tests := []struct {
+               name      string
+               eventName MetricName
+               handler   func(*MetadataMetricCollector, *MetadataMetricEvent)
+               prefix    string
+       }{
+               {
+                       name:      "register",
+                       eventName: MetadataMappingRegister,
+                       handler:   
(*MetadataMetricCollector).handleMetadataMappingRegister,
+                       prefix:    "dubbo_metadata_mapping_register",
+               },
+               {
+                       name:      "get",
+                       eventName: MetadataMappingGet,
+                       handler:   
(*MetadataMetricCollector).handleMetadataMappingGet,
+                       prefix:    "dubbo_metadata_mapping_get",
+               },
+               {
+                       name:      "listen",
+                       eventName: MetadataMappingListen,
+                       handler:   
(*MetadataMetricCollector).handleMetadataMappingListen,
+                       prefix:    "dubbo_metadata_mapping_listen",
+               },
+               {
+                       name:      "remove",
+                       eventName: MetadataMappingRemove,
+                       handler:   
(*MetadataMetricCollector).handleMetadataMappingRemove,
+                       prefix:    "dubbo_metadata_mapping_remove",
+               },
+       }
+
+       for _, tt := range tests {
+               for _, succ := range []bool{true, false} {
+                       t.Run(fmt.Sprintf("%s/succ=%v", tt.name, succ), func(t 
*testing.T) {
+                               registry := newMockMetricRegistry()
+                               collector := 
&MetadataMetricCollector{BaseCollector: metrics.BaseCollector{R: registry}}
+                               event := 
NewMetadataMetricTimeEvent(tt.eventName)
+                               event.End = event.Start.Add(10 * 
time.Millisecond)
+                               event.Succ = succ
+                               event.Attachment[constant.InterfaceKey] = 
"interfaceName"
+                               event.Attachment[constant.GroupKey] = "group"
+                               event.Attachment[constant.ApplicationKey] = 
"application"
+
+                               tt.handler(collector, event)
+
+                               assert.InDelta(t, 1.0, 
registry.counters[tt.prefix+"_num_total"], 0.000001)
+                               if succ {
+                                       assert.InDelta(t, 1.0, 
registry.counters[tt.prefix+"_num_succeed_total"], 0.000001)
+                                       assert.NotContains(t, 
registry.counters, tt.prefix+"_num_failed_total")
+                               } else {
+                                       assert.InDelta(t, 1.0, 
registry.counters[tt.prefix+"_num_failed_total"], 0.000001)
+                                       assert.NotContains(t, 
registry.counters, tt.prefix+"_num_succeed_total")
+                               }
+                               assert.Equal(t, []float64{10.0}, 
registry.rts[tt.prefix+"_rt_milliseconds"])
+
+                               id := registry.ids[tt.prefix+"_num_total"]
+                               assert.Equal(t, "interfaceName", 
id.Tags[constant.TagInterface])
+                               assert.Equal(t, "group", 
id.Tags[constant.TagGroup])
+                               assert.Equal(t, "application", 
id.Tags[constant.TagApplicationName])
+                       })
+               }
+
+               // a partial success (some reports failed, some succeeded) must 
be counted
+               // as failed, not succeed
+               for _, tt := range []struct {
+                       name      string
+                       eventName MetricName
+                       handler   func(*MetadataMetricCollector, 
*MetadataMetricEvent)
+                       prefix    string
+               }{
+                       {"get", MetadataMappingGet, 
(*MetadataMetricCollector).handleMetadataMappingGet, 
"dubbo_metadata_mapping_get"},
+                       {"listen", MetadataMappingListen, 
(*MetadataMetricCollector).handleMetadataMappingListen, 
"dubbo_metadata_mapping_listen"},
+               } {
+                       t.Run(fmt.Sprintf("%s/partial", tt.name), func(t 
*testing.T) {
+                               registry := newMockMetricRegistry()
+                               collector := 
&MetadataMetricCollector{BaseCollector: metrics.BaseCollector{R: registry}}
+                               event := 
NewMetadataMetricTimeEvent(tt.eventName)
+                               event.End = event.Start.Add(10 * 
time.Millisecond)
+                               event.Succ = true
+                               event.Partial = true
+                               event.Attachment[constant.InterfaceKey] = 
"interfaceName"
+                               event.Attachment[constant.GroupKey] = "group"
+                               event.Attachment[constant.ApplicationKey] = 
"application"
+
+                               tt.handler(collector, event)
+
+                               assert.InDelta(t, 1.0, 
registry.counters[tt.prefix+"_num_total"], 0.000001)
+                               assert.InDelta(t, 1.0, 
registry.counters[tt.prefix+"_num_failed_total"], 0.000001)
+                               assert.NotContains(t, registry.counters, 
tt.prefix+"_num_succeed_total")
+                               assert.Equal(t, []float64{10.0}, 
registry.rts[tt.prefix+"_rt_milliseconds"])
+                       })
+               }
+       }
+}
+
+type mockMetricRegistry struct {
+       mu       sync.Mutex
+       counters map[string]float64
+       rts      map[string][]float64
+       ids      map[string]*metrics.MetricId
+}
+
+func newMockMetricRegistry() *mockMetricRegistry {
+       return &mockMetricRegistry{
+               counters: make(map[string]float64),
+               rts:      make(map[string][]float64),
+               ids:      make(map[string]*metrics.MetricId),
+       }
+}
+
+func (m *mockMetricRegistry) Counter(id *metrics.MetricId) 
metrics.CounterMetric {
+       m.mu.Lock()
+       defer m.mu.Unlock()
+       m.ids[id.Name] = id
+       return &mockCounterMetric{m: m, name: id.Name}
+}
+
+func (m *mockMetricRegistry) Rt(id *metrics.MetricId, _ *metrics.RtOpts) 
metrics.ObservableMetric {
+       m.mu.Lock()
+       defer m.mu.Unlock()
+       m.ids[id.Name] = id
+       return &mockRtMetric{m: m, name: id.Name}
+}
+
+func (m *mockMetricRegistry) Gauge(id *metrics.MetricId) metrics.GaugeMetric {
+       return nil
+}
+
+func (m *mockMetricRegistry) Histogram(id *metrics.MetricId) 
metrics.ObservableMetric {
+       return nil
+}
+
+func (m *mockMetricRegistry) Summary(id *metrics.MetricId) 
metrics.ObservableMetric {
+       return nil
+}
+
+func (m *mockMetricRegistry) Export() {}
+
+type mockCounterMetric struct {
+       m    *mockMetricRegistry
+       name string
+}
+
+func (c *mockCounterMetric) Inc() {
+       c.m.mu.Lock()
+       defer c.m.mu.Unlock()
+       c.m.counters[c.name]++
+}
+
+func (c *mockCounterMetric) Add(v float64) {
+       c.m.mu.Lock()
+       defer c.m.mu.Unlock()
+       c.m.counters[c.name] += v
+}
+
+type mockRtMetric struct {
+       m    *mockMetricRegistry
+       name string
+}
+
+func (r *mockRtMetric) Observe(v float64) {
+       r.m.mu.Lock()
+       defer r.m.mu.Unlock()
+       r.m.rts[r.name] = append(r.m.rts[r.name], v)
+}
+
+// TestMetadataMetricCollectorPublishChain covers the production dispatch path:
+// a started collector subscribes to the event bus, and events published via
+// metrics.Publish must reach the registry. Removing any of the mapping switch
+// cases in start() must make this test fail.
+func TestMetadataMetricCollectorPublishChain(t *testing.T) {
+       // start() subscribes the package-level channel. Swap in a per-test
+       // channel so the teardown does not permanently close the production
+       // channel: Unsubscribe closes the registered channel, and 
re-subscribing
+       // a closed channel on the next run (go test -count=2) would panic on
+       // publish. The original channel is restored afterwards.
+       originalCh := ch
+       testCh := make(chan metrics.MetricsEvent, 10)
+       ch = testCh
+       defer func() {
+               metrics.Unsubscribe(constant.MetricsMetadata)
+               ch = originalCh
+       }()
+
+       registry := newMockMetricRegistry()
+       collector := &MetadataMetricCollector{BaseCollector: 
metrics.BaseCollector{R: registry}}
+       collector.start()
+
+       publish := func(name MetricName, succ bool) {
+               event := NewMetadataMetricTimeEvent(name)
+               event.End = event.Start.Add(10 * time.Millisecond)
+               event.Succ = succ
+               event.Attachment[constant.InterfaceKey] = "interfaceName"
+               event.Attachment[constant.GroupKey] = "group"
+               event.Attachment[constant.ApplicationKey] = "application"
+               metrics.Publish(event)
+       }
+
+       publish(MetadataMappingRegister, true)
+       publish(MetadataMappingGet, true)
+       publish(MetadataMappingListen, false)
+       publish(MetadataMappingRemove, true)
+
+       prefixes := []string{
+               "dubbo_metadata_mapping_register",
+               "dubbo_metadata_mapping_get",
+               "dubbo_metadata_mapping_listen",
+               "dubbo_metadata_mapping_remove",
+       }
+       require.Eventually(t, func() bool {
+               registry.mu.Lock()
+               defer registry.mu.Unlock()
+               if 
registry.counters["dubbo_metadata_mapping_register_num_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_get_num_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_listen_num_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_remove_num_total"] != 1 {
+                       return false
+               }
+               if 
registry.counters["dubbo_metadata_mapping_register_num_succeed_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_get_num_succeed_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_listen_num_failed_total"] != 1 ||
+                       
registry.counters["dubbo_metadata_mapping_remove_num_succeed_total"] != 1 {
+                       return false
+               }
+               for _, p := range prefixes {
+                       if len(registry.rts[p+"_rt_milliseconds"]) == 0 {
+                               return false
+                       }
+               }
+               return true
+       }, 5*time.Second, 10*time.Millisecond)
+}
diff --git a/metrics/metadata/metric_set.go b/metrics/metadata/metric_set.go
index 20d76a36f..c760e201c 100644
--- a/metrics/metadata/metric_set.go
+++ b/metrics/metadata/metric_set.go
@@ -31,16 +31,28 @@ const (
        // SubscribeRt
        // StoreProviderInterfaceRt
        SubscribeServiceRt
+       MetadataMappingRegister
+       MetadataMappingGet
+       MetadataMappingListen
+       MetadataMappingRemove
 )
 
 const (
-       dubboMetadataPush             = "dubbo_metadata_push_num"
-       dubboPushRt                   = "dubbo_push_rt_milliseconds"
-       dubboMetadataSubscribe        = "dubbo_metadata_subscribe_num"
-       dubboSubscribeRt              = "dubbo_subscribe_rt_milliseconds"
-       dubboMetadataStoreProvider    = "dubbo_metadata_store_provider"
-       dubboStoreProviderInterfaceRt = 
"dubbo_store_provider_interface_rt_milliseconds"
-       dubboSubscribeServiceRt       = 
"dubbo_subscribe_service_rt_milliseconds"
+       dubboMetadataPush              = "dubbo_metadata_push_num"
+       dubboPushRt                    = "dubbo_push_rt_milliseconds"
+       dubboMetadataSubscribe         = "dubbo_metadata_subscribe_num"
+       dubboSubscribeRt               = "dubbo_subscribe_rt_milliseconds"
+       dubboMetadataStoreProvider     = "dubbo_metadata_store_provider"
+       dubboStoreProviderInterfaceRt  = 
"dubbo_store_provider_interface_rt_milliseconds"
+       dubboSubscribeServiceRt        = 
"dubbo_subscribe_service_rt_milliseconds"
+       dubboMetadataMappingRegister   = "dubbo_metadata_mapping_register_num"
+       dubboMetadataMappingRegisterRt = 
"dubbo_metadata_mapping_register_rt_milliseconds"
+       dubboMetadataMappingGet        = "dubbo_metadata_mapping_get_num"
+       dubboMetadataMappingGetRt      = 
"dubbo_metadata_mapping_get_rt_milliseconds"
+       dubboMetadataMappingListen     = "dubbo_metadata_mapping_listen_num"
+       dubboMetadataMappingListenRt   = 
"dubbo_metadata_mapping_listen_rt_milliseconds"
+       dubboMetadataMappingRemove     = "dubbo_metadata_mapping_remove_num"
+       dubboMetadataMappingRemoveRt   = 
"dubbo_metadata_mapping_remove_rt_milliseconds"
 )
 
 const (
@@ -84,4 +96,24 @@ var (
        storeProviderInterfaceRt = 
metrics.NewMetricKey(dubboStoreProviderInterfaceRt, "Store Provider Interface 
Time")
 
        subscribeServiceRt = metrics.NewMetricKey(dubboSubscribeServiceRt, 
"Subscribe Service Time")
+
+       metadataMappingRegisterNum     = 
metrics.NewMetricKey(dubboMetadataMappingRegister+totalSuffix, "Total Metadata 
Mapping Register Num")
+       metadataMappingRegisterSucceed = 
metrics.NewMetricKey(dubboMetadataMappingRegister+succSuffix, "Succeed Metadata 
Mapping Register Num")
+       metadataMappingRegisterFailed  = 
metrics.NewMetricKey(dubboMetadataMappingRegister+failedSuffix, "Failed 
Metadata Mapping Register Num")
+       metadataMappingRegisterRt      = 
metrics.NewMetricKey(dubboMetadataMappingRegisterRt, "Metadata Mapping Register 
Time")
+
+       metadataMappingGetNum     = 
metrics.NewMetricKey(dubboMetadataMappingGet+totalSuffix, "Total Metadata 
Mapping Get Num")
+       metadataMappingGetSucceed = 
metrics.NewMetricKey(dubboMetadataMappingGet+succSuffix, "Succeed Metadata 
Mapping Get Num")
+       metadataMappingGetFailed  = 
metrics.NewMetricKey(dubboMetadataMappingGet+failedSuffix, "Failed Metadata 
Mapping Get Num")
+       metadataMappingGetRt      = 
metrics.NewMetricKey(dubboMetadataMappingGetRt, "Metadata Mapping Get Time")
+
+       metadataMappingListenNum     = 
metrics.NewMetricKey(dubboMetadataMappingListen+totalSuffix, "Total Metadata 
Mapping Listen Num")
+       metadataMappingListenSucceed = 
metrics.NewMetricKey(dubboMetadataMappingListen+succSuffix, "Succeed Metadata 
Mapping Listen Num")
+       metadataMappingListenFailed  = 
metrics.NewMetricKey(dubboMetadataMappingListen+failedSuffix, "Failed Metadata 
Mapping Listen Num")
+       metadataMappingListenRt      = 
metrics.NewMetricKey(dubboMetadataMappingListenRt, "Metadata Mapping Listen 
Time")
+
+       metadataMappingRemoveNum     = 
metrics.NewMetricKey(dubboMetadataMappingRemove+totalSuffix, "Total Metadata 
Mapping Remove Num")
+       metadataMappingRemoveSucceed = 
metrics.NewMetricKey(dubboMetadataMappingRemove+succSuffix, "Succeed Metadata 
Mapping Remove Num")
+       metadataMappingRemoveFailed  = 
metrics.NewMetricKey(dubboMetadataMappingRemove+failedSuffix, "Failed Metadata 
Mapping Remove Num")
+       metadataMappingRemoveRt      = 
metrics.NewMetricKey(dubboMetadataMappingRemoveRt, "Metadata Mapping Remove 
Time")
 )

Reply via email to