Migrate from libkv to valkeyrie library

This commit is contained in:
NicoMen 2018-01-24 17:52:03 +01:00 committed by Traefiker
parent a91080b060
commit 563a0bd274
33 changed files with 109 additions and 115 deletions

37
Gopkg.lock generated
View file

@ -126,6 +126,20 @@
revision = "65b0cdae8d7fe5c05c7430e055938ef6d24a66c9" revision = "65b0cdae8d7fe5c05c7430e055938ef6d24a66c9"
source = "github.com/containous/go-http-auth" source = "github.com/containous/go-http-auth"
[[projects]]
branch = "master"
name = "github.com/abronan/valkeyrie"
packages = [
".",
"store",
"store/boltdb",
"store/consul",
"store/etcd/v2",
"store/etcd/v3",
"store/zookeeper"
]
revision = "063d875e3c5fd734fa2aa12fac83829f62acfc70"
[[projects]] [[projects]]
name = "github.com/aokoli/goutils" name = "github.com/aokoli/goutils"
packages = ["."] packages = ["."]
@ -235,8 +249,8 @@
[[projects]] [[projects]]
name = "github.com/containous/staert" name = "github.com/containous/staert"
packages = ["."] packages = ["."]
revision = "af517d5b70db9c4b0505e0144fcc62b054057d2a" revision = "dbb7c840f31daec0d863ada6829d075a8dbb7469"
version = "v2.0.0" version = "v3.0.0"
[[projects]] [[projects]]
name = "github.com/containous/traefik-extra-service-fabric" name = "github.com/containous/traefik-extra-service-fabric"
@ -427,7 +441,7 @@
branch = "master" branch = "master"
name = "github.com/docker/leadership" name = "github.com/docker/leadership"
packages = ["."] packages = ["."]
revision = "af20da7d3e62be9259835e93261acf931b5adecf" revision = "a2e096d9fe0af5b4c37dd37aea719bc9c2e5eec6"
source = "github.com/containous/leadership" source = "github.com/containous/leadership"
[[projects]] [[projects]]
@ -457,21 +471,6 @@
] ]
revision = "57bd716502dcbe1799f026148016022b0f3b989c" revision = "57bd716502dcbe1799f026148016022b0f3b989c"
[[projects]]
branch = "master"
name = "github.com/docker/libkv"
packages = [
".",
"store",
"store/boltdb",
"store/consul",
"store/etcd/v2",
"store/etcd/v3",
"store/zookeeper"
]
revision = "5e4bb288a9a74320bb03f5c18d6bdbab0d8049de"
source = "github.com/abronan/libkv"
[[projects]] [[projects]]
name = "github.com/donovanhide/eventsource" name = "github.com/donovanhide/eventsource"
packages = ["."] packages = ["."]
@ -1513,6 +1512,6 @@
[solve-meta] [solve-meta]
analyzer-name = "dep" analyzer-name = "dep"
analyzer-version = 1 analyzer-version = 1
inputs-digest = "abd30ec16123cf6de5a869dcaba3e3247907785b4802f0eebaa0cf4d16ef444b" inputs-digest = "0d3d2cd01e06cfcc86cbd261395d59187d4d27d88398bd7aeb1298676becd3e7"
solver-name = "gps-cdcl" solver-name = "gps-cdcl"
solver-version = 1 solver-version = 1

View file

