docker.go
⎇
Raw
1package ci
2
3import (
4 "archive/tar"
5 "bytes"
6 "context"
7 "crypto/sha256"
8 "encoding/base64"
9 "encoding/binary"
10 "encoding/json"
11 "errors"
12 "fmt"
13 "io"
14 "log"
15 "net"
16 "net/http"
17 "net/url"
18 "os"
19 "strings"
20 "time"
21)
22
23// dockerAPIBase is the version-prefixed base URL. The host part is ignored:
24// every request goes to the unix socket.
25const dockerAPIBase = "http://localhost/v1.47"
26
27// errNoSocket is reported as a skipped run, not a failure.
28var errNoSocket = errors.New("no Docker/Podman socket found; set CI_DOCKER_SOCKET")
29
30// socketPath resolves and caches the engine socket. The order matches CI.md.
31func (r *Runner) socketPath() (string, error) {
32 r.socketMu.Lock()
33 defer r.socketMu.Unlock()
34 if r.socket != "" {
35 return r.socket, nil
36 }
37 candidates := []string{r.cfg.CIDockerSocket}
38 if r.cfg.CIDockerSocket == "" {
39 candidates = []string{
40 "/var/run/docker.sock",
41 "/run/podman/podman.sock",
42 fmt.Sprintf("/run/user/%d/podman/podman.sock", os.Getuid()),
43 }
44 }
45 for _, s := range candidates {
46 if _, err := os.Stat(s); err == nil {
47 r.socket = s
48 r.client = &http.Client{
49 Transport: &http.Transport{
50 DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
51 return (&net.Dialer{}).DialContext(ctx, "unix", s)
52 },
53 },
54 }
55 return s, nil
56 }
57 }
58 return "", errNoSocket
59}
60
61// do sends one request to the engine. The caller closes the body.
62func (r *Runner) do(ctx context.Context, method, endpoint string, body io.Reader, contentType string) (*http.Response, error) {
63 if _, err := r.socketPath(); err != nil {
64 return nil, err
65 }
66 r.socketMu.Lock()
67 client := r.client
68 r.socketMu.Unlock()
69 req, err := http.NewRequestWithContext(ctx, method, dockerAPIBase+endpoint, body)
70 if err != nil {
71 return nil, err
72 }
73 if contentType != "" {
74 req.Header.Set("Content-Type", contentType)
75 }
76 return client.Do(req)
77}
78
79// doJSON sends a JSON body and decodes nothing. It drains the response.
80func (r *Runner) doJSON(ctx context.Context, method, endpoint string, payload any) (*http.Response, []byte, error) {
81 var body io.Reader
82 ct := ""
83 if payload != nil {
84 buf, err := json.Marshal(payload)
85 if err != nil {
86 return nil, nil, err
87 }
88 body = bytes.NewReader(buf)
89 ct = "application/json"
90 }
91 resp, err := r.do(ctx, method, endpoint, body, ct)
92 if err != nil {
93 return nil, nil, err
94 }
95 defer resp.Body.Close()
96 data, err := io.ReadAll(resp.Body)
97 return resp, data, err
98}
99
100// discard drains and closes a response so the connection can be reused.
101func discard(resp *http.Response) {
102 if resp == nil {
103 return
104 }
105 io.Copy(io.Discard, resp.Body)
106 resp.Body.Close()
107}
108
109// --- Images ---
110
111// splitImageRef separates the image name from its tag. Only the part after
112// the last slash can hold a tag, because a registry host may carry a port.
113// A digest ref pins the image itself, so it is passed through whole with no
114// tag.
115func splitImageRef(image string) (name, tag string) {
116 last := image[strings.LastIndex(image, "/")+1:]
117 if strings.Contains(last, "@") {
118 return image, ""
119 }
120 i := strings.LastIndex(last, ":")
121 if i < 0 {
122 return image, "latest"
123 }
124 return image[:len(image)-len(last)+i], last[i+1:]
125}
126
127func (r *Runner) pullImage(ctx context.Context, image string) error {
128 name, tag := splitImageRef(image)
129 endpoint := "/images/create?fromImage=" + url.QueryEscape(name)
130 if tag != "" {
131 endpoint += "&tag=" + url.QueryEscape(tag)
132 }
133 resp, body, err := r.doJSON(ctx, http.MethodPost, endpoint, nil)
134 if err != nil {
135 return err
136 }
137 if resp.StatusCode >= 300 {
138 detail := strings.TrimSpace(string(body))
139 return fmt.Errorf("failed to pull image %s: HTTP %d %s", image, resp.StatusCode, detail)
140 }
141 // The stream must be read to the end: the pull only completes when the
142 // stream does. Each line is a JSON progress object, and a trailing
143 // {"error": …} means the pull failed despite the HTTP 200.
144 for _, line := range strings.Split(string(body), "\n") {
145 line = strings.TrimSpace(line)
146 if line == "" {
147 continue
148 }
149 var obj struct {
150 Error string `json:"error"`
151 }
152 if json.Unmarshal([]byte(line), &obj) != nil {
153 continue
154 }
155 if obj.Error != "" {
156 return fmt.Errorf("failed to pull image %s: %s", image, obj.Error)
157 }
158 }
159 return nil
160}
161
162// --- Volumes ---
163
164// cacheVolumeName hashes the path instead of encoding it: a truncated
165// encoding of two long paths can collide and silently share one volume.
166func cacheVolumeName(repoName, cachePath string) string {
167 sum := sha256.Sum256([]byte(repoName + ":" + cachePath))
168 return "hearthforge-ci-cache-" + base64.RawURLEncoding.EncodeToString(sum[:])[:24]
169}
170
171func (r *Runner) ensureVolume(ctx context.Context, volName, repoName, cachePath string) error {
172 resp, body, err := r.doJSON(ctx, http.MethodPost, "/volumes/create", map[string]any{
173 "Name": volName,
174 "Labels": map[string]string{
175 "com.hearthforge.repo": repoName,
176 // The name is a digest, so without this nothing can say which
177 // path a volume belongs to.
178 "com.hearthforge.cache-path": cachePath,
179 },
180 })
181 if err != nil {
182 return err
183 }
184 // 201 on success. An ignored error leaves an unlabelled volume that the
185 // cache prune never finds.
186 if resp.StatusCode >= 300 {
187 return fmt.Errorf("failed to create volume %s: %d %s",
188 volName, resp.StatusCode, strings.TrimSpace(string(body)))
189 }
190 return nil
191}
192
193func repoVolumeFilter(repoName string) string {
194 filters, _ := json.Marshal(map[string][]string{"label": {"com.hearthforge.repo=" + repoName}})
195 return "/volumes?filters=" + url.QueryEscape(string(filters))
196}
197
198// listRepoVolumes returns the names of this repo's cache volumes.
199func (r *Runner) listRepoVolumes(ctx context.Context, repoName string) ([]string, error) {
200 resp, body, err := r.doJSON(ctx, http.MethodGet, repoVolumeFilter(repoName), nil)
201 if err != nil {
202 return nil, err
203 }
204 if resp.StatusCode >= 300 {
205 return nil, fmt.Errorf("docker returned %d listing volumes", resp.StatusCode)
206 }
207 var data struct {
208 Volumes []struct{ Name string }
209 }
210 if err := json.Unmarshal(body, &data); err != nil {
211 return nil, err
212 }
213 names := make([]string, 0, len(data.Volumes))
214 for _, v := range data.Volumes {
215 names = append(names, v.Name)
216 }
217 return names, nil
218}
219
220// removeVolume deletes one volume. It reports whether the engine accepted it.
221func (r *Runner) removeVolume(ctx context.Context, name string) bool {
222 resp, _, err := r.doJSON(ctx, http.MethodDelete, "/volumes/"+name, nil)
223 return err == nil && resp.StatusCode < 300
224}
225
226// pruneStaleCaches deletes this repo's cache volumes the config no longer
227// mentions. Only safe for a default-branch run: the config is read per
228// commit, so pruning off a feature branch would delete the default branch's
229// volumes and the two would rebuild each other's caches forever.
230func (r *Runner) pruneStaleCaches(ctx context.Context, repoName string, cache []Cache) {
231 keep := map[string]bool{}
232 for _, c := range cache {
233 keep[cacheVolumeName(repoName, c.Path)] = true
234 }
235 names, err := r.listRepoVolumes(ctx, repoName)
236 if err != nil {
237 return
238 }
239 for _, name := range names {
240 if keep[name] {
241 continue
242 }
243 // A volume a concurrent run still holds cannot be deleted. That is fine.
244 // The next default-branch run deletes it.
245 r.removeVolume(ctx, name)
246 }
247}
248
249// droppedCache records one cache volume deleted for being over its cap.
250type droppedCache struct {
251 path string
252 size int64
253 maxSize int64
254}
255
256// enforceCacheLimits deletes cache volumes that outgrew their max_size.
257//
258// Sizes come from the daemon, not from a `du` in the CI container: the image
259// need not ship one, and a timed-out run has no container left to exec in. A
260// cache is only dropped on a size the daemon actually reported.
261func (r *Runner) enforceCacheLimits(ctx context.Context, repoName string, cache []Cache) []droppedCache {
262 capped := map[string]Cache{}
263 for _, c := range cache {
264 if c.MaxSize > 0 {
265 capped[cacheVolumeName(repoName, c.Path)] = c
266 }
267 }
268 if len(capped) == 0 {
269 return nil
270 }
271 resp, body, err := r.doJSON(ctx, http.MethodGet, "/system/df", nil)
272 if err != nil || resp.StatusCode >= 300 {
273 return nil
274 }
275 var data struct {
276 Volumes []struct {
277 Name string
278 UsageData *struct {
279 Size int64
280 RefCount int
281 }
282 }
283 }
284 if json.Unmarshal(body, &data) != nil {
285 return nil
286 }
287 var dropped []droppedCache
288 for _, vol := range data.Volumes {
289 cap, ok := capped[vol.Name]
290 if !ok || vol.UsageData == nil {
291 continue
292 }
293 if vol.UsageData.Size < 0 || vol.UsageData.Size <= cap.MaxSize {
294 continue
295 }
296 // A concurrent run still has it mounted, so the delete would fail.
297 if vol.UsageData.RefCount > 0 {
298 continue
299 }
300 if r.removeVolume(ctx, vol.Name) {
301 dropped = append(dropped, droppedCache{path: cap.Path, size: vol.UsageData.Size, maxSize: cap.MaxSize})
302 }
303 }
304 return dropped
305}
306
307// --- Containers ---
308
309// engineSocketInContainer is where an engine_socket step finds the socket.
310const engineSocketInContainer = "/run/hearthforge/engine.sock"
311
312func (r *Runner) createContainer(ctx context.Context, runID int64, repoName string, cfg *Config, envVars []string) (string, error) {
313 // No bind for the repo: the daemon resolves bind sources on the host,
314 // where HearthForge's own paths need not exist. uploadCheckout copies it in.
315 binds := []string{}
316 // A step that asked for the socket on a server that forbids it fails in
317 // the step loop instead.
318 if r.cfg.CIEngineSocket && cfg.WantsEngineSocket() {
319 // The socket is a host path the engine itself listens on, so the
320 // engine can resolve it even when Hearthforge runs in a container.
321 sock, err := r.socketPath()
322 if err != nil {
323 return "", err
324 }
325 binds = append(binds, sock+":"+engineSocketInContainer)
326 }
327 for _, c := range cfg.Cache {
328 volName := cacheVolumeName(repoName, c.Path)
329 if err := r.ensureVolume(ctx, volName, repoName, c.Path); err != nil {
330 return "", err
331 }
332 binds = append(binds, volName+":"+c.Path)
333 }
334
335 hostConfig := map[string]any{"Binds": binds}
336 // Empty means the engine's default network. A named network is the way
337 // to give CI containers IPv6 or an isolated segment.
338 if r.cfg.CINetwork != "" {
339 hostConfig["NetworkMode"] = r.cfg.CINetwork
340 }
341 if cfg.CPULimit > 0 {
342 hostConfig["NanoCpus"] = int64(cfg.CPULimit * 1e9)
343 }
344 if cfg.MemoryLimit != "" {
345 hostConfig["Memory"] = parseMemoryBytes(cfg.MemoryLimit)
346 }
347 workDir := cfg.WorkDir
348 if workDir == "" {
349 workDir = "/"
350 }
351 resp, body, err := r.doJSON(ctx, http.MethodPost,
352 fmt.Sprintf("/containers/create?name=hearthforge-ci-%d", runID),
353 map[string]any{
354 "Image": cfg.Image,
355 "Cmd": []string{"sleep", "infinity"},
356 "Env": envVars,
357 "WorkingDir": workDir,
358 "HostConfig": hostConfig,
359 })
360 if err != nil {
361 return "", err
362 }
363 if resp.StatusCode >= 300 {
364 return "", fmt.Errorf("failed to create container: %d %s", resp.StatusCode, string(body))
365 }
366 var data struct{ Id string }
367 if err := json.Unmarshal(body, &data); err != nil {
368 return "", err
369 }
370 return data.Id, nil
371}
372
373func (r *Runner) startContainer(ctx context.Context, containerID string) error {
374 resp, _, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/start", nil)
375 if err != nil {
376 return err
377 }
378 if resp.StatusCode >= 300 && resp.StatusCode != http.StatusNotModified {
379 return fmt.Errorf("failed to start container: %d", resp.StatusCode)
380 }
381 return nil
382}
383
384// removeContainer force-removes a container. Cleanup is best effort.
385func (r *Runner) removeContainer(ctx context.Context, containerID string) {
386 if _, _, err := r.doJSON(ctx, http.MethodDelete, "/containers/"+containerID+"?force=true", nil); err != nil {
387 log.Printf("remove container %s: %v", containerID, err)
388 }
389}
390
391// --- Exec ---
392
393type execResult struct {
394 log string
395 exitCode int
396}
397
398// readMuxFrames reads the Docker multiplexed stream and returns the payload
399// bytes. onPartial is called at most every two seconds with the log so far.
400func readMuxFrames(rd io.Reader, onPartial func(string)) string {
401 var out logBuffer
402 header := make([]byte, 8)
403 lastSave := time.Now()
404 for {
405 if _, err := io.ReadFull(rd, header); err != nil {
406 break
407 }
408 size := binary.BigEndian.Uint32(header[4:8])
409 if size == 0 {
410 continue
411 }
412 payload := make([]byte, size)
413 n, err := io.ReadFull(rd, payload)
414 if err != nil {
415 // A cut stream still carries what arrived before the cut.
416 out.append(string(payload[:n]))
417 break
418 }
419 out.append(string(payload))
420 // Once truncated the log no longer changes, so stop re-flushing it.
421 if onPartial != nil && !out.truncated && time.Since(lastSave) >= 2*time.Second {
422 onPartial(out.String())
423 lastSave = time.Now()
424 }
425 }
426 return out.String()
427}
428
429// exec runs one command in the container and returns its output and exit code.
430func (r *Runner) exec(ctx context.Context, containerID string, cmd []string, workDir string, envVars []string, onPartial func(string)) (execResult, error) {
431 payload := map[string]any{
432 "Cmd": cmd,
433 "AttachStdout": true,
434 "AttachStderr": true,
435 }
436 if workDir != "" {
437 payload["WorkingDir"] = workDir
438 }
439 if envVars != nil {
440 payload["Env"] = envVars
441 }
442 resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/exec", payload)
443 if err != nil {
444 return execResult{}, err
445 }
446 if resp.StatusCode >= 300 {
447 return execResult{}, fmt.Errorf("failed to create exec: %d %s", resp.StatusCode, string(body))
448 }
449 var created struct{ Id string }
450 if err := json.Unmarshal(body, &created); err != nil {
451 return execResult{}, err
452 }
453
454 startBody, _ := json.Marshal(map[string]any{"Detach": false, "Tty": false})
455 startResp, err := r.do(ctx, http.MethodPost, "/exec/"+created.Id+"/start",
456 bytes.NewReader(startBody), "application/json")
457 if err != nil {
458 return execResult{}, err
459 }
460 if startResp.StatusCode >= 300 {
461 detail, _ := io.ReadAll(io.LimitReader(startResp.Body, 4096))
462 startResp.Body.Close()
463 return execResult{}, fmt.Errorf("failed to start exec: %d %s",
464 startResp.StatusCode, strings.TrimSpace(string(detail)))
465 }
466 log := readMuxFrames(startResp.Body, onPartial)
467 startResp.Body.Close()
468
469 // The exit code needs a live request, so a cancelled context cannot be
470 // used here. The caller treats a cancelled exec as a failure anyway.
471 inspectCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
472 defer cancel()
473 _, inspectBody, err := r.doJSON(inspectCtx, http.MethodGet, "/exec/"+created.Id+"/json", nil)
474 if err != nil {
475 return execResult{log: log, exitCode: 1}, err
476 }
477 var inspect struct{ ExitCode *int }
478 if json.Unmarshal(inspectBody, &inspect) != nil || inspect.ExitCode == nil {
479 return execResult{log: log, exitCode: 1}, ctx.Err()
480 }
481 return execResult{log: log, exitCode: *inspect.ExitCode}, ctx.Err()
482}
483
484// logBuffer collects step output and stops at ciMaxLogBytes, so a runaway
485// step cannot exhaust RAM or make every partial flush rewrite a huge row.
486type logBuffer struct {
487 b strings.Builder
488 truncated bool
489}
490
491func (l *logBuffer) append(text string) {
492 if l.truncated || text == "" {
493 return
494 }
495 room := ciMaxLogBytes - l.b.Len()
496 if len(text) <= room {
497 l.b.WriteString(text)
498 return
499 }
500 if room > 0 {
501 l.b.WriteString(text[:room])
502 }
503 fmt.Fprintf(&l.b, "\n[log truncated at %d bytes]\n", ciMaxLogBytes)
504 l.truncated = true
505}
506
507func (l *logBuffer) String() string { return l.b.String() }
508
509// --- Archive transfer ---
510
511// copyFromImage lifts a file or directory out of another image, like
512// `COPY --from`. The source container is created but never started: the
513// archive endpoint reads its filesystem either way.
514//
515// `to` is a directory and the source basename is preserved. Renaming would
516// mean rewriting tar headers while streaming.
517func (r *Runner) copyFromImage(ctx context.Context, containerID string, spec Copy) error {
518 if err := r.pullImage(ctx, spec.Image); err != nil {
519 return err
520 }
521 // Unnamed on purpose. The startup sweep deletes by exact name, and a
522 // retry reuses the run id, so a leaked named container would fail every
523 // retry with a name conflict.
524 resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/create",
525 map[string]any{"Image": spec.Image, "Cmd": []string{"true"}})
526 if err != nil {
527 return err
528 }
529 if resp.StatusCode >= 300 {
530 return fmt.Errorf("copy: cannot create a container from %s: HTTP %d %s",
531 spec.Image, resp.StatusCode, strings.TrimSpace(string(body)))
532 }
533 var created struct{ Id string }
534 if err := json.Unmarshal(body, &created); err != nil {
535 return err
536 }
537 // The cleanup must still run when the run is cancelled mid-copy.
538 defer r.removeContainer(context.WithoutCancel(ctx), created.Id)
539
540 // Docker rejects the upload unless the destination already exists.
541 if err := r.mkdirInContainer(ctx, containerID, spec.To); err != nil {
542 return fmt.Errorf("copy: %w", err)
543 }
544
545 get, err := r.do(ctx, http.MethodGet,
546 "/containers/"+created.Id+"/archive?path="+url.QueryEscape(spec.From), nil, "")
547 if err != nil {
548 return err
549 }
550 defer discard(get)
551 if get.StatusCode >= 300 {
552 return fmt.Errorf("copy: cannot read %s from %s: HTTP %d", spec.From, spec.Image, get.StatusCode)
553 }
554 put, putBody, err := r.putArchive(ctx, containerID, spec.To, get.Body)
555 if err != nil {
556 return err
557 }
558 if put.StatusCode >= 300 {
559 return fmt.Errorf("copy: cannot write %s to %s: HTTP %d %s",
560 spec.From, spec.To, put.StatusCode, strings.TrimSpace(string(putBody)))
561 }
562 return nil
563}
564
565// putArchive uploads a tar stream into the container at destPath.
566func (r *Runner) putArchive(ctx context.Context, containerID, destPath string, body io.Reader) (*http.Response, []byte, error) {
567 resp, err := r.do(ctx, http.MethodPut,
568 "/containers/"+containerID+"/archive?path="+url.QueryEscape(destPath),
569 body, "application/x-tar")
570 if err != nil {
571 return nil, nil, err
572 }
573 defer resp.Body.Close()
574 data, _ := io.ReadAll(resp.Body)
575 return resp, data, nil
576}
577
578// copyFileFromContainer streams one file out of the container into destPath
579// and returns its size. The archive endpoint always answers with a tar, so
580// the body is read as a stream and never held in memory.
581func (r *Runner) copyFileFromContainer(ctx context.Context, containerID, containerPath, destPath string) (int64, error) {
582 resp, err := r.do(ctx, http.MethodGet,
583 "/containers/"+containerID+"/archive?path="+url.QueryEscape(containerPath), nil, "")
584 if err != nil {
585 return 0, err
586 }
587 defer discard(resp)
588 if resp.StatusCode >= 300 {
589 return 0, fmt.Errorf("cannot read %s: HTTP %d", containerPath, resp.StatusCode)
590 }
591 return writeFirstTarFile(resp.Body, destPath, r.cfg.CIMaxArtifactBytes)
592}
593
594// writeFirstTarFile copies the first regular file of a tar stream to destPath.
595// maxBytes caps the copy, so a container cannot fill the host disk.
596func writeFirstTarFile(src io.Reader, destPath string, maxBytes int64) (int64, error) {
597 tr := tar.NewReader(src)
598 for {
599 hdr, err := tr.Next()
600 if errors.Is(err, io.EOF) {
601 return 0, errors.New("archive holds no regular file")
602 }
603 if err != nil {
604 return 0, err
605 }
606 if hdr.Typeflag != tar.TypeReg {
607 continue
608 }
609 if hdr.Size > maxBytes {
610 return 0, fmt.Errorf("%s is %d bytes, over the %d byte limit", hdr.Name, hdr.Size, maxBytes)
611 }
612 f, err := os.Create(destPath)
613 if err != nil {
614 return 0, err
615 }
616 n, err := io.Copy(f, io.LimitReader(tr, maxBytes))
617 if cerr := f.Close(); err == nil {
618 err = cerr
619 }
620 if err != nil {
621 os.Remove(destPath)
622 return 0, err
623 }
624 return n, nil
625 }
626}
627
628// containerName is the name every run's container carries. The startup sweep
629// removes leftovers by this name.
630func containerName(runID int64) string {
631 return fmt.Sprintf("hearthforge-ci-%d", runID)
632}
633
634// basename and dirname keep the container-side (always POSIX) paths intact.
635