Attach a BuildKit session to builds and demultiplex container logs (closes #251) #256
@@ -20,6 +20,12 @@ regress.
|
|||||||
|
|
||||||
# Completed Steps
|
# 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
|
- 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
|
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
|
app's container starts and to leave an existing one alone, so nobody has to
|
||||||
|
|||||||
@@ -3,12 +3,14 @@ package docker
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bufio"
|
"bufio"
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
|
"net"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
"regexp"
|
"regexp"
|
||||||
@@ -25,9 +27,11 @@ import (
|
|||||||
"github.com/docker/docker/client"
|
"github.com/docker/docker/client"
|
||||||
"github.com/docker/docker/pkg/archive"
|
"github.com/docker/docker/pkg/archive"
|
||||||
"github.com/docker/docker/pkg/jsonmessage"
|
"github.com/docker/docker/pkg/jsonmessage"
|
||||||
|
"github.com/docker/docker/pkg/stdcopy"
|
||||||
"github.com/docker/go-connections/nat"
|
"github.com/docker/go-connections/nat"
|
||||||
controlapi "github.com/moby/buildkit/api/services/control"
|
controlapi "github.com/moby/buildkit/api/services/control"
|
||||||
buildkitclient "github.com/moby/buildkit/client"
|
buildkitclient "github.com/moby/buildkit/client"
|
||||||
|
"github.com/moby/buildkit/session"
|
||||||
"github.com/moby/buildkit/util/progress/progressui"
|
"github.com/moby/buildkit/util/progress/progressui"
|
||||||
"go.uber.org/fx"
|
"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 {
|
if err != nil {
|
||||||
return "", fmt.Errorf("failed to read container logs: %w", err)
|
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.
|
// 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
|
// 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
|
// its build cache, which Docker limits on its own, instead of being
|
||||||
// left behind as untagged images.
|
// left behind as untagged images.
|
||||||
resp, err := c.docker.ImageBuild(ctx, tarArchive, dockertypes.ImageBuildOptions{
|
resp, err := c.docker.ImageBuild(ctx, tarArchive, dockertypes.ImageBuildOptions{
|
||||||
Version: dockertypes.BuilderBuildKit,
|
Version: dockertypes.BuilderBuildKit,
|
||||||
|
SessionID: buildSession.ID(),
|
||||||
Dockerfile: opts.DockerfilePath,
|
Dockerfile: opts.DockerfilePath,
|
||||||
Tags: opts.Tags,
|
Tags: opts.Tags,
|
||||||
Remove: true,
|
Remove: true,
|
||||||
@@ -677,6 +699,34 @@ func (c *Client) performBuild(
|
|||||||
return "", nil
|
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.
|
// scannerInitialBufferSize is the initial buffer size for the build log scanner.
|
||||||
const scannerInitialBufferSize = 64 * 1024 // 64KB
|
const scannerInitialBufferSize = 64 * 1024 // 64KB
|
||||||
|
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
@@ -16,6 +17,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/docker/docker/client"
|
"github.com/docker/docker/client"
|
||||||
|
"github.com/docker/docker/pkg/stdcopy"
|
||||||
controlapi "github.com/moby/buildkit/api/services/control"
|
controlapi "github.com/moby/buildkit/api/services/control"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -273,6 +275,8 @@ func TestPerformBuildUsesBuildKit(t *testing.T) {
|
|||||||
switch {
|
switch {
|
||||||
case strings.HasSuffix(r.URL.Path, "/version"):
|
case strings.HasSuffix(r.URL.Path, "/version"):
|
||||||
_, _ = w.Write([]byte(`{"Version":"27.3.1","ApiVersion":"1.47"}`))
|
_, _ = 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"):
|
case strings.HasSuffix(r.URL.Path, "/build"):
|
||||||
if r.URL.Query().Get("version") != "2" {
|
if r.URL.Query().Get("version") != "2" {
|
||||||
http.Error(w, "not a BuildKit build", http.StatusBadRequest)
|
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"):
|
case strings.HasSuffix(r.URL.Path, "/version"):
|
||||||
_, _ = fmt.Fprintf(w, `{"Version":%q,"ApiVersion":%q}`,
|
_, _ = fmt.Fprintf(w, `{"Version":%q,"ApiVersion":%q}`,
|
||||||
tt.engine, tt.apiVersion)
|
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"):
|
case strings.HasSuffix(r.URL.Path, "/build"):
|
||||||
_, _ = w.Write([]byte(tt.buildOutput))
|
_, _ = w.Write([]byte(tt.buildOutput))
|
||||||
default:
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user