Files
weave-scope/vendor/github.com/sercand/kuberesolver/resolver.go
Marc CARRE 4439e6f717 Update weaveworks/common to latest.
```
$ gvt delete github.com/weaveworks/common
$ gvt fetch -revision 15746cbb5831e70c58d520634b09393b7c3cfd56 github.com/weaveworks/common
2017/10/08 XX:XX:XX Fetching: github.com/weaveworks/common
2017/10/08 XX:XX:XX · Skipping (existing): github.com/sirupsen/logrus
2017/10/08 XX:XX:XX · Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX · Skipping (existing): github.com/mgutz/ansi
2017/10/08 XX:XX:XX · Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX · Skipping (existing): github.com/opentracing/opentracing-go/ext
2017/10/08 XX:XX:XX · Skipping (existing): github.com/golang/protobuf/ptypes/empty
2017/10/08 XX:XX:XX · Fetching recursive dependency: github.com/opentracing-contrib/go-stdlib/nethttp
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2017/10/08 XX:XX:XX · Skipping (existing): github.com/aws/aws-sdk-go/aws/credentials
2017/10/08 XX:XX:XX · Skipping (existing): github.com/weaveworks/promrus
2017/10/08 XX:XX:XX · Skipping (existing): github.com/golang/protobuf/ptypes
2017/10/08 XX:XX:XX · Fetching recursive dependency: github.com/grpc-ecosystem/grpc-opentracing/go/otgrpc
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/metadata
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/codes
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/status
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2017/10/08 XX:XX:XX ·· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX · Skipping (existing): golang.org/x/tools/cover
2017/10/08 XX:XX:XX · Skipping (existing): github.com/aws/aws-sdk-go/aws
2017/10/08 XX:XX:XX · Skipping (existing): github.com/armon/go-socks5
2017/10/08 XX:XX:XX · Fetching recursive dependency: github.com/weaveworks-experiments/loki/pkg/client
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/apache/thrift/lib/go/thrift
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/golang.org/x/net/context
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/openzipkin/zipkin-go-opentracing
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/gogo/protobuf/proto
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/Shopify/sarama
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/eapache/queue
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/eapache/go-xerial-snappy
2017/10/08 XX:XX:XX ····· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/golang/snappy
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/pierrec/lz4
2017/10/08 XX:XX:XX ····· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/pierrec/xxHash/xxHash32
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/eapache/go-resiliency/breaker
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/klauspost/crc32
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/davecgh/go-spew/spew
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/rcrowley/go-metrics
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/go-logfmt/logfmt
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/weaveworks-experiments/loki/vendor/github.com/kr/logfmt
2017/10/08 XX:XX:XX · Skipping (existing): google.golang.org/grpc/metadata
2017/10/08 XX:XX:XX · Fetching recursive dependency: github.com/mwitkow/go-grpc-middleware
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/grpclog
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: go.uber.org/zap/zapcore
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/zap/internal/exit
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/multierr
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: go.uber.org/atomic
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/zap/internal/color
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/zap/buffer
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/zap/internal/bufferpool
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/metadata
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware/testing/testproto
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ··· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ··· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/log
2017/10/08 XX:XX:XX ·· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/sirupsen/logrus
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/golang/protobuf/jsonpb
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/ptypes/any
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/ptypes/struct
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/ptypes/duration
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/ptypes/timestamp
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/ptypes/wrappers
2017/10/08 XX:XX:XX ·· Skipping (existing): golang.org/x/net/trace
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/credentials
2017/10/08 XX:XX:XX ·· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware/logging
2017/10/08 XX:XX:XX ··· Skipping (existing): google.golang.org/grpc/grpclog
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: go.uber.org/zap
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/sirupsen/logrus
2017/10/08 XX:XX:XX ··· Skipping (existing): google.golang.org/grpc/codes
2017/10/08 XX:XX:XX ··· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware/tags
2017/10/08 XX:XX:XX ···· Skipping (existing): google.golang.org/grpc/peer
2017/10/08 XX:XX:XX ···· Fetching recursive dependency: github.com/grpc-ecosystem/go-grpc-middleware
2017/10/08 XX:XX:XX ····· Skipping (existing): github.com/opentracing/opentracing-go/log
2017/10/08 XX:XX:XX ····· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX ····· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc/metadata
2017/10/08 XX:XX:XX ····· Fetching recursive dependency: github.com/stretchr/testify/suite
2017/10/08 XX:XX:XX ······ Skipping (existing): github.com/stretchr/testify/assert
2017/10/08 XX:XX:XX ······ Fetching recursive dependency: github.com/stretchr/testify/require
2017/10/08 XX:XX:XX ······· Skipping (existing): github.com/stretchr/testify/assert
2017/10/08 XX:XX:XX ····· Skipping (existing): github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX ····· Skipping (existing): github.com/sirupsen/logrus
2017/10/08 XX:XX:XX ····· Fetching recursive dependency: github.com/gogo/protobuf/gogoproto
2017/10/08 XX:XX:XX ······ Fetching recursive dependency: github.com/gogo/protobuf/proto
2017/10/08 XX:XX:XX ······· Fetching recursive dependency: github.com/gogo/protobuf/types
2017/10/08 XX:XX:XX ········ Fetching recursive dependency: github.com/gogo/protobuf/sortkeys
2017/10/08 XX:XX:XX ······ Fetching recursive dependency: github.com/gogo/protobuf/protoc-gen-gogo/descriptor
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc/codes
2017/10/08 XX:XX:XX ····· Skipping (existing): github.com/opentracing/opentracing-go/ext
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc/grpclog
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc/peer
2017/10/08 XX:XX:XX ····· Skipping (existing): golang.org/x/net/trace
2017/10/08 XX:XX:XX ····· Skipping (existing): google.golang.org/grpc/credentials
2017/10/08 XX:XX:XX ···· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ···· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX ··· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ··· Skipping (existing): golang.org/x/net/context
2017/10/08 XX:XX:XX ··· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/opentracing/opentracing-go/ext
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/codes
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/peer
2017/10/08 XX:XX:XX · Skipping (existing): github.com/davecgh/go-spew/spew
2017/10/08 XX:XX:XX · Skipping (existing): google.golang.org/grpc/status
2017/10/08 XX:XX:XX · Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX · Skipping (existing): github.com/golang/protobuf/ptypes/any
2017/10/08 XX:XX:XX · Skipping (existing): github.com/opentracing/opentracing-go
2017/10/08 XX:XX:XX · Skipping (existing): github.com/prometheus/client_golang/prometheus
2017/10/08 XX:XX:XX · Fetching recursive dependency: google.golang.org/genproto/googleapis/rpc/status
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/golang/protobuf/proto
2017/10/08 XX:XX:XX ·· Skipping (existing): github.com/golang/protobuf/ptypes/any
2017/10/08 XX:XX:XX · Skipping (existing): github.com/pmezard/go-difflib/difflib
2017/10/08 XX:XX:XX · Fetching recursive dependency: github.com/sercand/kuberesolver
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/naming
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc/grpclog
2017/10/08 XX:XX:XX ·· Skipping (existing): google.golang.org/grpc
2017/10/08 XX:XX:XX · Skipping (existing): github.com/gorilla/mux
```
2017-10-08 18:37:28 +01:00

