package agent import ( "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", Mergeable: true, 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, }, { // A single mergeable=false poll is debounced: Gitea often reports // this transiently right after a push. name: "mergeable true to false for one poll is benign", mutate: func(s *PRState) { s.Mergeable = false }, wantChange: false, }, { // mergeable=false persisting into a second consecutive poll is a // real conflict and alerts. name: "mergeable false persisting a second poll alerts", mutatePrev: func(s *PRState) { s.Mergeable = false }, mutate: func(s *PRState) { s.Mergeable = false }, wantChange: true, }, { // mergeable recovered (false then true) must not alert. name: "mergeable recovered false to true is benign", mutatePrev: func(s *PRState) { s.Mergeable = false }, mutate: func(s *PRState) {}, 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 baselineFired := false ticks := make(chan time.Time, 1) ticks <- time.Now() res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", ticks, func() { baselineFired = true }, nil) if err != nil { t.Fatalf("Watch: %v", err) } if !baselineFired { t.Errorf("onBaseline should fire for an open baseline") } 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) } }