Files
upaas/internal/docker/client.go
T
clawbot 7cf7059d3a
Check / check (pull_request) Successful in 3m47s
Leave out of the build context what the app's .dockerignore names (closes #274)
Docker does not apply an app's `.dockerignore` to a build context sent as a tar, which is how upaas sends it, so every file in the clone, `.git/config` included, reached the build.

upaas now reads the ignore file as `docker build` does, with the `ignorefile` reader of `github.com/moby/patternmatcher`, and leaves those files out of the tar: an ignore file named after the Dockerfile and next to it, such as `Dockerfile.dockerignore`, otherwise `.dockerignore` at the root of the clone. The Dockerfile path is read as a path inside the clone, and the Dockerfile (or the lowercase `dockerfile` Docker builds when `Dockerfile` is missing) and the ignore file always stay in. An app without an ignore file builds as before.

Not handled: a Dockerfile that is a symlink to an ignored file fails the build.

Model: opus-5-5
Co-authored-by: clawbot <sneak+clawbot@sneak.cloud>
2026-10-03 03:34:39 +02:00

1093 lines
29 KiB
Go

// Package docker provides Docker client functionality.
package docker
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net"
"os"
"path/filepath"
"regexp"
"strconv"
"strings"
dockertypes "github.com/docker/docker/api/types"
"github.com/docker/docker/api/types/container"
"github.com/docker/docker/api/types/filters"
"github.com/docker/docker/api/types/image"
"github.com/docker/docker/api/types/mount"
"github.com/docker/docker/api/types/network"
"github.com/docker/docker/api/types/versions"
"github.com/docker/docker/client"
"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"
"sneak.berlin/go/upaas/internal/config"
"sneak.berlin/go/upaas/internal/logger"
)
// sshKeyPermissions is the file permission for SSH private keys.
const sshKeyPermissions = 0o600
// workDirPermissions is the file permission for the work directory.
const workDirPermissions = 0o750
// stopTimeoutSeconds is the timeout for stopping containers.
const stopTimeoutSeconds = 10
// gitImage is the Docker image used for git operations.
// alpine/git v2.47.2 - pulled 2025-12-30
const gitImage = "alpine/git@sha256:" +
"d86f367afb53d022acc4377741e7334bc20add161bb10234272b91b459b4b7d8"
// ErrNotConnected is returned when Docker client is not connected.
var ErrNotConnected = errors.New("docker client not connected")
// ErrGitCloneFailed is returned when git clone fails.
var ErrGitCloneFailed = errors.New("git clone failed")
// ErrInvalidBranch is returned when a branch name contains invalid characters.
var ErrInvalidBranch = errors.New("invalid branch name")
// ErrInvalidCommitSHA is returned when a commit SHA is not a valid hex string.
var ErrInvalidCommitSHA = errors.New("invalid commit SHA")
// ErrBuildKitUnavailable is returned when the Docker daemon is too old to
// build with BuildKit.
var ErrBuildKitUnavailable = errors.New("BuildKit is unavailable on the Docker daemon")
// minBuildKitAPIVersion is the API version of Docker Engine 18.09, the first
// that builds with BuildKit when asked to without experimental mode. Older
// daemons refuse the request or silently use the legacy builder.
const minBuildKitAPIVersion = "1.39"
// validBranchRe matches safe git branch names.
var validBranchRe = regexp.MustCompile(`^[a-zA-Z0-9._/\-]+$`)
// validCommitSHARe matches a full-length hex commit SHA.
var validCommitSHARe = regexp.MustCompile(`^[0-9a-f]{40}$`)
// Params contains dependencies for Client.
type Params struct {
fx.In
Logger *logger.Logger
Config *config.Config
}
// Client wraps the Docker client.
type Client struct {
docker *client.Client
log *slog.Logger
params *Params
}
// New creates a new Docker Client.
func New(lifecycle fx.Lifecycle, params Params) (*Client, error) {
dockerClient := &Client{
log: params.Logger.Get(),
params: &params,
}
// For testing, if lifecycle is nil, skip connection (tests mock Docker)
if lifecycle == nil {
return dockerClient, nil
}
lifecycle.Append(fx.Hook{
OnStart: func(ctx context.Context) error {
return dockerClient.connect(ctx)
},
OnStop: func(_ context.Context) error {
return dockerClient.close()
},
})
return dockerClient, nil
}
// IsConnected returns true if the Docker client is connected.
func (c *Client) IsConnected() bool {
return c.docker != nil
}
// BuildImageOptions contains options for building an image.
type BuildImageOptions struct {
ContextDir string
DockerfilePath string
Tags []string
LogWriter io.Writer // Optional writer for build output
}
// BuildImage builds a Docker image from a context directory.
func (c *Client) BuildImage(
ctx context.Context,
opts BuildImageOptions,
) (ImageID, error) {
if c.docker == nil {
return "", ErrNotConnected
}
c.log.Info(
"building docker image",
"context", opts.ContextDir,
"dockerfile", opts.DockerfilePath,
)
imageID, err := c.performBuild(ctx, opts)
if err != nil {
return "", err
}
return imageID, nil
}
// CreateContainerOptions contains options for creating a container.
type CreateContainerOptions struct {
Name string
Image string
Env map[string]string
Labels map[string]string
Volumes []VolumeMount
Ports []PortMapping
Network string
CPULimit float64 // CPU cores (0.5 = half a core). 0 means unlimited.
MemoryLimit int64 // Memory in bytes. 0 means unlimited.
}
// VolumeMount represents a volume mount.
type VolumeMount struct {
HostPath string
ContainerPath string
ReadOnly bool
}
// PortMapping represents a port mapping.
type PortMapping struct {
HostPort int
ContainerPort int
Protocol string // "tcp" or "udp"
}
// nanoCPUsPerCPU is the number of NanoCPUs per CPU core.
const nanoCPUsPerCPU = 1e9
// cpuLimitToNanoCPUs converts a CPU limit (e.g. 0.5 cores) to Docker NanoCPUs.
func cpuLimitToNanoCPUs(cpuLimit float64) int64 {
return int64(cpuLimit * nanoCPUsPerCPU)
}
// buildPortConfig converts port mappings to Docker port configuration.
func buildPortConfig(ports []PortMapping) (nat.PortSet, nat.PortMap) {
exposedPorts := make(nat.PortSet)
portBindings := make(nat.PortMap)
for _, p := range ports {
proto := p.Protocol
if proto == "" {
proto = "tcp"
}
containerPort := nat.Port(fmt.Sprintf("%d/%s", p.ContainerPort, proto))
exposedPorts[containerPort] = struct{}{}
portBindings[containerPort] = []nat.PortBinding{
{
HostIP: "0.0.0.0",
HostPort: strconv.Itoa(p.HostPort),
},
}
}
return exposedPorts, portBindings
}
// buildEnvSlice converts an env map to a Docker-compatible env slice.
func buildEnvSlice(env map[string]string) []string {
envSlice := make([]string, 0, len(env))
for key, val := range env {
envSlice = append(envSlice, key+"="+val)
}
return envSlice
}
// buildMounts converts volume mounts to Docker mount configuration. Docker,
// not upaas, creates a missing host path, because upaas runs in a container
// and cannot see the host's filesystem. An existing path is left alone.
func buildMounts(volumes []VolumeMount) []mount.Mount {
mounts := make([]mount.Mount, 0, len(volumes))
for _, vol := range volumes {
mounts = append(mounts, mount.Mount{
Type: mount.TypeBind,
Source: vol.HostPath,
Target: vol.ContainerPath,
ReadOnly: vol.ReadOnly,
BindOptions: &mount.BindOptions{
CreateMountpoint: true,
},
})
}
return mounts
}
// buildResources builds Docker resource constraints from container options.
func buildResources(opts CreateContainerOptions) container.Resources {
resources := container.Resources{}
if opts.CPULimit > 0 {
resources.NanoCPUs = cpuLimitToNanoCPUs(opts.CPULimit)
}
if opts.MemoryLimit > 0 {
resources.Memory = opts.MemoryLimit
}
return resources
}
// CreateContainer creates a new container.
func (c *Client) CreateContainer(
ctx context.Context,
opts CreateContainerOptions,
) (ContainerID, error) {
if c.docker == nil {
return "", ErrNotConnected
}
c.log.Info("creating container", "name", opts.Name, "image", opts.Image)
exposedPorts, portBindings := buildPortConfig(opts.Ports)
resp, err := c.docker.ContainerCreate(ctx,
&container.Config{
Image: opts.Image,
Env: buildEnvSlice(opts.Env),
Labels: opts.Labels,
ExposedPorts: exposedPorts,
},
&container.HostConfig{
Mounts: buildMounts(opts.Volumes),
PortBindings: portBindings,
NetworkMode: container.NetworkMode(opts.Network),
Resources: buildResources(opts),
RestartPolicy: container.RestartPolicy{
Name: container.RestartPolicyUnlessStopped,
},
},
&network.NetworkingConfig{},
nil,
opts.Name,
)
if err != nil {
return "", fmt.Errorf("failed to create container: %w", err)
}
return ContainerID(resp.ID), nil
}
// StartContainer starts a container.
func (c *Client) StartContainer(ctx context.Context, containerID ContainerID) error {
if c.docker == nil {
return ErrNotConnected
}
c.log.Info("starting container", "id", containerID)
err := c.docker.ContainerStart(ctx, containerID.String(), container.StartOptions{})
if err != nil {
return fmt.Errorf("failed to start container: %w", err)
}
return nil
}
// StopContainer stops a container.
func (c *Client) StopContainer(ctx context.Context, containerID ContainerID) error {
if c.docker == nil {
return ErrNotConnected
}
c.log.Info("stopping container", "id", containerID)
timeout := stopTimeoutSeconds
err := c.docker.ContainerStop(
ctx,
containerID.String(),
container.StopOptions{Timeout: &timeout},
)
if err != nil {
return fmt.Errorf("failed to stop container: %w", err)
}
return nil
}
// RemoveContainer removes a container.
func (c *Client) RemoveContainer(
ctx context.Context,
containerID ContainerID,
force bool,
) error {
if c.docker == nil {
return ErrNotConnected
}
c.log.Info("removing container", "id", containerID, "force", force)
err := c.docker.ContainerRemove(
ctx,
containerID.String(),
container.RemoveOptions{Force: force},
)
if err != nil {
return fmt.Errorf("failed to remove container: %w", err)
}
return nil
}
// ContainerLogs returns the logs for a container.
func (c *Client) ContainerLogs(
ctx context.Context,
containerID ContainerID,
tail string,
) (string, error) {
if c.docker == nil {
return "", ErrNotConnected
}
opts := container.LogsOptions{
ShowStdout: true,
ShowStderr: true,
Tail: tail,
}
reader, err := c.docker.ContainerLogs(ctx, containerID.String(), opts)
if err != nil {
return "", fmt.Errorf("failed to get container logs: %w", err)
}
defer func() {
closeErr := reader.Close()
if closeErr != nil {
c.log.Error("failed to close log reader", "error", closeErr)
}
}()
// 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 logs.String(), nil
}
// IsContainerRunning checks if a container is running.
func (c *Client) IsContainerRunning(
ctx context.Context,
containerID ContainerID,
) (bool, error) {
if c.docker == nil {
return false, ErrNotConnected
}
inspect, err := c.docker.ContainerInspect(ctx, containerID.String())
if err != nil {
return false, fmt.Errorf("failed to inspect container: %w", err)
}
return inspect.State.Running, nil
}
// IsContainerHealthy checks if a container is healthy.
func (c *Client) IsContainerHealthy(
ctx context.Context,
containerID ContainerID,
) (bool, error) {
if c.docker == nil {
return false, ErrNotConnected
}
inspect, err := c.docker.ContainerInspect(ctx, containerID.String())
if err != nil {
return false, fmt.Errorf("failed to inspect container: %w", err)
}
// If no health check defined, consider running as healthy
if inspect.State.Health == nil {
return inspect.State.Running, nil
}
return inspect.State.Health.Status == "healthy", nil
}
// LabelUpaasID is the Docker label key used to identify containers managed by upaas.
const LabelUpaasID = "upaas.id"
// ContainerInfo contains basic information about a container.
type ContainerInfo struct {
ID ContainerID
Running bool
}
// FindContainerByAppID finds a container by the upaas.id label.
// Returns nil if no container is found.
//
//nolint:nilnil // returning nil,nil is idiomatic for "not found"
func (c *Client) FindContainerByAppID(
ctx context.Context,
appID string,
) (*ContainerInfo, error) {
if c.docker == nil {
return nil, ErrNotConnected
}
filterArgs := filters.NewArgs()
filterArgs.Add("label", LabelUpaasID+"="+appID)
containers, err := c.docker.ContainerList(ctx, container.ListOptions{
All: true,
Filters: filterArgs,
})
if err != nil {
return nil, fmt.Errorf("failed to list containers: %w", err)
}
if len(containers) == 0 {
return nil, nil
}
// Return the first matching container
ctr := containers[0]
return &ContainerInfo{
ID: ContainerID(ctr.ID),
Running: ctr.State == "running",
}, nil
}
// cloneConfig holds configuration for a git clone operation.
type cloneConfig struct {
repoURL string
branch string
commitSHA string // Optional: specific commit to checkout
sshPrivateKey string
containerDir string // Path inside the upaas container (for file operations)
hostDir string // Path on the Docker host (for bind mounts)
keyFile string // Container path to SSH key file
hostKeyFile string // Host path to SSH key file
}
// CloneResult contains the result of a git clone operation.
type CloneResult struct {
Output string // Combined stdout/stderr from git clone
CommitSHA string // The HEAD commit SHA after clone/checkout
ShortSHA string // git's short form of CommitSHA (git rev-parse --short)
}
// CloneRepo clones a git repository using SSH and optionally checks out a
// specific commit.
// containerDir is the path inside the upaas container (for writing files).
// hostDir is the corresponding path on the Docker host (for bind mounts).
// If commitSHA is provided, that specific commit will be checked out.
func (c *Client) CloneRepo(
ctx context.Context,
repoURL, branch, commitSHA, sshPrivateKey, containerDir, hostDir string,
) (*CloneResult, error) {
// Validate inputs to prevent shell injection
if !validBranchRe.MatchString(branch) {
return nil, fmt.Errorf("%w: %q", ErrInvalidBranch, branch)
}
if commitSHA != "" && !validCommitSHARe.MatchString(commitSHA) {
return nil, fmt.Errorf("%w: %q", ErrInvalidCommitSHA, commitSHA)
}
if c.docker == nil {
return nil, ErrNotConnected
}
c.log.Info("cloning repository",
"url", repoURL,
"branch", branch,
"commit", commitSHA,
"containerDir", containerDir,
"hostDir", hostDir,
)
// Clone to 'work' subdirectory, SSH key stays in build directory root
cfg := &cloneConfig{
repoURL: repoURL,
branch: branch,
commitSHA: commitSHA,
sshPrivateKey: sshPrivateKey,
containerDir: filepath.Join(containerDir, "work"),
hostDir: filepath.Join(hostDir, "work"),
keyFile: filepath.Join(containerDir, "deploy_key"),
hostKeyFile: filepath.Join(hostDir, "deploy_key"),
}
return c.performClone(ctx, cfg)
}
// RemoveImage removes a Docker image by ID or tag.
// It returns nil if the image was successfully removed or does not exist.
func (c *Client) RemoveImage(ctx context.Context, imageID ImageID) error {
_, err := c.docker.ImageRemove(ctx, imageID.String(), image.RemoveOptions{
Force: true,
PruneChildren: true,
})
if err != nil && !client.IsErrNotFound(err) {
return fmt.Errorf("failed to remove image %s: %w", imageID, err)
}
return nil
}
// ListImageTags returns the tags in the given repository, such as
// "upaas-myapp:1a2b3c4" in "upaas-myapp", each with the ID of its image.
// Tags the same image has in other repositories are left out.
func (c *Client) ListImageTags(
ctx context.Context,
repository string,
) (map[string]ImageID, error) {
if c.docker == nil {
return nil, ErrNotConnected
}
images, err := c.docker.ImageList(ctx, image.ListOptions{
Filters: filters.NewArgs(filters.Arg("reference", repository)),
})
if err != nil {
return nil, fmt.Errorf("failed to list images: %w", err)
}
tags := make(map[string]ImageID)
for _, img := range images {
for _, tag := range img.RepoTags {
if strings.HasPrefix(tag, repository+":") {
tags[tag] = ImageID(img.ID)
}
}
}
return tags, nil
}
// ListUntaggedImages returns the IDs of the images that have no tag, such
// as one whose tag a later build gave to the image it built.
func (c *Client) ListUntaggedImages(ctx context.Context) ([]ImageID, error) {
if c.docker == nil {
return nil, ErrNotConnected
}
images, err := c.docker.ImageList(ctx, image.ListOptions{
Filters: filters.NewArgs(filters.Arg("dangling", "true")),
})
if err != nil {
return nil, fmt.Errorf("failed to list untagged images: %w", err)
}
imageIDs := make([]ImageID, 0, len(images))
for _, img := range images {
imageIDs = append(imageIDs, ImageID(img.ID))
}
return imageIDs, nil
}
// RemoveImageTag removes the tag name, such as "upaas-myapp:1a2b3c4",
// without force. Docker then deletes the image, and the untagged images it
// was built on, only if no other tag and no container still uses it. If name
// is instead the ID of an untagged image, it removes that image unless a
// container uses it.
func (c *Client) RemoveImageTag(ctx context.Context, name string) error {
if c.docker == nil {
return ErrNotConnected
}
_, err := c.docker.ImageRemove(ctx, name, image.RemoveOptions{
PruneChildren: true,
})
if err != nil && !client.IsErrNotFound(err) {
return fmt.Errorf("failed to remove image %s: %w", name, err)
}
return nil
}
func (c *Client) performBuild(
ctx context.Context,
opts BuildImageOptions,
) (ImageID, error) {
server, err := c.docker.ServerVersion(ctx)
if err != nil {
return "", fmt.Errorf("failed to get Docker version: %w", err)
}
if versions.LessThan(server.APIVersion, minBuildKitAPIVersion) {
return "", fmt.Errorf(
"%w: Docker Engine %s (API %s) is older than 18.09 (API %s); "+
"upgrade Docker Engine",
ErrBuildKitUnavailable, server.Version, server.APIVersion,
minBuildKitAPIVersion,
)
}
// Create tar archive of build context
tarArchive, err := tarBuildContext(opts.ContextDir, opts.DockerfilePath)
if err != nil {
return "", fmt.Errorf("failed to create build context: %w", err)
}
defer func() {
closeErr := tarArchive.Close()
if closeErr != nil {
c.log.Error("failed to close tar archive", "error", closeErr)
}
}()
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,
NoCache: false,
})
if err != nil {
return "", fmt.Errorf("failed to build image: %w", err)
}
defer func() {
closeErr := resp.Body.Close()
if closeErr != nil {
c.log.Error("failed to close response body", "error", closeErr)
}
}()
// Stream build output line by line for real-time log updates
err = c.streamBuildOutput(ctx, resp.Body, opts.LogWriter)
if err != nil {
return "", err
}
// Get image ID
if len(opts.Tags) > 0 {
inspect, _, inspectErr := c.docker.ImageInspectWithRaw(ctx, opts.Tags[0])
if inspectErr != nil {
return "", fmt.Errorf("failed to inspect image: %w", inspectErr)
}
return ImageID(inspect.ID), 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.
const scannerInitialBufferSize = 64 * 1024 // 64KB
// scannerMaxBufferSize is the max buffer size for build log lines
// (base64 layers can be large).
const scannerMaxBufferSize = 1024 * 1024 // 1MB
// streamBuildOutput reads Docker build output line by line and writes it to
// stdout and the optional log writer as it arrives. Docker sends
// newline-delimited JSON. BuildKit's progress arrives encoded in
// "moby.buildkit.trace" messages; these are decoded and written as plain
// text, as "docker build --progress=plain" shows it. Other lines, such as
// build errors, are written unchanged. Docker ends a failed build with a line
// carrying the error; it is returned once the output is written.
func (c *Client) streamBuildOutput(
ctx context.Context,
body io.Reader,
logWriter io.Writer,
) error {
out := io.Writer(os.Stdout)
if logWriter != nil {
out = io.MultiWriter(os.Stdout, logWriter)
}
display, err := progressui.NewDisplay(out, progressui.PlainMode)
if err != nil {
return fmt.Errorf("failed to create build progress display: %w", err)
}
statuses := make(chan *buildkitclient.SolveStatus)
displayDone := make(chan struct{})
go func() {
defer close(displayDone)
// The display must keep reading until statuses is closed, even
// after ctx is cancelled, or the loop below would block.
_, _ = display.UpdateFrom(context.WithoutCancel(ctx), statuses)
}()
scanner := bufio.NewScanner(body)
buf := make([]byte, 0, scannerInitialBufferSize)
scanner.Buffer(buf, scannerMaxBufferSize)
var buildErr error
for scanner.Scan() {
line := scanner.Bytes()
var msg jsonmessage.JSONMessage
err = json.Unmarshal(line, &msg)
if err == nil && msg.ID == "moby.buildkit.trace" && msg.Aux != nil {
var data []byte
var status controlapi.StatusResponse
if json.Unmarshal(*msg.Aux, &data) == nil && status.Unmarshal(data) == nil {
statuses <- buildkitclient.NewSolveStatus(&status)
}
continue
}
if err == nil && msg.Error != nil {
buildErr = msg.Error
}
// One write per line, so it is not split by the display's output.
_, _ = fmt.Fprintf(out, "%s\n", line)
}
close(statuses)
<-displayDone
scanErr := scanner.Err()
if scanErr != nil {
return fmt.Errorf("failed to read build output: %w", scanErr)
}
return buildErr
}
func (c *Client) performClone(
ctx context.Context,
cfg *cloneConfig,
) (*CloneResult, error) {
// Create work directory for clone destination
err := os.MkdirAll(cfg.containerDir, workDirPermissions)
if err != nil {
return nil, fmt.Errorf("failed to create work dir: %w", err)
}
// Write SSH key to temp file
err = os.WriteFile(cfg.keyFile, []byte(cfg.sshPrivateKey), sshKeyPermissions)
if err != nil {
return nil, fmt.Errorf("failed to write SSH key: %w", err)
}
defer func() {
removeErr := os.Remove(cfg.keyFile)
if removeErr != nil {
c.log.Error("failed to remove SSH key file", "error", removeErr)
}
}()
gitContainerID, err := c.createGitContainer(ctx, cfg)
if err != nil {
return nil, err
}
// The git image declares a volume, so Docker gives each clone container
// an anonymous volume; remove it with the container. The removal must
// still run when the deploy is cancelled.
defer func() {
_ = c.docker.ContainerRemove(
context.WithoutCancel(ctx),
gitContainerID.String(),
container.RemoveOptions{Force: true, RemoveVolumes: true},
)
}()
return c.runGitClone(ctx, gitContainerID)
}
// ensureImage pulls ref if it is not already present locally. The pinned
// digest is preserved: a pull of an image already present is a no-op, and a
// missing one is fetched before it is used to create a container.
func (c *Client) ensureImage(ctx context.Context, ref string) error {
_, _, err := c.docker.ImageInspectWithRaw(ctx, ref)
if err == nil {
return nil
}
if !client.IsErrNotFound(err) {
return fmt.Errorf("failed to inspect image %s: %w", ref, err)
}
c.log.Info("pulling image", "image", ref)
reader, err := c.docker.ImagePull(ctx, ref, image.PullOptions{})
if err != nil {
return fmt.Errorf("failed to pull image %s: %w", ref, err)
}
defer func() {
closeErr := reader.Close()
if closeErr != nil {
c.log.Error("failed to close image pull reader", "error", closeErr)
}
}()
// The pull only completes once its response stream is fully drained.
_, err = io.Copy(io.Discard, reader)
if err != nil {
return fmt.Errorf("failed to pull image %s: %w", ref, err)
}
return nil
}
func (c *Client) createGitContainer(
ctx context.Context,
cfg *cloneConfig,
) (ContainerID, error) {
err := c.ensureImage(ctx, gitImage)
if err != nil {
return "", err
}
gitSSHCmd := "ssh -i /keys/deploy_key -o StrictHostKeyChecking=no"
// Build the git command using environment variables to avoid shell injection.
// Arguments are passed via env vars and quoted in the shell script.
var script string
if cfg.commitSHA != "" {
// Clone without depth limit so we can checkout any commit, then checkout specific SHA
script = `git clone --branch "$CLONE_BRANCH" "$CLONE_URL" /repo` +
` && cd /repo && git checkout "$CLONE_SHA"` +
` && echo COMMIT:$(git rev-parse HEAD)` +
` && echo SHORT_SHA:$(git rev-parse --short HEAD)`
} else {
// Shallow clone of branch HEAD, then output commit SHA
script = `git clone --depth 1 --branch "$CLONE_BRANCH" "$CLONE_URL" /repo` +
` && cd /repo && echo COMMIT:$(git rev-parse HEAD)` +
` && echo SHORT_SHA:$(git rev-parse --short HEAD)`
}
env := []string{
"GIT_SSH_COMMAND=" + gitSSHCmd,
"CLONE_URL=" + cfg.repoURL,
"CLONE_BRANCH=" + cfg.branch,
}
if cfg.commitSHA != "" {
env = append(env, "CLONE_SHA="+cfg.commitSHA)
}
entrypoint := []string{}
cmd := []string{"sh", "-c", script}
// Use host paths for Docker bind mounts
// (Docker runs on the host, not in our container)
resp, err := c.docker.ContainerCreate(ctx,
&container.Config{
Image: gitImage,
Entrypoint: entrypoint,
Cmd: cmd,
Env: env,
WorkingDir: "/",
},
&container.HostConfig{
Mounts: []mount.Mount{
{Type: mount.TypeBind, Source: cfg.hostDir, Target: "/repo"},
{
Type: mount.TypeBind,
Source: cfg.hostKeyFile,
Target: "/keys/deploy_key",
ReadOnly: true,
},
},
},
nil,
nil,
"",
)
if err != nil {
return "", fmt.Errorf("failed to create git container: %w", err)
}
return ContainerID(resp.ID), nil
}
func (c *Client) runGitClone(
ctx context.Context,
containerID ContainerID,
) (*CloneResult, error) {
err := c.docker.ContainerStart(ctx, containerID.String(), container.StartOptions{})
if err != nil {
return nil, fmt.Errorf("failed to start git container: %w", err)
}
statusCh, errCh := c.docker.ContainerWait(
ctx,
containerID.String(),
container.WaitConditionNotRunning,
)
select {
case err := <-errCh:
return nil, fmt.Errorf("error waiting for git container: %w", err)
case status := <-statusCh:
// Always capture logs for the result
logs, _ := c.ContainerLogs(ctx, containerID, "100")
if status.StatusCode != 0 {
return nil, fmt.Errorf(
"%w with status %d: %s",
ErrGitCloneFailed,
status.StatusCode,
logs,
)
}
// Parse the commit from the "COMMIT:" and "SHORT_SHA:" lines.
result := &CloneResult{
Output: logs,
CommitSHA: parseCommitSHA(logs, commitMarker),
ShortSHA: parseCommitSHA(logs, shortSHAMarker),
}
// The short hash names the image the deploy builds.
if result.ShortSHA == "" {
return nil, fmt.Errorf("%w: no short commit hash in its output: %s",
ErrGitCloneFailed, logs)
}
return result, nil
}
}
// Prefixes of the lines in the clone output that carry the commit checked
// out, in full and in git's short form.
const (
commitMarker = "COMMIT:"
shortSHAMarker = "SHORT_SHA:"
)
// parseCommitSHA extracts a commit SHA from git clone output.
// It looks for a line starting with marker and returns the SHA after it.
func parseCommitSHA(output, marker string) string {
for line := range strings.SplitSeq(output, "\n") {
line = strings.TrimSpace(line)
sha, found := strings.CutPrefix(line, marker)
if found {
return strings.TrimSpace(sha)
}
}
return ""
}
func (c *Client) connect(ctx context.Context) error {
opts := []client.Opt{
client.FromEnv,
client.WithAPIVersionNegotiation(),
}
if c.params.Config.DockerHost != "" {
opts = append(opts, client.WithHost(c.params.Config.DockerHost))
}
docker, err := client.NewClientWithOpts(opts...)
if err != nil {
return fmt.Errorf("failed to create Docker client: %w", err)
}
// Test connection
_, err = docker.Ping(ctx)
if err != nil {
return fmt.Errorf("failed to ping Docker: %w", err)
}
c.docker = docker
c.log.Info("docker client connected")
return nil
}
func (c *Client) close() error {
if c.docker != nil {
err := c.docker.Close()
if err != nil {
return fmt.Errorf("failed to close docker client: %w", err)
}
}
return nil
}