Files
weave-scope/vendor/github.com/sercand/kuberesolver/stream.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

106 lines
2.5 KiB
Go

package kuberesolver
import (
"encoding/json"
"fmt"
"io"
"sync"
"google.golang.org/grpc/grpclog"
)
// Interface can be implemented by anything that knows how to watch and report changes.
type watchInterface interface {
// Stops watching. Will close the channel returned by ResultChan(). Releases
// any resources used by the watch.
Stop()
// Returns a chan which will receive all the events. If an error occurs
// or Stop() is called, this channel will be closed, in which case the
// watch should be completely cleaned up.
ResultChan() <-chan Event
}
// StreamWatcher turns any stream for which you can write a Decoder interface
// into a watch.Interface.
type streamWatcher struct {
result chan Event
r io.ReadCloser
decoder *json.Decoder
sync.Mutex
stopped bool
}
// NewStreamWatcher creates a StreamWatcher from the given io.ReadClosers.
func newStreamWatcher(r io.ReadCloser) watchInterface {
sw := &streamWatcher{
r: r,
decoder: json.NewDecoder(r),
result: make(chan Event),
}
go sw.receive()
return sw
}
// ResultChan implements Interface.
func (sw *streamWatcher) ResultChan() <-chan Event {
return sw.result
}
// Stop implements Interface.
func (sw *streamWatcher) Stop() {
sw.Lock()
defer sw.Unlock()
if !sw.stopped {
sw.stopped = true
sw.r.Close()
}
}
// stopping returns true if Stop() was called previously.
func (sw *streamWatcher) stopping() bool {
sw.Lock()
defer sw.Unlock()
return sw.stopped
}
// receive reads result from the decoder in a loop and sends down the result channel.
func (sw *streamWatcher) receive() {
defer close(sw.result)
defer sw.Stop()
for {
obj, err := sw.Decode()
if err != nil {
// Ignore expected error.
if sw.stopping() {
return
}
switch err {
case io.EOF:
// watch closed normally
case io.ErrUnexpectedEOF:
grpclog.Printf("kuberesolver: Unexpected EOF during watch stream event decoding: %v", err)
default:
grpclog.Printf("kuberesolver: Unable to decode an event from the watch stream: %v", err)
}
return
}
sw.result <- obj
}
}
// Decode blocks until it can return the next object in the writer. Returns an error
// if the writer is closed or an object can't be decoded.
func (sw *streamWatcher) Decode() (Event, error) {
var got Event
if err := sw.decoder.Decode(&got); err != nil {
return Event{}, err
}
switch got.Type {
case Added, Modified, Deleted, Error:
return got, nil
default:
return Event{}, fmt.Errorf("got invalid watch event type: %v", got.Type)
}
}