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 if cfg.CPULimit > 0 {
337 hostConfig["NanoCpus"] = int64(cfg.CPULimit * 1e9)
338 }
339 if cfg.MemoryLimit != "" {
340 hostConfig["Memory"] = parseMemoryBytes(cfg.MemoryLimit)
341 }
342 workDir := cfg.WorkDir
343 if workDir == "" {
344 workDir = "/"
345 }
346 resp, body, err := r.doJSON(ctx, http.MethodPost,
347 fmt.Sprintf("/containers/create?name=hearthforge-ci-%d", runID),
348 map[string]any{
349 "Image": cfg.Image,
350 "Cmd": []string{"sleep", "infinity"},
351 "Env": envVars,
352 "WorkingDir": workDir,
353 "HostConfig": hostConfig,
354 })
355 if err != nil {
356 return "", err
357 }
358 if resp.StatusCode >= 300 {
359 return "", fmt.Errorf("failed to create container: %d %s", resp.StatusCode, string(body))
360 }
361 var data struct{ Id string }
362 if err := json.Unmarshal(body, &data); err != nil {
363 return "", err
364 }
365 return data.Id, nil
366}
367
368func (r *Runner) startContainer(ctx context.Context, containerID string) error {
369 resp, _, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/start", nil)
370 if err != nil {
371 return err
372 }
373 if resp.StatusCode >= 300 && resp.StatusCode != http.StatusNotModified {
374 return fmt.Errorf("failed to start container: %d", resp.StatusCode)
375 }
376 return nil
377}
378
379// removeContainer force-removes a container. Cleanup is best effort.
380func (r *Runner) removeContainer(ctx context.Context, containerID string) {
381 if _, _, err := r.doJSON(ctx, http.MethodDelete, "/containers/"+containerID+"?force=true", nil); err != nil {
382 log.Printf("remove container %s: %v", containerID, err)
383 }
384}
385
386// --- Exec ---
387
388type execResult struct {
389 log string
390 exitCode int
391}
392
393// readMuxFrames reads the Docker multiplexed stream and returns the payload
394// bytes. onPartial is called at most every two seconds with the log so far.
395func readMuxFrames(rd io.Reader, onPartial func(string)) string {
396 var out logBuffer
397 header := make([]byte, 8)
398 lastSave := time.Now()
399 for {
400 if _, err := io.ReadFull(rd, header); err != nil {
401 break
402 }
403 size := binary.BigEndian.Uint32(header[4:8])
404 if size == 0 {
405 continue
406 }
407 payload := make([]byte, size)
408 n, err := io.ReadFull(rd, payload)
409 if err != nil {
410 // A cut stream still carries what arrived before the cut.
411 out.append(string(payload[:n]))
412 break
413 }
414 out.append(string(payload))
415 // Once truncated the log no longer changes, so stop re-flushing it.
416 if onPartial != nil && !out.truncated && time.Since(lastSave) >= 2*time.Second {
417 onPartial(out.String())
418 lastSave = time.Now()
419 }
420 }
421 return out.String()
422}
423
424// exec runs one command in the container and returns its output and exit code.
425func (r *Runner) exec(ctx context.Context, containerID string, cmd []string, workDir string, envVars []string, onPartial func(string)) (execResult, error) {
426 payload := map[string]any{
427 "Cmd": cmd,
428 "AttachStdout": true,
429 "AttachStderr": true,
430 }
431 if workDir != "" {
432 payload["WorkingDir"] = workDir
433 }
434 if envVars != nil {
435 payload["Env"] = envVars
436 }
437 resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/exec", payload)
438 if err != nil {
439 return execResult{}, err
440 }
441 if resp.StatusCode >= 300 {
442 return execResult{}, fmt.Errorf("failed to create exec: %d %s", resp.StatusCode, string(body))
443 }
444 var created struct{ Id string }
445 if err := json.Unmarshal(body, &created); err != nil {
446 return execResult{}, err
447 }
448
449 startBody, _ := json.Marshal(map[string]any{"Detach": false, "Tty": false})
450 startResp, err := r.do(ctx, http.MethodPost, "/exec/"+created.Id+"/start",
451 bytes.NewReader(startBody), "application/json")
452 if err != nil {
453 return execResult{}, err
454 }
455 if startResp.StatusCode >= 300 {
456 detail, _ := io.ReadAll(io.LimitReader(startResp.Body, 4096))
457 startResp.Body.Close()
458 return execResult{}, fmt.Errorf("failed to start exec: %d %s",
459 startResp.StatusCode, strings.TrimSpace(string(detail)))
460 }
461 log := readMuxFrames(startResp.Body, onPartial)
462 startResp.Body.Close()
463
464 // The exit code needs a live request, so a cancelled context cannot be
465 // used here. The caller treats a cancelled exec as a failure anyway.
466 inspectCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
467 defer cancel()
468 _, inspectBody, err := r.doJSON(inspectCtx, http.MethodGet, "/exec/"+created.Id+"/json", nil)
469 if err != nil {
470 return execResult{log: log, exitCode: 1}, err
471 }
472 var inspect struct{ ExitCode *int }
473 if json.Unmarshal(inspectBody, &inspect) != nil || inspect.ExitCode == nil {
474 return execResult{log: log, exitCode: 1}, ctx.Err()
475 }
476 return execResult{log: log, exitCode: *inspect.ExitCode}, ctx.Err()
477}
478
479// logBuffer collects step output and stops at ciMaxLogBytes, so a runaway
480// step cannot exhaust RAM or make every partial flush rewrite a huge row.
481type logBuffer struct {
482 b strings.Builder
483 truncated bool
484}
485
486func (l *logBuffer) append(text string) {
487 if l.truncated || text == "" {
488 return
489 }
490 room := ciMaxLogBytes - l.b.Len()
491 if len(text) <= room {
492 l.b.WriteString(text)
493 return
494 }
495 if room > 0 {
496 l.b.WriteString(text[:room])
497 }
498 fmt.Fprintf(&l.b, "\n[log truncated at %d bytes]\n", ciMaxLogBytes)
499 l.truncated = true
500}
501
502func (l *logBuffer) String() string { return l.b.String() }
503
504// --- Archive transfer ---
505
506// copyFromImage lifts a file or directory out of another image, like
507// `COPY --from`. The source container is created but never started: the
508// archive endpoint reads its filesystem either way.
509//
510// `to` is a directory and the source basename is preserved. Renaming would
511// mean rewriting tar headers while streaming.
512func (r *Runner) copyFromImage(ctx context.Context, containerID string, spec Copy) error {
513 if err := r.pullImage(ctx, spec.Image); err != nil {
514 return err
515 }
516 // Unnamed on purpose. The startup sweep deletes by exact name, and a
517 // retry reuses the run id, so a leaked named container would fail every
518 // retry with a name conflict.
519 resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/create",
520 map[string]any{"Image": spec.Image, "Cmd": []string{"true"}})
521 if err != nil {
522 return err
523 }
524 if resp.StatusCode >= 300 {
525 return fmt.Errorf("copy: cannot create a container from %s: HTTP %d %s",
526 spec.Image, resp.StatusCode, strings.TrimSpace(string(body)))
527 }
528 var created struct{ Id string }
529 if err := json.Unmarshal(body, &created); err != nil {
530 return err
531 }
532 // The cleanup must still run when the run is cancelled mid-copy.
533 defer r.removeContainer(context.WithoutCancel(ctx), created.Id)
534
535 // Docker rejects the upload unless the destination already exists.
536 if err := r.mkdirInContainer(ctx, containerID, spec.To); err != nil {
537 return fmt.Errorf("copy: %w", err)
538 }
539
540 get, err := r.do(ctx, http.MethodGet,
541 "/containers/"+created.Id+"/archive?path="+url.QueryEscape(spec.From), nil, "")
542 if err != nil {
543 return err
544 }
545 defer discard(get)
546 if get.StatusCode >= 300 {
547 return fmt.Errorf("copy: cannot read %s from %s: HTTP %d", spec.From, spec.Image, get.StatusCode)
548 }
549 put, putBody, err := r.putArchive(ctx, containerID, spec.To, get.Body)
550 if err != nil {
551 return err
552 }
553 if put.StatusCode >= 300 {
554 return fmt.Errorf("copy: cannot write %s to %s: HTTP %d %s",
555 spec.From, spec.To, put.StatusCode, strings.TrimSpace(string(putBody)))
556 }
557 return nil
558}
559
560// putArchive uploads a tar stream into the container at destPath.
561func (r *Runner) putArchive(ctx context.Context, containerID, destPath string, body io.Reader) (*http.Response, []byte, error) {
562 resp, err := r.do(ctx, http.MethodPut,
563 "/containers/"+containerID+"/archive?path="+url.QueryEscape(destPath),
564 body, "application/x-tar")
565 if err != nil {
566 return nil, nil, err
567 }
568 defer resp.Body.Close()
569 data, _ := io.ReadAll(resp.Body)
570 return resp, data, nil
571}
572
573// copyFileFromContainer streams one file out of the container into destPath
574// and returns its size. The archive endpoint always answers with a tar, so
575// the body is read as a stream and never held in memory.
576func (r *Runner) copyFileFromContainer(ctx context.Context, containerID, containerPath, destPath string) (int64, error) {
577 resp, err := r.do(ctx, http.MethodGet,
578 "/containers/"+containerID+"/archive?path="+url.QueryEscape(containerPath), nil, "")
579 if err != nil {
580 return 0, err
581 }
582 defer discard(resp)
583 if resp.StatusCode >= 300 {
584 return 0, fmt.Errorf("cannot read %s: HTTP %d", containerPath, resp.StatusCode)
585 }
586 return writeFirstTarFile(resp.Body, destPath, r.cfg.CIMaxArtifactBytes)
587}
588
589// writeFirstTarFile copies the first regular file of a tar stream to destPath.
590// maxBytes caps the copy, so a container cannot fill the host disk.
591func writeFirstTarFile(src io.Reader, destPath string, maxBytes int64) (int64, error) {
592 tr := tar.NewReader(src)
593 for {
594 hdr, err := tr.Next()
595 if errors.Is(err, io.EOF) {
596 return 0, errors.New("archive holds no regular file")
597 }
598 if err != nil {
599 return 0, err
600 }
601 if hdr.Typeflag != tar.TypeReg {
602 continue
603 }
604 if hdr.Size > maxBytes {
605 return 0, fmt.Errorf("%s is %d bytes, over the %d byte limit", hdr.Name, hdr.Size, maxBytes)
606 }
607 f, err := os.Create(destPath)
608 if err != nil {
609 return 0, err
610 }
611 n, err := io.Copy(f, io.LimitReader(tr, maxBytes))
612 if cerr := f.Close(); err == nil {
613 err = cerr
614 }
615 if err != nil {
616 os.Remove(destPath)
617 return 0, err
618 }
619 return n, nil
620 }
621}
622
623// containerName is the name every run's container carries. The startup sweep
624// removes leftovers by this name.
625func containerName(runID int64) string {
626 return fmt.Sprintf("hearthforge-ci-%d", runID)
627}
628
629// basename and dirname keep the container-side (always POSIX) paths intact.
630