Files
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

96 lines
2.1 KiB
Go

package kuberesolver
import (
"net"
"strconv"
"sync"
"google.golang.org/grpc/grpclog"
"google.golang.org/grpc/naming"
)
type watchResult struct {
ep *Event
err error
}
// A Watcher provides name resolution updates from Kubernetes endpoints
// identified by name.
type watcher struct {
target targetInfo
endpoints map[string]interface{}
stopCh chan struct{}
result chan watchResult
sync.Mutex
stopped bool
}
// Close closes the watcher, cleaning up any open connections.
func (w *watcher) Close() {
close(w.stopCh)
}
// Next updates the endpoints for the name being watched.
func (w *watcher) Next() ([]*naming.Update, error) {
updates := make([]*naming.Update, 0)
updatedEndpoints := make(map[string]interface{})
var ep Event
select {
case <-w.stopCh:
w.Lock()
if !w.stopped {
w.stopped = true
}
w.Unlock()
return updates, nil
case r := <-w.result:
if r.err == nil {
ep = *r.ep
} else {
return updates, r.err
}
}
for _, subset := range ep.Object.Subsets {
port := ""
if w.target.useFirstPort {
port = strconv.Itoa(subset.Ports[0].Port)
} else if w.target.resolveByPortName {
for _, p := range subset.Ports {
if p.Name == w.target.port {
port = strconv.Itoa(p.Port)
break
}
}
} else {
port = w.target.port
}
if len(port) == 0 {
port = strconv.Itoa(subset.Ports[0].Port)
}
for _, address := range subset.Addresses {
endpoint := net.JoinHostPort(address.IP, port)
updatedEndpoints[endpoint] = nil
}
}
// Create updates to add new endpoints.
for addr, md := range updatedEndpoints {
if _, ok := w.endpoints[addr]; !ok {
updates = append(updates, &naming.Update{naming.Add, addr, md})
grpclog.Printf("kuberesolver: %s ADDED to %s", addr, w.target.target)
}
}
// Create updates to delete old endpoints.
for addr := range w.endpoints {
if _, ok := updatedEndpoints[addr]; !ok {
updates = append(updates, &naming.Update{naming.Delete, addr, nil})
grpclog.Printf("kuberesolver: %s DELETED from %s", addr, w.target.target)
}
}
w.endpoints = updatedEndpoints
return updates, nil
}