diff --git a/internal/machine/cluster.go b/internal/machine/cluster.go index 1b5ae2a0..6ec44dd4 100644 --- a/internal/machine/cluster.go +++ b/internal/machine/cluster.go @@ -151,7 +151,7 @@ func (cc *clusterController) Run(ctx context.Context) error { slog.Info("Corrosion service started.") } - // Apply the seed to finish Corrosion migrations from 0.x to 2026.5.14 (upstream v1.0.0) if applicable. + // Apply the seed to finish Corrosion migrations from 0.x to 2026.x.x (upstream v1.0.0) if applicable. if err := corromigrate.ApplySeedIfPresent(ctx, cc.corrosionDir, cc.store); err != nil { return fmt.Errorf("apply corrosion migration seed: %w", err) } @@ -173,7 +173,9 @@ func (cc *clusterController) Run(ctx context.Context) error { return nil }) - // Start the network API server. Assume the management IP can't be changed when the network is running. + // Start the network API server before waiting for the store sync so the machine is reachable on the mesh + // during the sync and can serve requests that don't depend on the store. + // Assume the management IP can't be changed when the network is running. apiAddr := net.JoinHostPort(cc.state.Network.ManagementIP.String(), strconv.Itoa(constants.MachineAPIPort)) listener, err := net.Listen("tcp", apiAddr) if err != nil { @@ -197,12 +199,7 @@ func (cc *clusterController) Run(ctx context.Context) error { // Check if waitStoreSync exited because the context was cancelled. Return early in that case. if ctx.Err() != nil { cc.stopAPIServer() - - err := errGroup.Wait() - if corroErr := cc.stopCorrosion(); corroErr != nil { - err = errors.Join(err, corroErr) - } - return err + return errGroup.Wait() } errGroup.Go(func() error { @@ -289,15 +286,9 @@ func (cc *clusterController) Run(ctx context.Context) error { slog.Info("Unregistry server stopped.") } - // Wait for all controllers to finish. - err = errGroup.Wait() - - // Stop Corrosion after all controllers depending on it and API server are stopped. - if corroErr := cc.stopCorrosion(); corroErr != nil { - err = errors.Join(err, corroErr) - } - - return err + // Wait for all controllers to finish. The Corrosion service shutdown is handled by the machine after stopping all + // local API servers that may still serve requests depending on the store. + return errGroup.Wait() } // stopAPIServer gracefully stops the network API server with a timeout. @@ -323,19 +314,6 @@ func (cc *clusterController) stopAPIServer() { slog.Info("Network API server stopped.") } -// stopCorrosion stops the Corrosion service with a timeout. -func (cc *clusterController) stopCorrosion() error { - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - - if err := cc.corroService.Stop(ctx); err != nil { - return fmt.Errorf("stop corrosion service: %w", err) - } - slog.Info("Corrosion service stopped.") - - return nil -} - // ensureDockerNetwork ensures that the Docker network is configured and ready for containers. func (cc *clusterController) ensureDockerNetwork(ctx context.Context) error { if err := cc.dockerCtrl.WaitDaemonReady(ctx); err != nil { diff --git a/internal/machine/machine.go b/internal/machine/machine.go index 2594ceed..a56590f8 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -378,6 +378,18 @@ func (m *Machine) Run(ctx context.Context) error { if err := docker.WaitDaemonReady(ctx, m.config.DockerClient); err != nil { return fmt.Errorf("wait for Docker daemon: %w", err) } + defer m.config.DockerClient.Close() + + // Bind the local API listeners before starting the dependencies (e.g. corrosion) to not deal with the teardown + // on failure. + machineListener, err := listenUnixSocket(m.config.MachineSockPath) + if err != nil { + return fmt.Errorf("listen machine API unix socket %q: %w", m.config.MachineSockPath, err) + } + proxyListener, err := listenUnixSocket(m.config.UncloudSockPath) + if err != nil { + return fmt.Errorf("listen API proxy unix socket %q: %w", m.config.UncloudSockPath, err) + } // Configure and start the corrosion service on the loopback if the machine is not initialised as a cluster // member. This provides the store required for the machine to initialise a new cluster on it. Once the machine @@ -393,7 +405,7 @@ func (m *Machine) Run(ctx context.Context) error { } slog.Info("Corrosion service started.") } else { - // Migrate the on-disk Corrosion store to 2026.5.14 (v1.0.0 upstream) if a v0.x store.db is detected, + // Migrate the on-disk Corrosion store to 2026.x.x (v1.0.0 upstream) if a v0.x store.db is detected, // before any Corrosion start attempt. The legacy systemd unit (if installed) is stopped here too // so we own the data dir exclusively. if err := corromigrate.MigrateIfNeeded(ctx, m.config.CorrosionDataDir, m.config.CorrosionUser); err != nil { @@ -405,10 +417,6 @@ func (m *Machine) Run(ctx context.Context) error { errGroup, ctx := errgroup.WithContext(ctx) // Start the local machine API server. - machineListener, err := listenUnixSocket(m.config.MachineSockPath) - if err != nil { - return fmt.Errorf("listen machine API unix socket %q: %w", m.config.MachineSockPath, err) - } errGroup.Go(func() error { slog.Info("Starting local machine API server.", "path", m.config.MachineSockPath) if err := m.localMachineServer.Serve(machineListener); err != nil { @@ -418,10 +426,6 @@ func (m *Machine) Run(ctx context.Context) error { }) // Start the local API proxy server. - proxyListener, err := listenUnixSocket(m.config.UncloudSockPath) - if err != nil { - return fmt.Errorf("listen API proxy unix socket %q: %w", m.config.UncloudSockPath, err) - } errGroup.Go(func() error { slog.Info("Starting local API proxy server.", "path", m.config.UncloudSockPath) if err := m.localProxyServer.Serve(proxyListener); err != nil { @@ -551,8 +555,6 @@ func (m *Machine) Run(ctx context.Context) error { // Shutdown goroutine. errGroup.Go(func() error { - var err error - <-ctx.Done() slog.Info("Stopping local machine API server.") // TODO: implement timeout for graceful shutdown. @@ -566,28 +568,31 @@ func (m *Machine) Run(ctx context.Context) error { m.proxyDirector.Close() slog.Info("Local API proxy server stopped.") - // Stop the corrosion container so this node stops gossiping its membership as "Up" while the - // gRPC API is gone. Use a fresh context because ctx is already cancelled here. - slog.Info("Stopping corrosion service.") - stopCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) - if stopErr := m.config.CorrosionService.Stop(stopCtx); stopErr != nil { - slog.Error("Failed to stop corrosion service.", "err", stopErr) - } - cancel() - - // Clean up the machine data and resources if the machine shutdown was initiated by a reset. - if m.resetting { - slog.Info("Cleaning up machine data and resources.") - if err = m.cleanup(); err != nil { - slog.Error("Failed to clean up machine data and resources.", "err", err) - } - } - - m.config.DockerClient.Close() - return err + return nil }) - return errGroup.Wait() + err = errGroup.Wait() + + // Stop the corrosion container only after the API servers, cluster controller, and all components depending on the + // store have stopped, so this machine keeps serving the store until then. Use a fresh context because ctx + // is already cancelled here. + slog.Info("Stopping corrosion service.") + stopCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + if stopErr := m.config.CorrosionService.Stop(stopCtx); stopErr != nil { + slog.Error("Failed to stop corrosion service.", "err", stopErr) + } + cancel() + slog.Info("Corrosion service stopped.") + + // Clean up the machine data and resources if the machine shutdown was initiated by a reset. + if m.resetting { + slog.Info("Cleaning up machine data and resources.") + if cleanupErr := m.cleanup(); cleanupErr != nil { + slog.Error("Failed to clean up machine data and resources.", "err", cleanupErr) + } + } + + return err } // listenUnixSocket creates a new Unix socket listener with the specified path. The socket file is created with 0660