2021-02-22 17:27:13 +05:30
|
|
|
package gitaly
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"strings"
|
|
|
|
"sync"
|
|
|
|
|
2021-10-27 15:23:28 +05:30
|
|
|
"github.com/golang/protobuf/jsonpb" //lint:ignore SA1019 https://gitlab.com/gitlab-org/gitlab/-/issues/324868
|
|
|
|
"github.com/golang/protobuf/proto" //lint:ignore SA1019 https://gitlab.com/gitlab-org/gitlab/-/issues/324868
|
2021-02-22 17:27:13 +05:30
|
|
|
grpc_middleware "github.com/grpc-ecosystem/go-grpc-middleware"
|
|
|
|
grpc_prometheus "github.com/grpc-ecosystem/go-grpc-prometheus"
|
|
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
|
|
"github.com/prometheus/client_golang/prometheus/promauto"
|
2021-11-18 22:05:49 +05:30
|
|
|
"github.com/sirupsen/logrus"
|
2021-10-27 15:23:28 +05:30
|
|
|
"google.golang.org/grpc"
|
|
|
|
"google.golang.org/grpc/metadata"
|
|
|
|
|
2021-09-04 01:27:46 +05:30
|
|
|
gitalyauth "gitlab.com/gitlab-org/gitaly/v14/auth"
|
|
|
|
gitalyclient "gitlab.com/gitlab-org/gitaly/v14/client"
|
|
|
|
"gitlab.com/gitlab-org/gitaly/v14/proto/go/gitalypb"
|
2021-02-22 17:27:13 +05:30
|
|
|
|
|
|
|
grpccorrelation "gitlab.com/gitlab-org/labkit/correlation/grpc"
|
|
|
|
grpctracing "gitlab.com/gitlab-org/labkit/tracing/grpc"
|
|
|
|
)
|
|
|
|
|
|
|
|
type Server struct {
|
2021-11-18 22:05:49 +05:30
|
|
|
Address string `json:"address"`
|
|
|
|
Token string `json:"token"`
|
|
|
|
Features map[string]string `json:"features"`
|
|
|
|
Sidechannel bool `json:"sidechannel"`
|
2021-02-22 17:27:13 +05:30
|
|
|
}
|
|
|
|
|
2021-11-18 22:05:49 +05:30
|
|
|
type cacheKey struct {
|
|
|
|
address, token string
|
|
|
|
sidechannel bool
|
|
|
|
}
|
2021-02-22 17:27:13 +05:30
|
|
|
|
|
|
|
func (server Server) cacheKey() cacheKey {
|
2021-11-18 22:05:49 +05:30
|
|
|
return cacheKey{address: server.Address, token: server.Token, sidechannel: server.Sidechannel}
|
2021-02-22 17:27:13 +05:30
|
|
|
}
|
|
|
|
|
|
|
|
type connectionsCache struct {
|
|
|
|
sync.RWMutex
|
|
|
|
connections map[cacheKey]*grpc.ClientConn
|
|
|
|
}
|
|
|
|
|
|
|
|
var (
|
|
|
|
jsonUnMarshaler = jsonpb.Unmarshaler{AllowUnknownFields: true}
|
2021-11-18 22:05:49 +05:30
|
|
|
// This connection cache map contains two types of connections:
|
|
|
|
// - Normal gRPC connections
|
|
|
|
// - Sidechannel connections. When client dials to the Gitaly server, the
|
|
|
|
// server multiplexes the connection using Yamux. In the future, the server
|
|
|
|
// can open another stream to transfer data without gRPC. Besides, we apply
|
|
|
|
// a framing protocol to add the half-close capability to Yamux streams.
|
|
|
|
// Hence, we cannot use those connections interchangeably.
|
|
|
|
cache = connectionsCache{
|
2021-02-22 17:27:13 +05:30
|
|
|
connections: make(map[cacheKey]*grpc.ClientConn),
|
|
|
|
}
|
2021-11-18 22:05:49 +05:30
|
|
|
sidechannelRegistry *gitalyclient.SidechannelRegistry
|
2021-02-22 17:27:13 +05:30
|
|
|
|
|
|
|
connectionsTotal = promauto.NewCounterVec(
|
|
|
|
prometheus.CounterOpts{
|
|
|
|
Name: "gitlab_workhorse_gitaly_connections_total",
|
|
|
|
Help: "Number of Gitaly connections that have been established",
|
|
|
|
},
|
|
|
|
[]string{"status"},
|
|
|
|
)
|
|
|
|
)
|
|
|
|
|
2021-11-18 22:05:49 +05:30
|
|
|
func InitializeSidechannelRegistry(logger *logrus.Logger) {
|
|
|
|
if sidechannelRegistry == nil {
|
|
|
|
sidechannelRegistry = gitalyclient.NewSidechannelRegistry(logrus.NewEntry(logger))
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-02-22 17:27:13 +05:30
|
|
|
func withOutgoingMetadata(ctx context.Context, features map[string]string) context.Context {
|
|
|
|
md := metadata.New(nil)
|
|
|
|
for k, v := range features {
|
|
|
|
if !strings.HasPrefix(k, "gitaly-feature-") {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
md.Append(k, v)
|
|
|
|
}
|
|
|
|
|
|
|
|
return metadata.NewOutgoingContext(ctx, md)
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewSmartHTTPClient(ctx context.Context, server Server) (context.Context, *SmartHTTPClient, error) {
|
|
|
|
conn, err := getOrCreateConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
grpcClient := gitalypb.NewSmartHTTPServiceClient(conn)
|
2021-11-18 22:05:49 +05:30
|
|
|
smartHTTPClient := &SmartHTTPClient{
|
|
|
|
SmartHTTPServiceClient: grpcClient,
|
|
|
|
sidechannelRegistry: sidechannelRegistry,
|
|
|
|
useSidechannel: server.Sidechannel,
|
|
|
|
}
|
|
|
|
return withOutgoingMetadata(ctx, server.Features), smartHTTPClient, nil
|
2021-02-22 17:27:13 +05:30
|
|
|
}
|
|
|
|
|
|
|
|
func NewBlobClient(ctx context.Context, server Server) (context.Context, *BlobClient, error) {
|
|
|
|
conn, err := getOrCreateConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
grpcClient := gitalypb.NewBlobServiceClient(conn)
|
|
|
|
return withOutgoingMetadata(ctx, server.Features), &BlobClient{grpcClient}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewRepositoryClient(ctx context.Context, server Server) (context.Context, *RepositoryClient, error) {
|
|
|
|
conn, err := getOrCreateConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
grpcClient := gitalypb.NewRepositoryServiceClient(conn)
|
|
|
|
return withOutgoingMetadata(ctx, server.Features), &RepositoryClient{grpcClient}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// NewNamespaceClient is only used by the Gitaly integration tests at present
|
|
|
|
func NewNamespaceClient(ctx context.Context, server Server) (context.Context, *NamespaceClient, error) {
|
|
|
|
conn, err := getOrCreateConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
grpcClient := gitalypb.NewNamespaceServiceClient(conn)
|
|
|
|
return withOutgoingMetadata(ctx, server.Features), &NamespaceClient{grpcClient}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func NewDiffClient(ctx context.Context, server Server) (context.Context, *DiffClient, error) {
|
|
|
|
conn, err := getOrCreateConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, nil, err
|
|
|
|
}
|
|
|
|
grpcClient := gitalypb.NewDiffServiceClient(conn)
|
|
|
|
return withOutgoingMetadata(ctx, server.Features), &DiffClient{grpcClient}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func getOrCreateConnection(server Server) (*grpc.ClientConn, error) {
|
|
|
|
key := server.cacheKey()
|
|
|
|
|
|
|
|
cache.RLock()
|
|
|
|
conn := cache.connections[key]
|
|
|
|
cache.RUnlock()
|
|
|
|
|
|
|
|
if conn != nil {
|
|
|
|
return conn, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
cache.Lock()
|
|
|
|
defer cache.Unlock()
|
|
|
|
|
|
|
|
if conn := cache.connections[key]; conn != nil {
|
|
|
|
return conn, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
conn, err := newConnection(server)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
cache.connections[key] = conn
|
|
|
|
|
|
|
|
return conn, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func CloseConnections() {
|
|
|
|
cache.Lock()
|
|
|
|
defer cache.Unlock()
|
|
|
|
|
|
|
|
for _, conn := range cache.connections {
|
|
|
|
conn.Close()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func newConnection(server Server) (*grpc.ClientConn, error) {
|
|
|
|
connOpts := append(gitalyclient.DefaultDialOpts,
|
|
|
|
grpc.WithPerRPCCredentials(gitalyauth.RPCCredentialsV2(server.Token)),
|
|
|
|
grpc.WithStreamInterceptor(
|
|
|
|
grpc_middleware.ChainStreamClient(
|
|
|
|
grpctracing.StreamClientTracingInterceptor(),
|
|
|
|
grpc_prometheus.StreamClientInterceptor,
|
|
|
|
grpccorrelation.StreamClientCorrelationInterceptor(
|
|
|
|
grpccorrelation.WithClientName("gitlab-workhorse"),
|
|
|
|
),
|
|
|
|
),
|
|
|
|
),
|
|
|
|
|
|
|
|
grpc.WithUnaryInterceptor(
|
|
|
|
grpc_middleware.ChainUnaryClient(
|
|
|
|
grpctracing.UnaryClientTracingInterceptor(),
|
|
|
|
grpc_prometheus.UnaryClientInterceptor,
|
|
|
|
grpccorrelation.UnaryClientCorrelationInterceptor(
|
|
|
|
grpccorrelation.WithClientName("gitlab-workhorse"),
|
|
|
|
),
|
|
|
|
),
|
|
|
|
),
|
|
|
|
)
|
|
|
|
|
2021-11-18 22:05:49 +05:30
|
|
|
var conn *grpc.ClientConn
|
|
|
|
var connErr error
|
|
|
|
if server.Sidechannel {
|
|
|
|
conn, connErr = gitalyclient.DialSidechannel(context.Background(), server.Address, sidechannelRegistry, connOpts) // lint:allow context.Background
|
|
|
|
} else {
|
|
|
|
conn, connErr = gitalyclient.Dial(server.Address, connOpts)
|
|
|
|
}
|
2021-02-22 17:27:13 +05:30
|
|
|
|
|
|
|
label := "ok"
|
|
|
|
if connErr != nil {
|
|
|
|
label = "fail"
|
|
|
|
}
|
|
|
|
connectionsTotal.WithLabelValues(label).Inc()
|
|
|
|
|
|
|
|
return conn, connErr
|
|
|
|
}
|
|
|
|
|
|
|
|
func UnmarshalJSON(s string, msg proto.Message) error {
|
|
|
|
return jsonUnMarshaler.Unmarshal(strings.NewReader(s), msg)
|
|
|
|
}
|