run.go
| 1 | package ci |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "encoding/json" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | "io" |
| 10 | "log" |
| 11 | "net/http" |
| 12 | "os" |
| 13 | "os/exec" |
| 14 | "path" |
| 15 | "path/filepath" |
| 16 | "slices" |
| 17 | "sort" |
| 18 | "strconv" |
| 19 | "strings" |
| 20 | "sync" |
| 21 | "time" |
| 22 | |
| 23 | "hearthforge/internal/config" |
| 24 | "hearthforge/internal/db" |
| 25 | "hearthforge/internal/gitcmd" |
| 26 | ) |
| 27 | |
| 28 | // ciMaxLogBytes caps the log stored per step. |
| 29 | const ciMaxLogBytes = 2 * 1024 * 1024 |
| 30 | |
| 31 | // TriggerOpts describes why a run starts and which commit it builds. |
| 32 | type TriggerOpts struct { |
| 33 | // TriggerSource is "push", "tag" or "manual". |
| 34 | TriggerSource string |
| 35 | CommitSha string |
| 36 | CommitBranch string |
| 37 | CommitTag string |
| 38 | TriggeredBy int64 |
| 39 | VariableOverrides map[string]string |
| 40 | } |
| 41 | |
| 42 | // task is one run that holds a concurrency slot. |
| 43 | type task struct { |
| 44 | cancel context.CancelFunc |
| 45 | containerID string |
| 46 | } |
| 47 | |
| 48 | // Runner owns the CI queue and talks to the container engine. |
| 49 | type Runner struct { |
| 50 | cfg *config.Config |
| 51 | db *db.DB |
| 52 | |
| 53 | // ImportImage stores an image from a build VM in the registry and |
| 54 | // returns its manifest digest. The web server provides it. |
| 55 | ImportImage func(ctx context.Context, repoName, image string, tags []string, src io.Reader) (string, error) |
| 56 | |
| 57 | // vmImageSlot serialises builds of the default build VM image. |
| 58 | vmImageSlot chan struct{} |
| 59 | |
| 60 | socketMu sync.Mutex |
| 61 | socket string |
| 62 | client *http.Client |
| 63 | |
| 64 | mu sync.Mutex |
| 65 | running map[int64]*task |
| 66 | // queue holds run ids whose DB row is "queued". pumpQueue is the sole |
| 67 | // writer of `running`, so the slot check and the reservation happen |
| 68 | // together under one lock. |
| 69 | queue []int64 |
| 70 | } |
| 71 | |
| 72 | func New(cfg *config.Config, database *db.DB) *Runner { |
| 73 | return &Runner{cfg: cfg, db: database, running: map[int64]*task{}, vmImageSlot: make(chan struct{}, 1)} |
| 74 | } |
| 75 | |
| 76 | func (r *Runner) repoPath(name string) string { |
| 77 | return filepath.Join(r.cfg.ReposDir(), name+".git") |
| 78 | } |
| 79 | |
| 80 | // ErrNoConfig means the commit carries no .hearthforge-ci.toml. Callers tell |
| 81 | // it apart from a parse error to explain why CI cannot run. |
| 82 | var ErrNoConfig = errors.New(".hearthforge-ci.toml not found at commit") |
| 83 | |
| 84 | // ConfigAt reads and parses .hearthforge-ci.toml at a commit or ref. |
| 85 | func (r *Runner) ConfigAt(ctx context.Context, repoName, ref string) (*Config, error) { |
| 86 | cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "show", |
| 87 | "--end-of-options", ref+":.hearthforge-ci.toml") |
| 88 | cmd.Env = gitcmd.Env() |
| 89 | out, err := cmd.Output() |
| 90 | if err != nil || len(out) == 0 { |
| 91 | return nil, ErrNoConfig |
| 92 | } |
| 93 | return ParseCiConfig(string(out)) |
| 94 | } |
| 95 | |
| 96 | // --- Queue --- |
| 97 | |
| 98 | // QueuePosition returns the 1-based place of a queued run, or 0. |
| 99 | func (r *Runner) QueuePosition(runID int64) int { |
| 100 | r.mu.Lock() |
| 101 | defer r.mu.Unlock() |
| 102 | for i, id := range r.queue { |
| 103 | if id == runID { |
| 104 | return i + 1 |
| 105 | } |
| 106 | } |
| 107 | return 0 |
| 108 | } |
| 109 | |
| 110 | // enqueue puts a run in the queue and starts it when a slot is free. |
| 111 | // The row is flipped to "queued" under the same lock that decides there is no |
| 112 | // slot, so two triggers that arrive together cannot both stay "pending". |
| 113 | func (r *Runner) enqueue(ctx context.Context, runID int64) { |
| 114 | r.mu.Lock() |
| 115 | r.queue = append(r.queue, runID) |
| 116 | if len(r.running)+len(r.queue) > r.cfg.CIMaxConcurrent { |
| 117 | // pumpQueue is the only starter and it needs the same lock, so this |
| 118 | // update cannot clobber a row that just went running. |
| 119 | r.execSQL(ctx, `UPDATE ci_runs SET status = 'queued' WHERE id = ?`, runID) |
| 120 | } |
| 121 | r.mu.Unlock() |
| 122 | r.pumpQueue() |
| 123 | } |
| 124 | |
| 125 | func (r *Runner) pumpQueue() { |
| 126 | for { |
| 127 | r.mu.Lock() |
| 128 | if len(r.queue) == 0 || len(r.running) >= r.cfg.CIMaxConcurrent { |
| 129 | r.mu.Unlock() |
| 130 | return |
| 131 | } |
| 132 | runID := r.queue[0] |
| 133 | r.queue = r.queue[1:] |
| 134 | // The run outlives the request that created it, so it gets its own |
| 135 | // context. CancelRun cancels it. |
| 136 | ctx, cancel := context.WithCancel(context.Background()) |
| 137 | t := &task{cancel: cancel} |
| 138 | r.running[runID] = t |
| 139 | r.mu.Unlock() |
| 140 | |
| 141 | go r.spawnRun(ctx, runID, t) |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | // spawnRun promotes the row to pending, runs it, and reconciles the row when |
| 146 | // executeRun itself fails. |
| 147 | func (r *Runner) spawnRun(ctx context.Context, runID int64, t *task) { |
| 148 | defer func() { |
| 149 | if rec := recover(); rec != nil { |
| 150 | log.Printf("[ci] run %d panicked: %v", runID, rec) |
| 151 | r.finishSlot(runID, t) |
| 152 | r.execSQL(context.Background(), |
| 153 | `UPDATE ci_runs SET status = 'failure', finished_at = ? |
| 154 | WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID) |
| 155 | r.pumpQueue() |
| 156 | } |
| 157 | }() |
| 158 | // Promote queued→pending before executing, so executeRun's failure path |
| 159 | // (which only matches pending/running/queued) can mark it failed. |
| 160 | r.execSQL(ctx, `UPDATE ci_runs SET status = 'pending' WHERE id = ?`, runID) |
| 161 | r.executeRun(ctx, runID, t) |
| 162 | } |
| 163 | |
| 164 | // finishSlot releases the run's concurrency slot. It only drops the slot this |
| 165 | // execution owns: a retry that started the same run again holds a different |
| 166 | // task, and deleting that one would leak its slot forever. |
| 167 | func (r *Runner) finishSlot(runID int64, t *task) { |
| 168 | r.mu.Lock() |
| 169 | if r.running[runID] == t { |
| 170 | delete(r.running, runID) |
| 171 | } |
| 172 | r.mu.Unlock() |
| 173 | } |
| 174 | |
| 175 | func (r *Runner) setContainer(t *task, containerID string) { |
| 176 | r.mu.Lock() |
| 177 | t.containerID = containerID |
| 178 | r.mu.Unlock() |
| 179 | } |
| 180 | |
| 181 | // --- Small DB helpers --- |
| 182 | |
| 183 | func (r *Runner) execSQL(ctx context.Context, query string, args ...any) { |
| 184 | if _, err := r.db.ExecContext(ctx, query, args...); err != nil { |
| 185 | log.Printf("[ci] sql failed: %v", err) |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | func (r *Runner) insertStep(ctx context.Context, runID int64, name, status string, started bool) int64 { |
| 190 | var startedAt any |
| 191 | if started { |
| 192 | startedAt = db.NowISO() |
| 193 | } |
| 194 | res, err := r.db.ExecContext(ctx, |
| 195 | `INSERT INTO ci_steps (run_id, name, status, started_at) VALUES (?, ?, ?, ?)`, |
| 196 | runID, name, status, startedAt) |
| 197 | if err != nil { |
| 198 | log.Printf("[ci] insert step failed: %v", err) |
| 199 | return 0 |
| 200 | } |
| 201 | id, _ := res.LastInsertId() |
| 202 | return id |
| 203 | } |
| 204 | |
| 205 | // --- Public API --- |
| 206 | |
| 207 | // TriggerRun records a run and starts it, or queues it when every slot is |
| 208 | // taken. It returns the run id. |
| 209 | func (r *Runner) TriggerRun(ctx context.Context, repoName string, opts TriggerOpts) (int64, error) { |
| 210 | var repoID int64 |
| 211 | err := r.db.QueryRowContext(ctx, `SELECT id FROM repositories WHERE name = ?`, repoName).Scan(&repoID) |
| 212 | if err != nil { |
| 213 | return 0, errors.New("repository not found") |
| 214 | } |
| 215 | |
| 216 | var overrides any |
| 217 | if len(opts.VariableOverrides) > 0 { |
| 218 | buf, err := json.Marshal(opts.VariableOverrides) |
| 219 | if err != nil { |
| 220 | return 0, err |
| 221 | } |
| 222 | overrides = string(buf) |
| 223 | } |
| 224 | var triggeredBy any |
| 225 | if opts.TriggeredBy > 0 { |
| 226 | triggeredBy = opts.TriggeredBy |
| 227 | } |
| 228 | res, err := r.db.ExecContext(ctx, |
| 229 | `INSERT INTO ci_runs (repo_id, triggered_by, trigger_source, commit_sha, |
| 230 | commit_branch, commit_tag, status, variable_overrides) |
| 231 | VALUES (?, ?, ?, ?, ?, ?, 'pending', ?)`, |
| 232 | repoID, triggeredBy, opts.TriggerSource, opts.CommitSha, |
| 233 | nullString(opts.CommitBranch), nullString(opts.CommitTag), overrides) |
| 234 | if err != nil { |
| 235 | return 0, err |
| 236 | } |
| 237 | runID, err := res.LastInsertId() |
| 238 | if err != nil { |
| 239 | return 0, err |
| 240 | } |
| 241 | |
| 242 | // The human-facing run number comes from a per-repo counter. One |
| 243 | // upsert-and-increment cannot collide under concurrent triggers, and it |
| 244 | // never reuses a number after pruneHistory shrinks the table. |
| 245 | var repoRunID int64 |
| 246 | err = r.db.QueryRowContext(ctx, |
| 247 | `INSERT INTO ci_run_counters (repo_id, last_run_id) VALUES (?, 1) |
| 248 | ON CONFLICT(repo_id) DO UPDATE SET last_run_id = last_run_id + 1 |
| 249 | RETURNING last_run_id`, repoID).Scan(&repoRunID) |
| 250 | if err != nil { |
| 251 | return 0, err |
| 252 | } |
| 253 | r.execSQL(ctx, `UPDATE ci_runs SET repo_run_id = ? WHERE id = ?`, repoRunID, runID) |
| 254 | |
| 255 | r.enqueue(ctx, runID) |
| 256 | return runID, nil |
| 257 | } |
| 258 | |
| 259 | // ErrRunNotFinished is returned when a retry targets a run that is still |
| 260 | // queued or executing. |
| 261 | var ErrRunNotFinished = errors.New("run is not finished") |
| 262 | |
| 263 | // Active reports whether a run still holds its slot. A run frees it only |
| 264 | // after its cleanup, so a finished status alone does not mean idle. |
| 265 | func (r *Runner) Active(runID int64) bool { |
| 266 | r.mu.Lock() |
| 267 | defer r.mu.Unlock() |
| 268 | return r.running[runID] != nil |
| 269 | } |
| 270 | |
| 271 | // RetryRun clears a finished run's steps and artifacts and runs it again. |
| 272 | // The reset only applies to a run in a terminal status, so a retry cannot |
| 273 | // hijack a run that is still executing. |
| 274 | func (r *Runner) RetryRun(ctx context.Context, runID, retriedBy int64) error { |
| 275 | var repoID int64 |
| 276 | if err := r.db.QueryRowContext(ctx, `SELECT repo_id FROM ci_runs WHERE id = ?`, runID).Scan(&repoID); err != nil { |
| 277 | return errors.New("run not found") |
| 278 | } |
| 279 | // CancelRun marks the row before the execution has cleaned up. Its |
| 280 | // cleanup would clobber the new attempt's rows and container. |
| 281 | if r.Active(runID) { |
| 282 | return ErrRunNotFinished |
| 283 | } |
| 284 | res, err := r.db.ExecContext(ctx, |
| 285 | `UPDATE ci_runs SET status = 'pending', triggered_by = ?, started_at = NULL, finished_at = NULL |
| 286 | WHERE id = ? AND status IN ('success','failure','warning','cancelled','skipped')`, |
| 287 | retriedBy, runID) |
| 288 | if err != nil { |
| 289 | return err |
| 290 | } |
| 291 | if n, err := res.RowsAffected(); err != nil { |
| 292 | return err |
| 293 | } else if n == 0 { |
| 294 | return ErrRunNotFinished |
| 295 | } |
| 296 | r.execSQL(ctx, `DELETE FROM ci_steps WHERE run_id = ?`, runID) |
| 297 | os.RemoveAll(r.artifactDir(runID)) |
| 298 | r.execSQL(ctx, `DELETE FROM ci_artifacts WHERE run_id = ?`, runID) |
| 299 | r.enqueue(ctx, runID) |
| 300 | return nil |
| 301 | } |
| 302 | |
| 303 | // CancelRun stops a running or queued run. |
| 304 | func (r *Runner) CancelRun(ctx context.Context, runID int64) error { |
| 305 | r.mu.Lock() |
| 306 | t := r.running[runID] |
| 307 | var containerID string |
| 308 | if t != nil { |
| 309 | containerID = t.containerID |
| 310 | } |
| 311 | for i, id := range r.queue { |
| 312 | if id == runID { |
| 313 | r.queue = append(r.queue[:i], r.queue[i+1:]...) |
| 314 | break |
| 315 | } |
| 316 | } |
| 317 | r.mu.Unlock() |
| 318 | |
| 319 | if t != nil { |
| 320 | t.cancel() |
| 321 | if containerID != "" { |
| 322 | // Docker has no per-exec kill, so removing the container is what |
| 323 | // stops the command running inside it. |
| 324 | r.removeContainer(ctx, containerID) |
| 325 | } |
| 326 | } |
| 327 | r.execSQL(ctx, |
| 328 | `UPDATE ci_runs SET status = 'cancelled', finished_at = ? |
| 329 | WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID) |
| 330 | return nil |
| 331 | } |
| 332 | |
| 333 | // PurgeRepoCaches deletes this repo's cache volumes and returns how many |
| 334 | // went. It fails when the engine is unreachable. |
| 335 | func (r *Runner) PurgeRepoCaches(ctx context.Context, repoName string) (int, error) { |
| 336 | names, err := r.listRepoVolumes(ctx, repoName) |
| 337 | if err != nil { |
| 338 | return 0, err |
| 339 | } |
| 340 | removed := 0 |
| 341 | for _, name := range names { |
| 342 | if r.removeVolume(ctx, name) { |
| 343 | removed++ |
| 344 | } |
| 345 | } |
| 346 | return removed, nil |
| 347 | } |
| 348 | |
| 349 | // StopRepo cancels the repo's active runs, waits for them to clean up, and |
| 350 | // deletes their artifacts and cache volumes. Call it before the repo row goes: |
| 351 | // the run ids come from the database. |
| 352 | func (r *Runner) StopRepo(ctx context.Context, repoID int64, repoName string) { |
| 353 | ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE repo_id = ?`, repoID) |
| 354 | if err != nil { |
| 355 | log.Printf("[ci] stop repo %s: %v", repoName, err) |
| 356 | return |
| 357 | } |
| 358 | for _, id := range ids { |
| 359 | if err := r.CancelRun(ctx, id); err != nil { |
| 360 | log.Printf("[ci] stop repo %s: cancel run %d: %v", repoName, id, err) |
| 361 | } |
| 362 | } |
| 363 | // A cancelled run removes its container before it frees its slot, and a |
| 364 | // volume cannot be deleted while a container holds it. |
| 365 | for deadline := time.Now().Add(30 * time.Second); time.Now().Before(deadline); { |
| 366 | if !slices.ContainsFunc(ids, r.Active) { |
| 367 | break |
| 368 | } |
| 369 | time.Sleep(100 * time.Millisecond) |
| 370 | } |
| 371 | for _, id := range ids { |
| 372 | os.RemoveAll(r.artifactDir(id)) |
| 373 | } |
| 374 | if _, err := r.PurgeRepoCaches(ctx, repoName); err != nil && !errors.Is(err, errNoSocket) { |
| 375 | log.Printf("[ci] purge caches of %s: %v", repoName, err) |
| 376 | } |
| 377 | } |
| 378 | |
| 379 | func (r *Runner) queryIDs(ctx context.Context, q string, args ...any) ([]int64, error) { |
| 380 | rows, err := r.db.QueryContext(ctx, q, args...) |
| 381 | if err != nil { |
| 382 | return nil, err |
| 383 | } |
| 384 | defer rows.Close() |
| 385 | var ids []int64 |
| 386 | for rows.Next() { |
| 387 | var id int64 |
| 388 | if err := rows.Scan(&id); err != nil { |
| 389 | return nil, err |
| 390 | } |
| 391 | ids = append(ids, id) |
| 392 | } |
| 393 | return ids, rows.Err() |
| 394 | } |
| 395 | |
| 396 | // CancelStaleRuns cleans up runs a restart interrupted. Their containers are |
| 397 | // force-removed and the rows read "cancelled", so they can be retried. |
| 398 | func (r *Runner) CancelStaleRuns(ctx context.Context) error { |
| 399 | ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE status IN ('pending','running','queued')`) |
| 400 | if err != nil { |
| 401 | return err |
| 402 | } |
| 403 | // By name too: containers from before the label existed carry none. |
| 404 | for _, id := range ids { |
| 405 | r.removeContainer(ctx, containerName(id)) |
| 406 | } |
| 407 | if err := r.removeLabelledContainers(ctx); err != nil && !errors.Is(err, errNoSocket) { |
| 408 | log.Printf("[ci] sweep containers: %v", err) |
| 409 | } |
| 410 | now := db.NowISO() |
| 411 | r.execSQL(ctx, `UPDATE ci_runs SET status = 'cancelled', finished_at = ? |
| 412 | WHERE status IN ('pending','running','queued')`, now) |
| 413 | r.execSQL(ctx, `UPDATE ci_steps SET status = 'cancelled', finished_at = ?, log = 'Skipped: run was cancelled' |
| 414 | WHERE status IN ('pending','running')`, now) |
| 415 | return nil |
| 416 | } |
| 417 | |
| 418 | func (r *Runner) artifactDir(runID int64) string { |
| 419 | return filepath.Join(r.cfg.CIArtifactsDir(), strconv.FormatInt(runID, 10)) |
| 420 | } |
| 421 | |
| 422 | // pruneHistory keeps only the CI_MAX_HISTORY most recently finished runs of |
| 423 | // a repository. By finish time, not id: a retried old run counts as recent. |
| 424 | func (r *Runner) pruneHistory(ctx context.Context, repoID int64) { |
| 425 | // Active runs never get pruned: their row is what the queue and the |
| 426 | // finish path update. |
| 427 | ids, err := r.queryIDs(ctx, |
| 428 | `SELECT id FROM ci_runs WHERE repo_id = ? |
| 429 | AND status NOT IN ('pending','running','queued') ORDER BY finished_at DESC, id DESC`, repoID) |
| 430 | if err != nil || len(ids) <= r.cfg.CIMaxHistory { |
| 431 | return |
| 432 | } |
| 433 | toDelete := ids[r.cfg.CIMaxHistory:] |
| 434 | args := make([]any, len(toDelete)) |
| 435 | for i, id := range toDelete { |
| 436 | os.RemoveAll(r.artifactDir(id)) |
| 437 | args[i] = id |
| 438 | } |
| 439 | placeholders := strings.TrimSuffix(strings.Repeat("?,", len(args)), ",") |
| 440 | r.execSQL(ctx, `DELETE FROM ci_runs WHERE id IN (`+placeholders+`)`, args...) |
| 441 | } |
| 442 | |
| 443 | // --- Environment --- |
| 444 | |
| 445 | // buildEnvVars assembles the container environment. Later sources win on a |
| 446 | // name clash: config defaults, then manual overrides, then secrets. |
| 447 | func buildEnvVars(runID int64, repoName, baseURL, host string, run runRow, cfg *Config, |
| 448 | overrides map[string]string, secrets []secret, |
| 449 | ) (envArray, secretValues []string) { |
| 450 | order := []string{} |
| 451 | vars := map[string]string{} |
| 452 | set := func(k, v string) { |
| 453 | if _, seen := vars[k]; !seen { |
| 454 | order = append(order, k) |
| 455 | } |
| 456 | vars[k] = v |
| 457 | } |
| 458 | |
| 459 | refName := run.CommitTag |
| 460 | if refName == "" { |
| 461 | refName = run.CommitBranch |
| 462 | } |
| 463 | shortSha := run.CommitSha |
| 464 | if len(shortSha) > 8 { |
| 465 | shortSha = shortSha[:8] |
| 466 | } |
| 467 | set("CI", "true") |
| 468 | set("CI_PIPELINE_ID", strconv.FormatInt(runID, 10)) |
| 469 | set("CI_REPO_NAME", repoName) |
| 470 | set("CI_SERVER_URL", baseURL) |
| 471 | set("CI_TRIGGER_SOURCE", run.TriggerSource) |
| 472 | set("CI_COMMIT_SHA", run.CommitSha) |
| 473 | set("CI_COMMIT_SHORT_SHA", shortSha) |
| 474 | set("CI_COMMIT_BRANCH", run.CommitBranch) |
| 475 | set("CI_COMMIT_TAG", run.CommitTag) |
| 476 | set("CI_COMMIT_REF_NAME", refName) |
| 477 | // Image names must be lowercase. The registry lowers the repo name too. |
| 478 | set("CI_REGISTRY", host+"/"+strings.ToLower(repoName)) |
| 479 | |
| 480 | for _, name := range cfg.VariableOrder { |
| 481 | if def, ok := cfg.Variables[name]; ok && def.Default != "" { |
| 482 | set(name, def.Default) |
| 483 | } |
| 484 | } |
| 485 | names := make([]string, 0, len(overrides)) |
| 486 | for name := range overrides { |
| 487 | names = append(names, name) |
| 488 | } |
| 489 | sort.Strings(names) |
| 490 | for _, name := range names { |
| 491 | set(name, overrides[name]) |
| 492 | } |
| 493 | for _, s := range secrets { |
| 494 | set(s.name, s.value) |
| 495 | secretValues = append(secretValues, s.value) |
| 496 | } |
| 497 | |
| 498 | for _, k := range order { |
| 499 | envArray = append(envArray, k+"="+vars[k]) |
| 500 | } |
| 501 | return envArray, secretValues |
| 502 | } |
| 503 | |
| 504 | // buildVars maps the run environment without the secrets, for expanding |
| 505 | // build_image tags and args. |
| 506 | func buildVars(envArray []string, secrets []secret) map[string]string { |
| 507 | vars := map[string]string{} |
| 508 | for _, kv := range envArray { |
| 509 | k, v, _ := strings.Cut(kv, "=") |
| 510 | vars[k] = v |
| 511 | } |
| 512 | for _, s := range secrets { |
| 513 | delete(vars, s.name) |
| 514 | } |
| 515 | return vars |
| 516 | } |
| 517 | |
| 518 | // maskSecrets replaces secret values in a log. It is a plain text match: a |
| 519 | // secret printed in encoded or split form is not caught. |
| 520 | func maskSecrets(text string, secrets []string) string { |
| 521 | // Longest first: a secret that contains a shorter one would otherwise |
| 522 | // be left partly visible. |
| 523 | sorted := slices.SortedFunc(slices.Values(secrets), func(a, b string) int { return len(b) - len(a) }) |
| 524 | for _, s := range sorted { |
| 525 | if s != "" { |
| 526 | text = strings.ReplaceAll(text, s, "[MASKED]") |
| 527 | } |
| 528 | } |
| 529 | return text |
| 530 | } |
| 531 | |
| 532 | // maskLog masks a step log. A partial log, or one cut at the size limit, may |
| 533 | // end in the first half of a secret, so a tail that starts one is dropped. |
| 534 | func maskLog(text string, secrets []string, partial bool) string { |
| 535 | body, truncated := strings.CutSuffix(text, logTruncatedNotice) |
| 536 | body = maskSecrets(body, secrets) |
| 537 | if !partial && !truncated { |
| 538 | return body |
| 539 | } |
| 540 | cut := 0 |
| 541 | for _, s := range secrets { |
| 542 | for k := min(len(s)-1, len(body)); k > cut; k-- { |
| 543 | if strings.HasSuffix(body, s[:k]) { |
| 544 | cut = k |
| 545 | break |
| 546 | } |
| 547 | } |
| 548 | } |
| 549 | body = body[:len(body)-cut] |
| 550 | if truncated { |
| 551 | body += logTruncatedNotice |
| 552 | } |
| 553 | return body |
| 554 | } |
| 555 | |
| 556 | type secret struct{ name, value string } |
| 557 | |
| 558 | type runRow struct { |
| 559 | RepoID int64 |
| 560 | TriggerSource string |
| 561 | CommitSha string |
| 562 | CommitBranch string |
| 563 | CommitTag string |
| 564 | Overrides string |
| 565 | } |
| 566 | |
| 567 | func nullString(s string) any { |
| 568 | if s == "" { |
| 569 | return nil |
| 570 | } |
| 571 | return s |
| 572 | } |
| 573 | |
| 574 | func (r *Runner) loadRun(ctx context.Context, runID int64) (runRow, error) { |
| 575 | var row runRow |
| 576 | var sha, branch, tag, overrides sql.NullString |
| 577 | err := r.db.QueryRowContext(ctx, |
| 578 | `SELECT repo_id, trigger_source, commit_sha, commit_branch, commit_tag, variable_overrides |
| 579 | FROM ci_runs WHERE id = ?`, runID). |
| 580 | Scan(&row.RepoID, &row.TriggerSource, &sha, &branch, &tag, &overrides) |
| 581 | if err != nil { |
| 582 | return row, errors.New("run not found") |
| 583 | } |
| 584 | row.CommitSha, row.CommitBranch = sha.String, branch.String |
| 585 | row.CommitTag, row.Overrides = tag.String, overrides.String |
| 586 | return row, nil |
| 587 | } |
| 588 | |
| 589 | // --- Checkout upload --- |
| 590 | |
| 591 | // uploadCheckout puts the commit's tree into the container at destPath. |
| 592 | // |
| 593 | // The checkout deliberately does not happen inside the container: that would |
| 594 | // force every CI image to carry git. It is not a bind mount either, because |
| 595 | // the daemon resolves bind sources in the host filesystem. |
| 596 | // |
| 597 | // `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and no |
| 598 | // entry for the archive root, so the destination keeps the ownership the |
| 599 | // container gave it. No .git reaches the container. |
| 600 | func (r *Runner) uploadCheckout(ctx context.Context, containerID, repoName, commitSha, destPath string) error { |
| 601 | // The archive carries destPath as its prefix and is extracted at /, so |
| 602 | // the engine creates the directories and the image needs no mkdir. |
| 603 | prefix := strings.TrimPrefix(path.Clean(destPath), "/") + "/" |
| 604 | // --end-of-options: commit_sha is unvalidated text from the push. |
| 605 | cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "archive", |
| 606 | "--format=tar", "--prefix="+prefix, "--end-of-options", commitSha) |
| 607 | cmd.Env = gitcmd.Env() |
| 608 | var stderr strings.Builder |
| 609 | cmd.Stderr = &stderr |
| 610 | stdout, err := cmd.StdoutPipe() |
| 611 | if err != nil { |
| 612 | return err |
| 613 | } |
| 614 | if err := cmd.Start(); err != nil { |
| 615 | return err |
| 616 | } |
| 617 | resp, body, putErr := r.putArchive(ctx, containerID, "/", stdout) |
| 618 | // The PUT stopped reading, so git would block writing into a full pipe. |
| 619 | stdout.Close() |
| 620 | waitErr := cmd.Wait() |
| 621 | |
| 622 | if putErr != nil { |
| 623 | return putErr |
| 624 | } |
| 625 | if resp.StatusCode >= 300 { |
| 626 | return fmt.Errorf("failed to upload the checkout: HTTP %d %s", |
| 627 | resp.StatusCode, strings.TrimSpace(string(body))) |
| 628 | } |
| 629 | if waitErr != nil { |
| 630 | return fmt.Errorf("failed to read %s: git archive failed %v %s", |
| 631 | commitSha, waitErr, strings.TrimSpace(stderr.String())) |
| 632 | } |
| 633 | return nil |
| 634 | } |
| 635 | |
| 636 | // --- Artifacts --- |
| 637 | |
| 638 | // collectArtifacts copies this step's published files out of the container. |
| 639 | // Archives are built on the host from the engine's tar stream, so the image |
| 640 | // needs no archive tools. |
| 641 | func (r *Runner) collectArtifacts(ctx context.Context, runID int64, containerID string, step Step) { |
| 642 | dir := r.artifactDir(runID) |
| 643 | if err := os.MkdirAll(dir, 0o755); err != nil { |
| 644 | return |
| 645 | } |
| 646 | |
| 647 | // store streams the container file straight into the artifact directory. |
| 648 | store := func(filename, srcPath string) { |
| 649 | dest := filepath.Join(dir, filename) |
| 650 | // The run directory is shared by every step. A second file of the |
| 651 | // same name would overwrite the first and still add a row. |
| 652 | if _, err := os.Stat(dest); err == nil { |
| 653 | log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath) |
| 654 | return |
| 655 | } |
| 656 | size, err := r.copyFileFromContainer(ctx, containerID, srcPath, dest) |
| 657 | if err != nil { |
| 658 | log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err) |
| 659 | return |
| 660 | } |
| 661 | r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`, |
| 662 | runID, filename, size) |
| 663 | } |
| 664 | |
| 665 | for _, srcPath := range step.PublishFile { |
| 666 | store(path.Base(srcPath), srcPath) |
| 667 | } |
| 668 | |
| 669 | for _, f := range archiveFormats { |
| 670 | for _, srcPath := range f.pick(step) { |
| 671 | filename := path.Base(srcPath) + f.ext |
| 672 | dest := filepath.Join(dir, filename) |
| 673 | if _, err := os.Stat(dest); err == nil { |
| 674 | log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath) |
| 675 | continue |
| 676 | } |
| 677 | size, err := r.storeArchive(ctx, containerID, srcPath, dest, f.ext, r.cfg.CIMaxArtifactBytes) |
| 678 | if err != nil { |
| 679 | log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err) |
| 680 | continue |
| 681 | } |
| 682 | r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`, |
| 683 | runID, filename, size) |
| 684 | } |
| 685 | } |
| 686 | } |
| 687 |