package ci import ( "archive/tar" "bytes" "context" "crypto/sha256" "encoding/base64" "encoding/binary" "encoding/json" "errors" "fmt" "io" "log" "net" "net/http" "net/url" "os" "strings" "time" ) // dockerAPIBase is the version-prefixed base URL. The host part is ignored: // every request goes to the unix socket. const dockerAPIBase = "http://localhost/v1.47" // errNoSocket is reported as a skipped run, not a failure. var errNoSocket = errors.New("no Docker/Podman socket found; set CI_DOCKER_SOCKET") // socketPath resolves and caches the engine socket. The order matches CI.md. func (r *Runner) socketPath() (string, error) { r.socketMu.Lock() defer r.socketMu.Unlock() if r.socket != "" { return r.socket, nil } candidates := []string{r.cfg.CIDockerSocket} if r.cfg.CIDockerSocket == "" { candidates = []string{ "/var/run/docker.sock", "/run/podman/podman.sock", fmt.Sprintf("/run/user/%d/podman/podman.sock", os.Getuid()), } } for _, s := range candidates { if _, err := os.Stat(s); err == nil { r.socket = s r.client = &http.Client{ Transport: &http.Transport{ DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) { return (&net.Dialer{}).DialContext(ctx, "unix", s) }, }, } return s, nil } } return "", errNoSocket } // do sends one request to the engine. The caller closes the body. func (r *Runner) do(ctx context.Context, method, endpoint string, body io.Reader, contentType string) (*http.Response, error) { if _, err := r.socketPath(); err != nil { return nil, err } r.socketMu.Lock() client := r.client r.socketMu.Unlock() req, err := http.NewRequestWithContext(ctx, method, dockerAPIBase+endpoint, body) if err != nil { return nil, err } if contentType != "" { req.Header.Set("Content-Type", contentType) } return client.Do(req) } // doJSON sends a JSON body and decodes nothing. It drains the response. func (r *Runner) doJSON(ctx context.Context, method, endpoint string, payload any) (*http.Response, []byte, error) { var body io.Reader ct := "" if payload != nil { buf, err := json.Marshal(payload) if err != nil { return nil, nil, err } body = bytes.NewReader(buf) ct = "application/json" } resp, err := r.do(ctx, method, endpoint, body, ct) if err != nil { return nil, nil, err } defer resp.Body.Close() data, err := io.ReadAll(resp.Body) return resp, data, err } // discard drains and closes a response so the connection can be reused. func discard(resp *http.Response) { if resp == nil { return } io.Copy(io.Discard, resp.Body) resp.Body.Close() } // --- Images --- // splitImageRef separates the image name from its tag. Only the part after // the last slash can hold a tag, because a registry host may carry a port. // A digest ref pins the image itself, so it is passed through whole with no // tag. func splitImageRef(image string) (name, tag string) { last := image[strings.LastIndex(image, "/")+1:] if strings.Contains(last, "@") { return image, "" } i := strings.LastIndex(last, ":") if i < 0 { return image, "latest" } return image[:len(image)-len(last)+i], last[i+1:] } func (r *Runner) pullImage(ctx context.Context, image string) error { name, tag := splitImageRef(image) endpoint := "/images/create?fromImage=" + url.QueryEscape(name) if tag != "" { endpoint += "&tag=" + url.QueryEscape(tag) } resp, body, err := r.doJSON(ctx, http.MethodPost, endpoint, nil) if err != nil { return err } if resp.StatusCode >= 300 { detail := strings.TrimSpace(string(body)) return fmt.Errorf("failed to pull image %s: HTTP %d %s", image, resp.StatusCode, detail) } // The stream must be read to the end: the pull only completes when the // stream does. Each line is a JSON progress object, and a trailing // {"error": …} means the pull failed despite the HTTP 200. for _, line := range strings.Split(string(body), "\n") { line = strings.TrimSpace(line) if line == "" { continue } var obj struct { Error string `json:"error"` } if json.Unmarshal([]byte(line), &obj) != nil { continue } if obj.Error != "" { return fmt.Errorf("failed to pull image %s: %s", image, obj.Error) } } return nil } // --- Volumes --- // cacheVolumeName hashes the path instead of encoding it: a truncated // encoding of two long paths can collide and silently share one volume. func cacheVolumeName(repoName, cachePath string) string { sum := sha256.Sum256([]byte(repoName + ":" + cachePath)) return "hearthforge-ci-cache-" + base64.RawURLEncoding.EncodeToString(sum[:])[:24] } func (r *Runner) ensureVolume(ctx context.Context, volName, repoName, cachePath string) error { resp, body, err := r.doJSON(ctx, http.MethodPost, "/volumes/create", map[string]any{ "Name": volName, "Labels": map[string]string{ "com.hearthforge.repo": repoName, // The name is a digest, so without this nothing can say which // path a volume belongs to. "com.hearthforge.cache-path": cachePath, }, }) if err != nil { return err } // 201 on success. An ignored error leaves an unlabelled volume that the // cache prune never finds. if resp.StatusCode >= 300 { return fmt.Errorf("failed to create volume %s: %d %s", volName, resp.StatusCode, strings.TrimSpace(string(body))) } return nil } func repoVolumeFilter(repoName string) string { filters, _ := json.Marshal(map[string][]string{"label": {"com.hearthforge.repo=" + repoName}}) return "/volumes?filters=" + url.QueryEscape(string(filters)) } // listRepoVolumes returns the names of this repo's cache volumes. func (r *Runner) listRepoVolumes(ctx context.Context, repoName string) ([]string, error) { resp, body, err := r.doJSON(ctx, http.MethodGet, repoVolumeFilter(repoName), nil) if err != nil { return nil, err } if resp.StatusCode >= 300 { return nil, fmt.Errorf("docker returned %d listing volumes", resp.StatusCode) } var data struct { Volumes []struct{ Name string } } if err := json.Unmarshal(body, &data); err != nil { return nil, err } names := make([]string, 0, len(data.Volumes)) for _, v := range data.Volumes { names = append(names, v.Name) } return names, nil } // removeVolume deletes one volume. It reports whether the engine accepted it. func (r *Runner) removeVolume(ctx context.Context, name string) bool { resp, _, err := r.doJSON(ctx, http.MethodDelete, "/volumes/"+name, nil) return err == nil && resp.StatusCode < 300 } // pruneStaleCaches deletes this repo's cache volumes the config no longer // mentions. Only safe for a default-branch run: the config is read per // commit, so pruning off a feature branch would delete the default branch's // volumes and the two would rebuild each other's caches forever. func (r *Runner) pruneStaleCaches(ctx context.Context, repoName string, cache []Cache) { keep := map[string]bool{} for _, c := range cache { keep[cacheVolumeName(repoName, c.Path)] = true } names, err := r.listRepoVolumes(ctx, repoName) if err != nil { return } for _, name := range names { if keep[name] { continue } // A volume a concurrent run still holds cannot be deleted. That is fine. // The next default-branch run deletes it. r.removeVolume(ctx, name) } } // droppedCache records one cache volume deleted for being over its cap. type droppedCache struct { path string size int64 maxSize int64 } // enforceCacheLimits deletes cache volumes that outgrew their max_size. // // Sizes come from the daemon, not from a `du` in the CI container: the image // need not ship one, and a timed-out run has no container left to exec in. A // cache is only dropped on a size the daemon actually reported. func (r *Runner) enforceCacheLimits(ctx context.Context, repoName string, cache []Cache) []droppedCache { capped := map[string]Cache{} for _, c := range cache { if c.MaxSize > 0 { capped[cacheVolumeName(repoName, c.Path)] = c } } if len(capped) == 0 { return nil } resp, body, err := r.doJSON(ctx, http.MethodGet, "/system/df", nil) if err != nil || resp.StatusCode >= 300 { return nil } var data struct { Volumes []struct { Name string UsageData *struct { Size int64 RefCount int } } } if json.Unmarshal(body, &data) != nil { return nil } var dropped []droppedCache for _, vol := range data.Volumes { cap, ok := capped[vol.Name] if !ok || vol.UsageData == nil { continue } if vol.UsageData.Size < 0 || vol.UsageData.Size <= cap.MaxSize { continue } // A concurrent run still has it mounted, so the delete would fail. if vol.UsageData.RefCount > 0 { continue } if r.removeVolume(ctx, vol.Name) { dropped = append(dropped, droppedCache{path: cap.Path, size: vol.UsageData.Size, maxSize: cap.MaxSize}) } } return dropped } // --- Containers --- // engineSocketInContainer is where an engine_socket step finds the socket. const engineSocketInContainer = "/run/hearthforge/engine.sock" func (r *Runner) createContainer(ctx context.Context, runID int64, repoName string, cfg *Config, envVars []string) (string, error) { // No bind for the repo: the daemon resolves bind sources on the host, // where HearthForge's own paths need not exist. uploadCheckout copies it in. binds := []string{} // A step that asked for the socket on a server that forbids it fails in // the step loop instead. if r.cfg.CIEngineSocket && cfg.WantsEngineSocket() { // The socket is a host path the engine itself listens on, so the // engine can resolve it even when Hearthforge runs in a container. sock, err := r.socketPath() if err != nil { return "", err } binds = append(binds, sock+":"+engineSocketInContainer) } for _, c := range cfg.Cache { volName := cacheVolumeName(repoName, c.Path) if err := r.ensureVolume(ctx, volName, repoName, c.Path); err != nil { return "", err } binds = append(binds, volName+":"+c.Path) } hostConfig := map[string]any{"Binds": binds} if cfg.CPULimit > 0 { hostConfig["NanoCpus"] = int64(cfg.CPULimit * 1e9) } if cfg.MemoryLimit != "" { hostConfig["Memory"] = parseMemoryBytes(cfg.MemoryLimit) } workDir := cfg.WorkDir if workDir == "" { workDir = "/" } resp, body, err := r.doJSON(ctx, http.MethodPost, fmt.Sprintf("/containers/create?name=hearthforge-ci-%d", runID), map[string]any{ "Image": cfg.Image, "Cmd": []string{"sleep", "infinity"}, "Env": envVars, "WorkingDir": workDir, "HostConfig": hostConfig, }) if err != nil { return "", err } if resp.StatusCode >= 300 { return "", fmt.Errorf("failed to create container: %d %s", resp.StatusCode, string(body)) } var data struct{ Id string } if err := json.Unmarshal(body, &data); err != nil { return "", err } return data.Id, nil } func (r *Runner) startContainer(ctx context.Context, containerID string) error { resp, _, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/start", nil) if err != nil { return err } if resp.StatusCode >= 300 && resp.StatusCode != http.StatusNotModified { return fmt.Errorf("failed to start container: %d", resp.StatusCode) } return nil } // removeContainer force-removes a container. Cleanup is best effort. func (r *Runner) removeContainer(ctx context.Context, containerID string) { if _, _, err := r.doJSON(ctx, http.MethodDelete, "/containers/"+containerID+"?force=true", nil); err != nil { log.Printf("remove container %s: %v", containerID, err) } } // --- Exec --- type execResult struct { log string exitCode int } // readMuxFrames reads the Docker multiplexed stream and returns the payload // bytes. onPartial is called at most every two seconds with the log so far. func readMuxFrames(rd io.Reader, onPartial func(string)) string { var out logBuffer header := make([]byte, 8) lastSave := time.Now() for { if _, err := io.ReadFull(rd, header); err != nil { break } size := binary.BigEndian.Uint32(header[4:8]) if size == 0 { continue } payload := make([]byte, size) n, err := io.ReadFull(rd, payload) if err != nil { // A cut stream still carries what arrived before the cut. out.append(string(payload[:n])) break } out.append(string(payload)) // Once truncated the log no longer changes, so stop re-flushing it. if onPartial != nil && !out.truncated && time.Since(lastSave) >= 2*time.Second { onPartial(out.String()) lastSave = time.Now() } } return out.String() } // exec runs one command in the container and returns its output and exit code. func (r *Runner) exec(ctx context.Context, containerID string, cmd []string, workDir string, envVars []string, onPartial func(string)) (execResult, error) { payload := map[string]any{ "Cmd": cmd, "AttachStdout": true, "AttachStderr": true, } if workDir != "" { payload["WorkingDir"] = workDir } if envVars != nil { payload["Env"] = envVars } resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/exec", payload) if err != nil { return execResult{}, err } if resp.StatusCode >= 300 { return execResult{}, fmt.Errorf("failed to create exec: %d %s", resp.StatusCode, string(body)) } var created struct{ Id string } if err := json.Unmarshal(body, &created); err != nil { return execResult{}, err } startBody, _ := json.Marshal(map[string]any{"Detach": false, "Tty": false}) startResp, err := r.do(ctx, http.MethodPost, "/exec/"+created.Id+"/start", bytes.NewReader(startBody), "application/json") if err != nil { return execResult{}, err } if startResp.StatusCode >= 300 { detail, _ := io.ReadAll(io.LimitReader(startResp.Body, 4096)) startResp.Body.Close() return execResult{}, fmt.Errorf("failed to start exec: %d %s", startResp.StatusCode, strings.TrimSpace(string(detail))) } log := readMuxFrames(startResp.Body, onPartial) startResp.Body.Close() // The exit code needs a live request, so a cancelled context cannot be // used here. The caller treats a cancelled exec as a failure anyway. inspectCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) defer cancel() _, inspectBody, err := r.doJSON(inspectCtx, http.MethodGet, "/exec/"+created.Id+"/json", nil) if err != nil { return execResult{log: log, exitCode: 1}, err } var inspect struct{ ExitCode *int } if json.Unmarshal(inspectBody, &inspect) != nil || inspect.ExitCode == nil { return execResult{log: log, exitCode: 1}, ctx.Err() } return execResult{log: log, exitCode: *inspect.ExitCode}, ctx.Err() } // logBuffer collects step output and stops at ciMaxLogBytes, so a runaway // step cannot exhaust RAM or make every partial flush rewrite a huge row. type logBuffer struct { b strings.Builder truncated bool } func (l *logBuffer) append(text string) { if l.truncated || text == "" { return } room := ciMaxLogBytes - l.b.Len() if len(text) <= room { l.b.WriteString(text) return } if room > 0 { l.b.WriteString(text[:room]) } fmt.Fprintf(&l.b, "\n[log truncated at %d bytes]\n", ciMaxLogBytes) l.truncated = true } func (l *logBuffer) String() string { return l.b.String() } // --- Archive transfer --- // copyFromImage lifts a file or directory out of another image, like // `COPY --from`. The source container is created but never started: the // archive endpoint reads its filesystem either way. // // `to` is a directory and the source basename is preserved. Renaming would // mean rewriting tar headers while streaming. func (r *Runner) copyFromImage(ctx context.Context, containerID string, spec Copy) error { if err := r.pullImage(ctx, spec.Image); err != nil { return err } // Unnamed on purpose. The startup sweep deletes by exact name, and a // retry reuses the run id, so a leaked named container would fail every // retry with a name conflict. resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/create", map[string]any{"Image": spec.Image, "Cmd": []string{"true"}}) if err != nil { return err } if resp.StatusCode >= 300 { return fmt.Errorf("copy: cannot create a container from %s: HTTP %d %s", spec.Image, resp.StatusCode, strings.TrimSpace(string(body))) } var created struct{ Id string } if err := json.Unmarshal(body, &created); err != nil { return err } // The cleanup must still run when the run is cancelled mid-copy. defer r.removeContainer(context.WithoutCancel(ctx), created.Id) // Docker rejects the upload unless the destination already exists. if err := r.mkdirInContainer(ctx, containerID, spec.To); err != nil { return fmt.Errorf("copy: %w", err) } get, err := r.do(ctx, http.MethodGet, "/containers/"+created.Id+"/archive?path="+url.QueryEscape(spec.From), nil, "") if err != nil { return err } defer discard(get) if get.StatusCode >= 300 { return fmt.Errorf("copy: cannot read %s from %s: HTTP %d", spec.From, spec.Image, get.StatusCode) } put, putBody, err := r.putArchive(ctx, containerID, spec.To, get.Body) if err != nil { return err } if put.StatusCode >= 300 { return fmt.Errorf("copy: cannot write %s to %s: HTTP %d %s", spec.From, spec.To, put.StatusCode, strings.TrimSpace(string(putBody))) } return nil } // putArchive uploads a tar stream into the container at destPath. func (r *Runner) putArchive(ctx context.Context, containerID, destPath string, body io.Reader) (*http.Response, []byte, error) { resp, err := r.do(ctx, http.MethodPut, "/containers/"+containerID+"/archive?path="+url.QueryEscape(destPath), body, "application/x-tar") if err != nil { return nil, nil, err } defer resp.Body.Close() data, _ := io.ReadAll(resp.Body) return resp, data, nil } // copyFileFromContainer streams one file out of the container into destPath // and returns its size. The archive endpoint always answers with a tar, so // the body is read as a stream and never held in memory. func (r *Runner) copyFileFromContainer(ctx context.Context, containerID, containerPath, destPath string) (int64, error) { resp, err := r.do(ctx, http.MethodGet, "/containers/"+containerID+"/archive?path="+url.QueryEscape(containerPath), nil, "") if err != nil { return 0, err } defer discard(resp) if resp.StatusCode >= 300 { return 0, fmt.Errorf("cannot read %s: HTTP %d", containerPath, resp.StatusCode) } return writeFirstTarFile(resp.Body, destPath, r.cfg.CIMaxArtifactBytes) } // writeFirstTarFile copies the first regular file of a tar stream to destPath. // maxBytes caps the copy, so a container cannot fill the host disk. func writeFirstTarFile(src io.Reader, destPath string, maxBytes int64) (int64, error) { tr := tar.NewReader(src) for { hdr, err := tr.Next() if errors.Is(err, io.EOF) { return 0, errors.New("archive holds no regular file") } if err != nil { return 0, err } if hdr.Typeflag != tar.TypeReg { continue } if hdr.Size > maxBytes { return 0, fmt.Errorf("%s is %d bytes, over the %d byte limit", hdr.Name, hdr.Size, maxBytes) } f, err := os.Create(destPath) if err != nil { return 0, err } n, err := io.Copy(f, io.LimitReader(tr, maxBytes)) if cerr := f.Close(); err == nil { err = cerr } if err != nil { os.Remove(destPath) return 0, err } return n, nil } } // containerName is the name every run's container carries. The startup sweep // removes leftovers by this name. func containerName(runID int64) string { return fmt.Sprintf("hearthforge-ci-%d", runID) } // basename and dirname keep the container-side (always POSIX) paths intact.