mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 12:27:52 +00:00
saving state of branch
This commit is contained in:
parent
20429bb9dc
commit
60274f7112
5 changed files with 76 additions and 15 deletions
33
cmd/add_node.go
Normal file
33
cmd/add_node.go
Normal file
|
|
@ -0,0 +1,33 @@
|
|||
// Copyright 2022 Molecula Corp. (DBA FeatureBase).
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/ctl"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
func newAddNodeCommand(logdest logger.Logger) *cobra.Command {
|
||||
c := ctl.NewParquetInfoCommand(logdest)
|
||||
cmd := &cobra.Command{
|
||||
Use: "add-node PATH|URL",
|
||||
Short: "Add node to FeatureBase cluster.",
|
||||
Long: `
|
||||
Displays schema and sample data from the specified file
|
||||
`,
|
||||
Args: func(cmd *cobra.Command, args []string) error {
|
||||
if len(args) == 0 {
|
||||
return fmt.Errorf("data directory path required")
|
||||
} else if len(args) > 1 {
|
||||
return fmt.Errorf("too many command line arguments")
|
||||
}
|
||||
c.Path = args[0]
|
||||
return nil
|
||||
},
|
||||
RunE: usageErrorWrapper(c),
|
||||
}
|
||||
return cmd
|
||||
}
|
||||
|
|
@ -109,6 +109,7 @@ at https://docs.featurebase.com/.
|
|||
rc.AddCommand(newDataframeCsvLoaderCommand(logdest))
|
||||
rc.AddCommand(newPreSortCommand(logdest))
|
||||
rc.AddCommand(newParquetInfoCommand(logdest))
|
||||
rc.AddCommand(newAddNodeCommand(logdest))
|
||||
|
||||
rc.SetOutput(stderr)
|
||||
return rc
|
||||
|
|
|
|||
|
|
@ -83,6 +83,9 @@ func serverFlagSet(srv *server.Config, prefix string) *pflag.FlagSet {
|
|||
flags.StringVar(&srv.Etcd.EtcdHosts, "etcd.etcd-hosts", srv.Etcd.EtcdHosts, "EXPERIMENTAL etcd server host:port comma separated list")
|
||||
flags.MarkHidden("etcd.etcd-hosts") // TODO (twg) expose when ready for public consumption
|
||||
|
||||
flags.StringVar(&srv.Etcd.InitClusterToken, pre("etcd.initial-cluster-token"), srv.Etcd.InitClusterToken, "EXPERIMENTAL unique token for etcd cluster.")
|
||||
flags.StringVar(&srv.Etcd.InitClusterState, pre("etcd.initial-cluster-state"), srv.Etcd.InitClusterState, "EXPERIMENTAL initial cluster state when this node is created.")
|
||||
|
||||
// External postgres database for ExternalLookup
|
||||
flags.StringVar(&srv.LookupDBDSN, pre("lookup-db-dsn"), "", "external (postgres) database DSN to use for ExternalLookup calls")
|
||||
|
||||
|
|
|
|||
|
|
@ -29,16 +29,19 @@ import (
|
|||
)
|
||||
|
||||
type Options struct {
|
||||
Name string `toml:"name"`
|
||||
Dir string `toml:"dir"`
|
||||
LClientURL string `toml:"listen-client-url"`
|
||||
AClientURL string `toml:"advertise-client-url"`
|
||||
LPeerURL string `toml:"listen-peer-url"`
|
||||
APeerURL string `toml:"advertise-peer-url"`
|
||||
ClusterURL string `toml:"cluster-url"`
|
||||
InitCluster string `toml:"initial-cluster"`
|
||||
ClusterName string `toml:"cluster-name"`
|
||||
HeartbeatTTL int64 `toml:"heartbeat-ttl"`
|
||||
Name string `toml:"name"`
|
||||
Dir string `toml:"dir"`
|
||||
LClientURL string `toml:"listen-client-url"`
|
||||
AClientURL string `toml:"advertise-client-url"`
|
||||
LPeerURL string `toml:"listen-peer-url"`
|
||||
APeerURL string `toml:"advertise-peer-url"`
|
||||
ClusterURL string `toml:"cluster-url"`
|
||||
InitCluster string `toml:"initial-cluster"`
|
||||
InitClusterToken string `toml:"initial-cluster-token"`
|
||||
InitClusterState string `toml:"initial-cluster-state"`
|
||||
ClusterName string `toml:"cluster-name"`
|
||||
HeartbeatTTL int64 `toml:"heartbeat-ttl"`
|
||||
|
||||
// TLS provided tls files
|
||||
TrustedCAFile string `toml:"tls-trusted-cafile"`
|
||||
ClientCertFile string `toml:"tls-cert-file"`
|
||||
|
|
@ -388,14 +391,21 @@ func (e *Etcd) parseOptions() (*embed.Config, error) {
|
|||
// %% end sonarcloud ignore %%
|
||||
}
|
||||
cfg.InitialCluster = e.options.InitCluster
|
||||
cfg.ClusterState = embed.ClusterStateFlagNew
|
||||
if e.options.InitClusterState != embed.ClusterStateFlagNew && e.options.InitClusterState != embed.ClusterStateFlagExisting {
|
||||
return nil, errors.New(fmt.Sprintf("embed: initial cluster state must be %s or %s", embed.ClusterStateFlagNew, embed.ClusterStateFlagExisting))
|
||||
} else {
|
||||
cfg.ClusterState = e.options.InitClusterState
|
||||
}
|
||||
} else {
|
||||
cfg.InitialCluster = cfg.Name + "=" + e.options.APeerURL
|
||||
}
|
||||
|
||||
if e.options.ClusterURL != "" {
|
||||
return nil, errors.New("joining an existing cluster is unsupported")
|
||||
}
|
||||
/*
|
||||
if e.options.ClusterURL != "" {
|
||||
return nil, errors.New("joining an existing cluster is unsupported")
|
||||
}
|
||||
*/
|
||||
|
||||
// can only use tls if not using pre-configured listeners
|
||||
cfg.ClientTLSInfo = transport.TLSInfo{
|
||||
TrustedCAFile: e.options.TrustedCAFile,
|
||||
|
|
@ -415,11 +425,20 @@ func (e *Etcd) parseOptions() (*embed.Config, error) {
|
|||
|
||||
// Start starts etcd and hearbeat
|
||||
func (e *Etcd) Start(ctx context.Context) (_ disco.InitialClusterState, err error) {
|
||||
// starting etcd service embeded in featurebase
|
||||
|
||||
opts, err := e.parseOptions()
|
||||
if err != nil {
|
||||
return disco.InitialClusterStateNew, err
|
||||
}
|
||||
|
||||
state := disco.InitialClusterState(opts.ClusterState)
|
||||
if state == embed.ClusterStateFlagNew {
|
||||
e.logger.Infof("Starting new cluster")
|
||||
} else {
|
||||
e.logger.Infof("Starting a node which will be added to an existing cluster")
|
||||
}
|
||||
|
||||
// create a context that can be used for our watch processes, etcetera.
|
||||
e.childContext, e.childCancel = context.WithCancel(context.Background())
|
||||
if e.options.EtcdHosts != "" {
|
||||
|
|
|
|||
|
|
@ -625,8 +625,13 @@ func (s *Server) Open() error {
|
|||
}
|
||||
// I'm pretty sure this can't happen, because the path that would have led to it
|
||||
// happening now generates an error already, but let's be careful.
|
||||
/*
|
||||
if initState == disco.InitialClusterStateExisting {
|
||||
return errors.New("disco reports existing cluster, but this is not supported")
|
||||
}
|
||||
*/
|
||||
if initState == disco.InitialClusterStateExisting {
|
||||
return errors.New("disco reports existing cluster, but this is not supported")
|
||||
s.logger.Infof("trying to join an existing cluster")
|
||||
}
|
||||
|
||||
// Set node ID.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue