From 21ebf82c99ef1968a278e0c5db8254db7727b2de Mon Sep 17 00:00:00 2001 From: "Jonathan A. Sternberg" Date: Wed, 4 Jun 2025 15:06:31 -0500 Subject: [PATCH] monitor: refactor how reload works The build now happens in a loop and the monitor is run after every build. The monitor can return `ErrReload` to signal to the main thread that it should reload the build result. This will be used in the future to move the monitor into a callback rather than as a separate existence. It allows the monitor to not control the build itself which now makes it possible to completely remove the controller. Signed-off-by: Jonathan A. Sternberg --- commands/build.go | 119 ++++++++++++++++++--------------- commands/debug/root.go | 2 +- controller/local/controller.go | 4 -- monitor/commands/reload.go | 45 +------------ monitor/monitor.go | 49 ++++++-------- monitor/types/types.go | 11 +-- util/ioset/mux.go | 32 +++++---- 7 files changed, 113 insertions(+), 149 deletions(-) diff --git a/commands/build.go b/commands/build.go index 93791a4ac..6e8052ba2 100644 --- a/commands/build.go +++ b/commands/build.go @@ -431,57 +431,75 @@ func runControllerBuild(ctx context.Context, dockerCli command.Cli, opts *cbuild return nil, nil, errors.Errorf("Dockerfile or context from stdin is not supported with invoke") } - c := local.NewController(ctx, dockerCli) - defer c.Close() - - var in io.ReadCloser + var ( + in io.ReadCloser + f *ioset.SingleForwarder + ) if options.invokeConfig == nil { in = dockerCli.In() - } - - resp, inputs, err := c.Build(ctx, opts, in, printer) - if err != nil { - var be *controllererrors.BuildError - if errors.As(err, &be) { - retErr = err - // We can proceed to monitor - } else { - return nil, nil, errors.Wrapf(err, "failed to build") - } - } - - if options.invokeConfig != nil && options.invokeConfig.needsDebug(retErr) { - // Print errors before launching monitor - if err := printError(retErr, printer); err != nil { - logrus.Warnf("failed to print error information: %v", err) - } - - pr, pw := io.Pipe() - - f := ioset.NewSingleForwarder() + } else { + f = ioset.NewSingleForwarder() f.SetReader(dockerCli.In()) - f.SetWriter(pw, func() io.WriteCloser { - pw.Close() // propagate EOF - return nil - }) - - // TODO: ref was never set to a value in the original code. Removed the variable to - // reduce confusion but it also probably means this call is wrong in some way. - // This area should be removed during the refactor anyway so it doesn't matter that much. - monitorBuildResult, err := options.invokeConfig.runDebug(ctx, "", opts, c, pr, os.Stdout, os.Stderr, printer) - if err := pw.Close(); err != nil { - logrus.Debug("failed to close monitor stdin pipe reader") - } - if err != nil { - logrus.Warnf("failed to run monitor: %v", err) - } - if monitorBuildResult != nil { - // Update return values with the last build result from monitor - resp, retErr = monitorBuildResult.Resp, monitorBuildResult.Err - } } - return resp, inputs, retErr + for { + c := local.NewController(ctx, dockerCli) + + resp, inputs, err := c.Build(ctx, opts, in, printer) + if err != nil { + var be *controllererrors.BuildError + if errors.As(err, &be) { + retErr = err + // We can proceed to monitor + } else { + c.Close() + return nil, nil, errors.Wrapf(err, "failed to build") + } + } + + if options.invokeConfig != nil { + if err := runMonitorIfNeeded(ctx, options.invokeConfig, retErr, c, f, os.Stdout, os.Stderr, printer); err != nil { + c.Close() + if errors.Is(err, monitor.ErrReload) { + retErr = nil + continue + } + logrus.Warnf("failed to run monitor: %v", err) + } + } + + c.Close() + return resp, inputs, err + } +} + +func runMonitorIfNeeded(ctx context.Context, cfg *invokeConfig, retErr error, c *local.Controller, stdin *ioset.SingleForwarder, stdout io.WriteCloser, stderr console.File, printer *progress.Printer) error { + if !cfg.needsDebug(retErr) { + return nil + } + + // Print errors before launching monitor + if err := printError(retErr, printer); err != nil { + logrus.Warnf("failed to print error information: %v", err) + } + + pr, pw := io.Pipe() + stdin.SetWriter(pw, func() io.WriteCloser { + pw.Close() // propagate EOF + return nil + }) + + con := console.Current() + if err := con.SetRaw(); err != nil { + return errors.Errorf("failed to configure terminal: %v", err) + } + defer con.Reset() + + monitorErr := monitor.RunMonitor(ctx, &cfg.InvokeConfig, c, pr, stdout, stderr, printer) + if err := pw.Close(); err != nil { + logrus.Debug("failed to close monitor stdin pipe reader") + } + return monitorErr } func printError(err error, printer *progress.Printer) error { @@ -970,15 +988,6 @@ func (cfg *invokeConfig) needsDebug(retErr error) bool { } } -func (cfg *invokeConfig) runDebug(ctx context.Context, ref string, options *cbuild.Options, c *local.Controller, stdin io.ReadCloser, stdout io.WriteCloser, stderr console.File, progress *progress.Printer) (*monitor.MonitorBuildResult, error) { - con := console.Current() - if err := con.SetRaw(); err != nil { - return nil, errors.Errorf("failed to configure terminal: %v", err) - } - defer con.Reset() - return monitor.RunMonitor(ctx, ref, options, &cfg.InvokeConfig, c, stdin, stdout, stderr, progress) -} - func (cfg *invokeConfig) parseInvokeConfig(invoke, on string) error { cfg.onFlag = on cfg.invokeFlag = invoke diff --git a/commands/debug/root.go b/commands/debug/root.go index 9d0d0bb4f..e9294f022 100644 --- a/commands/debug/root.go +++ b/commands/debug/root.go @@ -53,7 +53,7 @@ func RootCmd(dockerCli command.Cli, children ...DebuggableCmd) *cobra.Command { return errors.Errorf("failed to configure terminal: %v", err) } - _, err = monitor.RunMonitor(ctx, "", nil, &controllerapi.InvokeConfig{ + err = monitor.RunMonitor(ctx, &controllerapi.InvokeConfig{ Tty: true, }, c, dockerCli.In(), os.Stdout, os.Stderr, printer) con.Reset() diff --git a/controller/local/controller.go b/controller/local/controller.go index 671c4ba0e..aa741e280 100644 --- a/controller/local/controller.go +++ b/controller/local/controller.go @@ -88,10 +88,6 @@ func (b *Controller) Invoke(ctx context.Context, processes *processes.Manager, p } } -func (b *Controller) Inspect(ctx context.Context) *cbuild.Options { - return b.buildConfig.buildOptions -} - func (b *Controller) Close() error { if b.buildConfig.resultCtx != nil { b.buildConfig.resultCtx.Done() diff --git a/monitor/commands/reload.go b/monitor/commands/reload.go index 2559d002a..00a59f84f 100644 --- a/monitor/commands/reload.go +++ b/monitor/commands/reload.go @@ -2,30 +2,16 @@ package commands import ( "context" - "fmt" - "io" - cbuild "github.com/docker/buildx/controller/build" - controllererrors "github.com/docker/buildx/controller/errdefs" - controllerapi "github.com/docker/buildx/controller/pb" "github.com/docker/buildx/monitor/types" - "github.com/docker/buildx/util/progress" - "github.com/moby/buildkit/solver/errdefs" - "github.com/pkg/errors" ) type ReloadCmd struct { m types.Monitor - - stdout io.WriteCloser - progress *progress.Printer - - options *cbuild.Options - invokeConfig *controllerapi.InvokeConfig } -func NewReloadCmd(m types.Monitor, stdout io.WriteCloser, progress *progress.Printer, options *cbuild.Options, invokeConfig *controllerapi.InvokeConfig) types.Command { - return &ReloadCmd{m, stdout, progress, options, invokeConfig} +func NewReloadCmd(m types.Monitor) types.Command { + return &ReloadCmd{m: m} } func (cm *ReloadCmd) Info() types.CommandInfo { @@ -40,31 +26,6 @@ Usage: } func (cm *ReloadCmd) Exec(ctx context.Context, args []string) error { - bo := cm.m.Inspect(ctx) - - var resultUpdated bool - cm.progress.Unpause() - _, _, err := cm.m.Build(ctx, bo, nil, cm.progress) // TODO: support stdin, hold build ref - cm.progress.Pause() - if err != nil { - var be *controllererrors.BuildError - if errors.As(err, &be) { - resultUpdated = true - } else { - fmt.Printf("failed to reload: %v\n", err) - } - // report error - for _, s := range errdefs.Sources(err) { - s.Print(cm.stdout) - } - fmt.Fprintf(cm.stdout, "ERROR: %v\n", err) - } else { - resultUpdated = true - } - if resultUpdated { - // rollback the running container with the new result - id := cm.m.Rollback(ctx, cm.invokeConfig) - fmt.Fprintf(cm.stdout, "Interactive container was restarted with process %q. Press Ctrl-a-c to switch to the new container\n", id) - } + cm.m.Reload() return nil } diff --git a/monitor/monitor.go b/monitor/monitor.go index 50cabd672..f7710778a 100644 --- a/monitor/monitor.go +++ b/monitor/monitor.go @@ -10,8 +10,6 @@ import ( "text/tabwriter" "github.com/containerd/console" - "github.com/docker/buildx/build" - cbuild "github.com/docker/buildx/controller/build" "github.com/docker/buildx/controller/local" controllerapi "github.com/docker/buildx/controller/pb" "github.com/docker/buildx/controller/processes" @@ -20,25 +18,23 @@ import ( "github.com/docker/buildx/util/ioset" "github.com/docker/buildx/util/progress" "github.com/google/shlex" - "github.com/moby/buildkit/client" "github.com/moby/buildkit/identity" "github.com/pkg/errors" "github.com/sirupsen/logrus" "golang.org/x/term" ) -type MonitorBuildResult struct { - Resp *client.SolveResponse - Err error -} +var ErrReload = errors.New("monitor: reload") // RunMonitor provides an interactive session for running and managing containers via specified IO. -func RunMonitor(ctx context.Context, curRef string, options *cbuild.Options, invokeConfig *controllerapi.InvokeConfig, c *local.Controller, stdin io.ReadCloser, stdout io.WriteCloser, stderr console.File, progress *progress.Printer) (*MonitorBuildResult, error) { +func RunMonitor(ctx context.Context, invokeConfig *controllerapi.InvokeConfig, c *local.Controller, stdin io.ReadCloser, stdout io.WriteCloser, stderr console.File, progress *progress.Printer) error { if err := progress.Pause(); err != nil { - return nil, err + return err } defer progress.Unpause() + defer stdin.Close() + monitorIn, monitorOut := ioset.Pipe() defer func() { monitorIn.Close() @@ -80,13 +76,13 @@ func RunMonitor(ctx context.Context, curRef string, options *cbuild.Options, inv return "Switched IO\n" }), } + m.ctx, m.cancel = context.WithCancelCause(context.Background()) defer func() { if err := m.Close(); err != nil { logrus.Warnf("close error: %v", err) } }() - m.ref.Store(curRef) // Start container automatically fmt.Fprintf(stdout, "Launching interactive container. Press Ctrl-a-c to switch to monitor console\n") @@ -96,7 +92,7 @@ func RunMonitor(ctx context.Context, curRef string, options *cbuild.Options, inv fmt.Fprintf(stdout, "Interactive container was restarted with process %q. Press Ctrl-a-c to switch to the new container\n", id) availableCommands := []types.Command{ - commands.NewReloadCmd(m, stdout, progress, options, invokeConfig), + commands.NewReloadCmd(m), commands.NewRollbackCmd(m, invokeConfig, stdout), commands.NewAttachCmd(m, stdout), commands.NewExecCmd(m, invokeConfig, stdout), @@ -128,6 +124,11 @@ func RunMonitor(ctx context.Context, curRef string, options *cbuild.Options, inv }() t := term.NewTerminal(readWriter{in.Stdin, in.Stdout}, "(buildx) ") for { + if err := m.ctx.Err(); err != nil { + errCh <- context.Cause(m.ctx) + return + } + l, err := t.ReadLine() if err != nil { if err != io.EOF { @@ -176,10 +177,10 @@ func RunMonitor(ctx context.Context, curRef string, options *cbuild.Options, inv select { case <-doneCh: m.close() - return m.lastBuildResult, nil + return nil case err := <-errCh: m.close() - return m.lastBuildResult, err + return err case <-monitorDisableCh: } monitorForwarder.SetOut(nil) @@ -233,33 +234,23 @@ type readWriter struct { } type monitor struct { - c *local.Controller - ref atomic.Value + ctx context.Context + cancel context.CancelCauseFunc + + c *local.Controller muxIO *ioset.MuxIO invokeIO *ioset.Forwarder invokeCancel func() attachedPid atomic.Value - lastBuildResult *MonitorBuildResult - processes *processes.Manager } -func (m *monitor) Build(ctx context.Context, options *cbuild.Options, in io.ReadCloser, progress progress.Writer) (resp *client.SolveResponse, input *build.Inputs, err error) { - resp, _, err = m.c.Build(ctx, options, in, progress) - m.lastBuildResult = &MonitorBuildResult{Resp: resp, Err: err} // Record build result - return -} - func (m *monitor) Invoke(ctx context.Context, pid string, cfg *controllerapi.InvokeConfig, ioIn io.ReadCloser, ioOut io.WriteCloser, ioErr io.WriteCloser) error { return m.c.Invoke(ctx, m.processes, pid, cfg, ioIn, ioOut, ioErr) } -func (m *monitor) Inspect(ctx context.Context) *cbuild.Options { - return m.c.Inspect(ctx) -} - func (m *monitor) Rollback(ctx context.Context, cfg *controllerapi.InvokeConfig) string { pid := identity.NewID() cfg1 := cfg @@ -281,6 +272,10 @@ func (m *monitor) Detach() { } } +func (m *monitor) Reload() { + m.cancel(ErrReload) +} + func (m *monitor) AttachedPID() string { return m.attachedPid.Load().(string) } diff --git a/monitor/types/types.go b/monitor/types/types.go index da837f32c..fb17b52f7 100644 --- a/monitor/types/types.go +++ b/monitor/types/types.go @@ -4,20 +4,12 @@ import ( "context" "io" - "github.com/docker/buildx/build" - cbuild "github.com/docker/buildx/controller/build" controllerapi "github.com/docker/buildx/controller/pb" "github.com/docker/buildx/controller/processes" - "github.com/docker/buildx/util/progress" - "github.com/moby/buildkit/client" ) // Monitor provides APIs for attaching and controlling the buildx server. type Monitor interface { - Build(ctx context.Context, options *cbuild.Options, in io.ReadCloser, progress progress.Writer) (resp *client.SolveResponse, inputs *build.Inputs, err error) - - Inspect(ctx context.Context) *cbuild.Options - // Invoke starts an IO session into the specified process. // If pid doesn't match to any running processes, it starts a new process with the specified config. // If there is no container running or InvokeConfig.Rollback is specified, the process will start in a newly created container. @@ -43,6 +35,9 @@ type Monitor interface { // Detach detaches IO from the container. Detach() + // Reload will signal the monitor to be reloaded. + Reload() + io.Closer } diff --git a/util/ioset/mux.go b/util/ioset/mux.go index 1d3ec616a..143ce853a 100644 --- a/util/ioset/mux.go +++ b/util/ioset/mux.go @@ -229,20 +229,24 @@ func copyToFunc(r io.Reader, wFunc func() (io.Writer, error)) error { buf := make([]byte, 4096) for { n, readErr := r.Read(buf) - if readErr != nil && readErr != io.EOF { - return readErr - } - w, err := wFunc() - if err != nil { - return err - } - if w != nil { - if _, err := w.Write(buf[:n]); err != nil { - logrus.WithError(err).Debugf("failed to copy") + + if n > 0 { + w, err := wFunc() + if err != nil { + return err + } + if w != nil { + if _, err := w.Write(buf[:n]); err != nil { + logrus.WithError(err).Debugf("failed to copy") + } } } - if readErr == io.EOF { - return nil + + if readErr != nil { + if isReaderClosed(readErr) { + return nil + } + return readErr } } } @@ -255,3 +259,7 @@ type readerWithClose struct { func (r *readerWithClose) Close() error { return r.closeFunc() } + +func isReaderClosed(err error) bool { + return errors.Is(err, io.EOF) || errors.Is(err, io.ErrClosedPipe) +}