From dbf4d3a8caf71b48b9ec5082553c4a2f0153fec4 Mon Sep 17 00:00:00 2001 From: Brian Goff Date: Wed, 31 Jan 2018 17:32:40 -0500 Subject: [PATCH] Refresh containerd remotes on containerd restarted Before this patch, when containerd is restarted (due to a crash, or kill, whatever), the daemon would keep trying to process the event stream against the old socket handles. This would lead to a CPU spin due to the error handling when the client can't connect to containerd. This change makes sure the containerd remote client is updated for all registered libcontainerd clients. This is not neccessarily the ideal fix which would likely require a major refactor, but at least gets things to a working state with a minimal patch. Signed-off-by: Brian Goff (cherry picked from commit 400126f8698233099259da967378c0a76bc3ea31) Signed-off-by: Sebastiaan van Stijn --- .../engine/libcontainerd/client_daemon.go | 31 +++++-- .../engine/libcontainerd/remote_daemon.go | 86 +++++++++++++------ 2 files changed, 81 insertions(+), 36 deletions(-) diff --git a/components/engine/libcontainerd/client_daemon.go b/components/engine/libcontainerd/client_daemon.go index a9f7c11dd1..1b5ea2ab90 100644 --- a/components/engine/libcontainerd/client_daemon.go +++ b/components/engine/libcontainerd/client_daemon.go @@ -113,8 +113,21 @@ type client struct { containers map[string]*container } +func (c *client) setRemote(remote *containerd.Client) { + c.Lock() + c.remote = remote + c.Unlock() +} + +func (c *client) getRemote() *containerd.Client { + c.RLock() + remote := c.remote + c.RUnlock() + return remote +} + func (c *client) Version(ctx context.Context) (containerd.Version, error) { - return c.remote.Version(ctx) + return c.getRemote().Version(ctx) } func (c *client) Restore(ctx context.Context, id string, attachStdio StdioCallback) (alive bool, pid int, err error) { @@ -188,7 +201,7 @@ func (c *client) Create(ctx context.Context, id string, ociSpec *specs.Spec, run c.logger.WithField("bundle", bdir).WithField("root", ociSpec.Root.Path).Debug("bundle dir created") - cdCtr, err := c.remote.NewContainer(ctx, id, + cdCtr, err := c.getRemote().NewContainer(ctx, id, containerd.WithSpec(ociSpec), // TODO(mlaventure): when containerd support lcow, revisit runtime value containerd.WithRuntime(fmt.Sprintf("io.containerd.runtime.v1.%s", runtime.GOOS), runtimeOptions)) @@ -231,7 +244,7 @@ func (c *client) Start(ctx context.Context, id, checkpointDir string, withStdin // remove the checkpoint when we're done defer func() { if cp != nil { - err := c.remote.ContentStore().Delete(context.Background(), cp.Digest) + err := c.getRemote().ContentStore().Delete(context.Background(), cp.Digest) if err != nil { c.logger.WithError(err).WithFields(logrus.Fields{ "ref": checkpointDir, @@ -533,14 +546,14 @@ func (c *client) CreateCheckpoint(ctx context.Context, containerID, checkpointDi } // Whatever happens, delete the checkpoint from containerd defer func() { - err := c.remote.ImageService().Delete(context.Background(), img.Name()) + err := c.getRemote().ImageService().Delete(context.Background(), img.Name()) if err != nil { c.logger.WithError(err).WithField("digest", img.Target().Digest). Warnf("failed to delete checkpoint image") } }() - b, err := content.ReadBlob(ctx, c.remote.ContentStore(), img.Target().Digest) + b, err := content.ReadBlob(ctx, c.getRemote().ContentStore(), img.Target().Digest) if err != nil { return wrapSystemError(errors.Wrapf(err, "failed to retrieve checkpoint data")) } @@ -560,7 +573,7 @@ func (c *client) CreateCheckpoint(ctx context.Context, containerID, checkpointDi return wrapSystemError(errors.Wrapf(err, "invalid checkpoint")) } - rat, err := c.remote.ContentStore().ReaderAt(ctx, cpDesc.Digest) + rat, err := c.getRemote().ContentStore().ReaderAt(ctx, cpDesc.Digest) if err != nil { return wrapSystemError(errors.Wrapf(err, "failed to get checkpoint reader")) } @@ -713,7 +726,7 @@ func (c *client) processEventStream(ctx context.Context) { } }() - eventStream, err = c.remote.EventService().Subscribe(ctx, &eventsapi.SubscribeRequest{ + eventStream, err = c.getRemote().EventService().Subscribe(ctx, &eventsapi.SubscribeRequest{ Filters: []string{ "namespace==" + c.namespace, "topic~=/tasks/", @@ -723,6 +736,8 @@ func (c *client) processEventStream(ctx context.Context) { return } + c.logger.WithField("namespace", c.namespace).Debug("processing event stream") + var oomKilled bool for { ev, err = eventStream.Recv() @@ -826,7 +841,7 @@ func (c *client) processEventStream(ctx context.Context) { } func (c *client) writeContent(ctx context.Context, mediaType, ref string, r io.Reader) (*types.Descriptor, error) { - writer, err := c.remote.ContentStore().Writer(ctx, ref, 0, "") + writer, err := c.getRemote().ContentStore().Writer(ctx, ref, 0, "") if err != nil { return nil, err } diff --git a/components/engine/libcontainerd/remote_daemon.go b/components/engine/libcontainerd/remote_daemon.go index 609bcfba7a..7d0c5bb9e1 100644 --- a/components/engine/libcontainerd/remote_daemon.go +++ b/components/engine/libcontainerd/remote_daemon.go @@ -260,7 +260,7 @@ func (r *remote) startContainerd() error { return nil } -func (r *remote) monitorConnection(client *containerd.Client) { +func (r *remote) monitorConnection(monitor *containerd.Client) { var transientFailureCount = 0 ticker := time.NewTicker(500 * time.Millisecond) @@ -269,7 +269,7 @@ func (r *remote) monitorConnection(client *containerd.Client) { for { <-ticker.C ctx, cancel := context.WithTimeout(r.shutdownContext, healthCheckTimeout) - _, err := client.IsServing(ctx) + _, err := monitor.IsServing(ctx) cancel() if err == nil { transientFailureCount = 0 @@ -279,39 +279,69 @@ func (r *remote) monitorConnection(client *containerd.Client) { select { case <-r.shutdownContext.Done(): r.logger.Info("stopping healthcheck following graceful shutdown") - client.Close() + monitor.Close() return default: } r.logger.WithError(err).WithField("binary", binaryName).Debug("daemon is not responding") - if r.daemonPid != -1 { - transientFailureCount++ - if transientFailureCount >= maxConnectionRetryCount || !system.IsProcessAlive(r.daemonPid) { - transientFailureCount = 0 - if system.IsProcessAlive(r.daemonPid) { - r.logger.WithField("pid", r.daemonPid).Info("killing and restarting containerd") - // Try to get a stack trace - syscall.Kill(r.daemonPid, syscall.SIGUSR1) - <-time.After(100 * time.Millisecond) - system.KillProcess(r.daemonPid) + if r.daemonPid == -1 { + continue + } + + transientFailureCount++ + if transientFailureCount < maxConnectionRetryCount || system.IsProcessAlive(r.daemonPid) { + continue + } + + transientFailureCount = 0 + if system.IsProcessAlive(r.daemonPid) { + r.logger.WithField("pid", r.daemonPid).Info("killing and restarting containerd") + // Try to get a stack trace + syscall.Kill(r.daemonPid, syscall.SIGUSR1) + <-time.After(100 * time.Millisecond) + system.KillProcess(r.daemonPid) + } + <-r.daemonWaitCh + + monitor.Close() + os.Remove(r.GRPC.Address) + if err := r.startContainerd(); err != nil { + r.logger.WithError(err).Error("failed restarting containerd") + continue + } + + newMonitor, err := containerd.New(r.GRPC.Address) + if err != nil { + r.logger.WithError(err).Error("failed connect to containerd") + continue + } + + monitor = newMonitor + var wg sync.WaitGroup + + for _, c := range r.clients { + wg.Add(1) + + go func(c *client) { + defer wg.Done() + c.logger.WithField("namespace", c.namespace).Debug("creating new containerd remote client") + c.remote.Close() + + remote, err := containerd.New(r.GRPC.Address, containerd.WithDefaultNamespace(c.namespace)) + if err != nil { + r.logger.WithError(err).Error("failed to connect to containerd") + // TODO: Better way to handle this? + // This *shouldn't* happen, but this could wind up where the daemon + // is not able to communicate with an eventually up containerd + return } - <-r.daemonWaitCh - var err error - client.Close() - os.Remove(r.GRPC.Address) - if err = r.startContainerd(); err != nil { - r.logger.WithError(err).Error("failed restarting containerd") - } else { - newClient, err := containerd.New(r.GRPC.Address) - if err != nil { - r.logger.WithError(err).Error("failed connect to containerd") - } else { - client = newClient - } - } - } + + c.setRemote(remote) + }(c) + + wg.Wait() } } }