From 3c9abc6face8a943128cd5554486eeaf323c2ee8 Mon Sep 17 00:00:00 2001 From: CrazyMax <1951866+crazy-max@users.noreply.github.com> Date: Thu, 20 Aug 2026 13:04:59 +0200 Subject: [PATCH] docker-container: wait for BuildKit before dialing Signed-off-by: CrazyMax <1951866+crazy-max@users.noreply.github.com> --- driver/docker-container/driver.go | 19 +++++++--- driver/manager.go | 18 ++++++--- driver/manager_test.go | 36 ++++++++++++++++++ tests/build.go | 62 +++++++++++++++++++++++++++++++ 4 files changed, 123 insertions(+), 12 deletions(-) create mode 100644 driver/manager_test.go diff --git a/driver/docker-container/driver.go b/driver/docker-container/driver.go index a7688606995d..fd7049991dbf 100644 --- a/driver/docker-container/driver.go +++ b/driver/docker-container/driver.go @@ -330,12 +330,14 @@ func (d *Driver) wait(ctx context.Context, l progress.SubLogger) error { bufStderr := &bytes.Buffer{} if err := d.run(ctx, []string{"buildctl", "debug", "workers"}, bufStdout, bufStderr); err != nil { if try > 15 { - d.copyLogs(context.TODO(), l) - if bufStdout.Len() != 0 { - l.Log(1, bufStdout.Bytes()) - } - if bufStderr.Len() != 0 { - l.Log(2, bufStderr.Bytes()) + if l != nil { + d.copyLogs(context.TODO(), l) + if bufStdout.Len() != 0 { + l.Log(1, bufStdout.Bytes()) + } + if bufStderr.Len() != 0 { + l.Log(2, bufStderr.Bytes()) + } } return err } @@ -508,6 +510,11 @@ func (d *Driver) Rm(ctx context.Context, force, rmVolume, rmDaemon bool) error { } func (d *Driver) Dial(ctx context.Context) (net.Conn, error) { + // Docker marks the container running before buildkitd has necessarily + // bound its socket, so verify readiness before opening dial-stdio. + if err := d.wait(ctx, nil); err != nil { + return nil, err + } _, conn, err := d.exec(ctx, []string{"buildctl", "dial-stdio"}) if err != nil { return nil, err diff --git a/driver/manager.go b/driver/manager.go index b9ef2c275af5..881b3ac12217 100644 --- a/driver/manager.go +++ b/driver/manager.go @@ -127,17 +127,23 @@ func GetFactories(instanceRequired bool) []Factory { type DriverHandle struct { Driver client *client.Client - err error - once sync.Once + clientMu sync.Mutex historyAPISupportedOnce sync.Once historyAPISupported bool } func (d *DriverHandle) Client(ctx context.Context, opt ...client.ClientOpt) (*client.Client, error) { - d.once.Do(func() { - d.client, d.err = d.Driver.Client(ctx, append(d.getClientOptions(), opt...)...) - }) - return d.client, d.err + d.clientMu.Lock() + defer d.clientMu.Unlock() + if d.client != nil { + return d.client, nil + } + c, err := d.Driver.Client(ctx, append(d.getClientOptions(), opt...)...) + if err != nil { + return nil, err + } + d.client = c + return c, nil } func (d *DriverHandle) UncachedClient(ctx context.Context) (*client.Client, error) { diff --git a/driver/manager_test.go b/driver/manager_test.go new file mode 100644 index 000000000000..a9d552a7e47d --- /dev/null +++ b/driver/manager_test.go @@ -0,0 +1,36 @@ +package driver + +import ( + "context" + "testing" + + "github.com/moby/buildkit/client" + "github.com/stretchr/testify/require" +) + +func TestBootRetriesClientAfterErrNotRunning(t *testing.T) { + d := &retryDriver{client: &client.Client{}} + + c, err := Boot(context.Background(), context.Background(), &DriverHandle{Driver: d}, nil) + require.NoError(t, err) + require.Same(t, d.client, c) + require.Equal(t, 2, d.clientCalls) +} + +type retryDriver struct { + Driver + client *client.Client + clientCalls int +} + +func (d *retryDriver) Info(context.Context) (*Info, error) { + return &Info{Status: Running}, nil +} + +func (d *retryDriver) Client(context.Context, ...client.ClientOpt) (*client.Client, error) { + d.clientCalls++ + if d.clientCalls == 1 { + return nil, ErrNotRunning{} + } + return d.client, nil +} diff --git a/tests/build.go b/tests/build.go index 1b8c7a2504e3..b0b7febf17b3 100644 --- a/tests/build.go +++ b/tests/build.go @@ -14,12 +14,14 @@ import ( "regexp" "strings" "testing" + "time" "github.com/containerd/containerd/v2/core/content" "github.com/containerd/containerd/v2/plugins/content/local" "github.com/containerd/continuity/fs/fstest" "github.com/containerd/platforms" "github.com/creack/pty" + "github.com/docker/buildx/driver" "github.com/docker/buildx/localstate" "github.com/docker/buildx/util/confutil" "github.com/docker/buildx/util/gitutil" @@ -42,6 +44,7 @@ import ( "github.com/pkg/errors" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "golang.org/x/sync/errgroup" ) func buildCmd(sb integration.Sandbox, opts ...cmdOpt) (string, error) { @@ -94,6 +97,7 @@ var buildTests = []func(t *testing.T, sb integration.Sandbox){ testBuildCheckCallOutput, testBuildExtraHosts, testBuildIndexAnnotationsLoadDocker, + testBuildDockerContainerConcurrentFirstBuild, } func testBuild(t *testing.T, sb integration.Sandbox) { @@ -1831,6 +1835,64 @@ func testBuildIndexAnnotationsLoadDocker(t *testing.T, sb integration.Sandbox) { require.Contains(t, out, "index annotations not supported for single platform export") } +func testBuildDockerContainerConcurrentFirstBuild(t *testing.T, sb integration.Sandbox) { + if !isDockerContainerWorker(sb) { + t.Skip("only testing with docker-container worker") + } + integration.SkipOnPlatform(t, "windows") + + dir := tmpdir(t, + fstest.CreateFile("Dockerfile", []byte("FROM scratch\nCOPY marker /marker\n"), 0o600), + fstest.CreateFile("marker", []byte("hi"), 0o600), + ) + + builderName := "concurrent-" + identity.NewID() + out, err := createCmd(sb, withArgs( + "--name", builderName, + "--driver", "docker-container", + )) + require.NoError(t, err, out) + t.Cleanup(func() { + out, err := rmCmd(sb, withArgs("-f", builderName)) + require.NoError(t, err, out) + }) + + var eg errgroup.Group + var waited bool + defer func() { + if !waited { + require.NoError(t, eg.Wait()) + } + }() + startBuild := func(name string) { + eg.Go(func() error { + cmd := buildxCmd(sb, withArgs("build", "--progress=quiet", "--output=type=cacheonly", dir)) + cmd.Env = append(cmd.Env, "BUILDX_BUILDER="+builderName) + out, err := cmd.CombinedOutput() + if err != nil { + return errors.Errorf("%s failed: %v\n%s", name, err, string(out)) + } + return nil + }) + } + + startBuild("bootstrap") + containerName := fmt.Sprintf("%s0", driver.BuilderName(builderName)) + require.EventuallyWithT(t, func(c *assert.CollectT) { + cmd := dockerCmd(sb, withArgs("inspect", "-f", "{{.State.Running}}", containerName)) + out, err := cmd.CombinedOutput() + assert.NoError(c, err, string(out)) + assert.Equal(c, "true", strings.TrimSpace(string(out))) + }, 60*time.Second, 20*time.Millisecond, "container %s did not report running", containerName) + + for i := range 5 { + startBuild(fmt.Sprintf("sibling %d", i+1)) + } + err = eg.Wait() + waited = true + require.NoError(t, err) +} + func createTestProject(t *testing.T) string { dockerfile := []byte(` FROM busybox:latest AS base