@ -62,7 +62,7 @@
[[constraint]] [[constraint]]
name = "github.com/containous/staert" name = "github.com/containous/staert"
version = "2.0.0" version = "3.0.0"
[[constraint]] [[constraint]]
name = "github.com/containous/traefik-extra-service-fabric" name = "github.com/containous/traefik-extra-service-fabric"
@ -77,10 +77,6 @@
name = "github.com/docker/leadership" name = "github.com/docker/leadership"
source = "github.com/containous/leadership" source = "github.com/containous/leadership"
[[constraint]]
name = "github.com/docker/libkv"
source = "github.com/abronan/libkv"
[[constraint]] [[constraint]]
name = "github.com/eapache/channels" name = "github.com/eapache/channels"
version = "1.1.0" version = "1.1.0"
@ -115,6 +111,10 @@
branch = "master" branch = "master"
name = "github.com/jjcollinge/servicefabric" name = "github.com/jjcollinge/servicefabric"
[[constraint]]
branch = "master"
name = "github.com/abronan/valkeyrie"
[[constraint]] [[constraint]]
name = "github.com/mattn/go-shellwords" name = "github.com/mattn/go-shellwords"
version = "1.0.3" version = "1.0.3"

View file

@ -7,12 +7,12 @@ import (
"sync" "sync"
"time" "time"
"github.com/abronan/valkeyrie/store"
"github.com/cenk/backoff" "github.com/cenk/backoff"
"github.com/containous/staert" "github.com/containous/staert"
"github.com/containous/traefik/job" "github.com/containous/traefik/job"
"github.com/containous/traefik/log" "github.com/containous/traefik/log"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/docker/libkv/store"
"github.com/satori/go.uuid" "github.com/satori/go.uuid"
) )

View file

