runner: every phase of a build is cancellable !170

merged merged by cmc on 2026-09-02 04:25 UTC · krz/gitbay:cancel-robust into main

2 files changed, +47 −29

Layout: unified · split

cmd/gitbay-runner/main.go +42 −26
@@ -186,12 +186,45 @@ func (r *runner) run(j job) bool {
186 logExited := make(chan struct{}) 186 logExited := make(chan struct{})
187 go func() { 187 go func() {
188 defer close(logExited) 188 defer close(logExited)
189 if err := logCmd.Wait(); err != nil { 189 err := logCmd.Wait()
190 if ee, ok := err.(*exec.ExitError); ok && ee.ExitCode() == 3 { 190 if ee, ok := err.(*exec.ExitError); ok && ee.ExitCode() == 3 {
191 close(cancelled) 191 close(cancelled)
192 } 192 return
193 }
194 if err != nil {
195 log.Printf("build %d: log session ended: %v", j.ID, err)
193 } 196 }
194 }() 197 }()
198 // runStep starts cmd and waits for it, the cancel signal, or the
199 // deadline. Every phase goes through it, so a cancel during the clone
200 // lands as fast as one during a step.
201 runStep := func(cmd *exec.Cmd, deadline time.Time) (bool, string) {
202 select {
203 case <-cancelled:
204 return false, "cancelled"
205 default:
206 }
207 if err := cmd.Start(); err != nil {
208 return false, fmt.Sprintf("start: %v", err)
209 }
210 done := make(chan error, 1)
211 go func() { done <- cmd.Wait() }()
212 select {
213 case err := <-done:
214 if err != nil {
215 return false, fmt.Sprintf("step failed: %v", err)
216 }
217 return true, ""
218 case <-cancelled:
219 cmd.Process.Kill()
220 <-done
221 return false, "cancelled"
222 case <-time.After(time.Until(deadline)):
223 cmd.Process.Kill()
224 <-done
225 return false, fmt.Sprintf("build timed out after %s", r.timeout)
226 }
227 }
195 defer func() { 228 defer func() {
196 select { 229 select {
197 case <-cancelled: 230 case <-cancelled:
@@ -207,6 +240,7 @@ func (r *runner) run(j job) bool {
207 240
208 gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " ")) 241 gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " "))
209 cloneURL := r.cloneBase + "/" + j.Repo + ".git" 242 cloneURL := r.cloneBase + "/" + j.Repo + ".git"
243 deadline := time.Now().Add(r.timeout)
210 fmt.Fprintf(sink, "$ git clone %s (%.10s)\n", cloneURL, j.SHA) 244 fmt.Fprintf(sink, "$ git clone %s (%.10s)\n", cloneURL, j.SHA)
211 for _, args := range [][]string{ 245 for _, args := range [][]string{
212 {"clone", "-q", cloneURL, dir}, 246 {"clone", "-q", cloneURL, dir},
@@ -215,13 +249,12 @@ func (r *runner) run(j job) bool {
215 cmd := exec.Command("git", args...) 249 cmd := exec.Command("git", args...)
216 cmd.Env = append(os.Environ(), "GIT_SSH_COMMAND="+gitSSH, "GIT_TERMINAL_PROMPT=0") 250 cmd.Env = append(os.Environ(), "GIT_SSH_COMMAND="+gitSSH, "GIT_TERMINAL_PROMPT=0")
217 cmd.Stdout, cmd.Stderr = sink, sink 251 cmd.Stdout, cmd.Stderr = sink, sink
218 if err := cmd.Run(); err != nil { 252 if ok, why := runStep(cmd, deadline); !ok {
219 fmt.Fprintf(sink, "git %s: %v\n", args[0], err) 253 fmt.Fprintf(sink, "git %s: %s\n", args[0], why)
220 return false 254 return false
221 } 255 }
222 } 256 }
223 257
224 deadline := time.Now().Add(r.timeout)
225 for _, step := range j.Steps { 258 for _, step := range j.Steps {
226 fmt.Fprintf(sink, "$ %s\n", step) 259 fmt.Fprintf(sink, "$ %s\n", step)
227 cmd := exec.Command("sh", "-c", step) 260 cmd := exec.Command("sh", "-c", step)
@@ -232,25 +265,8 @@ func (r *runner) run(j job) bool {
232 cmd.Env = append(cmd.Env, name+"="+value) 265 cmd.Env = append(cmd.Env, name+"="+value)
233 } 266 }
234 cmd.Stdout, cmd.Stderr = sink, sink 267 cmd.Stdout, cmd.Stderr = sink, sink
235 if err := cmd.Start(); err != nil { 268 if ok, why := runStep(cmd, deadline); !ok {
236 fmt.Fprintf(sink, "start: %v\n", err) 269 fmt.Fprintf(sink, "%s\n", why)
237 return false
238 }
239 done := make(chan error, 1)
240 go func() { done <- cmd.Wait() }()
241 select {
242 case err := <-done:
243 if err != nil {
244 fmt.Fprintf(sink, "step failed: %v\n", err)
245 return false
246 }
247 case <-cancelled:
248 cmd.Process.Kill()
249 <-done
250 return false
251 case <-time.After(time.Until(deadline)):
252 cmd.Process.Kill()
253 fmt.Fprintf(sink, "build timed out after %s\n", r.timeout)
254 return false 270 return false
255 } 271 }
256 } 272 }
e2e/build_cancel_test.go +5 −3
@@ -144,10 +144,12 @@ func TestBuildCancelRunning(t *testing.T) {
144 } 144 }
145 select { 145 select {
146 case <-exited: 146 case <-exited:
147 case <-time.After(20 * time.Second): 147 case <-time.After(60 * time.Second):
148 t.Fatalf("runner still running 20s after cancel:\n%s", runnerOut.String()) 148 t.Fatalf("runner still running 60s after cancel:\n%s", runnerOut.String())
149 } 149 }
150 if took := time.Since(started); took > 15*time.Second { 150 // Well inside the step's 120s sleep: the runner stopped because it
151 // was told to, not because the step ended.
152 if took := time.Since(started); took > 45*time.Second {
151 t.Fatalf("runner took %s to stop", took) 153 t.Fatalf("runner took %s to stop", took)
152 } 154 }
153 if !strings.Contains(runnerOut.String(), "cancelled") { 155 if !strings.Contains(runnerOut.String(), "cancelled") {