Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions common/metrics/defs.go
Original file line number Diff line number Diff line change
Expand Up @@ -3283,11 +3283,11 @@ var MetricDefs = map[ServiceIdx]map[MetricIdx]metricDefinition{
PersistenceEmptyResponseCounter: {metricName: "persistence_empty_response", metricType: Counter},
PersistenceResponseRowSize: {metricName: "persistence_response_row_size", metricType: Histogram, buckets: ResponseRowSizeBuckets},
PersistenceResponsePayloadSize: {metricName: "persistence_response_payload_size", metricType: Histogram, buckets: ResponsePayloadSizeBuckets},
PersistenceRequestsPerDomain: {metricName: "persistence_requests_per_domain", metricRollupName: "persistence_requests", metricType: Counter},
PersistenceRequestsPerDomain: {metricName: "persistence_requests_per_domain", metricRollupName: "persistence_requests_rollup", metricType: Counter},
PersistenceRequestsPerShard: {metricName: "persistence_requests_per_shard", metricType: Counter},
PersistenceFailuresPerDomain: {metricName: "persistence_errors_per_domain", metricRollupName: "persistence_errors", metricType: Counter},
PersistenceLatencyPerDomain: {metricName: "persistence_latency_per_domain", metricRollupName: "persistence_latency", metricType: Timer},
PersistenceLatencyPerDomainHistogram: {metricName: "persistence_latency_per_domain_ns", metricRollupName: "persistence_latency_ns", metricType: Histogram, exponentialBuckets: Default1ms100s},
PersistenceLatencyPerDomain: {metricName: "persistence_latency_per_domain", metricRollupName: "persistence_latency_rollup", metricType: Timer},
PersistenceLatencyPerDomainHistogram: {metricName: "persistence_latency_per_domain_ns", metricRollupName: "persistence_latency_ns_rollup", metricType: Histogram, exponentialBuckets: Default1ms100s},
PersistenceLatencyPerShard: {metricName: "persistence_latency_per_shard", metricType: Timer},
PersistenceLatencyPerShardHistogram: {metricName: "persistence_latency_per_shard_ns", metricType: Histogram, exponentialBuckets: Low1ms100s},
PersistenceErrShardExistsCounterPerDomain: {metricName: "persistence_errors_shard_exists_per_domain", metricRollupName: "persistence_errors_shard_exists", metricType: Counter},
Expand Down
44 changes: 39 additions & 5 deletions common/persistence/wrappers/metered/base.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,35 @@ import (
"github.com/uber/cadence/common/types"
)

// domainTagKey is the metric-tag key produced by metrics.DomainTag.
// It is captured once at init so call() can decide between the overall
// (persistence_requests) and per-domain (persistence_requests_per_domain)
// metrics by inspecting tag content rather than tag count.
var domainTagKey = metrics.DomainTag("").Key()

// taskCategoryTagKey is the metric-tag key produced by metrics.TaskCategoryTag.
// Base persistence metrics use task_category="none" for operations where the
// concept does not apply, keeping label keys consistent with history-task calls.
var taskCategoryTagKey = metrics.TaskCategoryTag("").Key()

func hasDomainTag(tags []metrics.Tag) bool {
for _, t := range tags {
if t.Key() == domainTagKey {
return true
}
}
return false
}

func ensureTaskCategoryTag(tags []metrics.Tag) []metrics.Tag {
for _, t := range tags {
if t.Key() == taskCategoryTagKey {
return tags
}
}
return append(tags, metrics.TaskCategoryTag("none"))
}

type base struct {
metricClient metrics.Client
logger log.Logger
Expand Down Expand Up @@ -133,16 +162,20 @@ func (p *base) recordLatencyHistogram(scope metrics.ScopeIdx, duration time.Dura
}

func (p *base) call(scope metrics.ScopeIdx, op func() error, tags ...metrics.Tag) error {
perDomain := hasDomainTag(tags)
if !perDomain {
tags = ensureTaskCategoryTag(tags)
}
Comment thread
gitar-bot[bot] marked this conversation as resolved.
metricsScope := p.metricClient.Scope(scope, tags...)
if len(tags) > 0 {
if perDomain {
metricsScope.IncCounter(metrics.PersistenceRequestsPerDomain)
} else {
metricsScope.IncCounter(metrics.PersistenceRequests)
}
before := time.Now()
err := op()
duration := time.Since(before)
if len(tags) > 0 {
if perDomain {
metricsScope.RecordTimer(metrics.PersistenceLatencyPerDomain, duration)
metricsScope.ExponentialHistogram(metrics.PersistenceLatencyPerDomainHistogram, duration)
} else {
Expand All @@ -153,7 +186,7 @@ func (p *base) call(scope metrics.ScopeIdx, op func() error, tags ...metrics.Tag

logger := p.logger.Helper()
if err != nil {
if len(tags) > 0 {
if perDomain {
p.updateErrorMetricPerDomain(scope, err, metricsScope, logger)
} else {
p.updateErrorMetric(scope, err, metricsScope, logger)
Expand All @@ -163,6 +196,7 @@ func (p *base) call(scope metrics.ScopeIdx, op func() error, tags ...metrics.Tag
}

func (p *base) callWithoutDomainTag(scope metrics.ScopeIdx, op func() error, tags ...metrics.Tag) error {
tags = ensureTaskCategoryTag(tags)
metricsScope := p.metricClient.Scope(scope, tags...)
metricsScope.IncCounter(metrics.PersistenceRequests)
before := time.Now()
Expand All @@ -179,10 +213,10 @@ func (p *base) callWithoutDomainTag(scope metrics.ScopeIdx, op func() error, tag
}

func (p *base) callWithDomainAndShardScope(scope metrics.ScopeIdx, op func() error, domainTag metrics.Tag, shardIDTag metrics.Tag, additionalTags ...metrics.Tag) error {
overallScope := p.metricClient.Scope(scope)
overallScope := p.metricClient.Scope(scope, ensureTaskCategoryTag(additionalTags)...)
domainMetricsScope := p.metricClient.Scope(scope, append([]metrics.Tag{domainTag}, additionalTags...)...)
shardOperationsMetricsScope := p.metricClient.Scope(scope, append([]metrics.Tag{shardIDTag}, additionalTags...)...)
shardOverallMetricsScope := p.metricClient.Scope(metrics.PersistenceShardRequestCountScope, shardIDTag)
shardOverallMetricsScope := p.metricClient.Scope(metrics.PersistenceShardRequestCountScope, append([]metrics.Tag{shardIDTag}, additionalTags...)...)

domainMetricsScope.IncCounter(metrics.PersistenceRequestsPerDomain)
shardOperationsMetricsScope.IncCounter(metrics.PersistenceRequestsPerShard)
Expand Down

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

28 changes: 21 additions & 7 deletions common/persistence/wrappers/metered/domain_generated.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

32 changes: 24 additions & 8 deletions common/persistence/wrappers/metered/history_generated.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

73 changes: 73 additions & 0 deletions common/persistence/wrappers/metered/metered_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,11 @@ import (
"testing"
"time"

p8s "github.com/m3db/prometheus_client_golang/prometheus"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/uber-go/tally"
tallyp8s "github.com/uber-go/tally/prometheus"
"go.uber.org/mock/gomock"
"go.uber.org/zap"
"go.uber.org/zap/zaptest/observer"
Expand All @@ -51,6 +53,77 @@ var _staticMethods = map[string]bool{
"GetShardID": true,
}

func TestPersistenceMetricsLabelConsistency(t *testing.T) {
ctrl := gomock.NewController(t)
wrapped := persistence.NewMockExecutionManager(ctrl)
wrapped.EXPECT().GetShardID().Return(1).AnyTimes()
wrapped.EXPECT().GetWorkflowExecution(gomock.Any(), gomock.Any()).Return(&persistence.GetWorkflowExecutionResponse{}, nil).Times(2)
wrapped.EXPECT().GetHistoryTasks(gomock.Any(), gomock.Any()).Return(&persistence.GetHistoryTasksResponse{}, nil).Times(1)
wrapped.EXPECT().GetReplicationDLQSize(gomock.Any(), gomock.Any()).Return(&persistence.GetReplicationDLQSizeResponse{}, nil).Times(1)
historyStore := persistence.NewMockHistoryManager(ctrl)
historyStore.EXPECT().AppendHistoryNodes(gomock.Any(), gomock.Any()).Return(&persistence.AppendHistoryNodesResponse{}, nil).Times(1)

var registrationErrors []error
promCfg := &tallyp8s.Configuration{
OnError: "none",
TimerType: "histogram",
}
reporter, err := promCfg.NewReporter(tallyp8s.ConfigurationOptions{
Registry: p8s.NewRegistry(),
OnError: func(err error) {
registrationErrors = append(registrationErrors, err)
},
})
require.NoError(t, err)
rootScope, closer := tally.NewRootScope(tally.ScopeOptions{
Tags: map[string]string{
metrics.CadenceServiceTagName: "history",
},
CachedReporter: reporter,
Separator: tallyp8s.DefaultSeparator,
}, time.Second)
defer closer.Close()

metricsClient := metrics.NewClient(rootScope, metrics.History, metrics.MigrationConfig{})

shardMetricsManager := NewExecutionManager(
wrapped,
metricsClient,
log.NewNoop(),
&config.Persistence{EnablePersistenceLatencyHistogramMetrics: true},
dynamicproperties.GetBoolPropertyFn(true),
)
noShardMetricsManager := NewExecutionManager(
wrapped,
metricsClient,
log.NewNoop(),
&config.Persistence{EnablePersistenceLatencyHistogramMetrics: true},
dynamicproperties.GetBoolPropertyFn(false),
)
historyManager := NewHistoryManager(
historyStore,
metricsClient,
log.NewNoop(),
&config.Persistence{EnablePersistenceLatencyHistogramMetrics: true},
)

ctx := context.Background()
_, err = shardMetricsManager.GetWorkflowExecution(ctx, &persistence.GetWorkflowExecutionRequest{})
assert.NoError(t, err)
_, err = shardMetricsManager.GetHistoryTasks(ctx, &persistence.GetHistoryTasksRequest{
TaskCategory: persistence.HistoryTaskCategoryTransfer,
})
assert.NoError(t, err)
_, err = shardMetricsManager.GetReplicationDLQSize(ctx, &persistence.GetReplicationDLQSizeRequest{})
assert.NoError(t, err)
_, err = noShardMetricsManager.GetWorkflowExecution(ctx, &persistence.GetWorkflowExecutionRequest{})
assert.NoError(t, err)
_, err = historyManager.AppendHistoryNodes(ctx, &persistence.AppendHistoryNodesRequest{})
assert.NoError(t, err)

assert.Empty(t, registrationErrors, "Prometheus registration errors must not be emitted")
}

func TestGetRetryCountFromContext(t *testing.T) {
ctrl := gomock.NewController(t)
wrapped := persistence.NewMockExecutionManager(ctrl)
Expand Down
Loading
Loading