diff --git a/ctl/server.go b/ctl/server.go index cb317c6af..815ae827b 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -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") diff --git a/etcd/embed.go b/etcd/embed.go index 6b5f5fd3d..77ebdb5cb 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -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) diff --git a/etcd/external.go b/etcd/external.go new file mode 100644 index 000000000..b6d342a5a --- /dev/null +++ b/etcd/external.go @@ -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, + }, + ) +} diff --git a/etcd/external_test.go b/etcd/external_test.go new file mode 100644 index 000000000..a41ff19a0 --- /dev/null +++ b/etcd/external_test.go @@ -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) + } + }) + } +} diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 44cea686f..8ae542dce 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -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) } diff --git a/server/config.go b/server/config.go index caeb5314b..6446b2032 100644 --- a/server/config.go +++ b/server/config.go @@ -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 { diff --git a/server/server.go b/server/server.go index 52ae12d3b..1e862f486 100644 --- a/server/server.go +++ b/server/server.go @@ -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") }