mirror of
https://github.com/weaveworks/scope.git
synced 2026-07-26 16:52:25 +00:00
``` $ gvt delete github.com/weaveworks/common $ gvt fetch --revision 4d96fd8dcf2c7b417912c6219b310112cb4a4626 github.com/weaveworks/common 2018/07/23 15:31:11 Fetching: github.com/weaveworks/common 2018/07/23 15:31:14 · Skipping (existing): github.com/golang/protobuf/ptypes/any 2018/07/23 15:31:14 · Fetching recursive dependency: github.com/pkg/errors 2018/07/23 15:31:16 · Skipping (existing): github.com/aws/aws-sdk-go/aws 2018/07/23 15:31:16 · Fetching recursive dependency: github.com/sirupsen/logrus 2018/07/23 15:31:18 ·· Skipping (existing): golang.org/x/sys/unix 2018/07/23 15:31:18 ·· Skipping (existing): golang.org/x/crypto/ssh/terminal 2018/07/23 15:31:18 · Skipping (existing): google.golang.org/grpc/status 2018/07/23 15:31:18 · Skipping (existing): github.com/gorilla/mux 2018/07/23 15:31:18 · Fetching recursive dependency: github.com/opentracing-contrib/go-stdlib/nethttp 2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:31:20 ·· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:31:20 · Skipping (existing): github.com/prometheus/client_golang/prometheus 2018/07/23 15:31:20 · Skipping (existing): google.golang.org/grpc 2018/07/23 15:31:20 · Skipping (existing): github.com/pmezard/go-difflib/difflib 2018/07/23 15:31:20 · Fetching recursive dependency: github.com/go-kit/kit/log 2018/07/23 15:31:23 ·· Fetching recursive dependency: github.com/go-logfmt/logfmt 2018/07/23 15:31:25 ··· Fetching recursive dependency: github.com/kr/logfmt 2018/07/23 15:31:27 ·· Fetching recursive dependency: github.com/go-stack/stack 2018/07/23 15:31:29 · Fetching recursive dependency: google.golang.org/genproto/googleapis/rpc/status 2018/07/23 15:31:37 ·· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:31:37 ·· Skipping (existing): github.com/golang/protobuf/ptypes/any 2018/07/23 15:31:37 · Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:31:37 · Fetching recursive dependency: github.com/sercand/kuberesolver 2018/07/23 15:31:39 ·· Skipping (existing): google.golang.org/grpc/grpclog 2018/07/23 15:31:39 ·· Skipping (existing): google.golang.org/grpc/resolver 2018/07/23 15:31:39 ·· Skipping (existing): golang.org/x/net/context 2018/07/23 15:31:39 · Skipping (existing): google.golang.org/grpc/metadata 2018/07/23 15:31:39 · Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:31:39 · Skipping (existing): github.com/armon/go-socks5 2018/07/23 15:31:39 · Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:31:39 · Skipping (existing): github.com/davecgh/go-spew/spew 2018/07/23 15:31:39 · Skipping (existing): github.com/golang/protobuf/ptypes 2018/07/23 15:31:39 · Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:31:39 · Fetching recursive dependency: github.com/grpc-ecosystem/grpc-opentracing/go/otgrpc 2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:31:41 ·· Skipping (existing): golang.org/x/net/context 2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/codes 2018/07/23 15:31:41 ·· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:31:41 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc 2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/metadata 2018/07/23 15:31:41 ·· Skipping (existing): google.golang.org/grpc/status 2018/07/23 15:31:41 · Fetching recursive dependency: github.com/uber/jaeger-client-go/config 2018/07/23 15:31:44 ·· Fetching recursive dependency: github.com/uber/jaeger-client-go/internal/throttler/remote 2018/07/23 15:31:44 ··· Fetching recursive dependency: github.com/uber/jaeger-client-go/utils 2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift 2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/agent 2018/07/23 15:31:44 ····· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/jaeger 2018/07/23 15:31:44 ····· Fetching recursive dependency: github.com/uber/jaeger-client-go/thrift-gen/zipkincore 2018/07/23 15:31:44 ··· Fetching recursive dependency: github.com/uber/jaeger-client-go 2018/07/23 15:31:44 ···· Fetching recursive dependency: github.com/crossdock/crossdock-go 2018/07/23 15:31:46 ····· Skipping (existing): github.com/davecgh/go-spew/spew 2018/07/23 15:31:46 ····· Skipping (existing): golang.org/x/net/context/ctxhttp 2018/07/23 15:31:46 ····· Skipping (existing): golang.org/x/net/context 2018/07/23 15:31:46 ····· Skipping (existing): github.com/pmezard/go-difflib/difflib 2018/07/23 15:31:46 ···· Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:31:46 ···· Fetching recursive dependency: go.uber.org/zap/zapcore 2018/07/23 15:31:49 ····· Fetching recursive dependency: go.uber.org/atomic 2018/07/23 15:31:51 ····· Fetching recursive dependency: go.uber.org/zap/internal/bufferpool 2018/07/23 15:31:51 ······ Fetching recursive dependency: go.uber.org/zap/buffer 2018/07/23 15:31:51 ····· Fetching recursive dependency: go.uber.org/multierr 2018/07/23 15:31:54 ····· Fetching recursive dependency: go.uber.org/zap/internal/exit 2018/07/23 15:31:54 ····· Fetching recursive dependency: go.uber.org/zap/internal/color 2018/07/23 15:31:54 ···· Fetching recursive dependency: go.uber.org/zap 2018/07/23 15:31:54 ···· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:31:54 ···· Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:31:54 ···· Fetching recursive dependency: github.com/uber/jaeger-lib/metrics 2018/07/23 15:31:56 ····· Fetching recursive dependency: github.com/uber-go/tally 2018/07/23 15:31:58 ······ Fetching recursive dependency: github.com/m3db/prometheus_client_golang/prometheus/promhttp 2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus 2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/common/expfmt 2018/07/23 15:32:00 ······· Skipping (existing): github.com/prometheus/client_model/go 2018/07/23 15:32:00 ······ Fetching recursive dependency: gopkg.in/validator.v2 2018/07/23 15:32:06 ······ Fetching recursive dependency: github.com/cactus/go-statsd-client/statsd 2018/07/23 15:32:08 ······ Skipping (existing): gopkg.in/yaml.v2 2018/07/23 15:32:08 ······ Fetching recursive dependency: github.com/m3db/prometheus_client_golang/prometheus 2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/procfs 2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/client_model/go 2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/common/expfmt 2018/07/23 15:32:08 ······· Skipping (existing): golang.org/x/net/context 2018/07/23 15:32:08 ······· Skipping (existing): github.com/beorn7/perks/quantile 2018/07/23 15:32:08 ······· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/common/model 2018/07/23 15:32:08 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus 2018/07/23 15:32:08 ······ Fetching recursive dependency: github.com/apache/thrift/lib/go/thrift 2018/07/23 15:32:13 ····· Skipping (existing): github.com/stretchr/testify/assert 2018/07/23 15:32:13 ····· Fetching recursive dependency: github.com/go-kit/kit/metrics/influx 2018/07/23 15:32:13 ······ Fetching recursive dependency: github.com/influxdata/influxdb/client/v2 2018/07/23 15:32:17 ······· Fetching recursive dependency: github.com/influxdata/influxdb/models 2018/07/23 15:32:17 ········ Fetching recursive dependency: github.com/influxdata/influxdb/pkg/escape 2018/07/23 15:32:17 ······ Fetching recursive dependency: github.com/go-kit/kit/metrics 2018/07/23 15:32:17 ······· Fetching recursive dependency: github.com/performancecopilot/speed 2018/07/23 15:32:19 ······· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/aws 2018/07/23 15:32:28 ········ Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/internal/sdk 2018/07/23 15:32:28 ········ Skipping (existing): github.com/go-ini/ini 2018/07/23 15:32:28 ········ Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/service/sts 2018/07/23 15:32:28 ········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/private/protocol/query 2018/07/23 15:32:28 ·········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/private/protocol 2018/07/23 15:32:28 ········· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/internal/awsutil 2018/07/23 15:32:28 ·········· Skipping (existing): github.com/jmespath/go-jmespath 2018/07/23 15:32:28 ······· Fetching recursive dependency: github.com/aws/aws-sdk-go-v2/service/cloudwatch 2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/aws 2018/07/23 15:32:29 ······· Skipping (existing): github.com/prometheus/client_golang/prometheus 2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/service/cloudwatch 2018/07/23 15:32:29 ······· Skipping (existing): github.com/aws/aws-sdk-go/service/cloudwatch/cloudwatchiface 2018/07/23 15:32:29 ······· Fetching recursive dependency: golang.org/x/sync/errgroup 2018/07/23 15:32:31 ········ Skipping (existing): golang.org/x/net/context 2018/07/23 15:32:31 ······· Fetching recursive dependency: github.com/go-kit/kit/util/conn 2018/07/23 15:32:31 ······· Fetching recursive dependency: github.com/VividCortex/gohistogram 2018/07/23 15:32:33 ····· Skipping (existing): github.com/prometheus/client_golang/prometheus 2018/07/23 15:32:33 ····· Fetching recursive dependency: github.com/codahale/hdrhistogram 2018/07/23 15:32:35 ·· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:32:35 · Fetching recursive dependency: github.com/mwitkow/go-grpc-middleware 2018/07/23 15:32:37 ·· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware/logging 2018/07/23 15:32:39 ··· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware 2018/07/23 15:32:39 ···· Fetching recursive dependency: github.com/golang/protobuf/jsonpb 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/timestamp 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/duration 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/any 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/struct 2018/07/23 15:32:42 ····· Skipping (existing): github.com/golang/protobuf/ptypes/wrappers 2018/07/23 15:32:42 ···· Skipping (existing): google.golang.org/grpc/metadata 2018/07/23 15:32:42 ···· Fetching recursive dependency: github.com/stretchr/testify/suite 2018/07/23 15:32:45 ····· Skipping (existing): github.com/stretchr/testify/assert 2018/07/23 15:32:45 ····· Fetching recursive dependency: github.com/stretchr/testify/require 2018/07/23 15:32:45 ······ Skipping (existing): github.com/stretchr/testify/assert 2018/07/23 15:32:45 ···· Skipping (existing): google.golang.org/grpc/peer 2018/07/23 15:32:45 ···· Skipping (existing): golang.org/x/net/context 2018/07/23 15:32:45 ···· Skipping (existing): golang.org/x/net/trace 2018/07/23 15:32:45 ···· Fetching recursive dependency: github.com/gogo/protobuf/gogoproto 2018/07/23 15:32:48 ····· Fetching recursive dependency: github.com/gogo/protobuf/protoc-gen-gogo/descriptor 2018/07/23 15:32:48 ······ Skipping (existing): github.com/gogo/protobuf/proto 2018/07/23 15:32:48 ····· Skipping (existing): github.com/gogo/protobuf/proto 2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/credentials 2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc 2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/codes 2018/07/23 15:32:48 ···· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:32:48 ···· Skipping (existing): google.golang.org/grpc/grpclog 2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:32:48 ···· Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:32:48 ··· Skipping (existing): golang.org/x/net/context 2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc 2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc/grpclog 2018/07/23 15:32:48 ··· Skipping (existing): google.golang.org/grpc/codes 2018/07/23 15:32:48 ··· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc 2018/07/23 15:32:48 ·· Skipping (existing): golang.org/x/net/context 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/codes 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/grpclog 2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go/log 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/metadata 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/peer 2018/07/23 15:32:48 ·· Skipping (existing): google.golang.org/grpc/credentials 2018/07/23 15:32:48 ·· Skipping (existing): github.com/golang/protobuf/proto 2018/07/23 15:32:48 ·· Skipping (existing): golang.org/x/net/trace 2018/07/23 15:32:48 ·· Skipping (existing): github.com/opentracing/opentracing-go/ext 2018/07/23 15:32:48 · Fetching recursive dependency: github.com/weaveworks/promrus 2018/07/23 15:32:53 ·· Skipping (existing): gopkg.in/yaml.v2 2018/07/23 15:32:53 ·· Skipping (existing): golang.org/x/net/context/ctxhttp 2018/07/23 15:32:53 ·· Fetching recursive dependency: github.com/stretchr/objx 2018/07/23 15:32:55 ·· Fetching recursive dependency: gopkg.in/alecthomas/kingpin.v2 2018/07/23 15:32:58 ··· Fetching recursive dependency: github.com/alecthomas/units 2018/07/23 15:33:00 ··· Fetching recursive dependency: github.com/alecthomas/template 2018/07/23 15:33:02 ·· Fetching recursive dependency: github.com/julienschmidt/httprouter 2018/07/23 15:33:05 ·· Skipping (existing): golang.org/x/net/context 2018/07/23 15:33:05 · Skipping (existing): github.com/aws/aws-sdk-go/aws/credentials 2018/07/23 15:33:05 · Skipping (existing): github.com/golang/protobuf/ptypes/empty 2018/07/23 15:33:05 · Skipping (existing): golang.org/x/net/context 2018/07/23 15:33:05 · Skipping (existing): golang.org/x/tools/cover 2018/07/23 15:33:05 · Skipping (existing): github.com/mgutz/ansi ```
297 lines
9.7 KiB
Go
297 lines
9.7 KiB
Go
// Copyright 2016 Michal Witkowski. All Rights Reserved.
|
|
// See LICENSE for licensing terms.
|
|
|
|
package grpc_retry
|
|
|
|
import (
|
|
"fmt"
|
|
"io"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/grpc-ecosystem/go-grpc-middleware/util/metautils"
|
|
"golang.org/x/net/context"
|
|
"golang.org/x/net/trace"
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/metadata"
|
|
)
|
|
|
|
const (
|
|
AttemptMetadataKey = "x-retry-attempty"
|
|
)
|
|
|
|
// UnaryClientInterceptor returns a new retrying unary client interceptor.
|
|
//
|
|
// The default configuration of the interceptor is to not retry *at all*. This behaviour can be
|
|
// changed through options (e.g. WithMax) on creation of the interceptor or on call (through grpc.CallOptions).
|
|
func UnaryClientInterceptor(optFuncs ...CallOption) grpc.UnaryClientInterceptor {
|
|
intOpts := reuseOrNewWithCallOptions(defaultOptions, optFuncs)
|
|
return func(parentCtx context.Context, method string, req, reply interface{}, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
|
|
grpcOpts, retryOpts := filterCallOptions(opts)
|
|
callOpts := reuseOrNewWithCallOptions(intOpts, retryOpts)
|
|
// short circuit for simplicity, and avoiding allocations.
|
|
if callOpts.max == 0 {
|
|
return invoker(parentCtx, method, req, reply, cc, grpcOpts...)
|
|
}
|
|
var lastErr error
|
|
for attempt := uint(0); attempt < callOpts.max; attempt++ {
|
|
if err := waitRetryBackoff(attempt, parentCtx, callOpts); err != nil {
|
|
return err
|
|
}
|
|
callCtx := perCallContext(parentCtx, callOpts, attempt)
|
|
lastErr = invoker(callCtx, method, req, reply, cc, grpcOpts...)
|
|
// TODO(mwitkow): Maybe dial and transport errors should be retriable?
|
|
if lastErr == nil {
|
|
return nil
|
|
}
|
|
logTrace(parentCtx, "grpc_retry attempt: %d, got err: %v", attempt, lastErr)
|
|
if isContextError(lastErr) {
|
|
if parentCtx.Err() != nil {
|
|
logTrace(parentCtx, "grpc_retry attempt: %d, parent context error: %v", attempt, parentCtx.Err())
|
|
// its the parent context deadline or cancellation.
|
|
return lastErr
|
|
} else {
|
|
logTrace(parentCtx, "grpc_retry attempt: %d, context error from retry call", attempt)
|
|
// its the callCtx deadline or cancellation, in which case try again.
|
|
continue
|
|
}
|
|
}
|
|
if !isRetriable(lastErr, callOpts) {
|
|
return lastErr
|
|
}
|
|
}
|
|
return lastErr
|
|
}
|
|
}
|
|
|
|
// StreamClientInterceptor returns a new retrying stream client interceptor for server side streaming calls.
|
|
//
|
|
// The default configuration of the interceptor is to not retry *at all*. This behaviour can be
|
|
// changed through options (e.g. WithMax) on creation of the interceptor or on call (through grpc.CallOptions).
|
|
//
|
|
// Retry logic is available *only for ServerStreams*, i.e. 1:n streams, as the internal logic needs
|
|
// to buffer the messages sent by the client. If retry is enabled on any other streams (ClientStreams,
|
|
// BidiStreams), the retry interceptor will fail the call.
|
|
func StreamClientInterceptor(optFuncs ...CallOption) grpc.StreamClientInterceptor {
|
|
intOpts := reuseOrNewWithCallOptions(defaultOptions, optFuncs)
|
|
return func(parentCtx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn, method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) {
|
|
grpcOpts, retryOpts := filterCallOptions(opts)
|
|
callOpts := reuseOrNewWithCallOptions(intOpts, retryOpts)
|
|
// short circuit for simplicity, and avoiding allocations.
|
|
if callOpts.max == 0 {
|
|
return streamer(parentCtx, desc, cc, method, grpcOpts...)
|
|
}
|
|
if desc.ClientStreams {
|
|
return nil, grpc.Errorf(codes.Unimplemented, "grpc_retry: cannot retry on ClientStreams, set grpc_retry.Disable()")
|
|
}
|
|
logTrace(parentCtx, "grpc_retry attempt: %d, no backoff for this call", 0)
|
|
callCtx := perCallContext(parentCtx, callOpts, 0)
|
|
newStreamer, err := streamer(callCtx, desc, cc, method, grpcOpts...)
|
|
if err != nil {
|
|
// TODO(mwitkow): Maybe dial and transport errors should be retriable?
|
|
return nil, err
|
|
}
|
|
retryingStreamer := &serverStreamingRetryingStream{
|
|
ClientStream: newStreamer,
|
|
callOpts: callOpts,
|
|
parentCtx: parentCtx,
|
|
streamerCall: func(ctx context.Context) (grpc.ClientStream, error) {
|
|
return streamer(ctx, desc, cc, method, grpcOpts...)
|
|
},
|
|
}
|
|
return retryingStreamer, nil
|
|
}
|
|
}
|
|
|
|
// type serverStreamingRetryingStream is the implementation of grpc.ClientStream that acts as a
|
|
// proxy to the underlying call. If any of the RecvMsg() calls fail, it will try to reestablish
|
|
// a new ClientStream according to the retry policy.
|
|
type serverStreamingRetryingStream struct {
|
|
grpc.ClientStream
|
|
bufferedSends []interface{} // single messsage that the client can sen
|
|
receivedGood bool // indicates whether any prior receives were successful
|
|
wasClosedSend bool // indicates that CloseSend was closed
|
|
parentCtx context.Context
|
|
callOpts *options
|
|
streamerCall func(ctx context.Context) (grpc.ClientStream, error)
|
|
mu sync.RWMutex
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) setStream(clientStream grpc.ClientStream) {
|
|
s.mu.Lock()
|
|
s.ClientStream = clientStream
|
|
s.mu.Unlock()
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) getStream() grpc.ClientStream {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.ClientStream
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) SendMsg(m interface{}) error {
|
|
s.mu.Lock()
|
|
s.bufferedSends = append(s.bufferedSends, m)
|
|
s.mu.Unlock()
|
|
return s.getStream().SendMsg(m)
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) CloseSend() error {
|
|
s.mu.Lock()
|
|
s.wasClosedSend = true
|
|
s.mu.Unlock()
|
|
return s.getStream().CloseSend()
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) Header() (metadata.MD, error) {
|
|
return s.getStream().Header()
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) Trailer() metadata.MD {
|
|
return s.getStream().Trailer()
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) RecvMsg(m interface{}) error {
|
|
attemptRetry, lastErr := s.receiveMsgAndIndicateRetry(m)
|
|
if !attemptRetry {
|
|
return lastErr // success or hard failure
|
|
}
|
|
// We start off from attempt 1, because zeroth was already made on normal SendMsg().
|
|
for attempt := uint(1); attempt < s.callOpts.max; attempt++ {
|
|
if err := waitRetryBackoff(attempt, s.parentCtx, s.callOpts); err != nil {
|
|
return err
|
|
}
|
|
callCtx := perCallContext(s.parentCtx, s.callOpts, attempt)
|
|
newStream, err := s.reestablishStreamAndResendBuffer(callCtx)
|
|
if err != nil {
|
|
// TODO(mwitkow): Maybe dial and transport errors should be retriable?
|
|
return err
|
|
}
|
|
s.setStream(newStream)
|
|
attemptRetry, lastErr = s.receiveMsgAndIndicateRetry(m)
|
|
//fmt.Printf("Received message and indicate: %v %v\n", attemptRetry, lastErr)
|
|
if !attemptRetry {
|
|
return lastErr
|
|
}
|
|
}
|
|
return lastErr
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) receiveMsgAndIndicateRetry(m interface{}) (bool, error) {
|
|
s.mu.RLock()
|
|
wasGood := s.receivedGood
|
|
s.mu.RUnlock()
|
|
err := s.getStream().RecvMsg(m)
|
|
if err == nil || err == io.EOF {
|
|
s.mu.Lock()
|
|
s.receivedGood = true
|
|
s.mu.Unlock()
|
|
return false, err
|
|
} else if wasGood {
|
|
// previous RecvMsg in the stream succeeded, no retry logic should interfere
|
|
return false, err
|
|
}
|
|
if isContextError(err) {
|
|
if s.parentCtx.Err() != nil {
|
|
logTrace(s.parentCtx, "grpc_retry parent context error: %v", s.parentCtx.Err())
|
|
return false, err
|
|
} else {
|
|
logTrace(s.parentCtx, "grpc_retry context error from retry call")
|
|
// its the callCtx deadline or cancellation, in which case try again.
|
|
return true, err
|
|
}
|
|
}
|
|
return isRetriable(err, s.callOpts), err
|
|
|
|
}
|
|
|
|
func (s *serverStreamingRetryingStream) reestablishStreamAndResendBuffer(callCtx context.Context) (grpc.ClientStream, error) {
|
|
s.mu.RLock()
|
|
bufferedSends := s.bufferedSends
|
|
s.mu.RUnlock()
|
|
newStream, err := s.streamerCall(callCtx)
|
|
if err != nil {
|
|
logTrace(callCtx, "grpc_retry failed redialing new stream: %v", err)
|
|
return nil, err
|
|
}
|
|
for _, msg := range bufferedSends {
|
|
if err := newStream.SendMsg(msg); err != nil {
|
|
logTrace(callCtx, "grpc_retry failed resending message: %v", err)
|
|
return nil, err
|
|
}
|
|
}
|
|
if err := newStream.CloseSend(); err != nil {
|
|
logTrace(callCtx, "grpc_retry failed CloseSend on new stream %v", err)
|
|
return nil, err
|
|
}
|
|
return newStream, nil
|
|
}
|
|
|
|
func waitRetryBackoff(attempt uint, parentCtx context.Context, callOpts *options) error {
|
|
var waitTime time.Duration = 0
|
|
if attempt > 0 {
|
|
waitTime = callOpts.backoffFunc(attempt)
|
|
}
|
|
if waitTime > 0 {
|
|
logTrace(parentCtx, "grpc_retry attempt: %d, backoff for %v", attempt, waitTime)
|
|
timer := time.NewTimer(waitTime)
|
|
select {
|
|
case <-parentCtx.Done():
|
|
timer.Stop()
|
|
return contextErrToGrpcErr(parentCtx.Err())
|
|
case <-timer.C:
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func isRetriable(err error, callOpts *options) bool {
|
|
errCode := grpc.Code(err)
|
|
if isContextError(err) {
|
|
// context errors are not retriable based on user settings.
|
|
return false
|
|
}
|
|
for _, code := range callOpts.codes {
|
|
if code == errCode {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func isContextError(err error) bool {
|
|
return grpc.Code(err) == codes.DeadlineExceeded || grpc.Code(err) == codes.Canceled
|
|
}
|
|
|
|
func perCallContext(parentCtx context.Context, callOpts *options, attempt uint) context.Context {
|
|
ctx := parentCtx
|
|
if callOpts.perCallTimeout != 0 {
|
|
ctx, _ = context.WithTimeout(ctx, callOpts.perCallTimeout)
|
|
}
|
|
if attempt > 0 && callOpts.includeHeader {
|
|
mdClone := metautils.ExtractOutgoing(ctx).Clone().Set(AttemptMetadataKey, fmt.Sprintf("%d", attempt))
|
|
ctx = mdClone.ToOutgoing(ctx)
|
|
}
|
|
return ctx
|
|
}
|
|
|
|
func contextErrToGrpcErr(err error) error {
|
|
switch err {
|
|
case context.DeadlineExceeded:
|
|
return grpc.Errorf(codes.DeadlineExceeded, err.Error())
|
|
case context.Canceled:
|
|
return grpc.Errorf(codes.Canceled, err.Error())
|
|
default:
|
|
return grpc.Errorf(codes.Unknown, err.Error())
|
|
}
|
|
}
|
|
|
|
func logTrace(ctx context.Context, format string, a ...interface{}) {
|
|
tr, ok := trace.FromContext(ctx)
|
|
if !ok {
|
|
return
|
|
}
|
|
tr.LazyPrintf(format, a...)
|
|
}
|