Files
agent-tools/internal/agent/watch.go
T
unkin-agent 54c1d868d3 watchpr: arm the conflict rule on a new merge computation
A mergeable=false baseline meant a conflict could never be reported, since
Gitea sends false while it recomputes after a push. Arm on a head or base SHA
change as well as on a mergeable poll, break the conflict run on unknown and
failed polls, and print the baseline with the conditions it suppresses.
2026-09-26 20:48:31 +10:00

318 lines
10 KiB
Go

package agent
import (
"errors"
"fmt"
"strings"
"time"
)
// errPRGone marks a 404 from the PR lookup itself. A 404 from any other endpoint
// can be a proxy or ingress blip and is left to the ordinary failure cap.
var errPRGone = errors.New("PR no longer visible")
// IsPRGone reports whether err is a 404 from the PR lookup, meaning the PR is no
// longer visible rather than one endpoint being briefly unreachable.
func IsPRGone(err error) bool {
return errors.Is(err, errPRGone)
}
// Mergeability is Gitea's mergeable flag. Gitea 1.26 always sends a plain bool,
// and sends false both for a real conflict and while it recomputes the merge
// base after a push, so false on its own decides nothing (prWatch resolves it).
// Unknown covers what the bool cannot carry: an absent or null flag from another
// Gitea build, and a snapshot no successful poll ever filled in.
type Mergeability int
const (
MergeUnknown Mergeability = iota
MergeYes
MergeNo
)
func (m Mergeability) String() string {
switch m {
case MergeYes:
return "true"
case MergeNo:
return "false"
}
return "unknown"
}
func (m Mergeability) MarshalJSON() ([]byte, error) {
switch m {
case MergeYes:
return []byte("true"), nil
case MergeNo:
return []byte("false"), nil
}
return []byte("null"), nil
}
func (m *Mergeability) UnmarshalJSON(b []byte) error {
switch strings.TrimSpace(string(b)) {
case "true":
*m = MergeYes
case "false":
*m = MergeNo
case "null":
*m = MergeUnknown
default:
return fmt.Errorf("mergeable: unexpected value %s", b)
}
return nil
}
// PRState is a point-in-time snapshot of the PR attributes watchpr tracks.
type PRState struct {
Ref PRRef `json:"ref"`
State string `json:"state"` // open / closed
Merged bool `json:"merged"`
HeadSHA string `json:"head_sha"`
BaseSHA string `json:"base_sha"`
Mergeable Mergeability `json:"mergeable"`
CIStatus string `json:"ci_status"` // success / pending / failure / error / ""
NonAgentComments int `json:"non_agent_comments"`
Title string `json:"title"`
URL string `json:"url"`
}
// FetchState builds a PRState for the given ref. agentLogin's comments are
// excluded from the non-agent comment count.
func FetchState(c *GiteaClient, ref PRRef, agentLogin string) (PRState, error) {
pr, err := c.GetPR(ref.RepoPath(), ref.Number)
if err != nil {
if IsNotFound(err) {
return PRState{}, fmt.Errorf("%w: %w", errPRGone, err)
}
return PRState{}, err
}
// A 404 here means the head commit is gone (branch deleted after a squash/
// rebase merge); the PR object is still authoritative, so treat CI as absent
// rather than discarding the merge signal and hanging the watch loop.
ci, err := c.CommitStatus(ref.RepoPath(), pr.Head.Sha)
if err != nil && !IsNotFound(err) {
return PRState{}, err
}
comments, err := c.ListComments(ref.RepoPath(), ref.Number)
if err != nil {
return PRState{}, err
}
return PRState{
Ref: ref,
State: pr.State,
Merged: pr.Merged,
HeadSHA: pr.Head.Sha,
BaseSHA: pr.Base.Sha,
Mergeable: pr.Mergeable,
CIStatus: ci,
NonAgentComments: countNonAgentComments(comments, agentLogin),
Title: pr.Title,
URL: pr.HTMLURL,
}, nil
}
// StateFetcher fetches the current PRState for a ref. *GiteaClient satisfies it
// via its FetchState method; tests inject fakes.
type StateFetcher interface {
FetchState(ref PRRef, agentLogin string) (PRState, error)
}
// FetchState makes *GiteaClient a StateFetcher.
func (c *GiteaClient) FetchState(ref PRRef, agentLogin string) (PRState, error) {
return FetchState(c, ref, agentLogin)
}
// WatchResult is the change that ended a watch.
type WatchResult struct {
Ref PRRef
Reason string
State PRState
}
// terminalState reports whether a PR has reached a final state from which no
// further meaningful change is possible, with a human-readable reason. Unlike a
// transition (see MeaningfulChange) this holds for a single snapshot, so it also
// catches a PR that is already merged/closed the moment watchpr starts.
func terminalState(st PRState) (bool, string) {
if st.Merged {
return true, "PR merged"
}
if st.State == "closed" {
return true, "PR closed without merging"
}
return false, ""
}
// conflictPolls is how many consecutive non-mergeable polls of an unchanged
// merge computation confirm a real conflict.
const conflictPolls = 2
// prWatch tracks one PR across polls, because mergeability needs more memory
// than the previous snapshot. Gitea reports mergeable=false while it recomputes
// the merge base after a push, so a false is only trusted once this watch has
// seen a merge computation start: the PR was mergeable at some point, or its
// head or base SHA moved. A bare false inherited from the baseline says nothing
// -- it is equally a conflict the operator is already waiting on and a recompute
// in flight -- so it arms nothing.
type prWatch struct {
prev PRState
armed bool
conflicts int
}
func newPRWatch(baseline PRState) *prWatch {
return &prWatch{prev: baseline, armed: baseline.Mergeable == MergeYes}
}
// mergeInputsChanged reports whether the commits Gitea merges have moved, which
// starts a fresh merge computation whose result is attributable to this watch.
func mergeInputsChanged(prev, cur PRState) bool {
return cur.HeadSHA != prev.HeadSHA || cur.BaseSHA != prev.BaseSHA
}
// track folds one snapshot's mergeability into the run of observations. The
// polls that confirm a conflict must be adjacent, so anything but another
// non-mergeable observation breaks the run.
func (w *prWatch) track(st PRState) {
switch st.Mergeable {
case MergeNo:
w.conflicts++
case MergeYes:
w.armed = true
w.conflicts = 0
default:
w.conflicts = 0
}
}
// missed records a poll that never produced a snapshot; the run of adjacent
// non-mergeable observations does not survive the gap.
func (w *prWatch) missed() {
w.conflicts = 0
}
// observe folds in the newest snapshot and reports whether the watch should end.
func (w *prWatch) observe(cur PRState) (bool, string) {
changed, reason := MeaningfulChange(w.prev, cur)
if mergeInputsChanged(w.prev, cur) {
w.armed = true
w.conflicts = 0
}
w.prev = cur
w.track(cur)
if changed {
return true, reason
}
if w.armed && w.conflicts >= conflictPolls && cur.State == "open" {
return true, "PR lost mergeability (conflict)"
}
return false, ""
}
// MaxPollFailures is how many consecutive failed polls of the same PR are
// tolerated before Watch gives up. The abort fires on the 20th failed tick, so
// at watchpr's default 60s interval a watch rides out ~19 minutes of failure.
const MaxPollFailures = 20
// Watch establishes a baseline for each ref, then polls on every tick until a
// tracked PR changes meaningfully, returning the first such change. A PR that is
// already terminal (merged/closed) at baseline is reported immediately rather
// than polled forever. Transient poll errors are handed to onError and the loop
// continues, but never blindly: a baseline fetch error, an authentication
// failure surviving a token re-mint, a 404 from the PR lookup itself (the repo
// is gone, renamed, or no longer visible), and MaxPollFailures consecutive
// failures of one PR all abort, because a watcher that sees nothing must not
// look healthy.
// onBaseline, if set, receives every captured baseline once, before the first
// tick, so a caller can show what state the watch started from.
func Watch(f StateFetcher, refs []PRRef, agentLogin string, ticks <-chan time.Time, onBaseline func([]PRState), onError func(PRRef, error)) (WatchResult, error) {
watches := make(map[string]*prWatch, len(refs))
baselines := make([]PRState, 0, len(refs))
for _, ref := range refs {
st, err := f.FetchState(ref, agentLogin)
if err != nil {
return WatchResult{}, err
}
if terminal, reason := terminalState(st); terminal {
return WatchResult{Ref: ref, Reason: reason, State: st}, nil
}
watches[ref.String()] = newPRWatch(st)
baselines = append(baselines, st)
}
if onBaseline != nil {
onBaseline(baselines)
}
fails := make(map[string]int, len(refs))
for range ticks {
for _, ref := range refs {
key := ref.String()
cur, err := f.FetchState(ref, agentLogin)
if err != nil {
if IsAuthError(err) || IsPRGone(err) {
return WatchResult{}, fmt.Errorf("polling %s: %w", key, err)
}
fails[key]++
watches[key].missed()
if onError != nil {
onError(ref, err)
}
if fails[key] >= MaxPollFailures {
return WatchResult{}, fmt.Errorf("polling %s: giving up after %d consecutive failures: %w", key, fails[key], err)
}
continue
}
fails[key] = 0
if changed, reason := watches[key].observe(cur); changed {
return WatchResult{Ref: ref, Reason: reason, State: cur}, nil
}
}
}
return WatchResult{}, nil
}
// countNonAgentComments counts comments authored by anyone other than agentLogin.
func countNonAgentComments(comments []Comment, agentLogin string) int {
n := 0
for _, cm := range comments {
if cm.User.Login != agentLogin {
n++
}
}
return n
}
// isFailedCI reports whether a combined CI state is a terminal failure.
func isFailedCI(state string) bool {
return state == "failure" || state == "error"
}
// MeaningfulChange compares a previous state to the current one and reports
// whether a change warrants alerting the operator, with a human-readable
// reason. Benign transitions (CI pending→success, the agent's own comments, a
// new head or base commit, an unchanged snapshot) return false. Mergeability is
// not decided here: it takes a whole run of observations, which prWatch keeps.
//
// Alerting conditions:
// - the PR merged
// - the PR closed without merging
// - a new comment from someone other than the agent
// - CI transitioned into failure/error
func MeaningfulChange(prev, cur PRState) (bool, string) {
if !prev.Merged && cur.Merged {
return true, "PR merged"
}
// Closed (not merged): only alert on the open→closed edge.
if prev.State == "open" && cur.State == "closed" && !cur.Merged {
return true, "PR closed without merging"
}
if cur.NonAgentComments > prev.NonAgentComments {
return true, "new comment from a non-agent user"
}
if isFailedCI(cur.CIStatus) && !isFailedCI(prev.CIStatus) {
return true, "CI failed (" + cur.CIStatus + ")"
}
return false, ""
}