Attach a BuildKit session to builds and demultiplex container logs (closes #251) #256

Merged
clawbot merged 1 commits from issue-251-build-session into next 2026-10-01 21:15:07 +02:00
3 changed files with 225 additions and 2 deletions
+6
View File
@@ -20,6 +20,12 @@ regress.
# Completed Steps
- 2026-10-01: Builds attach a BuildKit session, as the docker command line does,
so a base image that is not on the host is pulled instead of the build failing
with "no active sessions" on Docker Engine 27. Container logs, and so the
clone output in the build log and the app logs, no longer carry Docker's
stream frame headers, and the commit is now read from the clone output (#251).
- 2026-10-01: An app's first deploy no longer fails when a volume's host path
does not exist yet: upaas asks Docker to create a missing host path when the
app's container starts and to leave an existing one alone, so nobody has to
+52 -2
View File
@@ -3,12 +3,14 @@ package docker
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net"
"os"
"path/filepath"
"regexp"
@@ -25,9 +27,11 @@ import (
"github.com/docker/docker/client"
"github.com/docker/docker/pkg/archive"
"github.com/docker/docker/pkg/jsonmessage"
"github.com/docker/docker/pkg/stdcopy"
"github.com/docker/go-connections/nat"
controlapi "github.com/moby/buildkit/api/services/control"
buildkitclient "github.com/moby/buildkit/client"
"github.com/moby/buildkit/session"
"github.com/moby/buildkit/util/progress/progressui"
"go.uber.org/fx"
@@ -388,12 +392,17 @@ func (c *Client) ContainerLogs(
}
}()
logs, err := io.ReadAll(reader)
// A container without a terminal, as all of upaas's are, sends its
// output in frames, each with a header naming stdout or stderr. Both
// go to one buffer, in the order they were written.
var logs bytes.Buffer
_, err = stdcopy.StdCopy(&logs, &logs, reader)
if err != nil {
return "", fmt.Errorf("failed to read container logs: %w", err)
}
return string(logs), nil
return logs.String(), nil
}
// IsContainerRunning checks if a container is running.
@@ -637,11 +646,24 @@ func (c *Client) performBuild(
}
}()
buildSession, err := c.startBuildSession(ctx)
if err != nil {
return "", err
}
defer func() {
closeErr := buildSession.Close()
if closeErr != nil {
c.log.Error("failed to close build session", "error", closeErr)
}
}()
// Build with BuildKit: the stages of a multi-stage build are kept in
// its build cache, which Docker limits on its own, instead of being
// left behind as untagged images.
resp, err := c.docker.ImageBuild(ctx, tarArchive, dockertypes.ImageBuildOptions{
Version: dockertypes.BuilderBuildKit,
SessionID: buildSession.ID(),
Dockerfile: opts.DockerfilePath,
Tags: opts.Tags,
Remove: true,
@@ -677,6 +699,34 @@ func (c *Client) performBuild(
return "", nil
}
// startBuildSession attaches a BuildKit session to the daemon, as the docker
// command line does for a build. BuildKit asks the client, over the session,
// for registry access to fetch a base image that is not on the host; without
// a session, Docker Engine 27 fails the build with "no active sessions". The
// shared key is only used for a build context sent over the session; upaas
// sends the context with the build request. The caller closes the session.
func (c *Client) startBuildSession(ctx context.Context) (*session.Session, error) {
buildSession, err := session.NewSession(ctx, "")
if err != nil {
return nil, fmt.Errorf("failed to create build session: %w", err)
}
go func() {
runErr := buildSession.Run(ctx, func(
ctx context.Context,
proto string,
meta map[string][]string,
) (net.Conn, error) {
return c.docker.DialHijack(ctx, "/session", proto, meta)
})
if runErr != nil {
c.log.Error("build session failed", "error", runErr)
}
}()
return buildSession, nil
}
// scannerInitialBufferSize is the initial buffer size for the build log scanner.
const scannerInitialBufferSize = 64 * 1024 // 64KB
+167
View File
@@ -6,6 +6,7 @@ import (
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
@@ -16,6 +17,7 @@ import (
"time"
"github.com/docker/docker/client"
"github.com/docker/docker/pkg/stdcopy"
controlapi "github.com/moby/buildkit/api/services/control"
)
@@ -273,6 +275,8 @@ func TestPerformBuildUsesBuildKit(t *testing.T) {
switch {
case strings.HasSuffix(r.URL.Path, "/version"):
_, _ = w.Write([]byte(`{"Version":"27.3.1","ApiVersion":"1.47"}`))
case strings.HasSuffix(r.URL.Path, "/session"):
serveSession(t, w, r, make(chan string, 1))
case strings.HasSuffix(r.URL.Path, "/build"):
if r.URL.Query().Get("version") != "2" {
http.Error(w, "not a BuildKit build", http.StatusBadRequest)
@@ -367,6 +371,8 @@ func TestPerformBuildFails(t *testing.T) {
case strings.HasSuffix(r.URL.Path, "/version"):
_, _ = fmt.Fprintf(w, `{"Version":%q,"ApiVersion":%q}`,
tt.engine, tt.apiVersion)
case strings.HasSuffix(r.URL.Path, "/session"):
serveSession(t, w, r, make(chan string, 1))
case strings.HasSuffix(r.URL.Path, "/build"):
_, _ = w.Write([]byte(tt.buildOutput))
default:
@@ -395,3 +401,164 @@ func TestPerformBuildFails(t *testing.T) {
})
}
}
// TestPerformBuildAttachesSession runs a build against a fake Docker API
// that, like the real daemon, fails the build with "no active sessions"
// unless the build names a session the client attached over the session
// endpoint. It also checks that the session is closed when the build ends.
func TestPerformBuildAttachesSession(t *testing.T) {
t.Parallel()
attached := make(chan string, 1)
sessionClosed := make(chan struct{})
srv := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
switch {
case strings.HasSuffix(r.URL.Path, "/version"):
_, _ = w.Write([]byte(`{"Version":"27.3.1","ApiVersion":"1.47"}`))
case strings.HasSuffix(r.URL.Path, "/session"):
serveSession(t, w, r, attached)
close(sessionClosed)
case strings.HasSuffix(r.URL.Path, "/build"):
id := r.URL.Query().Get("session")
attachedID := ""
// The daemon waits a few seconds for the build's session
// to attach.
if id != "" {
select {
case attachedID = <-attached:
case <-time.After(5 * time.Second):
}
}
if id == "" || attachedID != id {
_, _ = w.Write([]byte(`{"errorDetail":{"message":"no active sessions"},` +
`"error":"no active sessions"}`))
}
default:
t.Errorf("unexpected request to %s", r.URL.Path)
}
},
))
t.Cleanup(srv.Close)
dockerAPI, err := client.NewClientWithOpts(
client.WithHost("tcp://" + srv.Listener.Addr().String()),
)
if err != nil {
t.Fatal(err)
}
c := &Client{docker: dockerAPI, log: slog.Default()}
_, err = c.performBuild(t.Context(), BuildImageOptions{ContextDir: t.TempDir()})
if err != nil {
t.Fatal(err)
}
select {
case <-sessionClosed:
case <-time.After(5 * time.Second):
t.Error("the build's session was not closed when the build ended")
}
}
// serveSession answers a request to attach a session as the Docker daemon
// does: it switches the connection over to the session, sends the session's
// ID on attached, and holds the connection until the client closes it.
func serveSession(
t *testing.T,
w http.ResponseWriter,
r *http.Request,
attached chan<- string,
) {
t.Helper()
conn, _, err := http.NewResponseController(w).Hijack()
if err != nil {
t.Error(err)
return
}
defer func() { _ = conn.Close() }()
_, err = io.WriteString(conn, "HTTP/1.1 101 Switching Protocols\r\n"+
"Connection: Upgrade\r\nUpgrade: h2c\r\n\r\n")
if err != nil {
t.Error(err)
return
}
attached <- r.Header.Get("X-Docker-Expose-Session-Uuid")
_, _ = io.Copy(io.Discard, conn)
}
// TestPerformCloneReadsFramedLogs runs a clone against a fake Docker API that
// sends the clone container's output in frames, as Docker does for a
// container without a terminal, and checks that the output comes back as
// plain text and that the commit is read from it.
func TestPerformCloneReadsFramedLogs(t *testing.T) {
t.Parallel()
const commit = "1647b43aa6b211686719313bc6372c3693c54ca9"
srv := httptest.NewServer(http.HandlerFunc(
func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch {
case strings.HasSuffix(r.URL.Path, "/containers/create"):
_, _ = w.Write([]byte(`{"Id":"gitcontainer"}`))
case strings.HasSuffix(r.URL.Path, "/wait"):
_, _ = w.Write([]byte(`{"StatusCode":0}`))
case strings.HasSuffix(r.URL.Path, "/logs"):
_, _ = stdcopy.NewStdWriter(w, stdcopy.Stderr).
Write([]byte("Cloning into '/repo'...\n"))
_, _ = stdcopy.NewStdWriter(w, stdcopy.Stdout).
Write([]byte("COMMIT:" + commit + "\n"))
default:
_, _ = w.Write([]byte(`{}`))
}
},
))
t.Cleanup(srv.Close)
dockerAPI, err := client.NewClientWithOpts(
client.WithHost("tcp://" + srv.Listener.Addr().String()),
)
if err != nil {
t.Fatal(err)
}
c := &Client{docker: dockerAPI, log: slog.Default()}
dir := t.TempDir()
cfg := &cloneConfig{
repoURL: "git@example.com:repo.git",
branch: mainBranch,
sshPrivateKey: "fake-key",
containerDir: filepath.Join(dir, "repo"),
hostDir: filepath.Join(dir, "repo"),
keyFile: filepath.Join(dir, "deploy_key"),
hostKeyFile: filepath.Join(dir, "deploy_key"),
}
result, err := c.performClone(t.Context(), cfg)
if err != nil {
t.Fatal(err)
}
want := "Cloning into '/repo'...\nCOMMIT:" + commit + "\n"
if result.Output != want {
t.Errorf("got clone output %q, want %q", result.Output, want)
}
if result.CommitSHA != commit {
t.Errorf("got commit %q, want %q", result.CommitSHA, commit)
}
}