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