@ -5,11 +5,11 @@ import (
"fmt" "fmt"
stdlog "log" stdlog "log"
"github.com/abronan/valkeyrie/store"
"github.com/containous/flaeg" "github.com/containous/flaeg"
"github.com/containous/staert" "github.com/containous/staert"
"github.com/containous/traefik/acme" "github.com/containous/traefik/acme"
"github.com/containous/traefik/cluster" "github.com/containous/traefik/cluster"
"github.com/docker/libkv/store"
) )
func newStoreConfigCmd(traefikConfiguration *TraefikConfiguration, traefikPointersConfiguration *TraefikConfiguration) *flaeg.Command { func newStoreConfigCmd(traefikConfiguration *TraefikConfiguration, traefikPointersConfiguration *TraefikConfiguration) *flaeg.Command {

View file

@ -10,13 +10,13 @@ import (
"sync" "sync"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/consul"
"github.com/containous/staert" "github.com/containous/staert"
"github.com/containous/traefik/cluster" "github.com/containous/traefik/cluster"
"github.com/containous/traefik/integration/try" "github.com/containous/traefik/integration/try"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/consul"
"github.com/go-check/check" "github.com/go-check/check"
checker "github.com/vdemeester/shakers" checker "github.com/vdemeester/shakers"
) )
@ -32,7 +32,7 @@ func (s *ConsulSuite) setupConsul(c *check.C) {
s.composeProject.Start(c) s.composeProject.Start(c)
consul.Register() consul.Register()
kv, err := libkv.NewStore( kv, err := valkeyrie.NewStore(
store.CONSUL, store.CONSUL,
[]string{s.composeProject.Container(c, "consul").NetworkSettings.IPAddress + ":8500"}, []string{s.composeProject.Container(c, "consul").NetworkSettings.IPAddress + ":8500"},
&store.Config{ &store.Config{
@ -63,7 +63,7 @@ func (s *ConsulSuite) setupConsulTLS(c *check.C) {
TLSConfig, err := clientTLS.CreateTLSConfig() TLSConfig, err := clientTLS.CreateTLSConfig()
c.Assert(err, checker.IsNil) c.Assert(err, checker.IsNil)
kv, err := libkv.NewStore( kv, err := valkeyrie.NewStore(
store.CONSUL, store.CONSUL,
[]string{s.composeProject.Container(c, "consul").NetworkSettings.IPAddress + ":8585"}, []string{s.composeProject.Container(c, "consul").NetworkSettings.IPAddress + ":8585"},
&store.Config{ &store.Config{

View file

@ -8,10 +8,10 @@ import (
"strings" "strings"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/etcd/v3"
"github.com/containous/traefik/integration/try" "github.com/containous/traefik/integration/try"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/etcd/v3"
"github.com/go-check/check" "github.com/go-check/check"
checker "github.com/vdemeester/shakers" checker "github.com/vdemeester/shakers"
@ -41,7 +41,7 @@ func (s *Etcd3Suite) SetUpTest(c *check.C) {
etcdv3.Register() etcdv3.Register()
url := ipEtcd + ":2379" url := ipEtcd + ":2379"
kv, err := libkv.NewStore( kv, err := valkeyrie.NewStore(
store.ETCDV3, store.ETCDV3,
[]string{url}, []string{url},
&store.Config{ &store.Config{

View file

@ -8,10 +8,10 @@ import (
"strings" "strings"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/etcd/v2"
"github.com/containous/traefik/integration/try" "github.com/containous/traefik/integration/try"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/etcd/v2"
"github.com/go-check/check" "github.com/go-check/check"
checker "github.com/vdemeester/shakers" checker "github.com/vdemeester/shakers"
@ -29,7 +29,7 @@ func (s *EtcdSuite) SetUpTest(c *check.C) {
etcd.Register() etcd.Register()
url := s.composeProject.Container(c, "etcd").NetworkSettings.IPAddress + ":2379" url := s.composeProject.Container(c, "etcd").NetworkSettings.IPAddress + ":2379"
kv, err := libkv.NewStore( kv, err := valkeyrie.NewStore(
store.ETCD, store.ETCD,
[]string{url}, []string{url},
&store.Config{ &store.Config{

View file

@ -7,7 +7,7 @@ import (
"net/http" "net/http"
"strings" "strings"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
) )
// ResponseCondition is a retry condition function. // ResponseCondition is a retry condition function.

View file

@ -3,12 +3,12 @@ package boltdb
import ( import (
"fmt" "fmt"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/boltdb"
"github.com/containous/traefik/provider" "github.com/containous/traefik/provider"
"github.com/containous/traefik/provider/kv" "github.com/containous/traefik/provider/kv"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/boltdb"
) )
var _ provider.Provider = (*Provider)(nil) var _ provider.Provider = (*Provider)(nil)
@ -23,7 +23,7 @@ type Provider struct {
func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error { func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error {
store, err := p.CreateStore() store, err := p.CreateStore()
if err != nil { if err != nil {
return fmt.Errorf("Failed to Connect to KV store: %v", err) return fmt.Errorf("failed to Connect to KV store: %v", err)
} }
p.SetKVClient(store) p.SetKVClient(store)
return p.Provider.Provide(configurationChan, pool, constraints) return p.Provider.Provide(configurationChan, pool, constraints)

View file

@ -3,12 +3,12 @@ package consul
import ( import (
"fmt" "fmt"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/consul"
"github.com/containous/traefik/provider" "github.com/containous/traefik/provider"
"github.com/containous/traefik/provider/kv" "github.com/containous/traefik/provider/kv"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/consul"
) )
var _ provider.Provider = (*Provider)(nil) var _ provider.Provider = (*Provider)(nil)
@ -23,7 +23,7 @@ type Provider struct {
func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error { func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error {
store, err := p.CreateStore() store, err := p.CreateStore()
if err != nil { if err != nil {
return fmt.Errorf("Failed to Connect to KV store: %v", err) return fmt.Errorf("failed to Connect to KV store: %v", err)
} }
p.SetKVClient(store) p.SetKVClient(store)
return p.Provider.Provide(configurationChan, pool, constraints) return p.Provider.Provide(configurationChan, pool, constraints)

View file

@ -3,14 +3,14 @@ package etcd
import ( import (
"fmt" "fmt"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/etcd/v2"
"github.com/abronan/valkeyrie/store/etcd/v3"
"github.com/containous/traefik/log" "github.com/containous/traefik/log"
"github.com/containous/traefik/provider" "github.com/containous/traefik/provider"
"github.com/containous/traefik/provider/kv" "github.com/containous/traefik/provider/kv"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/etcd/v2"
"github.com/docker/libkv/store/etcd/v3"
) )
var _ provider.Provider = (*Provider)(nil) var _ provider.Provider = (*Provider)(nil)
@ -26,7 +26,7 @@ type Provider struct {
func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error { func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error {
store, err := p.CreateStore() store, err := p.CreateStore()
if err != nil { if err != nil {
return fmt.Errorf("Failed to Connect to KV store: %v", err) return fmt.Errorf("failed to Connect to KV store: %v", err)
} }
p.SetKVClient(store) p.SetKVClient(store)
return p.Provider.Provide(configurationChan, pool, constraints) return p.Provider.Provide(configurationChan, pool, constraints)

View file

@ -5,7 +5,7 @@ import (
"strings" "strings"
"testing" "testing"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
) )

View file

@ -6,14 +6,14 @@ import (
"strings" "strings"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/cenk/backoff" "github.com/cenk/backoff"
"github.com/containous/traefik/job" "github.com/containous/traefik/job"
"github.com/containous/traefik/log" "github.com/containous/traefik/log"
"github.com/containous/traefik/provider" "github.com/containous/traefik/provider"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
) )
// Provider holds common configurations of key-value providers. // Provider holds common configurations of key-value providers.
@ -44,7 +44,7 @@ func (p *Provider) CreateStore() (store.Store, error) {
return nil, err return nil, err
} }
} }
return libkv.NewStore( return valkeyrie.NewStore(
p.storeType, p.storeType,
strings.Split(p.Endpoint, ","), strings.Split(p.Endpoint, ","),
storeConfig, storeConfig,

View file

@ -10,12 +10,12 @@ import (
"text/template" "text/template"
"github.com/BurntSushi/ty/fun" "github.com/BurntSushi/ty/fun"
"github.com/abronan/valkeyrie/store"
"github.com/containous/flaeg" "github.com/containous/flaeg"
"github.com/containous/traefik/log" "github.com/containous/traefik/log"
"github.com/containous/traefik/provider/label" "github.com/containous/traefik/provider/label"
"github.com/containous/traefik/tls" "github.com/containous/traefik/tls"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
) )
func (p *Provider) buildConfiguration() *types.Configuration { func (p *Provider) buildConfiguration() *types.Configuration {

View file

@ -6,11 +6,11 @@ import (
"testing" "testing"
"time" "time"
"github.com/abronan/valkeyrie/store"
"github.com/containous/flaeg" "github.com/containous/flaeg"
"github.com/containous/traefik/provider/label" "github.com/containous/traefik/provider/label"
"github.com/containous/traefik/tls" "github.com/containous/traefik/tls"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
) )

View file

@ -4,7 +4,7 @@ import (
"errors" "errors"
"strings" "strings"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
) )
func newProviderMock(kvPairs []*store.KVPair) *Provider { func newProviderMock(kvPairs []*store.KVPair) *Provider {

View file

@ -4,8 +4,8 @@ import (
"testing" "testing"
"time" "time"
"github.com/abronan/valkeyrie/store"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
) )
func TestKvWatchTree(t *testing.T) { func TestKvWatchTree(t *testing.T) {

View file

@ -3,12 +3,12 @@ package zk
import ( import (
"fmt" "fmt"
"github.com/abronan/valkeyrie/store"
"github.com/abronan/valkeyrie/store/zookeeper"
"github.com/containous/traefik/provider" "github.com/containous/traefik/provider"
"github.com/containous/traefik/provider/kv" "github.com/containous/traefik/provider/kv"
"github.com/containous/traefik/safe" "github.com/containous/traefik/safe"
"github.com/containous/traefik/types" "github.com/containous/traefik/types"
"github.com/docker/libkv/store"
"github.com/docker/libkv/store/zookeeper"
) )
var _ provider.Provider = (*Provider)(nil) var _ provider.Provider = (*Provider)(nil)
@ -23,7 +23,7 @@ type Provider struct {
func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error { func (p *Provider) Provide(configurationChan chan<- types.ConfigMessage, pool *safe.Pool, constraints types.Constraints) error {
store, err := p.CreateStore() store, err := p.CreateStore()
if err != nil { if err != nil {
return fmt.Errorf("Failed to Connect to KV store: %v", err) return fmt.Errorf("failed to Connect to KV store: %v", err)
} }
p.SetKVClient(store) p.SetKVClient(store)
return p.Provider.Provide(configurationChan, pool, constraints) return p.Provider.Provide(configurationChan, pool, constraints)

View file

@ -11,10 +11,10 @@ import (
"strconv" "strconv"
"strings" "strings"
"github.com/abronan/valkeyrie/store"
"github.com/containous/flaeg" "github.com/containous/flaeg"
"github.com/containous/traefik/log" "github.com/containous/traefik/log"
traefikTls "github.com/containous/traefik/tls" traefikTls "github.com/containous/traefik/tls"
"github.com/docker/libkv/store"
"github.com/ryanuber/go-glob" "github.com/ryanuber/go-glob"
) )

View file

@ -10,9 +10,9 @@ import (
"sync/atomic" "sync/atomic"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/coreos/bbolt" "github.com/coreos/bbolt"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
) )
var ( var (
@ -34,7 +34,7 @@ type BoltDB struct {
dbIndex uint64 dbIndex uint64
path string path string
timeout time.Duration timeout time.Duration
// By default libkv opens and closes the bolt DB connection for every // By default valkeyrie opens and closes the bolt DB connection for every
// get/put operation. This allows multiple apps to use a Bolt DB at the // get/put operation. This allows multiple apps to use a Bolt DB at the
// same time. // same time.
// PersistConnection flag provides an option to override ths behavior. // PersistConnection flag provides an option to override ths behavior.
@ -44,13 +44,13 @@ type BoltDB struct {
} }
const ( const (
libkvmetadatalen = 8 metadatalen = 8
transientTimeout = time.Duration(10) * time.Second transientTimeout = time.Duration(10) * time.Second
) )
// Register registers boltdb to libkv // Register registers boltdb to valkeyrie
func Register() { func Register() {
libkv.AddStore(store.BOLTDB, New) valkeyrie.AddStore(store.BOLTDB, New)
} }
// New opens a new BoltDB connection to the specified path and bucket // New opens a new BoltDB connection to the specified path and bucket
@ -126,7 +126,7 @@ func (b *BoltDB) releaseDBhandle() {
} }
// Get the value at "key". BoltDB doesn't provide an inbuilt last modified index with every kv pair. Its implemented by // Get the value at "key". BoltDB doesn't provide an inbuilt last modified index with every kv pair. Its implemented by
// by a atomic counter maintained by the libkv and appened to the value passed by the client. // by a atomic counter maintained by the valkeyrie and appened to the value passed by the client.
func (b *BoltDB) Get(key string, opts *store.ReadOptions) (*store.KVPair, error) { func (b *BoltDB) Get(key string, opts *store.ReadOptions) (*store.KVPair, error) {
var ( var (
val []byte val []byte
@ -161,8 +161,8 @@ func (b *BoltDB) Get(key string, opts *store.ReadOptions) (*store.KVPair, error)
return nil, err return nil, err
} }
dbIndex := binary.LittleEndian.Uint64(val[:libkvmetadatalen]) dbIndex := binary.LittleEndian.Uint64(val[:metadatalen])
val = val[libkvmetadatalen:] val = val[metadatalen:]
return &store.KVPair{Key: key, Value: val, LastIndex: (dbIndex)}, nil return &store.KVPair{Key: key, Value: val, LastIndex: (dbIndex)}, nil
} }
@ -177,7 +177,7 @@ func (b *BoltDB) Put(key string, value []byte, opts *store.WriteOptions) error {
b.Lock() b.Lock()
defer b.Unlock() defer b.Unlock()
dbval := make([]byte, libkvmetadatalen) dbval := make([]byte, metadatalen)
if db, err = b.getDBhandle(); err != nil { if db, err = b.getDBhandle(); err != nil {
return err return err
@ -287,8 +287,8 @@ func (b *BoltDB) List(keyPrefix string, opts *store.ReadOptions) ([]*store.KVPai
for key, v := cursor.Seek(prefix); bytes.HasPrefix(key, prefix); key, v = cursor.Next() { for key, v := cursor.Seek(prefix); bytes.HasPrefix(key, prefix); key, v = cursor.Next() {
hasResult = true hasResult = true
dbIndex := binary.LittleEndian.Uint64(v[:libkvmetadatalen]) dbIndex := binary.LittleEndian.Uint64(v[:metadatalen])
v = v[libkvmetadatalen:] v = v[metadatalen:]
val := make([]byte, len(v)) val := make([]byte, len(v))
copy(val, v) copy(val, v)
@ -338,7 +338,7 @@ func (b *BoltDB) AtomicDelete(key string, previous *store.KVPair) (bool, error)
if val == nil { if val == nil {
return store.ErrKeyNotFound return store.ErrKeyNotFound
} }
dbIndex := binary.LittleEndian.Uint64(val[:libkvmetadatalen]) dbIndex := binary.LittleEndian.Uint64(val[:metadatalen])
if dbIndex != previous.LastIndex { if dbIndex != previous.LastIndex {
return store.ErrKeyModified return store.ErrKeyModified
} }
@ -363,7 +363,7 @@ func (b *BoltDB) AtomicPut(key string, value []byte, previous *store.KVPair, opt
b.Lock() b.Lock()
defer b.Unlock() defer b.Unlock()
dbval := make([]byte, libkvmetadatalen) dbval := make([]byte, metadatalen)
if db, err = b.getDBhandle(); err != nil { if db, err = b.getDBhandle(); err != nil {
return false, nil, err return false, nil, err
@ -392,7 +392,7 @@ func (b *BoltDB) AtomicPut(key string, value []byte, previous *store.KVPair, opt
if len(val) == 0 { if len(val) == 0 {
return store.ErrKeyNotFound return store.ErrKeyNotFound
} }
dbIndex = binary.LittleEndian.Uint64(val[:libkvmetadatalen]) dbIndex = binary.LittleEndian.Uint64(val[:metadatalen])
if dbIndex != previous.LastIndex { if dbIndex != previous.LastIndex {
return store.ErrKeyModified return store.ErrKeyModified
} }

View file

@ -8,8 +8,8 @@ import (
"sync" "sync"
"time" "time"
"github.com/docker/libkv" "github.com/abronan/valkeyrie"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
api "github.com/hashicorp/consul/api" api "github.com/hashicorp/consul/api"
) )
@ -55,9 +55,9 @@ type consulLock struct {
renewCh chan struct{} renewCh chan struct{}
} }
// Register registers consul to libkv // Register registers consul to valkeyrie
func Register() { func Register() {
libkv.AddStore(store.CONSUL, New) valkeyrie.AddStore(store.CONSUL, New)
} }
// New creates a new Consul client given a list // New creates a new Consul client given a list

View file

@ -12,9 +12,9 @@ import (
"golang.org/x/net/context" "golang.org/x/net/context"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
etcd "github.com/coreos/etcd/client" etcd "github.com/coreos/etcd/client"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
) )
const ( const (
@ -53,9 +53,9 @@ const (
defaultUpdateTime = 5 * time.Second defaultUpdateTime = 5 * time.Second
) )
// Register registers etcd to libkv // Register registers etcd to valkeyrie
func Register() { func Register() {
libkv.AddStore(store.ETCD, New) valkeyrie.AddStore(store.ETCD, New)
} }
// New creates a new Etcd client given a list // New creates a new Etcd client given a list

View file

@ -7,10 +7,10 @@ import (
"sync" "sync"
"time" "time"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
etcd "github.com/coreos/etcd/clientv3" etcd "github.com/coreos/etcd/clientv3"
"github.com/coreos/etcd/clientv3/concurrency" "github.com/coreos/etcd/clientv3/concurrency"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
) )
const ( const (
@ -38,9 +38,9 @@ type etcdLock struct {
ttl time.Duration ttl time.Duration
} }
// Register registers etcd to libkv // Register registers etcd to valkeyrie
func Register() { func Register() {
libkv.AddStore(store.ETCDV3, New) valkeyrie.AddStore(store.ETCDV3, New)
} }
// New creates a new Etcd client given a list // New creates a new Etcd client given a list

View file

@ -25,7 +25,7 @@ const (
) )
var ( var (
// ErrBackendNotSupported is thrown when the backend k/v store is not supported by libkv // ErrBackendNotSupported is thrown when the backend k/v store is not supported by valkeyrie
ErrBackendNotSupported = errors.New("Backend storage not supported yet, please choose one of") ErrBackendNotSupported = errors.New("Backend storage not supported yet, please choose one of")
// ErrCallNotSupported is thrown when a method is not implemented/supported by the current backend // ErrCallNotSupported is thrown when a method is not implemented/supported by the current backend
ErrCallNotSupported = errors.New("The current call is not supported with this backend") ErrCallNotSupported = errors.New("The current call is not supported with this backend")
@ -66,7 +66,7 @@ type ClientTLSConfig struct {
// Store represents the backend K/V storage // Store represents the backend K/V storage
// Each store should support every call listed // Each store should support every call listed
// here. Or it couldn't be implemented as a K/V // here. Or it couldn't be implemented as a K/V
// backend for libkv // backend for valkeyrie
type Store interface { type Store interface {
// Put a value at the specified key // Put a value at the specified key
Put(key string, value []byte, options *WriteOptions) error Put(key string, value []byte, options *WriteOptions) error

View file

@ -4,8 +4,8 @@ import (
"strings" "strings"
"time" "time"
"github.com/docker/libkv" "github.com/abronan/valkeyrie"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
zk "github.com/samuel/go-zookeeper/zk" zk "github.com/samuel/go-zookeeper/zk"
) )
@ -32,9 +32,9 @@ type zookeeperLock struct {
value []byte value []byte
} }
// Register registers zookeeper to libkv // Register registers zookeeper to valkeyrie
func Register() { func Register() {
libkv.AddStore(store.ZK, New) valkeyrie.AddStore(store.ZK, New)
} }
// New creates a new Zookeeper client given a // New creates a new Zookeeper client given a
@ -485,7 +485,7 @@ func (s *Zookeeper) get(key string) ([]byte, *zk.Stat, error) {
var meta *zk.Stat var meta *zk.Stat
var err error var err error
// To guard against older versions of libkv // To guard against older versions of valkeyrie
// creating and writing to znodes non-atomically, // creating and writing to znodes non-atomically,
// We try to resync few times if we read SOH or // We try to resync few times if we read SOH or
// an empty string // an empty string
@ -518,7 +518,7 @@ func (s *Zookeeper) getW(key string) ([]byte, *zk.Stat, <-chan zk.Event, error)
var ech <-chan zk.Event var ech <-chan zk.Event
var err error var err error
// To guard against older versions of libkv // To guard against older versions of valkeyrie
// creating and writing to znodes non-atomically, // creating and writing to znodes non-atomically,
// We try to resync few times if we read SOH or // We try to resync few times if we read SOH or
// an empty string // an empty string

View file

@ -1,11 +1,11 @@
package libkv package valkeyrie
import ( import (
"fmt" "fmt"
"sort" "sort"
"strings" "strings"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
) )
// Initialize creates a new Store object, initializing the client // Initialize creates a new Store object, initializing the client
@ -34,7 +34,7 @@ func NewStore(backend store.Backend, addrs []string, options *store.Config) (sto
return nil, fmt.Errorf("%s %s", store.ErrBackendNotSupported.Error(), supportedBackend) return nil, fmt.Errorf("%s %s", store.ErrBackendNotSupported.Error(), supportedBackend)
} }
// AddStore adds a new store backend to libkv // AddStore adds a new store backend to valkeyrie
func AddStore(store store.Backend, init Initialize) { func AddStore(store store.Backend, init Initialize) {
initializers[store] = init initializers[store] = init
} }

View file

@ -10,9 +10,9 @@ import (
"strconv" "strconv"
"strings" "strings"
"github.com/abronan/valkeyrie"
"github.com/abronan/valkeyrie/store"
"github.com/containous/flaeg" "github.com/containous/flaeg"
"github.com/docker/libkv"
"github.com/docker/libkv/store"
"github.com/mitchellh/mapstructure" "github.com/mitchellh/mapstructure"
) )
@ -27,11 +27,11 @@ type KvSource struct {
// NewKvSource creates a new KvSource // NewKvSource creates a new KvSource
func NewKvSource(backend store.Backend, addrs []string, options *store.Config, prefix string) (*KvSource, error) { func NewKvSource(backend store.Backend, addrs []string, options *store.Config, prefix string) (*KvSource, error) {
kvStore, err := libkv.NewStore(backend, addrs, options) kvStore, err := valkeyrie.NewStore(backend, addrs, options)
return &KvSource{Store: kvStore, Prefix: prefix}, err return &KvSource{Store: kvStore, Prefix: prefix}, err
} }
// Parse uses libkv and mapstructure to fill the structure // Parse uses valkeyrie and mapstructure to fill the structure
func (kv *KvSource) Parse(cmd *flaeg.Command) (*flaeg.Command, error) { func (kv *KvSource) Parse(cmd *flaeg.Command) (*flaeg.Command, error) {
err := kv.LoadConfig(cmd.Config) err := kv.LoadConfig(cmd.Config)
if err != nil { if err != nil {

View file

@ -105,14 +105,8 @@ func (ts *TomlSource) ConfigFileUsed() string {
func preprocessDir(dirIn string) (string, error) { func preprocessDir(dirIn string) (string, error) {
dirOut := dirIn dirOut := dirIn
if strings.HasPrefix(dirIn, "$") { expanded := os.ExpandEnv(dirIn)
end := strings.Index(dirIn, string(os.PathSeparator)) dirOut, err := filepath.Abs(expanded)
if end == -1 {
end = len(dirIn)
}
dirOut = os.Getenv(dirIn[1:end]) + dirIn[end:]
}
dirOut, err := filepath.Abs(dirOut)
return dirOut, err return dirOut, err
} }
@ -123,7 +117,8 @@ func findFile(filename string, dirNfile []string) string {
if fileInfo, err := os.Stat(fullPath); err == nil && !fileInfo.IsDir() { if fileInfo, err := os.Stat(fullPath); err == nil && !fileInfo.IsDir() {
return fullPath return fullPath
} }
fullPath = fullPath + "/" + filename + ".toml"
fullPath = filepath.Join(fullPath, filename+".toml")
if fileInfo, err := os.Stat(fullPath); err == nil && !fileInfo.IsDir() { if fileInfo, err := os.Stat(fullPath); err == nil && !fileInfo.IsDir() {
return fullPath return fullPath
} }

View file

@ -4,7 +4,7 @@ import (
"sync" "sync"
"time" "time"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
) )
const ( const (

View file

@ -3,7 +3,7 @@ package leadership
import ( import (
"errors" "errors"
"github.com/docker/libkv/store" "github.com/abronan/valkeyrie/store"
) )
// Follower can follow an election in real-time and push notifications whenever // Follower can follow an election in real-time and push notifications whenever