docker.go
| 1 | package ci |
| 2 | |
| 3 | import ( |
| 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. |
| 25 | const dockerAPIBase = "http://localhost/v1.47" |
| 26 | |
| 27 | // errNoSocket is reported as a skipped run, not a failure. |
| 28 | var 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. |
| 31 | func (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. |
| 62 | func (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. |
| 80 | func (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. |
| 101 | func 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. |
| 115 | func 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 | |
| 127 | func (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. |
| 166 | func cacheVolumeName(repoName, cachePath string) string { |
| 167 | sum := sha256.Sum256([]byte(repoName + ":" + cachePath)) |
| 168 | return "hearthforge-ci-cache-" + base64.RawURLEncoding.EncodeToString(sum[:])[:24] |
| 169 | } |
| 170 | |
| 171 | func (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 | |
| 193 | func 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. |
| 199 | func (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. |
| 221 | func (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. |
| 230 | func (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. |
| 250 | type 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. |
| 261 | func (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. |
| 310 | const engineSocketInContainer = "/run/hearthforge/engine.sock" |
| 311 | |
| 312 | func (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 | // A retry reuses the name. An earlier attempt cancelled mid-create can |
| 352 | // leave a container the runner never learned the id of. |
| 353 | r.removeContainer(ctx, containerName(runID)) |
| 354 | resp, body, err := r.doJSON(ctx, http.MethodPost, |
| 355 | "/containers/create?name="+containerName(runID), |
| 356 | map[string]any{ |
| 357 | "Image": cfg.Image, |
| 358 | "Cmd": []string{"sleep", "infinity"}, |
| 359 | "Env": envVars, |
| 360 | "WorkingDir": workDir, |
| 361 | "HostConfig": hostConfig, |
| 362 | "Labels": ciContainerLabels, |
| 363 | }) |
| 364 | if err != nil { |
| 365 | return "", err |
| 366 | } |
| 367 | if resp.StatusCode >= 300 { |
| 368 | return "", fmt.Errorf("failed to create container: %d %s", resp.StatusCode, string(body)) |
| 369 | } |
| 370 | var data struct{ Id string } |
| 371 | if err := json.Unmarshal(body, &data); err != nil { |
| 372 | return "", err |
| 373 | } |
| 374 | return data.Id, nil |
| 375 | } |
| 376 | |
| 377 | func (r *Runner) startContainer(ctx context.Context, containerID string) error { |
| 378 | resp, _, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/start", nil) |
| 379 | if err != nil { |
| 380 | return err |
| 381 | } |
| 382 | if resp.StatusCode >= 300 && resp.StatusCode != http.StatusNotModified { |
| 383 | return fmt.Errorf("failed to start container: %d", resp.StatusCode) |
| 384 | } |
| 385 | return nil |
| 386 | } |
| 387 | |
| 388 | // ciContainerLabels marks every container the runner creates, so the startup |
| 389 | // sweep finds the ones a crash left behind. |
| 390 | var ciContainerLabels = map[string]string{"com.hearthforge.ci": "1"} |
| 391 | |
| 392 | // cleanupTimeout bounds each engine call made after a run ends or is cancelled. |
| 393 | const cleanupTimeout = 2 * time.Minute |
| 394 | |
| 395 | // removeContainer force-removes a container by id or name. Cleanup is best |
| 396 | // effort. It still runs when ctx is cancelled. |
| 397 | func (r *Runner) removeContainer(ctx context.Context, containerID string) { |
| 398 | ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cleanupTimeout) |
| 399 | defer cancel() |
| 400 | resp, body, err := r.doJSON(ctx, http.MethodDelete, "/containers/"+containerID+"?force=true", nil) |
| 401 | switch { |
| 402 | case errors.Is(err, errNoSocket): |
| 403 | case err != nil: |
| 404 | log.Printf("[ci] remove container %s: %v", containerID, err) |
| 405 | case resp.StatusCode >= 300 && resp.StatusCode != http.StatusNotFound: |
| 406 | log.Printf("[ci] remove container %s: HTTP %d %s", containerID, resp.StatusCode, strings.TrimSpace(string(body))) |
| 407 | } |
| 408 | } |
| 409 | |
| 410 | // removeLabelledContainers force-removes every container the runner created. |
| 411 | // Only safe while no run executes. |
| 412 | func (r *Runner) removeLabelledContainers(ctx context.Context) error { |
| 413 | filters, _ := json.Marshal(map[string][]string{"label": {"com.hearthforge.ci"}}) |
| 414 | resp, body, err := r.doJSON(ctx, http.MethodGet, |
| 415 | "/containers/json?all=true&filters="+url.QueryEscape(string(filters)), nil) |
| 416 | if err != nil { |
| 417 | return err |
| 418 | } |
| 419 | if resp.StatusCode >= 300 { |
| 420 | return fmt.Errorf("docker returned %d listing containers", resp.StatusCode) |
| 421 | } |
| 422 | var list []struct{ Id string } |
| 423 | if err := json.Unmarshal(body, &list); err != nil { |
| 424 | return err |
| 425 | } |
| 426 | for _, c := range list { |
| 427 | r.removeContainer(ctx, c.Id) |
| 428 | } |
| 429 | return nil |
| 430 | } |
| 431 | |
| 432 | // --- Exec --- |
| 433 | |
| 434 | type execResult struct { |
| 435 | log string |
| 436 | exitCode int |
| 437 | } |
| 438 | |
| 439 | // readMuxFrames reads the Docker multiplexed stream and returns the payload |
| 440 | // bytes. onPartial is called at most every two seconds with the log so far. |
| 441 | func readMuxFrames(rd io.Reader, onPartial func(string)) string { |
| 442 | var out logBuffer |
| 443 | header := make([]byte, 8) |
| 444 | lastSave := time.Now() |
| 445 | for { |
| 446 | if _, err := io.ReadFull(rd, header); err != nil { |
| 447 | break |
| 448 | } |
| 449 | size := binary.BigEndian.Uint32(header[4:8]) |
| 450 | if size == 0 { |
| 451 | continue |
| 452 | } |
| 453 | payload := make([]byte, size) |
| 454 | n, err := io.ReadFull(rd, payload) |
| 455 | if err != nil { |
| 456 | // A cut stream still carries what arrived before the cut. |
| 457 | out.append(string(payload[:n])) |
| 458 | break |
| 459 | } |
| 460 | out.append(string(payload)) |
| 461 | // Once truncated the log no longer changes, so stop re-flushing it. |
| 462 | if onPartial != nil && !out.truncated && time.Since(lastSave) >= 2*time.Second { |
| 463 | onPartial(out.String()) |
| 464 | lastSave = time.Now() |
| 465 | } |
| 466 | } |
| 467 | return out.String() |
| 468 | } |
| 469 | |
| 470 | // exec runs one command in the container and returns its output and exit code. |
| 471 | func (r *Runner) exec(ctx context.Context, containerID string, cmd []string, workDir string, envVars []string, onPartial func(string)) (execResult, error) { |
| 472 | payload := map[string]any{ |
| 473 | "Cmd": cmd, |
| 474 | "AttachStdout": true, |
| 475 | "AttachStderr": true, |
| 476 | } |
| 477 | if workDir != "" { |
| 478 | payload["WorkingDir"] = workDir |
| 479 | } |
| 480 | if envVars != nil { |
| 481 | payload["Env"] = envVars |
| 482 | } |
| 483 | resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/"+containerID+"/exec", payload) |
| 484 | if err != nil { |
| 485 | return execResult{}, err |
| 486 | } |
| 487 | if resp.StatusCode >= 300 { |
| 488 | return execResult{}, fmt.Errorf("failed to create exec: %d %s", resp.StatusCode, string(body)) |
| 489 | } |
| 490 | var created struct{ Id string } |
| 491 | if err := json.Unmarshal(body, &created); err != nil { |
| 492 | return execResult{}, err |
| 493 | } |
| 494 | |
| 495 | startBody, _ := json.Marshal(map[string]any{"Detach": false, "Tty": false}) |
| 496 | startResp, err := r.do(ctx, http.MethodPost, "/exec/"+created.Id+"/start", |
| 497 | bytes.NewReader(startBody), "application/json") |
| 498 | if err != nil { |
| 499 | return execResult{}, err |
| 500 | } |
| 501 | if startResp.StatusCode >= 300 { |
| 502 | detail, _ := io.ReadAll(io.LimitReader(startResp.Body, 4096)) |
| 503 | startResp.Body.Close() |
| 504 | return execResult{}, fmt.Errorf("failed to start exec: %d %s", |
| 505 | startResp.StatusCode, strings.TrimSpace(string(detail))) |
| 506 | } |
| 507 | log := readMuxFrames(startResp.Body, onPartial) |
| 508 | startResp.Body.Close() |
| 509 | |
| 510 | // The exit code needs a live request, so a cancelled context cannot be |
| 511 | // used here. The caller treats a cancelled exec as a failure anyway. |
| 512 | inspectCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second) |
| 513 | defer cancel() |
| 514 | _, inspectBody, err := r.doJSON(inspectCtx, http.MethodGet, "/exec/"+created.Id+"/json", nil) |
| 515 | if err != nil { |
| 516 | return execResult{log: log, exitCode: 1}, err |
| 517 | } |
| 518 | var inspect struct{ ExitCode *int } |
| 519 | if json.Unmarshal(inspectBody, &inspect) != nil || inspect.ExitCode == nil { |
| 520 | return execResult{log: log, exitCode: 1}, ctx.Err() |
| 521 | } |
| 522 | return execResult{log: log, exitCode: *inspect.ExitCode}, ctx.Err() |
| 523 | } |
| 524 | |
| 525 | // logBuffer collects step output and stops at ciMaxLogBytes, so a runaway |
| 526 | // step cannot exhaust RAM or make every partial flush rewrite a huge row. |
| 527 | type logBuffer struct { |
| 528 | b strings.Builder |
| 529 | truncated bool |
| 530 | } |
| 531 | |
| 532 | func (l *logBuffer) append(text string) { |
| 533 | if l.truncated || text == "" { |
| 534 | return |
| 535 | } |
| 536 | room := ciMaxLogBytes - l.b.Len() |
| 537 | if len(text) <= room { |
| 538 | l.b.WriteString(text) |
| 539 | return |
| 540 | } |
| 541 | if room > 0 { |
| 542 | l.b.WriteString(text[:room]) |
| 543 | } |
| 544 | l.b.WriteString(logTruncatedNotice) |
| 545 | l.truncated = true |
| 546 | } |
| 547 | |
| 548 | var logTruncatedNotice = fmt.Sprintf("\n[log truncated at %d bytes]\n", ciMaxLogBytes) |
| 549 | |
| 550 | func (l *logBuffer) String() string { return l.b.String() } |
| 551 | |
| 552 | // --- Archive transfer --- |
| 553 | |
| 554 | // copyFromImage lifts a file or directory out of another image, like |
| 555 | // `COPY --from`. The source container is created but never started: the |
| 556 | // archive endpoint reads its filesystem either way. |
| 557 | // |
| 558 | // `to` is a directory and the source basename is preserved. Renaming would |
| 559 | // mean rewriting tar headers while streaming. |
| 560 | func (r *Runner) copyFromImage(ctx context.Context, containerID string, spec Copy) error { |
| 561 | if err := r.pullImage(ctx, spec.Image); err != nil { |
| 562 | return err |
| 563 | } |
| 564 | // Unnamed: one run creates several helpers. |
| 565 | resp, body, err := r.doJSON(ctx, http.MethodPost, "/containers/create", |
| 566 | map[string]any{"Image": spec.Image, "Cmd": []string{"true"}, "Labels": ciContainerLabels}) |
| 567 | if err != nil { |
| 568 | return err |
| 569 | } |
| 570 | if resp.StatusCode >= 300 { |
| 571 | return fmt.Errorf("copy: cannot create a container from %s: HTTP %d %s", |
| 572 | spec.Image, resp.StatusCode, strings.TrimSpace(string(body))) |
| 573 | } |
| 574 | var created struct{ Id string } |
| 575 | if err := json.Unmarshal(body, &created); err != nil { |
| 576 | return err |
| 577 | } |
| 578 | defer r.removeContainer(ctx, created.Id) |
| 579 | |
| 580 | // Docker rejects the upload unless the destination already exists. |
| 581 | if err := r.mkdirInContainer(ctx, containerID, spec.To); err != nil { |
| 582 | return fmt.Errorf("copy: %w", err) |
| 583 | } |
| 584 | |
| 585 | get, err := r.do(ctx, http.MethodGet, |
| 586 | "/containers/"+created.Id+"/archive?path="+url.QueryEscape(spec.From), nil, "") |
| 587 | if err != nil { |
| 588 | return err |
| 589 | } |
| 590 | defer discard(get) |
| 591 | if get.StatusCode >= 300 { |
| 592 | return fmt.Errorf("copy: cannot read %s from %s: HTTP %d", spec.From, spec.Image, get.StatusCode) |
| 593 | } |
| 594 | put, putBody, err := r.putArchive(ctx, containerID, spec.To, get.Body) |
| 595 | if err != nil { |
| 596 | return err |
| 597 | } |
| 598 | if put.StatusCode >= 300 { |
| 599 | return fmt.Errorf("copy: cannot write %s to %s: HTTP %d %s", |
| 600 | spec.From, spec.To, put.StatusCode, strings.TrimSpace(string(putBody))) |
| 601 | } |
| 602 | return nil |
| 603 | } |
| 604 | |
| 605 | // putArchive uploads a tar stream into the container at destPath. |
| 606 | func (r *Runner) putArchive(ctx context.Context, containerID, destPath string, body io.Reader) (*http.Response, []byte, error) { |
| 607 | resp, err := r.do(ctx, http.MethodPut, |
| 608 | "/containers/"+containerID+"/archive?path="+url.QueryEscape(destPath), |
| 609 | body, "application/x-tar") |
| 610 | if err != nil { |
| 611 | return nil, nil, err |
| 612 | } |
| 613 | defer resp.Body.Close() |
| 614 | data, _ := io.ReadAll(resp.Body) |
| 615 | return resp, data, nil |
| 616 | } |
| 617 | |
| 618 | // copyFileFromContainer streams one file out of the container into destPath |
| 619 | // and returns its size. The archive endpoint always answers with a tar, so |
| 620 | // the body is read as a stream and never held in memory. |
| 621 | func (r *Runner) copyFileFromContainer(ctx context.Context, containerID, containerPath, destPath string) (int64, error) { |
| 622 | resp, err := r.do(ctx, http.MethodGet, |
| 623 | "/containers/"+containerID+"/archive?path="+url.QueryEscape(containerPath), nil, "") |
| 624 | if err != nil { |
| 625 | return 0, err |
| 626 | } |
| 627 | defer discard(resp) |
| 628 | if resp.StatusCode >= 300 { |
| 629 | return 0, fmt.Errorf("cannot read %s: HTTP %d", containerPath, resp.StatusCode) |
| 630 | } |
| 631 | return writeFirstTarFile(resp.Body, destPath, r.cfg.CIMaxArtifactBytes) |
| 632 | } |
| 633 | |
| 634 | // writeFirstTarFile copies the first regular file of a tar stream to destPath. |
| 635 | // maxBytes caps the copy, so a container cannot fill the host disk. |
| 636 | func writeFirstTarFile(src io.Reader, destPath string, maxBytes int64) (int64, error) { |
| 637 | tr := tar.NewReader(src) |
| 638 | for { |
| 639 | hdr, err := tr.Next() |
| 640 | if errors.Is(err, io.EOF) { |
| 641 | return 0, errors.New("archive holds no regular file") |
| 642 | } |
| 643 | if err != nil { |
| 644 | return 0, err |
| 645 | } |
| 646 | if hdr.Typeflag != tar.TypeReg { |
| 647 | continue |
| 648 | } |
| 649 | if hdr.Size > maxBytes { |
| 650 | return 0, fmt.Errorf("%s is %d bytes, over the %d byte limit", hdr.Name, hdr.Size, maxBytes) |
| 651 | } |
| 652 | f, err := os.Create(destPath) |
| 653 | if err != nil { |
| 654 | return 0, err |
| 655 | } |
| 656 | n, err := io.Copy(f, io.LimitReader(tr, maxBytes)) |
| 657 | if cerr := f.Close(); err == nil { |
| 658 | err = cerr |
| 659 | } |
| 660 | if err != nil { |
| 661 | os.Remove(destPath) |
| 662 | return 0, err |
| 663 | } |
| 664 | return n, nil |
| 665 | } |
| 666 | } |
| 667 | |
| 668 | // containerName is the name every run's container carries. |
| 669 | func containerName(runID int64) string { |
| 670 | return fmt.Sprintf("hearthforge-ci-%d", runID) |
| 671 | } |
| 672 | |
| 673 | // basename and dirname keep the container-side (always POSIX) paths intact. |
| 674 |