| @@ -1,12 +1,14 @@ |
| 1 | 1 | package e2e |
| 2 | 2 | |
| 3 | 3 | import ( |
| 4 | "bytes" |
| 4 | 5 | "fmt" |
| 5 | 6 | "os" |
| 6 | 7 | "os/exec" |
| 7 | 8 | "path/filepath" |
| 8 | 9 | "strings" |
| 9 | 10 | "testing" |
| 11 | "time" |
| 10 | 12 | ) |
| 11 | 13 | |
| 12 | 14 | func buildRunner(t *testing.T) string { |
| @@ -268,3 +270,107 @@ func TestCI(t *testing.T) { |
| 268 | 270 | t.Fatalf("config failure status missing:\n%s", out) |
| 269 | 271 | } |
| 270 | 272 | } |
| 273 | |
| 274 | // runnerJobs runs a runner with -jobs n until it has nothing left to do, |
| 275 | // and returns its output. Unlike runnerOnce it is not bounded to one |
| 276 | // build, so it is stopped when the queue drains. |
| 277 | func (i *instance) runnerJobs(t *testing.T, key, repo string, jobs int) string { |
| 278 | t.Helper() |
| 279 | opts := fmt.Sprintf("-p %d -i %s -o IdentitiesOnly=yes -o StrictHostKeyChecking=no -o UserKnownHostsFile=%s -o BatchMode=yes", |
| 280 | i.port, key, filepath.Join(i.sshDir, "known_hosts")) |
| 281 | cmd := exec.Command(i.runner, |
| 282 | "-jobs", fmt.Sprint(jobs), |
| 283 | "-poll", "200ms", |
| 284 | "-remote", "git@127.0.0.1", |
| 285 | "-ssh-opts", opts, |
| 286 | "-clone-base", fmt.Sprintf("ssh://git@127.0.0.1:%d", i.port), |
| 287 | "-workdir", t.TempDir()) |
| 288 | cmd.Env = append(os.Environ(), "GIT_CONFIG_NOSYSTEM=1", "GIT_CONFIG_GLOBAL=/dev/null") |
| 289 | var buf bytes.Buffer |
| 290 | cmd.Stdout, cmd.Stderr = &buf, &buf |
| 291 | if err := cmd.Start(); err != nil { |
| 292 | t.Fatal(err) |
| 293 | } |
| 294 | defer func() { cmd.Process.Kill(); cmd.Wait() }() |
| 295 | |
| 296 | // Wait for every queued build to leave the pending and running states. |
| 297 | deadline := time.Now().Add(60 * time.Second) |
| 298 | for time.Now().Before(deadline) { |
| 299 | out, _, code := i.ssh(t, key, "", "build", "list", repo, "--json") |
| 300 | if code == 0 && !strings.Contains(out, `"status":"pending"`) && |
| 301 | !strings.Contains(out, `"status":"running"`) { |
| 302 | break |
| 303 | } |
| 304 | time.Sleep(200 * time.Millisecond) |
| 305 | } |
| 306 | return buf.String() |
| 307 | } |
| 308 | |
| 309 | // -jobs N runs N builds at once. ClaimBuild has always been a single |
| 310 | // transaction that selects and updates, so several workers claiming |
| 311 | // together is safe; the runner simply never used more than one (#115). |
| 312 | func TestRunnerConcurrentJobs(t *testing.T) { |
| 313 | inst := startInstance(t) |
| 314 | inst.runner = buildRunner(t) |
| 315 | aliceKey := inst.newKey(t, "alice") |
| 316 | inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub", "--admin") |
| 317 | if _, errOut, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/app"); code != 0 { |
| 318 | t.Fatalf("repo create: %s", errOut) |
| 319 | } |
| 320 | |
| 321 | // Four jobs, each sleeping longer than the poll interval, so serial |
| 322 | // execution and concurrent execution are distinguishable. |
| 323 | ci := "jobs:\n" |
| 324 | for _, name := range []string{"one", "two", "three", "four"} { |
| 325 | ci += fmt.Sprintf(" %s:\n steps:\n - sleep 1\n - echo done-%s\n", name, name) |
| 326 | } |
| 327 | env := inst.gitEnv(aliceKey) |
| 328 | work := t.TempDir() |
| 329 | mustGit(t, work, env, "clone", inst.sshURL("alice/app"), "w") |
| 330 | dir := filepath.Join(work, "w") |
| 331 | os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755) |
| 332 | os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte(ci), 0o644) |
| 333 | mustGit(t, dir, env, "checkout", "-q", "-b", "main") |
| 334 | mustGit(t, dir, env, "add", ".") |
| 335 | mustGit(t, dir, env, "commit", "-q", "-m", "ci") |
| 336 | mustGit(t, dir, env, "push", "-q", "origin", "main") |
| 337 | |
| 338 | out := inst.runnerJobs(t, aliceKey, "alice/app", 4) |
| 339 | |
| 340 | listing, _, code := inst.ssh(t, aliceKey, "", "build", "list", "alice/app", "--json") |
| 341 | if code != 0 { |
| 342 | t.Fatalf("build list failed:\n%s", out) |
| 343 | } |
| 344 | for _, name := range []string{"one", "two", "three", "four"} { |
| 345 | if !strings.Contains(listing, `"job":"`+name+`"`) { |
| 346 | t.Fatalf("%s never ran:\n%s\n%s", name, listing, out) |
| 347 | } |
| 348 | } |
| 349 | if strings.Contains(listing, `"status":"pending"`) || strings.Contains(listing, `"status":"running"`) { |
| 350 | t.Fatalf("builds did not finish:\n%s", listing) |
| 351 | } |
| 352 | // Overlap, asserted from the runner's own log rather than from wall |
| 353 | // clock: elapsed time also covers four concurrent clones and the |
| 354 | // polling this test does, and would make a slow machine look serial. |
| 355 | // Every build announces itself when it starts and again when it |
| 356 | // finishes, so concurrency is "a build started before the first one |
| 357 | // finished". |
| 358 | started, firstFinish := 0, -1 |
| 359 | for _, line := range strings.Split(out, "\n") { |
| 360 | switch { |
| 361 | case strings.Contains(line, "alice/app"): |
| 362 | started++ |
| 363 | case strings.Contains(line, ": success"), strings.Contains(line, ": failure"): |
| 364 | if firstFinish < 0 { |
| 365 | firstFinish = started |
| 366 | } |
| 367 | } |
| 368 | } |
| 369 | if started != 4 { |
| 370 | t.Fatalf("%d builds started, want 4:\n%s", started, out) |
| 371 | } |
| 372 | if firstFinish < 2 { |
| 373 | t.Errorf("only %d build(s) had started when the first finished; -jobs 4 ran them serially\n%s", |
| 374 | firstFinish, out) |
| 375 | } |
| 376 | } |