85 lines
1.9 KiB
Go

package kuberesolver
import (
"fmt"
"net/http"
"net/url"
"time"
"google.golang.org/grpc/grpclog"
"google.golang.org/grpc/naming"
)
// kubeResolver resolves service names using Kubernetes endpoints.
type kubeResolver struct {
k8sClient *k8sClient
namespace string
watcher *watcher
}
// NewResolver returns a new Kubernetes resolver.
func newResolver(client *k8sClient, namespace string) *kubeResolver {
if namespace == "" {
namespace = "default"
}
return &kubeResolver{client, namespace, nil}
}
// Resolve creates a Kubernetes watcher for the named target.
func (r *kubeResolver) Resolve(target string) (naming.Watcher, error) {
pt, err := parseTarget(target)
if err != nil {
return nil, err
}
resultChan := make(chan watchResult)
stopCh := make(chan struct{})
wtarget := pt.target
go until(func() {
err := r.watch(wtarget, stopCh, resultChan)
if err != nil {
grpclog.Printf("kuberesolver: watching ended with error='%v', will reconnect again", err)
}
}, time.Second, stopCh)
r.watcher = &watcher{
target: pt,
endpoints: make(map[string]interface{}),
stopCh: stopCh,
result: resultChan,
}
return r.watcher, nil
}
func (r *kubeResolver) watch(target string, stopCh <-chan struct{}, resultCh chan<- watchResult) error {
u, err := url.Parse(fmt.Sprintf("%s/api/v1/watch/namespaces/%s/endpoints/%s",
r.k8sClient.host, r.namespace, target))
if err != nil {
return err
}
req, err := r.k8sClient.getRequest(u.String())
if err != nil {
return err
}
resp, err := r.k8sClient.Do(req)
if err != nil {
return err
}
if resp.StatusCode != http.StatusOK {
defer resp.Body.Close()
return fmt.Errorf("invalid response code %d", resp.StatusCode)
}
sw := newStreamWatcher(resp.Body)
for {
select {
case <-stopCh:
return nil
case up, more := <-sw.ResultChan():
if more {
resultCh <- watchResult{err: nil, ep: &up}
} else {
return nil
}
}
}
}