107 lines
2 KiB
Go
107 lines
2 KiB
Go
package etcd
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"strings"
|
|
|
|
"git.giftfish.de/ston1th/haproxy-lb/pkg/cluster"
|
|
"github.com/go-logr/logr"
|
|
clientv3 "go.etcd.io/etcd/client/v3"
|
|
)
|
|
|
|
type Cluster struct {
|
|
log logr.Logger
|
|
cli *clientv3.Client
|
|
ctx context.Context
|
|
|
|
stepdown chan struct{}
|
|
stop chan struct{}
|
|
done chan struct{}
|
|
callbacks cluster.Callbacks
|
|
|
|
kvPrefix string
|
|
id string
|
|
keyID string
|
|
session bool
|
|
sessionCancel func()
|
|
fatal error
|
|
}
|
|
|
|
func NewCluster(log logr.Logger) cluster.Cluster {
|
|
c := &Cluster{log: log}
|
|
c.Reset()
|
|
return c
|
|
}
|
|
|
|
func (c *Cluster) Reset() {
|
|
c.ctx = context.Background()
|
|
c.stepdown = make(chan struct{}, 1)
|
|
c.stop = make(chan struct{}, 1)
|
|
c.done = make(chan struct{}, 1)
|
|
c.session = false
|
|
c.fatal = nil
|
|
}
|
|
|
|
func (c *Cluster) Get(k string) (v []byte, err error) {
|
|
r, err := c.cli.Get(c.ctx, c.kvPrefix+"/"+k)
|
|
if err != nil {
|
|
return
|
|
}
|
|
if len(r.Kvs) == 0 {
|
|
err = cluster.ErrKeyNotFound
|
|
return
|
|
}
|
|
v = r.Kvs[0].Value
|
|
return
|
|
}
|
|
|
|
func (c *Cluster) GetPrefix(pk string) (m map[string][]byte, err error) {
|
|
r, err := c.cli.Get(c.ctx, c.kvPrefix+"/"+pk, clientv3.WithPrefix())
|
|
if err != nil {
|
|
return
|
|
}
|
|
if len(r.Kvs) == 0 {
|
|
err = cluster.ErrPrefixNotFound
|
|
return
|
|
}
|
|
m = make(map[string][]byte)
|
|
for _, kv := range r.Kvs {
|
|
key := strings.TrimPrefix(string(kv.Key), c.kvPrefix+"/"+pk)
|
|
m[key] = kv.Value
|
|
}
|
|
return
|
|
}
|
|
|
|
func (c *Cluster) Set(k string, v []byte) (err error) {
|
|
_, err = c.cli.Put(c.ctx, c.kvPrefix+"/"+k, string(v))
|
|
return
|
|
}
|
|
|
|
func (c *Cluster) Delete(k string) (err error) {
|
|
_, err = c.cli.Delete(c.ctx, c.kvPrefix+"/"+k)
|
|
return
|
|
}
|
|
|
|
func (c *Cluster) SetCallbacks(callbacks cluster.Callbacks) error {
|
|
c.callbacks = callbacks
|
|
return c.checkConfig()
|
|
}
|
|
|
|
func (c *Cluster) Logger() logr.Logger {
|
|
return c.log
|
|
}
|
|
|
|
func (c *Cluster) ID() string {
|
|
return c.id
|
|
}
|
|
|
|
func (c *Cluster) checkConfig() error {
|
|
if c.callbacks.Leader == nil {
|
|
return errors.New("Leader callback is nil")
|
|
}
|
|
if c.callbacks.Follower == nil {
|
|
return errors.New("Follower callback is nil")
|
|
}
|
|
return nil
|
|
}
|