diff --git a/pkg/cli/server/server.go b/pkg/cli/server/server.go index cdada93cc2..ac4d2a7f4c 100644 --- a/pkg/cli/server/server.go +++ b/pkg/cli/server/server.go @@ -597,7 +597,7 @@ func run(app *cli.Context, cfg *cmds.Server, leaderControllers server.CustomCont go func() { if !serverConfig.ControlConfig.DisableETCD { - <-serverConfig.ControlConfig.Runtime.ETCDReady + <-executor.ETCDReadyChan() logrus.Info("ETCD server is now running") } if !serverConfig.ControlConfig.DisableAPIServer { diff --git a/pkg/cluster/cluster.go b/pkg/cluster/cluster.go index 69bdb97df7..a9e3447087 100644 --- a/pkg/cluster/cluster.go +++ b/pkg/cluster/cluster.go @@ -9,6 +9,7 @@ import ( "github.com/k3s-io/k3s/pkg/clientaccess" "github.com/k3s-io/k3s/pkg/cluster/managed" "github.com/k3s-io/k3s/pkg/daemons/config" + "github.com/k3s-io/k3s/pkg/daemons/executor" "github.com/k3s-io/k3s/pkg/etcd" "github.com/k3s-io/kine/pkg/endpoint" pkgerrors "github.com/pkg/errors" @@ -34,36 +35,36 @@ func (c *Cluster) ListenAndServe(ctx context.Context) error { return c.initClusterAndHTTPS(ctx) } -// Start handles writing/reading bootstrap data, and returns a channel -// that will be closed when datastore is ready. If embedded etcd is in use, +// Start handles writing/reading bootstrap data. If embedded etcd is in use, // a secondary call to Cluster.save is made. -func (c *Cluster) Start(ctx context.Context) (<-chan struct{}, error) { +func (c *Cluster) Start(ctx context.Context) error { + if c.config.DisableETCD || c.managedDB == nil { + // if etcd is disabled or we're using kine, perform a no-op start of etcd + // to close the etcd ready channel. When etcd is in use, this is handled by + // c.start() -> c.managedDB.Start() -> etcd.Start() -> executor.ETCD() + executor.ETCD(ctx, nil, nil, func(context.Context) error { return nil }) + } + if c.config.DisableETCD { - ready := make(chan struct{}) - defer close(ready) - return ready, nil + return nil } // start managed etcd database; when kine is in use this is a no-op. if err := c.start(ctx); err != nil { - return nil, pkgerrors.WithMessage(err, "start managed database") + return pkgerrors.WithMessage(err, "start managed database") } - // get the wait channel for testing etcd server readiness; when kine is in - // use the channel is closed immediately. - ready := c.testClusterDB(ctx) - // set c.config.Datastore and c.config.Runtime.EtcdConfig with values // necessary to build etcd clients, and start kine listener if necessary. if err := c.startStorage(ctx, false); err != nil { - return nil, err + return err } // if necessary, store bootstrap data to datastore. saveBootstrap is only set // when using kine, so this can be done before the ready channel has been closed. if c.saveBootstrap { if err := Save(ctx, c.config, false); err != nil { - return nil, err + return err } } @@ -71,7 +72,7 @@ func (c *Cluster) Start(ctx context.Context) (<-chan struct{}, error) { go func() { for { select { - case <-ready: + case <-executor.ETCDReadyChan(): // always save to managed etcd, to ensure that any file modified locally are in sync with the datastore. // this will panic if multiple keys exist, to prevent nodes from running with different bootstrap data. if err := Save(ctx, c.config, false); err != nil { @@ -104,7 +105,7 @@ func (c *Cluster) Start(ctx context.Context) (<-chan struct{}, error) { }() } - return ready, nil + return nil } // startEtcdProxy starts an etcd load-balancer proxy, for control-plane-only nodes diff --git a/pkg/cluster/managed.go b/pkg/cluster/managed.go index 79a3d13906..b8b0e10331 100644 --- a/pkg/cluster/managed.go +++ b/pkg/cluster/managed.go @@ -23,36 +23,6 @@ import ( "k8s.io/apimachinery/pkg/util/wait" ) -// testClusterDB returns a channel that will be closed when the datastore connection is available. -// The datastore is tested for readiness every 5 seconds until the test succeeds. -func (c *Cluster) testClusterDB(ctx context.Context) <-chan struct{} { - result := make(chan struct{}) - if c.managedDB == nil { - close(result) - return result - } - - go func() { - defer close(result) - for { - if err := c.managedDB.Test(ctx); err != nil { - logrus.Infof("Failed to test data store connection: %v", err) - } else { - logrus.Info(c.managedDB.EndpointName() + " data store connection OK") - return - } - - select { - case <-time.After(5 * time.Second): - case <-ctx.Done(): - return - } - } - }() - - return result -} - // start starts the database, unless a cluster reset has been requested, in which case // it does that instead. func (c *Cluster) start(ctx context.Context) error { diff --git a/pkg/cluster/managed/drivers.go b/pkg/cluster/managed/drivers.go index 421e7509f8..4a0515ba96 100644 --- a/pkg/cluster/managed/drivers.go +++ b/pkg/cluster/managed/drivers.go @@ -20,7 +20,6 @@ type Driver interface { IsReset() (bool, error) ResetFile() string Start(ctx context.Context, clientAccessInfo *clientaccess.Info) error - Test(ctx context.Context) error Restore(ctx context.Context) error EndpointName() string Snapshot(ctx context.Context) (*SnapshotResult, error) diff --git a/pkg/daemons/config/types.go b/pkg/daemons/config/types.go index 1085456081..d589220f6b 100644 --- a/pkg/daemons/config/types.go +++ b/pkg/daemons/config/types.go @@ -312,7 +312,6 @@ type ControlRuntimeBootstrap struct { type ControlRuntime struct { ControlRuntimeBootstrap - ETCDReady <-chan struct{} StartupHooksWg *sync.WaitGroup ClusterControllerStarts map[string]leader.Callback LeaderElectedClusterControllerStarts map[string]leader.Callback @@ -382,7 +381,7 @@ type ControlRuntime struct { type Cluster interface { Bootstrap(ctx context.Context, reset bool) error ListenAndServe(ctx context.Context) error - Start(ctx context.Context) (<-chan struct{}, error) + Start(ctx context.Context) error } type CoreFactory interface { diff --git a/pkg/daemons/control/server.go b/pkg/daemons/control/server.go index eddec8faba..97b71e8bde 100644 --- a/pkg/daemons/control/server.go +++ b/pkg/daemons/control/server.go @@ -70,10 +70,8 @@ func Prepare(ctx context.Context, cfg *config.Control) error { // Server starts the apiserver and whatever other control-plane components are // not disabled on this node. func Server(ctx context.Context, cfg *config.Control) error { - if ready, err := cfg.Cluster.Start(ctx); err != nil { + if err := cfg.Cluster.Start(ctx); err != nil { return pkgerrors.WithMessage(err, "failed to start cluster") - } else { - cfg.Runtime.ETCDReady = ready } if !cfg.DisableAPIServer { @@ -264,7 +262,7 @@ func apiServer(ctx context.Context, cfg *config.Control) error { logrus.Infof("Running kube-apiserver %s", config.ArgString(args)) - return executor.APIServer(ctx, runtime.ETCDReady, args) + return executor.APIServer(ctx, args) } func defaults(config *config.Control) { diff --git a/pkg/daemons/executor/embed.go b/pkg/daemons/executor/embed.go index 75ddd536f3..682d80f379 100644 --- a/pkg/daemons/executor/embed.go +++ b/pkg/daemons/executor/embed.go @@ -45,6 +45,7 @@ func init() { func (e *Embedded) Bootstrap(ctx context.Context, nodeConfig *daemonconfig.Node, cfg cmds.Agent) error { e.apiServerReady = util.APIServerReadyChan(ctx, nodeConfig.AgentConfig.KubeConfigK3sController, util.DefaultAPIServerReadyTimeout) + e.etcdReady = make(chan struct{}) e.criReady = make(chan struct{}) e.nodeConfig = nodeConfig @@ -119,12 +120,12 @@ func (*Embedded) APIServerHandlers(ctx context.Context) (authenticator.Request, return startupConfig.Authenticator, startupConfig.Handler, nil } -func (*Embedded) APIServer(ctx context.Context, etcdReady <-chan struct{}, args []string) error { +func (e *Embedded) APIServer(ctx context.Context, args []string) error { command := apiapp.NewAPIServerCommand(ctx.Done()) command.SetArgs(args) go func() { - <-etcdReady + <-e.ETCDReadyChan() defer func() { if err := recover(); err != nil { logrus.WithField("stack", string(debug.Stack())).Fatalf("apiserver panic: %v", err) @@ -262,6 +263,13 @@ func (e *Embedded) APIServerReadyChan() <-chan struct{} { return e.apiServerReady } +func (e *Embedded) ETCDReadyChan() <-chan struct{} { + if e.etcdReady == nil { + panic("executor not bootstrapped") + } + return e.etcdReady +} + func (e *Embedded) CRIReadyChan() <-chan struct{} { if e.criReady == nil { panic("executor not bootstrapped") diff --git a/pkg/daemons/executor/etcd.go b/pkg/daemons/executor/etcd.go index c0ba0c32bb..5ec3c7d4b9 100644 --- a/pkg/daemons/executor/etcd.go +++ b/pkg/daemons/executor/etcd.go @@ -5,6 +5,7 @@ import ( "errors" "os" "path/filepath" + "time" daemonconfig "github.com/k3s-io/k3s/pkg/daemons/config" "github.com/k3s-io/k3s/pkg/version" @@ -17,11 +18,40 @@ import ( // of the embedded execututor is disabled by build flags type Embedded struct { apiServerReady <-chan struct{} + etcdReady chan struct{} criReady chan struct{} nodeConfig *daemonconfig.Node } -func (e *Embedded) ETCD(ctx context.Context, args ETCDConfig, extraArgs []string) error { +func (e *Embedded) ETCD(ctx context.Context, args *ETCDConfig, extraArgs []string, test TestFunc) error { + // An unbootstrapped executor is used to start up a temporary embedded etcd when reconciling. + // This temporary executor doesn't have any ready channels set up, so don't bother testing. + if e.etcdReady != nil { + go func() { + defer close(e.etcdReady) + for { + if err := test(ctx); err != nil { + logrus.Infof("Failed to test etcd connection: %v", err) + } else { + logrus.Info("Connection to etcd is ready") + return + } + + select { + case <-time.After(5 * time.Second): + case <-ctx.Done(): + return + } + } + }() + } + + // nil args indicates a no-op start; all we need to do is wait for the test + // func to indicate readiness and close the channel. + if args == nil { + return nil + } + configFile, err := args.ToConfigFile(extraArgs) if err != nil { return err diff --git a/pkg/daemons/executor/executor.go b/pkg/daemons/executor/executor.go index 95735ff1db..2c6df8dd53 100644 --- a/pkg/daemons/executor/executor.go +++ b/pkg/daemons/executor/executor.go @@ -20,21 +20,25 @@ var ( executor Executor ) +// TestFunc is the signature of a function that returns nil error when the component is ready +type TestFunc func(context.Context) error + type Executor interface { Bootstrap(ctx context.Context, nodeConfig *daemonconfig.Node, cfg cmds.Agent) error Kubelet(ctx context.Context, args []string) error KubeProxy(ctx context.Context, args []string) error APIServerHandlers(ctx context.Context) (authenticator.Request, http.Handler, error) - APIServer(ctx context.Context, etcdReady <-chan struct{}, args []string) error + APIServer(ctx context.Context, args []string) error Scheduler(ctx context.Context, nodeReady <-chan struct{}, args []string) error ControllerManager(ctx context.Context, args []string) error CurrentETCDOptions() (InitialOptions, error) - ETCD(ctx context.Context, args ETCDConfig, extraArgs []string) error + ETCD(ctx context.Context, args *ETCDConfig, extraArgs []string, test TestFunc) error CloudControllerManager(ctx context.Context, ccmRBACReady <-chan struct{}, args []string) error Containerd(ctx context.Context, node *daemonconfig.Node) error Docker(ctx context.Context, node *daemonconfig.Node) error CRI(ctx context.Context, node *daemonconfig.Node) error APIServerReadyChan() <-chan struct{} + ETCDReadyChan() <-chan struct{} CRIReadyChan() <-chan struct{} } @@ -153,8 +157,8 @@ func APIServerHandlers(ctx context.Context) (authenticator.Request, http.Handler return executor.APIServerHandlers(ctx) } -func APIServer(ctx context.Context, etcdReady <-chan struct{}, args []string) error { - return executor.APIServer(ctx, etcdReady, args) +func APIServer(ctx context.Context, args []string) error { + return executor.APIServer(ctx, args) } func Scheduler(ctx context.Context, nodeReady <-chan struct{}, args []string) error { @@ -169,8 +173,8 @@ func CurrentETCDOptions() (InitialOptions, error) { return executor.CurrentETCDOptions() } -func ETCD(ctx context.Context, args ETCDConfig, extraArgs []string) error { - return executor.ETCD(ctx, args, extraArgs) +func ETCD(ctx context.Context, args *ETCDConfig, extraArgs []string, test TestFunc) error { + return executor.ETCD(ctx, args, extraArgs, test) } func CloudControllerManager(ctx context.Context, ccmRBACReady <-chan struct{}, args []string) error { @@ -193,6 +197,10 @@ func APIServerReadyChan() <-chan struct{} { return executor.APIServerReadyChan() } +func ETCDReadyChan() <-chan struct{} { + return executor.ETCDReadyChan() +} + func CRIReadyChan() <-chan struct{} { return executor.CRIReadyChan() } diff --git a/pkg/etcd/etcd.go b/pkg/etcd/etcd.go index 58ff7c485c..606ac126e8 100644 --- a/pkg/etcd/etcd.go +++ b/pkg/etcd/etcd.go @@ -1009,7 +1009,7 @@ func (e *ETCD) listenClientHTTPURLs() string { // cluster calls the executor to start etcd running with the provided configuration. func (e *ETCD) cluster(ctx context.Context, reset bool, options executor.InitialOptions) error { ctx, e.cancel = context.WithCancel(ctx) - return executor.ETCD(ctx, executor.ETCDConfig{ + return executor.ETCD(ctx, &executor.ETCDConfig{ Name: e.name, InitialOptions: options, ForceNewCluster: reset, @@ -1039,7 +1039,7 @@ func (e *ETCD) cluster(ctx context.Context, reset bool, options executor.Initial ExperimentalInitialCorruptCheck: true, ExperimentalWatchProgressNotifyInterval: e.config.Datastore.NotifyInterval, - }, e.config.ExtraEtcdArgs) + }, e.config.ExtraEtcdArgs, e.Test) } func (e *ETCD) StartEmbeddedTemporary(ctx context.Context) error { @@ -1089,7 +1089,7 @@ func (e *ETCD) StartEmbeddedTemporary(ctx context.Context) error { embedded := executor.Embedded{} ctx, e.cancel = context.WithCancel(ctx) - return embedded.ETCD(ctx, executor.ETCDConfig{ + return embedded.ETCD(ctx, &executor.ETCDConfig{ InitialOptions: executor.InitialOptions{AdvertisePeerURL: peerURL}, DataDir: tmpDataDir, ForceNewCluster: true, @@ -1106,7 +1106,7 @@ func (e *ETCD) StartEmbeddedTemporary(ctx context.Context) error { ExperimentalInitialCorruptCheck: true, ExperimentalWatchProgressNotifyInterval: e.config.Datastore.NotifyInterval, - }, append(e.config.ExtraEtcdArgs, "--max-snapshots=0", "--max-wals=0")) + }, append(e.config.ExtraEtcdArgs, "--max-snapshots=0", "--max-wals=0"), e.Test) } func addPort(address string, offset int) (string, error) { diff --git a/tests/mock/executor.go b/tests/mock/executor.go index 7f8c7587e2..92f612d5c4 100644 --- a/tests/mock/executor.go +++ b/tests/mock/executor.go @@ -30,7 +30,7 @@ func (e *Executor) APIServerHandlers(ctx context.Context) (authenticator.Request return nil, nil, errors.New("not implemented") } -func (e *Executor) APIServer(ctx context.Context, etcdReady <-chan struct{}, args []string) error { +func (e *Executor) APIServer(ctx context.Context, args []string) error { return errors.New("not implemented") } @@ -46,9 +46,9 @@ func (e *Executor) CurrentETCDOptions() (executor.InitialOptions, error) { return executor.InitialOptions{}, nil } -func (e *Executor) ETCD(ctx context.Context, args executor.ETCDConfig, extraArgs []string) error { +func (e *Executor) ETCD(ctx context.Context, args *executor.ETCDConfig, extraArgs []string, test executor.TestFunc) error { embed := &executor.Embedded{} - return embed.ETCD(ctx, args, extraArgs) + return embed.ETCD(ctx, args, extraArgs, test) } func (e *Executor) CloudControllerManager(ctx context.Context, ccmRBACReady <-chan struct{}, args []string) error { @@ -73,6 +73,12 @@ func (e *Executor) APIServerReadyChan() <-chan struct{} { return c } +func (e *Executor) ETCDReadyChan() <-chan struct{} { + c := make(chan struct{}) + close(c) + return c +} + func (e *Executor) CRIReadyChan() <-chan struct{} { c := make(chan struct{}) close(c) diff --git a/tests/unit.go b/tests/unit.go index 7f94c42be9..da00b2b06b 100644 --- a/tests/unit.go +++ b/tests/unit.go @@ -8,6 +8,7 @@ import ( "github.com/k3s-io/k3s/pkg/daemons/config" "github.com/k3s-io/k3s/pkg/daemons/control/deps" "github.com/k3s-io/k3s/pkg/daemons/executor" + "github.com/k3s-io/k3s/tests/mock" ) // GenerateDataDir creates a temporary directory at "/tmp/k3s//".