mirror of
https://github.com/psviderski/uncloud.git
synced 2026-08-26 11:03:34 +00:00
refactor(machine): improve startup and shuitdown logic, shutdown corrosion in one place
This commit is contained in:
@@ -151,7 +151,7 @@ func (cc *clusterController) Run(ctx context.Context) error {
|
|||||||
slog.Info("Corrosion service started.")
|
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 {
|
if err := corromigrate.ApplySeedIfPresent(ctx, cc.corrosionDir, cc.store); err != nil {
|
||||||
return fmt.Errorf("apply corrosion migration seed: %w", err)
|
return fmt.Errorf("apply corrosion migration seed: %w", err)
|
||||||
}
|
}
|
||||||
@@ -173,7 +173,9 @@ func (cc *clusterController) Run(ctx context.Context) error {
|
|||||||
return nil
|
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))
|
apiAddr := net.JoinHostPort(cc.state.Network.ManagementIP.String(), strconv.Itoa(constants.MachineAPIPort))
|
||||||
listener, err := net.Listen("tcp", apiAddr)
|
listener, err := net.Listen("tcp", apiAddr)
|
||||||
if err != nil {
|
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.
|
// Check if waitStoreSync exited because the context was cancelled. Return early in that case.
|
||||||
if ctx.Err() != nil {
|
if ctx.Err() != nil {
|
||||||
cc.stopAPIServer()
|
cc.stopAPIServer()
|
||||||
|
return errGroup.Wait()
|
||||||
err := errGroup.Wait()
|
|
||||||
if corroErr := cc.stopCorrosion(); corroErr != nil {
|
|
||||||
err = errors.Join(err, corroErr)
|
|
||||||
}
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
errGroup.Go(func() error {
|
errGroup.Go(func() error {
|
||||||
@@ -289,15 +286,9 @@ func (cc *clusterController) Run(ctx context.Context) error {
|
|||||||
slog.Info("Unregistry server stopped.")
|
slog.Info("Unregistry server stopped.")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Wait for all controllers to finish.
|
// Wait for all controllers to finish. The Corrosion service shutdown is handled by the machine after stopping all
|
||||||
err = errGroup.Wait()
|
// local API servers that may still serve requests depending on the store.
|
||||||
|
return 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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// stopAPIServer gracefully stops the network API server with a timeout.
|
// stopAPIServer gracefully stops the network API server with a timeout.
|
||||||
@@ -323,19 +314,6 @@ func (cc *clusterController) stopAPIServer() {
|
|||||||
slog.Info("Network API server stopped.")
|
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.
|
// ensureDockerNetwork ensures that the Docker network is configured and ready for containers.
|
||||||
func (cc *clusterController) ensureDockerNetwork(ctx context.Context) error {
|
func (cc *clusterController) ensureDockerNetwork(ctx context.Context) error {
|
||||||
if err := cc.dockerCtrl.WaitDaemonReady(ctx); err != nil {
|
if err := cc.dockerCtrl.WaitDaemonReady(ctx); err != nil {
|
||||||
|
|||||||
+36
-31
@@ -378,6 +378,18 @@ func (m *Machine) Run(ctx context.Context) error {
|
|||||||
if err := docker.WaitDaemonReady(ctx, m.config.DockerClient); err != nil {
|
if err := docker.WaitDaemonReady(ctx, m.config.DockerClient); err != nil {
|
||||||
return fmt.Errorf("wait for Docker daemon: %w", err)
|
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
|
// 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
|
// 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.")
|
slog.Info("Corrosion service started.")
|
||||||
} else {
|
} 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
|
// before any Corrosion start attempt. The legacy systemd unit (if installed) is stopped here too
|
||||||
// so we own the data dir exclusively.
|
// so we own the data dir exclusively.
|
||||||
if err := corromigrate.MigrateIfNeeded(ctx, m.config.CorrosionDataDir, m.config.CorrosionUser); err != nil {
|
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)
|
errGroup, ctx := errgroup.WithContext(ctx)
|
||||||
|
|
||||||
// Start the local machine API server.
|
// 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 {
|
errGroup.Go(func() error {
|
||||||
slog.Info("Starting local machine API server.", "path", m.config.MachineSockPath)
|
slog.Info("Starting local machine API server.", "path", m.config.MachineSockPath)
|
||||||
if err := m.localMachineServer.Serve(machineListener); err != nil {
|
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.
|
// 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 {
|
errGroup.Go(func() error {
|
||||||
slog.Info("Starting local API proxy server.", "path", m.config.UncloudSockPath)
|
slog.Info("Starting local API proxy server.", "path", m.config.UncloudSockPath)
|
||||||
if err := m.localProxyServer.Serve(proxyListener); err != nil {
|
if err := m.localProxyServer.Serve(proxyListener); err != nil {
|
||||||
@@ -551,8 +555,6 @@ func (m *Machine) Run(ctx context.Context) error {
|
|||||||
|
|
||||||
// Shutdown goroutine.
|
// Shutdown goroutine.
|
||||||
errGroup.Go(func() error {
|
errGroup.Go(func() error {
|
||||||
var err error
|
|
||||||
|
|
||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
slog.Info("Stopping local machine API server.")
|
slog.Info("Stopping local machine API server.")
|
||||||
// TODO: implement timeout for graceful shutdown.
|
// TODO: implement timeout for graceful shutdown.
|
||||||
@@ -566,28 +568,31 @@ func (m *Machine) Run(ctx context.Context) error {
|
|||||||
m.proxyDirector.Close()
|
m.proxyDirector.Close()
|
||||||
slog.Info("Local API proxy server stopped.")
|
slog.Info("Local API proxy server stopped.")
|
||||||
|
|
||||||
// Stop the corrosion container so this node stops gossiping its membership as "Up" while the
|
return nil
|
||||||
// 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 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
|
// listenUnixSocket creates a new Unix socket listener with the specified path. The socket file is created with 0660
|
||||||
|
|||||||
Reference in New Issue
Block a user