Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
3 changes: 0 additions & 3 deletions pkg/errs/errno.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,9 +86,6 @@ var (
ErrNotStarted = status.Error(codes.Unavailable, "server not started")
ErrEtcdNotStarted = status.Error(codes.Unavailable, "server is started, but etcd not started")
ErrFollowerHandlingNotAllowed = status.Error(codes.Unavailable, "not leader and follower handling not allowed")

// Unimplemented indicates operation is not implemented or not supported.
ErrRecvNotSupported = status.Error(codes.Unimplemented, "recv is not supported by this stream")
)

// common error in multiple packages
Expand Down
1 change: 0 additions & 1 deletion pkg/mcs/metastorage/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,6 @@ func (s *Service) checkServing() error {

// Watch watches the key with a given prefix and revision.
func (s *Service) Watch(req *meta_storagepb.WatchRequest, server meta_storagepb.MetaStorage_WatchServer) error {
server = newWatchMetricsStream(server)
if err := s.checkServing(); err != nil {
return err
}
Expand Down
38 changes: 0 additions & 38 deletions pkg/mcs/metastorage/server/metrics.go

This file was deleted.

1 change: 0 additions & 1 deletion pkg/mcs/resourcemanager/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -188,7 +188,6 @@ func (s *Service) ModifyResourceGroup(_ context.Context, req *rmpb.PutResourceGr

// AcquireTokenBuckets implements ResourceManagerServer.AcquireTokenBuckets.
func (s *Service) AcquireTokenBuckets(stream rmpb.ResourceManager_AcquireTokenBucketsServer) error {
stream = newAcquireTokenBucketsMetricsStream(stream)
for {
select {
case <-s.ctx.Done():
Expand Down
9 changes: 0 additions & 9 deletions pkg/mcs/resourcemanager/server/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,6 @@ import (
"github.com/prometheus/client_golang/prometheus"

rmpb "github.com/pingcap/kvproto/pkg/resource_manager"

"github.com/tikv/pd/pkg/utils/grpcutil"
)

const (
Expand Down Expand Up @@ -60,8 +58,6 @@ const (
)

var (
grpcStreamSendDuration = grpcutil.NewGRPCStreamSendDuration(namespace, serverSubsystem)

// RU cost metrics.
// `sum` is added to the name to maintain compatibility with the previous use of histogram.
readRequestUnitCost = prometheus.NewCounterVec(
Expand Down Expand Up @@ -313,7 +309,6 @@ type requestMetricsKey struct {
}

func init() {
prometheus.MustRegister(grpcStreamSendDuration)
prometheus.MustRegister(readRequestUnitCost)
prometheus.MustRegister(writeRequestUnitCost)
prometheus.MustRegister(activeRequestUnitCost)
Expand Down Expand Up @@ -1052,7 +1047,3 @@ func (t *maxPerSecCostTracker) flushMetrics() {
t.maxPerSecWRU = 0
}
}

func newAcquireTokenBucketsMetricsStream(stream rmpb.ResourceManager_AcquireTokenBucketsServer) rmpb.ResourceManager_AcquireTokenBucketsServer {
return grpcutil.NewMetricsStream(stream, stream.Send, stream.Recv, grpcStreamSendDuration, "acquire-token-buckets")
}
1 change: 0 additions & 1 deletion pkg/mcs/router/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,6 @@ func (s *Service) GetRegionByID(_ctx context.Context, request *pdpb.GetRegionByI

// QueryRegion implements the QueryRegion RPC method.
func (s *Service) QueryRegion(stream routerpb.Router_QueryRegionServer) error {
stream = newQueryRegionMetricsStream(stream)
for {
request, err := stream.Recv()
if err == io.EOF {
Expand Down
11 changes: 0 additions & 11 deletions pkg/mcs/router/server/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,6 @@ package server

import (
"github.com/prometheus/client_golang/prometheus"

"github.com/pingcap/kvproto/pkg/routerpb"

"github.com/tikv/pd/pkg/utils/grpcutil"
)

const (
Expand All @@ -28,8 +24,6 @@ const (
)

var (
grpcStreamSendDuration = grpcutil.NewGRPCStreamSendDuration(namespace, serverSubsystem)

queryRegionDuration = prometheus.NewHistogram(
prometheus.HistogramOpts{
Namespace: namespace,
Expand Down Expand Up @@ -57,12 +51,7 @@ var (
)

func init() {
prometheus.MustRegister(grpcStreamSendDuration)
prometheus.MustRegister(regionRequestCounter)
prometheus.MustRegister(queryRegionDuration)
prometheus.MustRegister(regionSyncerStatus)
}

func newQueryRegionMetricsStream(stream routerpb.Router_QueryRegionServer) routerpb.Router_QueryRegionServer {
return grpcutil.NewMetricsStream(stream, stream.Send, stream.Recv, grpcStreamSendDuration, "query-region")
}
4 changes: 1 addition & 3 deletions pkg/mcs/scheduling/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,9 +116,8 @@ func (s *heartbeatServer) recv() (*schedulingpb.RegionHeartbeatRequest, error) {

// RegionHeartbeat implements gRPC SchedulingServer.
func (s *Service) RegionHeartbeat(stream schedulingpb.Scheduling_RegionHeartbeatServer) error {
wrappedStream := newRegionHeartbeatMetricsStream(stream)
var (
server = &heartbeatServer{stream: wrappedStream}
server = &heartbeatServer{stream: stream}
cancel context.CancelFunc
lastBind time.Time
)
Expand Down Expand Up @@ -183,7 +182,6 @@ func (s *Service) RegionHeartbeat(stream schedulingpb.Scheduling_RegionHeartbeat

// RegionBuckets implements gRPC SchedulingServer.
func (s *Service) RegionBuckets(stream schedulingpb.Scheduling_RegionBucketsServer) error {
stream = newRegionBucketsMetricsStream(stream)
var cancel context.CancelFunc
defer func() {
// cancel the forward stream
Expand Down
19 changes: 1 addition & 18 deletions pkg/mcs/scheduling/server/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,22 +14,14 @@

package server

import (
"github.com/prometheus/client_golang/prometheus"

"github.com/pingcap/kvproto/pkg/schedulingpb"

"github.com/tikv/pd/pkg/utils/grpcutil"
)
import "github.com/prometheus/client_golang/prometheus"

const (
namespace = "scheduling"
serverSubsystem = "server"
)

var (
grpcStreamSendDuration = grpcutil.NewGRPCStreamSendDuration(namespace, serverSubsystem)

// Store heartbeat metrics
storeHeartbeatHandleDuration = prometheus.NewHistogramVec(
prometheus.HistogramOpts{
Expand Down Expand Up @@ -94,7 +86,6 @@ var (
)

func init() {
prometheus.MustRegister(grpcStreamSendDuration)
prometheus.MustRegister(storeHeartbeatHandleDuration)
prometheus.MustRegister(storeHeartbeatCounter)
prometheus.MustRegister(regionHeartbeatHandleDuration)
Expand All @@ -103,11 +94,3 @@ func init() {
prometheus.MustRegister(regionBucketsCounter)
prometheus.MustRegister(regionBucketsReportInterval)
}

func newRegionHeartbeatMetricsStream(stream schedulingpb.Scheduling_RegionHeartbeatServer) schedulingpb.Scheduling_RegionHeartbeatServer {
return grpcutil.NewMetricsStream(stream, stream.Send, stream.Recv, grpcStreamSendDuration, "region-heartbeat")
}

func newRegionBucketsMetricsStream(stream schedulingpb.Scheduling_RegionBucketsServer) schedulingpb.Scheduling_RegionBucketsServer {
return grpcutil.NewMetricsStream(stream, stream.Send, stream.Recv, grpcStreamSendDuration, "region-buckets")
}
1 change: 0 additions & 1 deletion pkg/mcs/tso/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,6 @@ func (s *Service) RegisterRESTHandler(userDefineHandlers map[string]http.Handler

// Tso returns a stream of timestamps
func (s *Service) Tso(stream tsopb.TSO_TsoServer) error {
stream = newTsoMetricsStream(stream)
ctx, cancel := context.WithCancel(stream.Context())
defer cancel()
for {
Expand Down
15 changes: 1 addition & 14 deletions pkg/mcs/tso/server/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,19 +14,11 @@

package server

import (
"github.com/prometheus/client_golang/prometheus"

"github.com/pingcap/kvproto/pkg/tsopb"

"github.com/tikv/pd/pkg/utils/grpcutil"
)
import "github.com/prometheus/client_golang/prometheus"

const namespace = "tso"

var (
grpcStreamSendDuration = grpcutil.NewGRPCStreamSendDuration(namespace, "server")

timeJumpBackCounter = prometheus.NewCounter(
prometheus.CounterOpts{
Namespace: namespace,
Expand Down Expand Up @@ -54,12 +46,7 @@ var (
)

func init() {
prometheus.MustRegister(grpcStreamSendDuration)
prometheus.MustRegister(timeJumpBackCounter)
prometheus.MustRegister(metaDataGauge)
prometheus.MustRegister(tsoHandleDuration)
}

func newTsoMetricsStream(stream tsopb.TSO_TsoServer) tsopb.TSO_TsoServer {
return grpcutil.NewMetricsStream(stream, stream.Send, stream.Recv, grpcStreamSendDuration, "tso")
}
126 changes: 0 additions & 126 deletions pkg/utils/grpcutil/stream.go

This file was deleted.

Loading
Loading