package agent import ( "encoding/json" "errors" "fmt" "net/http" "net/http/httptest" "strings" "sync/atomic" "testing" "time" ) func base() PRState { return PRState{ Ref: PRRef{Owner: "unkin", Repo: "repo", Number: 1}, State: "open", Merged: false, HeadSHA: "abc123", BaseSHA: "base000", Mergeable: MergeYes, CIStatus: "pending", NonAgentComments: 0, } } func TestMeaningfulChange(t *testing.T) { tests := []struct { name string mutatePrev func(s *PRState) mutate func(s *PRState) wantChange bool }{ { name: "no change", mutate: func(s *PRState) {}, wantChange: false, }, { name: "CI pending to success is benign", mutate: func(s *PRState) { s.CIStatus = "success" }, wantChange: false, }, { name: "open to merged alerts", mutate: func(s *PRState) { s.Merged = true; s.State = "closed" }, wantChange: true, }, { name: "open to closed without merge alerts", mutate: func(s *PRState) { s.State = "closed" }, wantChange: true, }, { name: "new non-agent comment alerts", mutate: func(s *PRState) { s.NonAgentComments = 1 }, wantChange: true, }, { name: "CI to failure alerts", mutate: func(s *PRState) { s.CIStatus = "failure" }, wantChange: true, }, { name: "CI to error alerts", mutate: func(s *PRState) { s.CIStatus = "error" }, wantChange: true, }, { // Mergeability takes a run of observations, so no pair of snapshots // decides it here; prWatch owns that rule. name: "mergeable false pair alone is not a pairwise change", mutatePrev: func(s *PRState) { s.Mergeable = MergeNo }, mutate: func(s *PRState) { s.Mergeable = MergeNo }, wantChange: false, }, { name: "new head sha alone is benign", mutate: func(s *PRState) { s.HeadSHA = "def456" }, wantChange: false, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { prev := base() if tt.mutatePrev != nil { tt.mutatePrev(&prev) } cur := base() tt.mutate(&cur) got, reason := MeaningfulChange(prev, cur) if got != tt.wantChange { t.Errorf("MeaningfulChange() = %v (%q), want %v", got, reason, tt.wantChange) } if got && reason == "" { t.Errorf("change reported without a reason") } }) } } // A comment that only the agent posts must not alert: the non-agent count is // unchanged, so MeaningfulChange sees nothing. func TestMeaningfulChangeAgentCommentIgnored(t *testing.T) { prev := base() cur := base() // agent commented, but NonAgentComments stayed 0 if got, _ := MeaningfulChange(prev, cur); got { t.Errorf("agent-only comment should not alert") } } // Once CI is already failing, staying failed must not re-alert. func TestMeaningfulChangeStaysFailed(t *testing.T) { prev := base() prev.CIStatus = "failure" cur := base() cur.CIStatus = "failure" if got, _ := MeaningfulChange(prev, cur); got { t.Errorf("CI staying failed should not re-alert") } } // fakeFetcher returns a scripted sequence of (state, error) results per call, // so tests can drive Watch across baseline and successive polls. type fakeFetcher struct { states []PRState errs []error calls int } func (f *fakeFetcher) FetchState(ref PRRef, agentLogin string) (PRState, error) { i := f.calls if i >= len(f.states) { i = len(f.states) - 1 } f.calls++ var err error if f.calls-1 < len(f.errs) { err = f.errs[f.calls-1] } return f.states[i], err } func TestTerminalState(t *testing.T) { open := base() open.State = "open" if term, _ := terminalState(open); term { t.Errorf("open PR should not be terminal") } merged := base() merged.State = "closed" merged.Merged = true if term, reason := terminalState(merged); !term || reason != "PR merged" { t.Errorf("merged PR: got (%v, %q), want (true, %q)", term, reason, "PR merged") } closed := base() closed.State = "closed" if term, reason := terminalState(closed); !term || reason != "PR closed without merging" { t.Errorf("closed PR: got (%v, %q), want (true, %q)", term, reason, "PR closed without merging") } } // The production hang: a PR that is already merged when watchpr starts must be // reported at baseline and exit, without ever consuming a tick. Before the fix, // Watch only reported transitions, so a terminal baseline was polled forever. func TestWatchExitsWhenAlreadyMergedAtBaseline(t *testing.T) { merged := base() merged.State = "closed" merged.Merged = true f := &fakeFetcher{states: []PRState{merged}} ticks := make(chan time.Time) // never fires; a hang would block here res, err := Watch(f, []PRRef{merged.Ref}, "unkin-agent", ticks, nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", res.Reason, "PR merged") } if f.calls != 1 { t.Errorf("fetch calls = %d, want 1 (baseline only)", f.calls) } } // A PR already closed-without-merge at baseline must also exit immediately. func TestWatchExitsWhenAlreadyClosedAtBaseline(t *testing.T) { closed := base() closed.State = "closed" f := &fakeFetcher{states: []PRState{closed}} ticks := make(chan time.Time) res, err := Watch(f, []PRRef{closed.Ref}, "unkin-agent", ticks, nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR closed without merging" { t.Errorf("reason = %q, want %q", res.Reason, "PR closed without merging") } } // An open→merged transition observed during polling must be detected and end // the watch. func TestWatchDetectsMergeAfterBaseline(t *testing.T) { open := base() merged := base() merged.State = "closed" merged.Merged = true f := &fakeFetcher{states: []PRState{open, merged}} // baseline open, then merged var baselines []PRState ticks := make(chan time.Time, 1) ticks <- time.Now() res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, func(sts []PRState) { baselines = sts }, nil) if err != nil { t.Fatalf("Watch: %v", err) } if len(baselines) != 1 || baselines[0].Ref != open.Ref { t.Errorf("onBaseline received %v, want the one open baseline", baselines) } if res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", res.Reason, "PR merged") } } // A transient poll error must be reported and the loop must keep polling; a // merge on the following tick still ends the watch. func TestWatchContinuesPastPollError(t *testing.T) { open := base() merged := base() merged.State = "closed" merged.Merged = true // baseline ok, first poll errors, second poll sees the merge. f := &fakeFetcher{ states: []PRState{open, open, merged}, errs: []error{nil, errors.New("HTTP 502"), nil}, } var gotErr error ticks := make(chan time.Time, 2) ticks <- time.Now() ticks <- time.Now() res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, nil, func(_ PRRef, e error) { gotErr = e }) if err != nil { t.Fatalf("Watch: %v", err) } if gotErr == nil { t.Errorf("onError should have received the transient poll error") } if res.Reason != "PR merged" { t.Errorf("reason = %q, want %q (loop must survive the error)", res.Reason, "PR merged") } } // A baseline fetch error aborts the watch (nothing to establish a baseline // from), unlike a mid-loop poll error. func TestWatchBaselineErrorAborts(t *testing.T) { f := &fakeFetcher{states: []PRState{base()}, errs: []error{errors.New("HTTP 500")}} ticks := make(chan time.Time) if _, err := Watch(f, []PRRef{base().Ref}, "unkin-agent", ticks, nil, nil); err == nil { t.Fatal("Watch should return the baseline fetch error") } } // The production hang, end to end: a watched PR stays open across several polls, // then is squash-merged and its branch deleted, so the commit-status endpoint // 404s. Driven through a real *GiteaClient, the watch loop must still detect the // merge on the poll it happens. Before the fix, FetchState returned an error on // that poll (the 404 masked the merge), so the loop reported only poll errors // and never exited -- exactly the 37-minute hang seen in production. func TestWatchDetectsMergeWhenCommitGone(t *testing.T) { const sha = "cafebabecafebabe" var polls atomic.Int32 // number of PR fetches so far mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { n := polls.Add(1) if n >= 4 { // baseline + two unchanged polls, then merged _, _ = fmt.Fprintf(w, `{"number":7,"state":"closed","merged":true,"mergeable":true,"head":{"sha":%q}}`, sha) return } _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { if polls.Load() >= 4 { // branch deleted post-merge: commit is gone w.WriteHeader(http.StatusNotFound) _, _ = fmt.Fprint(w, `{"message":"not found"}`) return } _, _ = fmt.Fprint(w, `{"state":"success"}`) }) mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} // A real ticker so the loop advances on its own; a hang (the bug) is caught // by the timeout below instead of blocking the suite. tk := time.NewTicker(5 * time.Millisecond) defer tk.Stop() var pollErr atomic.Pointer[error] type outcome struct { res WatchResult err error } done := make(chan outcome, 1) go func() { res, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(_ PRRef, e error) { pollErr.Store(&e) }) done <- outcome{res, err} }() select { case o := <-done: if o.err != nil { t.Fatalf("Watch: %v", o.err) } if o.res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", o.res.Reason, "PR merged") } if p := pollErr.Load(); p != nil { t.Errorf("no poll error expected once a 404 status is tolerated, got: %v", *p) } case <-time.After(3 * time.Second): var got error if p := pollErr.Load(); p != nil { got = *p } t.Fatalf("Watch hung: a merge with a gone head commit was never detected (last poll error: %v)", got) } } func TestCountNonAgentComments(t *testing.T) { comments := []Comment{ {User: User{Login: "unkin-agent"}}, {User: User{Login: "ben"}}, {User: User{Login: "unkin-agent"}}, {User: User{Login: "reviewer"}}, } if n := countNonAgentComments(comments, "unkin-agent"); n != 2 { t.Errorf("countNonAgentComments = %d, want 2", n) } } // The production failure: the Vault-minted token expired mid-watch and every // poll 401'd, which the loop logged as a warning and polled past forever. An // auth error that survived the client's re-mint must end the watch with an // error so watchpr exits non-zero instead of watching blind. func TestWatchAbortsOnAuthError(t *testing.T) { open := base() merged := base() merged.State = "closed" merged.Merged = true f := &fakeFetcher{ states: []PRState{open, open, merged}, errs: []error{nil, &APIError{Method: "GET", Path: "/p", StatusCode: 401, Body: "invalid token"}, nil}, } warned := 0 ticks := make(chan time.Time, 2) ticks <- time.Now() ticks <- time.Now() _, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, nil, func(PRRef, error) { warned++ }) if err == nil { t.Fatal("Watch should return the auth failure, not keep polling") } if !IsAuthError(err) { t.Errorf("Watch error = %v, want an auth error", err) } if warned != 0 { t.Errorf("auth failure was logged as a warning %d time(s); it must abort", warned) } if f.calls != 2 { t.Errorf("fetch calls = %d, want 2 (baseline + the failing poll)", f.calls) } } // A 5xx keeps its retry behaviour: warn and poll on. func TestWatchContinuesPastServerError(t *testing.T) { open := base() merged := base() merged.State = "closed" merged.Merged = true f := &fakeFetcher{ states: []PRState{open, open, merged}, errs: []error{nil, &APIError{Method: "GET", Path: "/p", StatusCode: 502, Body: "bad gateway"}, nil}, } warned := 0 ticks := make(chan time.Time, 2) ticks <- time.Now() ticks <- time.Now() res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, nil, func(PRRef, error) { warned++ }) if err != nil { t.Fatalf("Watch: %v", err) } if warned != 1 { t.Errorf("warnings = %d, want 1", warned) } if res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", res.Reason, "PR merged") } } // The production failure: a watched repo was renamed mid-watch, so every poll // 404'd (Gitea hides a repo the caller may not see rather than 403ing) and the // loop warned past it forever while reporting nothing. A 404 on a tracked PR // must end the watch with an error naming that PR. func TestWatchAbortsOnMidRunNotFound(t *testing.T) { const sha = "deadbeefdeadbeef" var polls atomic.Int32 mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { if polls.Add(1) > 1 { // repo renamed/made private after the baseline w.WriteHeader(http.StatusNotFound) _, _ = fmt.Fprint(w, `{"message":"Not Found"}`) return } _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `{"state":"success"}`) }) mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} tk := time.NewTicker(5 * time.Millisecond) defer tk.Stop() var warned atomic.Int32 done := make(chan error, 1) go func() { _, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(PRRef, error) { warned.Add(1) }) done <- err }() select { case err := <-done: if err == nil { t.Fatal("Watch should abort on a mid-run 404, not keep polling") } if !IsNotFound(err) { t.Errorf("Watch error = %v, want a 404", err) } if !strings.Contains(err.Error(), ref.String()) { t.Errorf("Watch error = %v, want it to name %s", err, ref.String()) } if n := warned.Load(); n != 0 { t.Errorf("404 was logged as a warning %d time(s); it must abort", n) } case <-time.After(3 * time.Second): t.Fatal("Watch hung: a vanished repo was warned past instead of aborting") } } // A 404 from a sub-resource is not proof the PR is gone: an ingress can serve // one during a Gitea rolling restart. Only the PR lookup itself is authoritative, // so a comments 404 must warn and keep polling like any other transient failure, // and still catch the merge that lands afterwards. func TestWatchSurvivesCommentsNotFound(t *testing.T) { const sha = "0badc0de0badc0de" var polls atomic.Int32 mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { if polls.Add(1) >= 4 { _, _ = fmt.Fprintf(w, `{"number":7,"state":"closed","merged":true,"mergeable":true,"head":{"sha":%q}}`, sha) return } _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `{"state":"success"}`) }) mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { if n := polls.Load(); n == 2 || n == 3 { // proxy blip across two polls w.WriteHeader(http.StatusNotFound) _, _ = fmt.Fprint(w, `{"message":"Not Found"}`) return } _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} tk := time.NewTicker(5 * time.Millisecond) defer tk.Stop() var warned atomic.Int32 type outcome struct { res WatchResult err error } done := make(chan outcome, 1) go func() { res, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(PRRef, error) { warned.Add(1) }) done <- outcome{res, err} }() select { case o := <-done: if o.err != nil { t.Fatalf("Watch: %v (a comments 404 must not be terminal)", o.err) } if o.res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", o.res.Reason, "PR merged") } if n := warned.Load(); n != 2 { t.Errorf("warnings = %d, want 2", n) } case <-time.After(3 * time.Second): t.Fatal("Watch hung: a comments 404 must warn and keep polling") } } // A comments 404 costs a poll from the same budget as any other failure: it must // not be free, and a permanently 404ing sub-resource must still end the watch. func TestWatchCommentsNotFoundCountsTowardCap(t *testing.T) { const sha = "1badc0de1badc0de" mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `{"state":"success"}`) }) var comments atomic.Int32 mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { if comments.Add(1) > 1 { // healthy at baseline, gone from the first poll on w.WriteHeader(http.StatusNotFound) _, _ = fmt.Fprint(w, `{"message":"Not Found"}`) return } _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} tk := time.NewTicker(time.Millisecond) defer tk.Stop() var warned atomic.Int32 done := make(chan error, 1) go func() { _, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(PRRef, error) { warned.Add(1) }) done <- err }() select { case err := <-done: if err == nil { t.Fatal("Watch should give up once the comments 404 stops being transient") } if !strings.Contains(err.Error(), "consecutive failures") { t.Errorf("Watch error = %v, want it to report the failure cap", err) } if n := warned.Load(); n != MaxPollFailures { t.Errorf("warnings = %d, want %d", n, MaxPollFailures) } case <-time.After(3 * time.Second): t.Fatal("Watch hung: a permanently 404ing comments endpoint must hit the cap") } } // The PR lookup is the call whose 404 means the PR is gone, so it aborts on the // very first occurrence rather than spending the failure budget. func TestWatchAbortsOnFirstPRLookupNotFound(t *testing.T) { const sha = "2badc0de2badc0de" var polls atomic.Int32 mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { if polls.Add(1) > 1 { w.WriteHeader(http.StatusNotFound) _, _ = fmt.Fprint(w, `{"message":"Not Found"}`) return } _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `{"state":"success"}`) }) mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} tk := time.NewTicker(5 * time.Millisecond) defer tk.Stop() done := make(chan error, 1) go func() { _, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(PRRef, error) { t.Errorf("a PR-lookup 404 must abort, not warn") }) done <- err }() select { case err := <-done: if !IsNotFound(err) { t.Fatalf("Watch error = %v, want a 404", err) } if n := polls.Load(); n != 2 { t.Errorf("PR fetches = %d, want 2 (baseline + the 404 that aborts)", n) } case <-time.After(3 * time.Second): t.Fatal("Watch hung: a vanished PR must abort") } } // A 5xx blip must not kill a long watch: it warns, keeps polling, and still // catches the merge that lands afterwards. func TestWatchSurvivesTransientServerError(t *testing.T) { const sha = "feedfacefeedface" var polls atomic.Int32 mux := http.NewServeMux() mux.HandleFunc("/api/v1/repos/unkin/repo/pulls/7", func(w http.ResponseWriter, r *http.Request) { switch n := polls.Add(1); { case n == 2 || n == 3: // gateway blip across two polls w.WriteHeader(http.StatusBadGateway) _, _ = fmt.Fprint(w, `bad gateway`) case n >= 4: _, _ = fmt.Fprintf(w, `{"number":7,"state":"closed","merged":true,"mergeable":true,"head":{"sha":%q}}`, sha) default: _, _ = fmt.Fprintf(w, `{"number":7,"state":"open","merged":false,"mergeable":true,"head":{"sha":%q}}`, sha) } }) mux.HandleFunc("/api/v1/repos/unkin/repo/commits/"+sha+"/status", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `{"state":"success"}`) }) mux.HandleFunc("/api/v1/repos/unkin/repo/issues/7/comments", func(w http.ResponseWriter, r *http.Request) { _, _ = fmt.Fprint(w, `[]`) }) srv := httptest.NewServer(mux) defer srv.Close() c := &GiteaClient{BaseURL: srv.URL, Token: "t", HTTP: srv.Client()} ref := PRRef{Owner: "unkin", Repo: "repo", Number: 7} tk := time.NewTicker(5 * time.Millisecond) defer tk.Stop() var warned atomic.Int32 type outcome struct { res WatchResult err error } done := make(chan outcome, 1) go func() { res, err := Watch(c, []PRRef{ref}, "unkin-agent", tk.C, nil, func(PRRef, error) { warned.Add(1) }) done <- outcome{res, err} }() select { case o := <-done: if o.err != nil { t.Fatalf("Watch: %v", o.err) } if o.res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", o.res.Reason, "PR merged") } if n := warned.Load(); n != 2 { t.Errorf("warnings = %d, want 2", n) } case <-time.After(3 * time.Second): t.Fatal("Watch hung: a transient 5xx must not stop the watch") } } // pollFailures scripts a fetcher whose polls fail with a 502 at the given call // indexes (0 is the baseline); the final call returns merged. func pollFailures(calls int, failAt map[int]bool) *fakeFetcher { open, merged := base(), base() merged.State = "closed" merged.Merged = true f := &fakeFetcher{states: make([]PRState, calls), errs: make([]error, calls)} for i := range calls { f.states[i] = open if failAt[i] { f.errs[i] = &APIError{Method: "GET", Path: "/p", StatusCode: 502, Body: "bad gateway"} } } f.states[calls-1] = merged return f } // A permanently wedged endpoint (5xx forever) must eventually give up instead of // warning on every tick for the life of the process. func TestWatchAbortsAfterConsecutiveFailures(t *testing.T) { failAt := map[int]bool{} for i := 1; i <= MaxPollFailures; i++ { failAt[i] = true } f := pollFailures(MaxPollFailures+1, failAt) warned := 0 ticks := make(chan time.Time, MaxPollFailures) for range MaxPollFailures { ticks <- time.Now() } close(ticks) _, err := Watch(f, []PRRef{base().Ref}, "unkin-agent", ticks, nil, func(PRRef, error) { warned++ }) if err == nil { t.Fatal("Watch should give up once the failures stop being transient") } if !strings.Contains(err.Error(), "consecutive failures") { t.Errorf("Watch error = %v, want it to report the failure cap", err) } if warned != MaxPollFailures { t.Errorf("warnings = %d, want %d", warned, MaxPollFailures) } if f.calls != MaxPollFailures+1 { t.Errorf("fetch calls = %d, want %d", f.calls, MaxPollFailures+1) } } // The cap counts consecutive failures only: a single successful poll clears it, // so an intermittent endpoint is watched indefinitely and the merge is caught. func TestWatchFailureCountResetsOnSuccess(t *testing.T) { const runs = MaxPollFailures - 1 failAt := map[int]bool{} for i := 1; i <= runs; i++ { // first run of failures failAt[i] = true } for i := runs + 2; i <= 2*runs+1; i++ { // second run, after one good poll failAt[i] = true } f := pollFailures(2*runs+3, failAt) ticks := make(chan time.Time, 2*runs+2) for range 2*runs + 2 { ticks <- time.Now() } close(ticks) res, err := Watch(f, []PRRef{base().Ref}, "unkin-agent", ticks, nil, func(PRRef, error) {}) if err != nil { t.Fatalf("Watch: %v (a successful poll must reset the failure count)", err) } if res.Reason != "PR merged" { t.Errorf("reason = %q, want %q", res.Reason, "PR merged") } } // Anonymous watching of a public repo must poll on without a credential in // sight: no token, no mint, no exit until something actually changes. func TestWatchAnonymousKeepsPolling(t *testing.T) { open := base() f := &fakeFetcher{states: []PRState{open}} ticks := make(chan time.Time, 2) ticks <- time.Now() ticks <- time.Now() close(ticks) res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, nil, func(_ PRRef, e error) { t.Errorf("unexpected poll error: %v", e) }) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "" { t.Errorf("reason = %q, want no change reported", res.Reason) } if f.calls != 3 { t.Errorf("fetch calls = %d, want 3 (baseline + two polls)", f.calls) } } // drainableTicks returns a channel holding n ticks and already closed, so Watch // polls exactly n times and then returns instead of blocking. The ticks are // spaced far past conflictWindow, so a run of non-mergeable polls confirms on // its second observation; spacedTicks drives the intervals where it must not. func drainableTicks(n int) <-chan time.Time { return spacedTicks(n, 10*time.Minute) } // spacedTicks is drainableTicks with the poll interval named. The tick carries // the time Watch measures the conflict window on, which makes every wall-clock // assertion in these tests exact and instant. func spacedTicks(n int, interval time.Duration) <-chan time.Time { start := time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC) ticks := make(chan time.Time, n) for i := 1; i <= n; i++ { ticks <- start.Add(time.Duration(i) * interval) } close(ticks) return ticks } // The production bug: a PR that was already conflicted (and already CI-failing) // when watching began must not be reported as having just changed. That state is // what the watcher is waiting to see resolved, so the loop keeps polling. func TestWatchIgnoresBaselineConflictAndFailure(t *testing.T) { stuck := base() stuck.Mergeable = MergeNo stuck.CIStatus = "failure" f := &fakeFetcher{states: []PRState{stuck}} res, err := Watch(f, []PRRef{stuck.Ref}, "unkin-agent", drainableTicks(5), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "" { t.Fatalf("Watch ended with %q; a conflict/failure predating the watch is not a change", res.Reason) } if f.calls != 6 { t.Errorf("fetch calls = %d, want 6 (baseline + 5 polls)", f.calls) } } // Mergeability lost after the baseline still alerts, on the second consecutive // conflicted poll. func TestWatchDetectsConflictAfterBaseline(t *testing.T) { ok := base() conflicted := base() conflicted.Mergeable = MergeNo f := &fakeFetcher{states: []PRState{ok, conflicted, conflicted}} res, err := Watch(f, []PRRef{ok.Ref}, "unkin-agent", drainableTicks(3), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR lost mergeability (conflict)" { t.Errorf("reason = %q, want the mergeability loss", res.Reason) } if f.calls != 3 { t.Errorf("fetch calls = %d, want 3 (baseline + the two conflicted polls)", f.calls) } } // CI that goes green→red during the watch alerts. func TestWatchDetectsCIFailureAfterBaseline(t *testing.T) { ok := base() ok.CIStatus = "pending" green := base() green.CIStatus = "success" red := base() red.CIStatus = "failure" f := &fakeFetcher{states: []PRState{ok, green, red}} res, err := Watch(f, []PRRef{ok.Ref}, "unkin-agent", drainableTicks(3), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "CI failed (failure)" { t.Errorf("reason = %q, want the CI failure (pending→success must pass silently)", res.Reason) } if f.calls != 3 { t.Errorf("fetch calls = %d, want 3 (the success poll must not end the watch)", f.calls) } } // The agent's own pushes and comments must not end a watch; a comment from // anyone else must. func TestWatchIgnoresAgentActivity(t *testing.T) { start := base() pushed := base() pushed.HeadSHA = "def456" // the agent pushed a fix; non-agent comments unchanged commented := pushed commented.NonAgentComments = 1 f := &fakeFetcher{states: []PRState{start, pushed, pushed, commented}} res, err := Watch(f, []PRRef{start.Ref}, "unkin-agent", drainableTicks(3), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "new comment from a non-agent user" { t.Errorf("reason = %q, want the non-agent comment", res.Reason) } if f.calls != 4 { t.Errorf("fetch calls = %d, want 4 (the agent's push and comment must not end the watch)", f.calls) } } func TestPRWatchMergeability(t *testing.T) { snap := func(m Mergeability) PRState { s := base() s.Mergeable = m return s } tests := []struct { name string baseline Mergeability polls []Mergeability wantPoll int // 1-based poll that ends the watch; 0 for none }{ {"conflicted before the watch never alerts", MergeNo, []Mergeability{MergeNo, MergeNo, MergeNo}, 0}, {"loss after a mergeable baseline alerts on the second poll", MergeYes, []Mergeability{MergeNo, MergeNo}, 2}, {"a lone conflicted poll is debounced", MergeYes, []Mergeability{MergeNo, MergeYes, MergeNo}, 0}, {"unknown at baseline does not arm the rule", MergeUnknown, []Mergeability{MergeNo, MergeNo, MergeNo}, 0}, {"unknown at baseline then a real loss alerts", MergeUnknown, []Mergeability{MergeYes, MergeNo, MergeNo}, 3}, {"a baseline conflict resolved then lost again alerts", MergeNo, []Mergeability{MergeYes, MergeNo, MergeNo}, 3}, {"unknown between conflicted polls breaks the run", MergeYes, []Mergeability{MergeNo, MergeUnknown, MergeNo}, 0}, {"the run restarts after an unknown", MergeYes, []Mergeability{MergeNo, MergeUnknown, MergeNo, MergeNo}, 4}, {"unknown alone is never a conflict", MergeYes, []Mergeability{MergeUnknown, MergeUnknown, MergeUnknown}, 0}, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { w := newPRWatch(snap(tt.baseline)) start := time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC) got := 0 for i, m := range tt.polls { changed, reason := w.observe(snap(m), start.Add(time.Duration(i+1)*10*time.Minute)) if !changed { continue } if reason != "PR lost mergeability (conflict)" { t.Fatalf("poll %d ended the watch with %q, want a mergeability loss", i+1, reason) } got = i + 1 break } if got != tt.wantPoll { t.Errorf("alerted on poll %d, want %d", got, tt.wantPoll) } }) } } // Gitea 1.26 always sends mergeable as a plain bool, so unknown is reserved for // what that bool cannot carry: an absent or null flag must decode as unknown // rather than as a conflict, and unknown must encode back as null. func TestMergeabilityDecoding(t *testing.T) { tests := map[string]Mergeability{ `{"number":7}`: MergeUnknown, `{"number":7,"mergeable":null}`: MergeUnknown, `{"number":7,"mergeable":true}`: MergeYes, `{"number":7,"mergeable":false}`: MergeNo, } for body, want := range tests { var pr PullRequest if err := json.Unmarshal([]byte(body), &pr); err != nil { t.Fatalf("Unmarshal(%s): %v", body, err) } if pr.Mergeable != want { t.Errorf("Unmarshal(%s) mergeable = %v, want %v", body, pr.Mergeable, want) } } out, err := json.Marshal(PRState{Mergeable: MergeUnknown}) if err != nil { t.Fatalf("Marshal: %v", err) } if !strings.Contains(string(out), `"mergeable":null`) { t.Errorf("unknown mergeability encoded as %s, want null", out) } } // The inverse of the bug this PR fixed: demanding a mergeable observation to arm // the rule left a PR that was non-mergeable at baseline silent forever, even // after a push gave Gitea a fresh merge computation to answer for. A head that // moves during the watch arms the rule, so the conflict the new head keeps // reporting is attributable to this watch and is reported. func TestWatchAlertsWhenAPushArmsABaselineConflict(t *testing.T) { baseline := base() baseline.Mergeable = MergeNo // still recomputing; the push is not visible yet pushed := base() pushed.Mergeable = MergeNo pushed.HeadSHA = "def456" f := &fakeFetcher{states: []PRState{baseline, pushed, pushed, pushed}} res, err := Watch(f, []PRRef{baseline.Ref}, "unkin-agent", drainableTicks(3), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR lost mergeability (conflict)" { t.Fatalf("reason = %q, want the mergeability loss; a conflict must not be unreportable because the baseline caught the recompute", res.Reason) } if f.calls != 3 { t.Errorf("fetch calls = %d, want 3 (baseline + the two conflicted polls on the new head)", f.calls) } } // A base-branch move restarts the same merge computation, so a conflict that // only becomes visible once someone merges into main alerts too. func TestWatchAlertsWhenBaseMovedIntoAConflict(t *testing.T) { baseline := base() baseline.Mergeable = MergeNo moved := base() moved.Mergeable = MergeNo moved.BaseSHA = "base111" f := &fakeFetcher{states: []PRState{baseline, moved, moved}} res, err := Watch(f, []PRRef{baseline.Ref}, "unkin-agent", drainableTicks(3), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR lost mergeability (conflict)" { t.Errorf("reason = %q, want the mergeability loss", res.Reason) } } // The fix this PR exists for: a conflict predating the watch, on a head and base // that never move, is the condition the operator is already waiting on and must // stay silent however long the watch runs. func TestWatchNeverAlertsOnAStableBaselineConflict(t *testing.T) { stuck := base() stuck.Mergeable = MergeNo f := &fakeFetcher{states: []PRState{stuck}} res, err := Watch(f, []PRRef{stuck.Ref}, "unkin-agent", drainableTicks(100), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "" { t.Fatalf("Watch ended with %q after 100 unchanged polls; a conflict that predates the watch is not a change", res.Reason) } } // Head and base movement arm the conflict rule but are not themselves alerts: a // push, or a base that moves under the PR, must not end a watch. func TestWatchArmingMovementDoesNotAlert(t *testing.T) { start := base() pushed := base() pushed.HeadSHA = "def456" rebased := base() rebased.HeadSHA = "789abc" rebased.BaseSHA = "base111" f := &fakeFetcher{states: []PRState{start, pushed, rebased}} res, err := Watch(f, []PRRef{start.Ref}, "unkin-agent", drainableTicks(5), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "" { t.Fatalf("Watch ended with %q; a new head or base is not a change worth alerting on", res.Reason) } } // The polls that confirm a conflict must be adjacent. A failed poll produces no // snapshot, so the run cannot span it and two non-adjacent falses do not fire. func TestWatchConflictRunDoesNotSpanAFailedPoll(t *testing.T) { ok := base() conflicted := base() conflicted.Mergeable = MergeNo f := &fakeFetcher{ states: []PRState{ok, conflicted, conflicted, conflicted, conflicted}, errs: []error{nil, nil, errors.New("HTTP 502"), nil, nil}, } res, err := Watch(f, []PRRef{ok.Ref}, "unkin-agent", drainableTicks(4), nil, func(PRRef, error) {}) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR lost mergeability (conflict)" { t.Fatalf("reason = %q, want the mergeability loss on the two adjacent polls", res.Reason) } if f.calls != 5 { t.Errorf("fetch calls = %d, want 5: the falses either side of the failed poll are not a run", f.calls) } } // perRefFetcher scripts a separate sequence per ref, so a multi-ref watch can be // driven with each PR doing something different. type perRefFetcher struct { states map[string][]PRState calls map[string]int } func (f *perRefFetcher) FetchState(ref PRRef, _ string) (PRState, error) { key := ref.String() seq := f.states[key] i := f.calls[key] if i >= len(seq) { i = len(seq) - 1 } f.calls[key]++ return seq[i], nil } // Each ref keeps its own run of observations: one PR's mergeability, pushes and // resolutions must neither arm nor disarm another's conflict rule. func TestWatchTracksRefsIndependently(t *testing.T) { refA := PRRef{Owner: "unkin", Repo: "repo", Number: 1} refB := PRRef{Owner: "unkin", Repo: "repo", Number: 2} snap := func(ref PRRef, m Mergeability, head string) PRState { s := base() s.Ref = ref s.Mergeable = m s.HeadSHA = head return s } tests := []struct { name string a, b []PRState want string // reason, "" for no alert }{ { name: "a mergeable neighbour does not arm a baseline conflict", a: []PRState{snap(refA, MergeYes, "aaa")}, b: []PRState{snap(refB, MergeNo, "bbb")}, want: "", }, { name: "a neighbour's push does not arm a baseline conflict", a: []PRState{snap(refA, MergeYes, "aaa"), snap(refA, MergeYes, "aa2"), snap(refA, MergeYes, "aa3")}, b: []PRState{snap(refB, MergeNo, "bbb")}, want: "", }, { name: "a neighbour's mergeable polls do not clear another's run", a: []PRState{snap(refA, MergeYes, "aaa")}, b: []PRState{snap(refB, MergeYes, "bbb"), snap(refB, MergeNo, "bbb"), snap(refB, MergeNo, "bbb")}, want: "PR lost mergeability (conflict)", }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { f := &perRefFetcher{ states: map[string][]PRState{refA.String(): tt.a, refB.String(): tt.b}, calls: map[string]int{}, } res, err := Watch(f, []PRRef{refA, refB}, "unkin-agent", drainableTicks(6), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != tt.want { t.Fatalf("reason = %q, want %q", res.Reason, tt.want) } if tt.want != "" && res.Ref != refB { t.Errorf("alert names %s, want %s", res.Ref, refB) } }) } } // A watch that will deliberately stay silent about a pre-existing conflict has // to hand its caller the baseline it is staying silent about. func TestWatchReportsBaselineStates(t *testing.T) { stuck := base() stuck.Mergeable = MergeNo stuck.CIStatus = "failure" f := &fakeFetcher{states: []PRState{stuck}} var got []PRState if _, err := Watch(f, []PRRef{stuck.Ref}, "unkin-agent", drainableTicks(1), func(sts []PRState) { got = sts }, nil); err != nil { t.Fatalf("Watch: %v", err) } if len(got) != 1 { t.Fatalf("onBaseline received %d state(s), want 1", len(got)) } if got[0].Mergeable != MergeNo || got[0].CIStatus != "failure" { t.Errorf("baseline = %+v, want the conflicted, CI-red snapshot the watch started from", got[0]) } } // The residual gap, pinned so it is a decision rather than an accident: the // agent pushes, the push conflicts, watchpr starts inside Gitea's recompute and // nothing moves again. Every poll answers false and none of them is // attributable to this watch, so no alert is ever sent -- the baseline line is // the only notice. Gitea's payload carries no field that separates this from a // merge check still running. func TestWatchNeverAlertsOnAConflictLandedByThePushBeforeTheWatch(t *testing.T) { pushed := base() pushed.Mergeable = MergeNo pushed.HeadSHA = "def456" f := &fakeFetcher{states: []PRState{pushed}} res, err := Watch(f, []PRRef{pushed.Ref}, "unkin-agent", drainableTicks(50), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "" { t.Fatalf("Watch ended with %q; nothing distinguishes this conflict from a merge check in flight, so it must stay silent", res.Reason) } } // The whole mergeability rule as a sequence table, driven through Watch at the // intervals that decide it. The debounce is wall-clock, so the same sequence of // snapshots must alert or stay silent according to how far apart the polls are. func TestWatchConflictSequences(t *testing.T) { snap := func(m Mergeability, head, bse string) PRState { st := base() st.Mergeable = m st.HeadSHA = head st.BaseSHA = bse return st } repeat := func(st PRState, n int) []PRState { out := make([]PRState, n) for i := range out { out[i] = st } return out } movingBase := func(m Mergeability, n int) []PRState { out := make([]PRState, n) for i := range out { out[i] = snap(m, "h1", fmt.Sprintf("b%d", i)) } return out } const ( quick = 5 * time.Second // well inside conflictWindow normal = 60 * time.Second // watchpr's default relaxed = 10 * time.Minute // past conflictWindow in a single gap ) const conflict = "PR lost mergeability (conflict)" yes := snap(MergeYes, "h1", "b1") no := snap(MergeNo, "h1", "b1") pushedNo := snap(MergeNo, "h2", "b1") movedNo := snap(MergeNo, "h1", "b2") tests := []struct { name string interval time.Duration baseline PRState polls []PRState errs []error want string }{ { name: "a push then the recompute's falses is not a conflict", interval: normal, baseline: yes, polls: repeat(pushedNo, 2), want: "", }, { name: "a fast poller rides out the whole recompute after a push", interval: quick, baseline: yes, polls: repeat(pushedNo, 20), want: "", }, { name: "a base move then the recompute's falses is not a conflict", interval: normal, baseline: no, polls: repeat(movedNo, 2), want: "", }, { name: "a base moving under every poll still confirms a conflict", interval: normal, baseline: yes, polls: movingBase(MergeNo, 20), want: conflict, }, { name: "a conflict that predates the watch stays silent forever", interval: relaxed, baseline: no, polls: repeat(no, 50), want: "", }, { name: "a sustained loss after a mergeable baseline alerts", interval: normal, baseline: yes, polls: repeat(no, 3), want: conflict, }, { name: "a mergeable poll breaks the run however long it ran", interval: relaxed, baseline: yes, polls: []PRState{no, yes, no}, want: "", }, { name: "an unknown poll breaks the run however long it ran", interval: relaxed, baseline: yes, polls: []PRState{no, snap(MergeUnknown, "h1", "b1"), no}, want: "", }, { name: "a baseline conflict that clears and returns alerts", interval: relaxed, baseline: no, polls: []PRState{yes, no, no}, want: conflict, }, { name: "a failed poll breaks the run", interval: relaxed, baseline: yes, polls: []PRState{no, no, no}, errs: []error{nil, nil, errors.New("HTTP 502"), nil}, want: "", }, { name: "the run restarts after a failed poll and still confirms", interval: relaxed, baseline: yes, polls: []PRState{no, no, no, no}, errs: []error{nil, nil, errors.New("HTTP 502"), nil, nil}, want: conflict, }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { f := &fakeFetcher{states: append([]PRState{tt.baseline}, tt.polls...), errs: tt.errs} res, err := Watch(f, []PRRef{tt.baseline.Ref}, "unkin-agent", spacedTicks(len(tt.polls), tt.interval), nil, func(PRRef, error) {}) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != tt.want { t.Fatalf("reason = %q, want %q", res.Reason, tt.want) } }) } } // Two refs watched together keep separate runs and separate clocks: a PR being // pushed to must neither delay nor trigger the conflict its neighbour is really // in, and the alert must name the conflicted one. func TestWatchConflictIsolatedFromANeighboursPushes(t *testing.T) { conflicted := PRRef{Owner: "unkin", Repo: "repo", Number: 1} pushing := PRRef{Owner: "unkin", Repo: "repo", Number: 2} snap := func(ref PRRef, m Mergeability, head string) PRState { st := base() st.Ref = ref st.Mergeable = m st.HeadSHA = head return st } f := &perRefFetcher{ states: map[string][]PRState{ conflicted.String(): {snap(conflicted, MergeYes, "a1"), snap(conflicted, MergeNo, "a1")}, pushing.String(): {snap(pushing, MergeYes, "b1"), snap(pushing, MergeYes, "b2"), snap(pushing, MergeYes, "b3")}, }, calls: map[string]int{}, } res, err := Watch(f, []PRRef{conflicted, pushing}, "unkin-agent", drainableTicks(4), nil, nil) if err != nil { t.Fatalf("Watch: %v", err) } if res.Reason != "PR lost mergeability (conflict)" { t.Fatalf("reason = %q, want the conflict on %s", res.Reason, conflicted) } if res.Ref != conflicted { t.Errorf("alert names %s, want %s", res.Ref, conflicted) } }