[FB-1598] Phase 0 of etcd extraction (#2169)

* Phase 0 of etcd extraction

* added test shell

* bare basic test

* address comments
This commit is contained in:
tgruben 2022-07-29 15:27:04 -05:00 • committed by GitHub
parent 6f1514933b
commit 0300cb071f
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 277 additions and 76 deletions

View file

@ -55,6 +55,11 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
flags.StringVar(&srv.Config.Etcd.InitCluster, "etcd.initial-cluster", srv.Config.Etcd.InitCluster, "Initial cluster name1=apurl1,name2=apurl2")
flags.Int64Var(&srv.Config.Etcd.HeartbeatTTL, "etcd.heartbeat-ttl", srv.Config.Etcd.HeartbeatTTL, "Timeout used to determine cluster status")
flags.StringVar(&srv.Config.Etcd.Cluster, "etcd.static-cluster", srv.Config.Etcd.Cluster, "EXPERIMENTAL static featurebase cluster name1=apurl1,name2=apurl2")
flags.MarkHidden("etcd.static-cluster")
flags.StringVar(&srv.Config.Etcd.EtcdHosts, "etcd.etcd-hosts", srv.Config.Etcd.EtcdHosts, "EXPERIMENTAL etcd server host:port comma separated list")
flags.MarkHidden("etcd.etcd-hosts") // TODO (twg) expose when ready for public consumption
// External postgres database for ExternalLookup
flags.StringVar(&srv.Config.LookupDBDSN, "lookup-db-dsn", "", "external (postgres) database DSN to use for ExternalLookup calls")

View file

@ -48,7 +48,10 @@ type Options struct {
LPeerSocket []*net.TCPListener
LClientSocket []*net.TCPListener
UnsafeNoFsync bool `toml:"no-fsync"`
UnsafeNoFsync bool `toml:"no-fsync"`
Cluster string `toml:"static-cluster"`
EtcdHosts string `toml:"etcd-hosts"`
Id string
}
var (
@ -68,9 +71,7 @@ const (
shardPrefix = "/shard/"
)
var (
errEtcdShuttingDown = errors.New("etcd shutting down")
)
var errEtcdShuttingDown = errors.New("etcd shutting down")
// nodeData is an internal tracker of the data we're keeping about
// nodes in etcd, which we update from data collected either directly
@ -90,13 +91,108 @@ type nodeData struct {
node *disco.Node
}
// EtcdServicer provides access to either an embeded etcd server or an external host
type EtcdServicer interface {
Shutdown()
Peers() []*disco.Peer
IsLeader() bool
Leader() *disco.Peer
ID() string
Startup(ctx context.Context, state disco.InitialClusterState) (disco.InitialClusterState, error)
NewClient() (*clientv3.Client, error)
}
type EmbeddedEtcd struct {
e *embed.Etcd
parent *Etcd
}
func NewEmbeddedEtcd(p *Etcd, opts *embed.Config) (*EmbeddedEtcd, error) {
e, err := embed.StartEtcd(opts)
if err != nil {
return nil, errors.Wrap(err, "starting etcd")
}
this := &EmbeddedEtcd{}
this.parent = p
this.e = e
p.e = e
return this, nil
}
func (e *EmbeddedEtcd) ID() string {
if e.e == nil || e.e.Server == nil {
return ""
}
return e.e.Server.ID().String()
}
func (e *EmbeddedEtcd) Shutdown() {
if e.e != nil {
e.e.Close()
<-e.e.Server.StopNotify()
}
}
func (e *EmbeddedEtcd) Peers() []*disco.Peer {
var peers []*disco.Peer
for _, member := range e.e.Server.Cluster().Members() {
peers = append(peers, &disco.Peer{ID: member.ID.String(), URL: member.PickPeerURL()})
}
return peers
}
func (e *EmbeddedEtcd) IsLeader() bool {
if e.e == nil || e.e.Server == nil {
return false
}
return e.e.Server.Leader() == e.e.Server.ID()
}
func (e *EmbeddedEtcd) Leader() *disco.Peer {
id := e.e.Server.Leader()
peer := &disco.Peer{ID: id.String()}
if m := e.e.Server.Cluster().Member(id); m != nil {
peer.URL = m.PickPeerURL()
}
return peer
}
func (e *EmbeddedEtcd) Startup(ctx context.Context, state disco.InitialClusterState) (disco.InitialClusterState, error) {
select {
case <-ctx.Done():
return state, ctx.Err()
case err := <-e.e.Err():
return state, err
case <-e.e.Server.ReadyNotify():
members := e.e.Server.Cluster().Members()
e.parent.nodeMu.Lock()
defer e.parent.nodeMu.Unlock()
// mark everything unknown so we show a state for nodes we haven't
// heard from yet.
for _, member := range members {
peerID := member.ID.String()
_ = e.parent.seeNode(peerID)
}
e.parent.nodesDirty = true
return state, e.parent.startHeartbeatAndWatcher(ctx)
}
}
func (e *EmbeddedEtcd) NewClient() (*clientv3.Client, error) {
return v3client.New(e.e.Server), nil
}
type Etcd struct {
options Options
replicas int
e *embed.Etcd
cli *clientv3.Client
cliMu sync.Mutex
e *embed.Etcd // TODO (twg) factor out
service EtcdServicer
cli *clientv3.Client
cliMu sync.Mutex
heartbeatLeasedKV *leasedKV
@ -149,10 +245,7 @@ func (e *Etcd) Close() error {
e.childCancel()
}
// shut down the server, if we have one.
if e.e != nil {
e.e.Close()
<-e.e.Server.StopNotify()
}
e.service.Shutdown()
// shut down the client, if we have one. if something's still
// using it, we anticipate the client failing its current call,
// and any retry will be checking for the child context being
@ -178,7 +271,7 @@ func (e *Etcd) newClient(cli *clientv3.Client) *clientv3.Client {
return cli
}
_ = cli.Close()
e.cli = v3client.New(e.e.Server)
e.cli, _ = e.service.NewClient()
return e.cli
}
@ -282,7 +375,7 @@ func (e *Etcd) parseOptions() (*embed.Config, error) {
monitor.InitErrorMonitor(e.version)
e.logger.Infof("Initializing Monitor: Capturing usage metrics")
//check for multiple nodes in the cluster and error if present
// check for multiple nodes in the cluster and error if present
nodes := strings.Split(e.options.InitCluster, ",")
if len(nodes) > 1 {
return nil, fmt.Errorf("multiple cluster nodes detected - this version of FeatureBase only supports single node. %+v", e.options.InitCluster)
@ -324,47 +417,37 @@ func (e *Etcd) Start(ctx context.Context) (_ disco.InitialClusterState, err erro
state := disco.InitialClusterState(opts.ClusterState)
// create a context that can be used for our watch processes, etcetera.
e.childContext, e.childCancel = context.WithCancel(context.Background())
e.e, err = embed.StartEtcd(opts)
if err != nil {
return state, errors.Wrap(err, "starting etcd")
}
// If we are returning an error, the caller won't be shutting us down
// later, so we have to stop the server ourselves.
defer func() {
if e.options.EtcdHosts != "" {
e.service, err = NewExternalEtcd(e, e.options)
if err != nil {
// shut down everything on our way out.
e.Close()
return state, errors.Wrap(err, "connecting etcd")
}
}()
e.cli = v3client.New(e.e.Server)
e.logger.Infof("using external etcd %v with fixed fb cluster %v", e.options.EtcdHosts, e.options.Cluster)
} else {
select {
case <-ctx.Done():
return state, ctx.Err()
case err := <-e.e.Err():
return state, err
case <-e.e.Server.ReadyNotify():
members := e.e.Server.Cluster().Members()
e.nodeMu.Lock()
defer e.nodeMu.Unlock()
// mark everything unknown so we show a state for nodes we haven't
// heard from yet.
for _, member := range members {
peerID := member.ID.String()
_ = e.seeNode(peerID)
e.service, err = NewEmbeddedEtcd(e, opts)
if err != nil {
return state, errors.Wrap(err, "starting etcd")
}
e.nodesDirty = true
return state, e.startHeartbeatAndWatcher(ctx)
// If we are returning an error, the caller won't be shutting us down
// later, so we have to stop the server ourselves.
defer func() {
if err != nil {
// shut down everything on our way out.
e.Close()
}
}()
}
e.cli, _ = e.service.NewClient()
return e.service.Startup(ctx, state)
}
// startHeartbeatAndWatcher spins up the heartbeat, and also a background
// watcher that watches for changes to events we care about.
func (e *Etcd) startHeartbeatAndWatcher(ctx context.Context) error {
key := heartbeatPrefix + e.e.Server.ID().String()
key := heartbeatPrefix + e.service.ID()
e.heartbeatLeasedKV = newLeasedKV(e, e.childContext, key, e.options.HeartbeatTTL)
if err := e.heartbeatLeasedKV.Start(string(disco.NodeStateStarting)); err != nil {
@ -383,36 +466,19 @@ func (e *Etcd) SetState(ctx context.Context, state disco.NodeState) (err error)
}
func (e *Etcd) ID() string {
if e.e == nil || e.e.Server == nil {
return ""
}
return e.e.Server.ID().String()
return e.service.ID()
}
func (e *Etcd) Peers() []*disco.Peer {
var peers []*disco.Peer
for _, member := range e.e.Server.Cluster().Members() {
peers = append(peers, &disco.Peer{ID: member.ID.String(), URL: member.PickPeerURL()})
}
return peers
return e.service.Peers()
}
func (e *Etcd) IsLeader() bool {
if e.e == nil || e.e.Server == nil {
return false
}
return e.e.Server.Leader() == e.e.Server.ID()
return e.service.IsLeader()
}
func (e *Etcd) Leader() *disco.Peer {
id := e.e.Server.Leader()
peer := &disco.Peer{ID: id.String()}
if m := e.e.Server.Cluster().Member(id); m != nil {
peer.URL = m.PickPeerURL()
}
return peer
return e.service.Leader()
}
func (e *Etcd) ClusterState(ctx context.Context) (out disco.ClusterState, err error) {
@ -737,7 +803,7 @@ func (e *Etcd) SetMetadata(ctx context.Context, node *disco.Node) error {
return errors.Wrap(err, "marshaling json metadata")
}
err = e.putKey(ctx, path.Join(metadataPrefix,
e.e.Server.ID().String()),
e.service.ID()),
string(data),
)
if err != nil {
@ -1011,7 +1077,7 @@ func (e *Etcd) Shards(ctx context.Context, index, field string) ([][]byte, error
// SetShards implements the Sharder interface.
func (e *Etcd) SetShards(ctx context.Context, index, field string, shards []byte) error {
key := path.Join(shardPrefix, index, field, e.e.Server.ID().String())
key := path.Join(shardPrefix, index, field, e.service.ID())
op := clientv3.OpPut(key, "")
op.WithValueBytes(shards)

94
etcd/external.go Normal file
View file

@ -0,0 +1,94 @@
package etcd
import (
"context"
"errors"
"strings"
"time"
"github.com/molecula/featurebase/v3/disco"
clientv3 "go.etcd.io/etcd/client/v3"
)
// TODO (twg) test coverage once we figure out how to setup external etcd in test
type ExternalEtcd struct {
parent *Etcd
id string
peers []*disco.Peer
etcdHosts []string
}
func NewExternalEtcd(p *Etcd, o Options) (*ExternalEtcd, error) {
if o.Cluster == "" {
return nil, errors.New("cluster required")
}
if o.Id == "" {
return nil, errors.New("nodeid required")
}
if o.EtcdHosts == "" {
return nil, errors.New("remote etcd host required")
}
// %% begin sonarcloud ignore %%
nodes := strings.Split(o.Cluster, ",") // id=url,id2=url2,....,idn=urln
peers := make([]*disco.Peer, len(nodes))
for i, node := range nodes {
parts := strings.Split(node, "=") // id=url pairs
peers[i] = &disco.Peer{ID: parts[0], URL: parts[1]}
}
this := &ExternalEtcd{}
this.parent = p
this.peers = peers
this.id = o.Id
// populate parent.KnownNodes
for _, n := range peers {
this.parent.seeNode(n.ID)
}
this.parent.nodesDirty = true
// populate parent.sortedNodes
hosts := strings.Split(o.EtcdHosts, ",")
this.etcdHosts = make([]string, len(hosts))
for i := range hosts {
this.etcdHosts[i] = hosts[i]
}
return this, nil
// %% end sonarcloud ignore %%
}
func (e *ExternalEtcd) ID() string {
return e.id
}
func (e *ExternalEtcd) Shutdown() {
}
func (e *ExternalEtcd) Peers() []*disco.Peer {
return e.peers
}
// TODO this is a nonsense implementation right now
// will fill out when understand the issue
func (e *ExternalEtcd) IsLeader() bool {
return e.Leader().String() == e.ID()
}
func (e *ExternalEtcd) Leader() *disco.Peer {
// TODO (twg) NOP?
return e.peers[0]
}
func (e *ExternalEtcd) Startup(ctx context.Context, state disco.InitialClusterState) (disco.InitialClusterState, error) {
// TODO (twg) trye disabled
return state, e.parent.startHeartbeatAndWatcher(ctx)
}
func (e *ExternalEtcd) NewClient() (*clientv3.Client, error) {
return clientv3.New(
clientv3.Config{
DialTimeout: 10 * time.Second,
Endpoints: e.etcdHosts,
},
)
}

41
etcd/external_test.go Normal file
View file

@ -0,0 +1,41 @@
package etcd
import (
"testing"
"github.com/molecula/featurebase/v3/logger"
)
func TestNewExternalEtcd(t *testing.T) {
logger := logger.NewLogfLogger(t)
type args struct {
p *Etcd
o Options
}
tests := []struct {
name string
args args
want string
wantErr bool
}{
{name: "constructor nil", args: args{nil, Options{}}, want: "", wantErr: true},
{name: "constructor", args: args{&Etcd{knownNodes: make(map[string]*nodeData), logger: logger}, Options{
Id: "f0",
Cluster: "f0=http://localhost:10101",
EtcdHosts: "0.0.0.0:2329",
}}, want: "f0", wantErr: false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := NewExternalEtcd(tt.args.p, tt.args.o)
if (err != nil) != tt.wantErr {
t.Errorf("NewExternalEtcd() error = %v, wantErr %v", err, tt.wantErr)
return
}
if tt.wantErr == false && got.ID() != tt.want {
t.Errorf("NewExternalEtcd() = %v, want %v", got.id, tt.want)
}
})
}
}

View file

@ -168,7 +168,6 @@ func (l *leasedKV) Stop() {
// So this is a low-interest message usually.
l.e.logger.Debugf("revoking lease during shutdown: %v", err)
}
}
// Set will change the specific value for this key.
@ -183,7 +182,6 @@ func (l *leasedKV) Set(ctx context.Context, value string) error {
return err
})
// l.e.logger.Printf("set key %q on %q value %q: err %v", l.key, l.e.options.Name, value, err)
if err != nil {
return errors.Wrapf(err, "creating key %s with value [%s]", l.key, l.value)
}

View file

@ -364,7 +364,7 @@ func NewConfig() *Config {
// Cluster config.
c.Cluster.Name = "cluster0"
c.Cluster.ReplicaN = 1
c.Cluster.LongQueryTime = toml.Duration(-time.Minute) //TODO remove this once cluster.longQueryTime is fully deprecated
c.Cluster.LongQueryTime = toml.Duration(-time.Minute) // TODO remove this once cluster.longQueryTime is fully deprecated
c.Cluster.PartitionToNodeAssignment = PartitionToNodeJmp
// AntiEntropy config.
@ -701,7 +701,6 @@ func (c *Config) ValidateAuth() (errors []error) {
}
func (c *Config) ValidatePermissions(permsFile io.Reader) (errors []error) {
var p authz.GroupPermissions
if err := p.ReadPermissionsFile(permsFile); err != nil {
return append(errors, err)
@ -737,14 +736,12 @@ func (c *Config) ValidatePermissions(permsFile io.Reader) (errors []error) {
if p.Admin == "" {
errors = append(errors, fmt.Errorf("empty string for admin in permissions file: %s", c.Auth.PermissionsFile))
}
return errors
}
func (c *Config) ValidatePermissionsFile() (err error) {
if c.Auth.PermissionsFile == "" {
return fmt.Errorf("empty string for auth config permissions file")
}
@ -757,7 +754,6 @@ func (c *Config) ValidatePermissionsFile() (err error) {
}
func (c *Config) MustValidateAuth() {
errorsAuth := c.ValidateAuth()
if len(errorsAuth) > 0 {
for _, e := range errorsAuth {

View file

@ -458,6 +458,7 @@ func (m *Command) SetupServer() error {
m.Config.Etcd.Dir = filepath.Join(path, pilosa.DiscoDir)
}
m.Config.Etcd.Id = m.Config.Name // TODO(twg) rethink this
e := petcd.NewEtcd(m.Config.Etcd, m.logger, m.Config.Cluster.ReplicaN, version)
serverOptions := []pilosa.ServerOption{
@ -635,12 +636,12 @@ func (m *Command) setupQueryLogger() error {
var err error
if m.Config.Auth.QueryLogPath == "" {
f, err = logger.NewFileWriterMode("queries/query.log", 0600)
f, err = logger.NewFileWriterMode("queries/query.log", 0o600)
if err != nil {
return errors.Wrap(err, "opening file")
}
} else {
f, err = logger.NewFileWriterMode(m.Config.Auth.QueryLogPath, 0600)
f, err = logger.NewFileWriterMode(m.Config.Auth.QueryLogPath, 0o600)
if err != nil {
return errors.Wrap(err, "opening file")
}