Let a runner name the repositories it builds !141
5 files changed, +150 −6
Layout: unified · split
cmd/gitbay-runner/main.go +11 −1
| @@ -41,6 +41,10 @@ type runner struct { | |||
| 41 | cloneBase string // e.g. ssh://git@gitbay.org | 41 | cloneBase string // e.g. ssh://git@gitbay.org |
| 42 | workdir string | 42 | workdir string |
| 43 | timeout time.Duration | 43 | timeout time.Duration |
| 44 | // repos limits which repositories this runner claims builds for. Empty | ||
| 45 | // means any, which is what a runner on the server itself wants; a runner | ||
| 46 | // somewhere that should not execute every repository's steps names them. | ||
| 47 | repos []string | ||
| 44 | } | 48 | } |
| 45 | 49 | ||
| 46 | func main() { | 50 | func main() { |
| @@ -51,6 +55,7 @@ func main() { | |||
| 51 | workdir = flag.String("workdir", filepath.Join(os.TempDir(), "gitbay-runner"), "build workspace root") | 55 | workdir = flag.String("workdir", filepath.Join(os.TempDir(), "gitbay-runner"), "build workspace root") |
| 52 | poll = flag.Duration("poll", 5*time.Second, "idle poll interval") | 56 | poll = flag.Duration("poll", 5*time.Second, "idle poll interval") |
| 53 | timeout = flag.Duration("timeout", 30*time.Minute, "per-build time limit") | 57 | timeout = flag.Duration("timeout", 30*time.Minute, "per-build time limit") |
| 58 | repos = flag.String("repos", "", "only claim builds for these repositories, comma-separated owner/name (default: any)") | ||
| 54 | once = flag.Bool("once", false, "process at most one build, then exit") | 59 | once = flag.Bool("once", false, "process at most one build, then exit") |
| 55 | version = flag.Bool("version", false, "print the commit this binary was built from, then exit") | 60 | version = flag.Bool("version", false, "print the commit this binary was built from, then exit") |
| 56 | ) | 61 | ) |
| @@ -71,6 +76,11 @@ func main() { | |||
| 71 | if *sshOpts != "" { | 76 | if *sshOpts != "" { |
| 72 | r.sshOpts = strings.Fields(*sshOpts) | 77 | r.sshOpts = strings.Fields(*sshOpts) |
| 73 | } | 78 | } |
| 79 | for _, name := range strings.Split(*repos, ",") { | ||
| 80 | if name = strings.TrimSpace(name); name != "" { | ||
| 81 | r.repos = append(r.repos, name) | ||
| 82 | } | ||
| 83 | } | ||
| 74 | if r.cloneBase == "" { | 84 | if r.cloneBase == "" { |
| 75 | r.cloneBase = "ssh://" + *remote | 85 | r.cloneBase = "ssh://" + *remote |
| 76 | } | 86 | } |
| @@ -94,7 +104,7 @@ func main() { | |||
| 94 | // step claims and executes at most one build. ran reports whether there was | 104 | // step claims and executes at most one build. ran reports whether there was |
| 95 | // one, so the caller knows when to idle. | 105 | // one, so the caller knows when to idle. |
| 96 | func (r *runner) step() (bool, error) { | 106 | func (r *runner) step() (bool, error) { |
| 97 | out, err := r.ssh(nil, "runner", "next", "--json") | 107 | out, err := r.ssh(nil, append([]string{"runner", "next"}, append(r.repos, "--json")...)...) |
| 98 | if err != nil { | 108 | if err != nil { |
| 99 | return false, fmt.Errorf("claiming build: %w (%s)", err, out) | 109 | return false, fmt.Errorf("claiming build: %w (%s)", err, out) |
| 100 | } | 110 | } |
e2e/runner_scope_test.go added +60
| @@ -0,0 +1,60 @@ | |||
| 1 | package e2e | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "os" | ||
| 5 | "path/filepath" | ||
| 6 | "strings" | ||
| 7 | "testing" | ||
| 8 | ) | ||
| 9 | |||
| 10 | // A runner names the repositories it will take builds for. Without that, any | ||
| 11 | // runner claims whatever is next in the global queue, so a runner on a machine | ||
| 12 | // that should only build one project ends up executing every repository's | ||
| 13 | // steps — including those of a repository it has nothing to do with. | ||
| 14 | func TestRunnerNextScopedToRepos(t *testing.T) { | ||
| 15 | inst := startInstance(t) | ||
| 16 | aliceKey := inst.newKey(t, "alice") | ||
| 17 | inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub") | ||
| 18 | runnerKey := inst.newKey(t, "ci") | ||
| 19 | inst.admin(t, "admin", "user", "create", "ci", "--key", runnerKey+".pub", "--admin") | ||
| 20 | |||
| 21 | // Two repositories, each with a build queued. "other" is pushed first, so | ||
| 22 | // an unscoped claim would take it. | ||
| 23 | for _, name := range []string{"other", "site"} { | ||
| 24 | if _, errOut, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/"+name); code != 0 { | ||
| 25 | t.Fatalf("repo create %s: %s", name, errOut) | ||
| 26 | } | ||
| 27 | work := t.TempDir() | ||
| 28 | env := inst.gitEnv(aliceKey) | ||
| 29 | mustGit(t, work, env, "clone", inst.sshURL("alice/"+name), "w") | ||
| 30 | dir := filepath.Join(work, "w") | ||
| 31 | os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755) | ||
| 32 | os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte( | ||
| 33 | "jobs:\n "+name+":\n steps:\n - echo hi\n"), 0o644) | ||
| 34 | mustGit(t, dir, env, "checkout", "-q", "-b", "main") | ||
| 35 | mustGit(t, dir, env, "add", ".") | ||
| 36 | mustGit(t, dir, env, "commit", "-q", "-m", "base") | ||
| 37 | mustGit(t, dir, env, "push", "-q", "origin", "main") | ||
| 38 | } | ||
| 39 | |||
| 40 | // Scoped to alice/site: takes the site build, not the older other build. | ||
| 41 | out, errOut, code := inst.ssh(t, runnerKey, "", "runner", "next", "alice/site", "--json") | ||
| 42 | if code != 0 { | ||
| 43 | t.Fatalf("runner next: %s", errOut) | ||
| 44 | } | ||
| 45 | if !strings.Contains(out, `"repo":"alice/site"`) { | ||
| 46 | t.Fatalf("scoped claim took the wrong repo:\n%s", out) | ||
| 47 | } | ||
| 48 | |||
| 49 | // That scope is now empty, though alice/other is still pending. | ||
| 50 | out, _, code = inst.ssh(t, runnerKey, "", "runner", "next", "alice/site", "--json") | ||
| 51 | if code != 0 || strings.Contains(out, `"repo":`) { | ||
| 52 | t.Fatalf("scoped claim took a build outside its scope:\n%s", out) | ||
| 53 | } | ||
| 54 | |||
| 55 | // An unscoped runner still takes it, so the default is unchanged. | ||
| 56 | out, _, code = inst.ssh(t, runnerKey, "", "runner", "next", "--json") | ||
| 57 | if code != 0 || !strings.Contains(out, `"repo":"alice/other"`) { | ||
| 58 | t.Fatalf("unscoped claim did not take the remaining build:\n%s", out) | ||
| 59 | } | ||
| 60 | } | ||
internal/control/build.go +13 −2
| @@ -55,7 +55,7 @@ func init() { | |||
| 55 | // instance operator's call. | 55 | // instance operator's call. |
| 56 | register(Command{Path: []string{"runner", "next"}, | 56 | register(Command{Path: []string{"runner", "next"}, |
| 57 | Summary: "claim the oldest pending build (runner protocol)", | 57 | Summary: "claim the oldest pending build (runner protocol)", |
| 58 | Usage: "runner next", SSHOnly: true, Run: runRunnerNext}) | 58 | Usage: "runner next [<owner/name>...]", SSHOnly: true, Run: runRunnerNext}) |
| 59 | register(Command{Path: []string{"runner", "log"}, | 59 | register(Command{Path: []string{"runner", "log"}, |
| 60 | Summary: "append a build's log from stdin", | 60 | Summary: "append a build's log from stdin", |
| 61 | Usage: "runner log <build-id>", SSHOnly: true, ReadsStdin: true, Run: runRunnerLog}) | 61 | Usage: "runner log <build-id>", SSHOnly: true, ReadsStdin: true, Run: runRunnerLog}) |
| @@ -329,7 +329,18 @@ func runRunnerNext(c *Ctx, args []string) int { | |||
| 329 | "build abandoned", url, c.User.ID) | 329 | "build abandoned", url, c.User.ID) |
| 330 | } | 330 | } |
| 331 | } | 331 | } |
| 332 | b, ok, err := c.Store.ClaimBuild() | 332 | // A runner may limit itself to named repositories. The operator chooses |
| 333 | // what a given runner executes by how they start it; this is scoping the | ||
| 334 | // runner asks for, not an ACL the server holds over it. | ||
| 335 | var repoIDs []int64 | ||
| 336 | for _, arg := range args { | ||
| 337 | repo, code := resolveRepo(c, arg, policy.CanRead) | ||
| 338 | if code >= 0 { | ||
| 339 | return code | ||
| 340 | } | ||
| 341 | repoIDs = append(repoIDs, repo.ID) | ||
| 342 | } | ||
| 343 | b, ok, err := c.Store.ClaimBuild(repoIDs) | ||
| 333 | if err != nil { | 344 | if err != nil { |
| 334 | return c.fail(protocol.ExitFailure, "%v", err) | 345 | return c.fail(protocol.ExitFailure, "%v", err) |
| 335 | } | 346 | } |
internal/store/builds.go +17 −2
| @@ -3,6 +3,7 @@ package store | |||
| 3 | import ( | 3 | import ( |
| 4 | "database/sql" | 4 | "database/sql" |
| 5 | "errors" | 5 | "errors" |
| 6 | "strings" | ||
| 6 | "time" | 7 | "time" |
| 7 | ) | 8 | ) |
| 8 | 9 | ||
| @@ -59,14 +60,28 @@ func scanBuild(row interface{ Scan(...any) error }) (Build, error) { | |||
| 59 | } | 60 | } |
| 60 | 61 | ||
| 61 | // ClaimBuild atomically hands the oldest pending build to a runner. | 62 | // ClaimBuild atomically hands the oldest pending build to a runner. |
| 62 | func (s *Store) ClaimBuild() (Build, bool, error) { | 63 | // ClaimBuild takes the oldest pending build and marks it running. A |
| 64 | // non-empty repoIDs restricts the claim to those repositories, which is how | ||
| 65 | // a runner on a machine that should not execute every repository's steps | ||
| 66 | // limits what it picks up. | ||
| 67 | func (s *Store) ClaimBuild(repoIDs []int64) (Build, bool, error) { | ||
| 63 | tx, err := s.DB.Begin() | 68 | tx, err := s.DB.Begin() |
| 64 | if err != nil { | 69 | if err != nil { |
| 65 | return Build{}, false, err | 70 | return Build{}, false, err |
| 66 | } | 71 | } |
| 67 | defer tx.Rollback() | 72 | defer tx.Rollback() |
| 73 | query := "SELECT id FROM builds WHERE status = 'pending' ORDER BY id LIMIT 1" | ||
| 74 | args := []any{} | ||
| 75 | if len(repoIDs) > 0 { | ||
| 76 | marks := strings.TrimSuffix(strings.Repeat("?,", len(repoIDs)), ",") | ||
| 77 | query = "SELECT id FROM builds WHERE status = 'pending' AND repo_id IN (" + | ||
| 78 | marks + ") ORDER BY id LIMIT 1" | ||
| 79 | for _, id := range repoIDs { | ||
| 80 | args = append(args, id) | ||
| 81 | } | ||
| 82 | } | ||
| 68 | var id int64 | 83 | var id int64 |
| 69 | err = tx.QueryRow("SELECT id FROM builds WHERE status = 'pending' ORDER BY id LIMIT 1").Scan(&id) | 84 | err = tx.QueryRow(query, args...).Scan(&id) |
| 70 | if errors.Is(err, sql.ErrNoRows) { | 85 | if errors.Is(err, sql.ErrNoRows) { |
| 71 | return Build{}, false, nil | 86 | return Build{}, false, nil |
| 72 | } | 87 | } |
internal/store/builds_test.go +49 −1
| @@ -32,7 +32,7 @@ func TestReapStaleBuilds(t *testing.T) { | |||
| 32 | 32 | ||
| 33 | // Claim both, then age only the first past the deadline. | 33 | // Claim both, then age only the first past the deadline. |
| 34 | for range 2 { | 34 | for range 2 { |
| 35 | if _, ok, err := s.ClaimBuild(); err != nil || !ok { | 35 | if _, ok, err := s.ClaimBuild(nil); err != nil || !ok { |
| 36 | t.Fatalf("claim: %v ok=%v", err, ok) | 36 | t.Fatalf("claim: %v ok=%v", err, ok) |
| 37 | } | 37 | } |
| 38 | } | 38 | } |
| @@ -118,3 +118,51 @@ func TestBuildsForCommitTiming(t *testing.T) { | |||
| 118 | t.Fatalf("unfinished build reported %s", d) | 118 | t.Fatalf("unfinished build reported %s", d) |
| 119 | } | 119 | } |
| 120 | } | 120 | } |
| 121 | |||
| 122 | // A runner that names repositories claims only their builds, so a runner on a | ||
| 123 | // machine that should not execute every repository's steps does not pick one | ||
| 124 | // up by being first to ask. | ||
| 125 | func TestClaimBuildScopedToRepos(t *testing.T) { | ||
| 126 | s := open(t) | ||
| 127 | if err := s.MigrateUp(); err != nil { | ||
| 128 | t.Fatal(err) | ||
| 129 | } | ||
| 130 | uid, err := s.CreateUser("cmc", true) | ||
| 131 | if err != nil { | ||
| 132 | t.Fatal(err) | ||
| 133 | } | ||
| 134 | mine, err := s.CreateRepo("user", uid, "site", "public") | ||
| 135 | if err != nil { | ||
| 136 | t.Fatal(err) | ||
| 137 | } | ||
| 138 | theirs, err := s.CreateRepo("user", uid, "stranger", "public") | ||
| 139 | if err != nil { | ||
| 140 | t.Fatal(err) | ||
| 141 | } | ||
| 142 | // Queued first, so an unscoped claim would take it. | ||
| 143 | if _, err := s.CreateBuild(theirs, "evil", "abc123", "main", `["true"]`); err != nil { | ||
| 144 | t.Fatal(err) | ||
| 145 | } | ||
| 146 | wanted, err := s.CreateBuild(mine, "deploy", "def456", "main", `["true"]`) | ||
| 147 | if err != nil { | ||
| 148 | t.Fatal(err) | ||
| 149 | } | ||
| 150 | |||
| 151 | b, ok, err := s.ClaimBuild([]int64{mine}) | ||
| 152 | if err != nil || !ok { | ||
| 153 | t.Fatalf("claim: %v ok=%v", err, ok) | ||
| 154 | } | ||
| 155 | if b.RepoID != mine || b.Number != wanted { | ||
| 156 | t.Fatalf("claimed repo %d build %d, want repo %d build %d", | ||
| 157 | b.RepoID, b.Number, mine, wanted) | ||
| 158 | } | ||
| 159 | |||
| 160 | // Nothing left for that scope, even though another repo's build is pending. | ||
| 161 | if _, ok, err := s.ClaimBuild([]int64{mine}); err != nil || ok { | ||
| 162 | t.Fatalf("second scoped claim: err=%v ok=%v, want no build", err, ok) | ||
| 163 | } | ||
| 164 | // An unscoped runner still takes it. | ||
| 165 | if b, ok, err := s.ClaimBuild(nil); err != nil || !ok || b.RepoID != theirs { | ||
| 166 | t.Fatalf("unscoped claim: err=%v ok=%v repo=%d", err, ok, b.RepoID) | ||
| 167 | } | ||
| 168 | } | ||