9 Commits

Author SHA1 Message Date
benvin 5b6e15b50f Merge pull request 'watchpr: add --max-wait with distinct timeout exit code' (#25) from benvin/watchpr-max-wait into main
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #25
2026-10-05 22:14:49 +11:00
unkin-agent 00ba3df49c watchpr: reject --max-wait with --once, test exit-code mapping
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-05 22:05:09 +11:00
unkin-agent 639bcbda49 watchpr: add --max-wait, exit 3 on timeout
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-05 22:03:13 +11:00
benvin 0f2132b236 Merge pull request 'ci: use container-rpmbuilder image' (#24) from benvin/rpmbuilder-image into main
Reviewed-on: #24
2026-10-05 21:44:32 +11:00
unkin-agent f98414518e ci: use container-rpmbuilder image
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-05 14:39:17 +11:00
benvin 4f95e2f59c Merge pull request 'ci: switch Go steps to gobuilder, enable S3 build cache' (#23) from benvin/gocache into main
Reviewed-on: #23
2026-10-02 23:51:56 +10:00
unkin-agent 0056be8a65 ci: switch Go steps to gobuilder, enable S3 build cache
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
golang:1.25 and the stale almalinux9-gobuilder pin are replaced by
gobuilder:0.1.2-alma9 on every step that invokes the Go toolchain, with
GOCACHEPROG wired to the baked-in go-cache-plugin against the shared
S3 cache bucket.

- build.yaml/test.yaml/release.yaml: build+test steps use gobuilder,
  cache env added
- pre-commit.yaml: switch pin, add cache (hooks run go vet/go test)
- lint step (test.yaml) and rpm/upload/release steps left untouched
2026-10-02 23:41:39 +10:00
benvin 215e4a96d3 Merge pull request 'agentws: refuse rm on an ambiguous branch name' (#22) from benvin/agentws-rm-ambiguous into main
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #22
2026-10-02 23:03:57 +10:00
unkin-agent 6c88d17736 agentws: refuse rm on a branch name shared by several repos
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-02 22:50:59 +10:00
12 changed files with 425 additions and 26 deletions
+12 -1
View File
@@ -3,7 +3,18 @@ when:
steps: steps:
- name: build - name: build
image: golang:1.25 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-agent-tools
GOCACHEPROG: go-cache-plugin --cache-dir=/tmp/gocache
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- make build - make build
backend_options: backend_options:
+12 -1
View File
@@ -3,7 +3,18 @@ when:
steps: steps:
- name: pre-commit - name: pre-commit
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-agent-tools
GOCACHEPROG: go-cache-plugin --cache-dir=/tmp/gocache
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- uvx pre-commit run --all-files - uvx pre-commit run --all-files
backend_options: backend_options:
+25 -3
View File
@@ -3,7 +3,18 @@ when:
steps: steps:
- name: test - name: test
image: golang:1.25 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-agent-tools
GOCACHEPROG: go-cache-plugin --cache-dir=/tmp/gocache
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- go test -race ./... - go test -race ./...
backend_options: backend_options:
@@ -21,7 +32,18 @@ steps:
# cross-platform binaries attached to the Gitea release. Each tool is a # cross-platform binaries attached to the Gitea release. Each tool is a
# separate main package, so they are built individually per os/arch. # separate main package, so they are built individually per os/arch.
- name: build - name: build
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-agent-tools
GOCACHEPROG: go-cache-plugin --cache-dir=/tmp/gocache
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- make build VERSION=${CI_COMMIT_TAG} - make build VERSION=${CI_COMMIT_TAG}
# Shell variables/expansions are escaped as $$ so Woodpecker leaves them # Shell variables/expansions are escaped as $$ so Woodpecker leaves them
@@ -51,7 +73,7 @@ steps:
# Package the built binaries + generated shell completions into an RPM. # Package the built binaries + generated shell completions into an RPM.
- name: package - name: package
image: git.unkin.net/unkin/almalinux9-rpmbuilder:latest image: artifactapi.k8s.syd1.au.unkin.net/docker-internal/rpmbuilder:0.1.0-alma9
commands: commands:
- ./scripts/build-rpm.sh ${CI_COMMIT_TAG} - ./scripts/build-rpm.sh ${CI_COMMIT_TAG}
depends_on: [build] depends_on: [build]
+12 -1
View File
@@ -18,7 +18,18 @@ steps:
cpu: 2 cpu: 2
- name: test - name: test
image: golang:1.25 image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
environment:
GOCACHE_S3_BUCKET: gocache
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-agent-tools
GOCACHEPROG: go-cache-plugin --cache-dir=/tmp/gocache
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
commands: commands:
- go test -v -race ./... - go test -v -race ./...
backend_options: backend_options:
+3
View File
@@ -178,6 +178,9 @@ wrapped per stage (login / read denied / write denied) via `ErrVaultDenied`.
## Gotchas ## Gotchas
- `watchpr` exits 0 with no output changes on `--once` (just prints state). - `watchpr` exits 0 with no output changes on `--once` (just prints state).
- Background commands are killed at a 2h cap, indistinguishable from a crash.
Orchestrators run `watchpr --max-wait 110m ...` and restart it on exit code 3
(no change within max-wait); any other non-zero exit is a real failure.
- Gitea tokens expire in ~1h, shorter than a watch: the client re-mints once when - Gitea tokens expire in ~1h, shorter than a watch: the client re-mints once when
the credential it sent was rejected and replays the request. If the fresh token the credential it sent was rejected and replays the request. If the fresh token
is rejected too, `watchpr` exits non-zero rather than polling blind. is rejected too, `watchpr` exits non-zero rather than polling blind.
+12 -1
View File
@@ -131,11 +131,20 @@ watchpr --interval 30 unkin/argocd-apps#42
# One-shot: print current state and exit 0 (great for scripts) # One-shot: print current state and exit 0 (great for scripts)
watchpr --once unkin/argocd-apps#42 watchpr --once unkin/argocd-apps#42
watchpr --once --json unkin/argocd-apps#42 watchpr --once --json unkin/argocd-apps#42
# Orchestrators: run inside a background command (2h hard cap) and restart on exit 3
watchpr --max-wait 110m unkin/argocd-apps#42
``` ```
On a meaningful change `watchpr` prints the reason and the PR's current state, On a meaningful change `watchpr` prints the reason and the PR's current state,
then exits 0. Use `--json` for machine-readable output. then exits 0. Use `--json` for machine-readable output.
`--max-wait` (same forms as `--interval`; default `0` = unlimited) bounds a
watch: when it elapses with no change, `watchpr` prints `timeout: no change
within <d>` and each PR's current state (`{"timeout":true,...}` under `--json`)
and exits **3**, so a timed-out watch is distinguishable from a crash and is
safe to restart.
### Exit behaviour ### Exit behaviour
A watcher that sees nothing must not look healthy, so every terminal failure A watcher that sees nothing must not look healthy, so every terminal failure
@@ -174,8 +183,10 @@ agentws new argocd-apps --branch benvin/hotfix --from release-1.2
# List managed worktrees (repo, branch, path) # List managed worktrees (repo, branch, path)
agentws list agentws list
# Remove a worktree (by path or branch); refreshes the source repo afterwards # Remove a worktree (by path or branch); refreshes the source repo afterwards.
# A branch name shared by several repos is refused; pass a path or --repo.
agentws rm benvin/my-change agentws rm benvin/my-change
agentws rm benvin/my-change --repo argocd-apps
agentws rm ~/.cache/agentws/argocd-apps__benvin-my-change --delete-branch agentws rm ~/.cache/agentws/argocd-apps__benvin-my-change --delete-branch
# Classify every worktree found; dry run unless --yes is given # Classify every worktree found; dry run unless --yes is given
+26 -6
View File
@@ -12,7 +12,7 @@
// //
// agentws new <repo> [--branch benvin/<name>] [--from <base-branch>] // agentws new <repo> [--branch benvin/<name>] [--from <base-branch>]
// agentws list // agentws list
// agentws rm <path-or-branch> [--delete-branch] // agentws rm <path-or-branch> [--repo <repo>] [--delete-branch]
// agentws prune [--yes] [--keep-branches] [--no-fetch] [--json] // agentws prune [--yes] [--keep-branches] [--no-fetch] [--json]
// [--include-unmanaged] [--include-keep] // [--include-unmanaged] [--include-keep]
// agentws clean // agentws clean
@@ -591,6 +591,7 @@ func underRoot(path, root string) bool {
func newRmCmd() *cobra.Command { func newRmCmd() *cobra.Command {
var deleteBranch bool var deleteBranch bool
var repo string
cmd := &cobra.Command{ cmd := &cobra.Command{
Use: "rm <path-or-branch>", Use: "rm <path-or-branch>",
Short: "Remove a managed worktree and refresh its source repo", Short: "Remove a managed worktree and refresh its source repo",
@@ -598,7 +599,7 @@ func newRmCmd() *cobra.Command {
SilenceUsage: true, SilenceUsage: true,
RunE: func(cmd *cobra.Command, args []string) error { RunE: func(cmd *cobra.Command, args []string) error {
target := strings.TrimSpace(args[0]) target := strings.TrimSpace(args[0])
wt, err := resolveWorktree(target) wt, err := resolveWorktree(target, repo)
if err != nil { if err != nil {
return err return err
} }
@@ -607,22 +608,41 @@ func newRmCmd() *cobra.Command {
}, },
} }
cmd.Flags().BoolVar(&deleteBranch, "delete-branch", false, "Also delete the local branch after removing the worktree") cmd.Flags().BoolVar(&deleteBranch, "delete-branch", false, "Also delete the local branch after removing the worktree")
cmd.Flags().StringVar(&repo, "repo", "", "Only match worktrees of this repo (disambiguates a branch name)")
return cmd return cmd
} }
// resolveWorktree finds a managed worktree by exact path or by branch name. // resolveWorktree finds a managed worktree by exact path or by branch name. A
func resolveWorktree(target string) (managedWt, error) { // branch name shared by several repos is refused unless repo narrows it to one.
func resolveWorktree(target, repo string) (managedWt, error) {
managed, err := managedWorktrees() managed, err := managedWorktrees()
if err != nil { if err != nil {
return managedWt{}, err return managedWt{}, err
} }
abs, _ := filepath.Abs(target) abs, _ := filepath.Abs(target)
var matches []managedWt
for _, w := range managed { for _, w := range managed {
if w.path == target || w.path == abs || (w.branch != "" && w.branch == target) { if repo != "" && w.repo != repo {
continue
}
if w.path == target || w.path == abs {
return w, nil return w, nil
} }
if w.branch != "" && w.branch == target {
matches = append(matches, w)
}
} }
return managedWt{}, fmt.Errorf("no managed worktree matching %q (try `agentws list`)", target) switch len(matches) {
case 0:
return managedWt{}, fmt.Errorf("no managed worktree matching %q (try `agentws list`)", target)
case 1:
return matches[0], nil
}
paths := make([]string, len(matches))
for i, w := range matches {
paths[i] = " " + w.path
}
return managedWt{}, fmt.Errorf("branch %q matches %d worktrees; pass a path or --repo:\n%s", target, len(matches), strings.Join(paths, "\n"))
} }
// removeWorktree removes a managed worktree and, when asked, its local branch. // removeWorktree removes a managed worktree and, when asked, its local branch.
+103
View File
@@ -0,0 +1,103 @@
package main
import (
"path/filepath"
"strings"
"testing"
"git.unkin.net/unkin/agent-tools/internal/agent"
)
// addOtherRepoWorktree clones a second source repo "other" and gives it a
// managed worktree on branch, so two repos share one branch name.
func addOtherRepoWorktree(t *testing.T, f *fixture, branch string) string {
t.Helper()
src := filepath.Join(f.root, "src", "other")
git(t, filepath.Join(f.root, "src"), "clone", f.bare, src)
path := filepath.Join(f.wtRoot, agent.WorktreeDirName("other", branch))
git(t, src, "worktree", "add", path, "-b", branch, "origin/main")
return path
}
func TestResolveUniqueBranch(t *testing.T) {
f := newFixture(t)
want := f.addWorktree(t, "benvin/one")
f.addWorktree(t, "benvin/two")
wt, err := resolveWorktree("benvin/one", "")
if err != nil {
t.Fatal(err)
}
if wt.path != want {
t.Errorf("path %q, want %q", wt.path, want)
}
}
func TestResolveAmbiguousBranchRefused(t *testing.T) {
f := newFixture(t)
a := f.addWorktree(t, "benvin/shared")
b := addOtherRepoWorktree(t, f, "benvin/shared")
_, err := resolveWorktree("benvin/shared", "")
if err == nil {
t.Fatal("ambiguous branch resolved, want refusal")
}
for _, p := range []string{a, b} {
if !strings.Contains(err.Error(), p) {
t.Errorf("error %q does not list candidate %s", err, p)
}
}
}
func TestRmAmbiguousBranchRemovesNothing(t *testing.T) {
f := newFixture(t)
a := f.addWorktree(t, "benvin/shared")
b := addOtherRepoWorktree(t, f, "benvin/shared")
cmd := newRootCmd()
cmd.SetArgs([]string{"rm", "benvin/shared"})
cmd.SetOut(new(strings.Builder))
cmd.SetErr(new(strings.Builder))
if err := cmd.Execute(); err == nil {
t.Fatal("rm succeeded on an ambiguous branch")
}
if !exists(a) || !exists(b) {
t.Error("ambiguous rm removed a worktree")
}
}
func TestResolveByPathDespiteSharedBranch(t *testing.T) {
f := newFixture(t)
f.addWorktree(t, "benvin/shared")
b := addOtherRepoWorktree(t, f, "benvin/shared")
wt, err := resolveWorktree(b, "")
if err != nil {
t.Fatal(err)
}
if wt.path != b || wt.repo != "other" {
t.Errorf("got %s (%s), want %s (other)", wt.path, wt.repo, b)
}
}
func TestResolveRepoDisambiguates(t *testing.T) {
f := newFixture(t)
a := f.addWorktree(t, "benvin/shared")
b := addOtherRepoWorktree(t, f, "benvin/shared")
for repo, want := range map[string]string{"repo": a, "other": b} {
wt, err := resolveWorktree("benvin/shared", repo)
if err != nil {
t.Fatalf("--repo %s: %v", repo, err)
}
if wt.path != want {
t.Errorf("--repo %s: path %q, want %q", repo, wt.path, want)
}
}
if _, err := resolveWorktree("benvin/shared", "missing"); err == nil {
t.Error("--repo missing resolved, want no match")
}
}
func TestResolvePathOutsideRepoFilterRefused(t *testing.T) {
f := newFixture(t)
a := f.addWorktree(t, "benvin/one")
if _, err := resolveWorktree(a, "other"); err == nil {
t.Error("path in repo resolved under --repo other")
}
}
+95 -7
View File
@@ -19,10 +19,12 @@
// watchpr --once --json owner/repo#12 // watchpr --once --json owner/repo#12
// watchpr --interval 30s owner/repo#12 // watchpr --interval 30s owner/repo#12
// watchpr --interval 30 owner/repo#12 // watchpr --interval 30 owner/repo#12
// watchpr --max-wait 110m owner/repo#12 # exits 3 if nothing changed
package main package main
import ( import (
"encoding/json" "encoding/json"
"errors"
"fmt" "fmt"
"io" "io"
"os" "os"
@@ -36,19 +38,34 @@ import (
var version = "dev" var version = "dev"
// exitTimedOut is the exit status when --max-wait elapses with no change, so a
// caller can tell a deliberate timeout from a crash and restart the watch.
const exitTimedOut = 3
var errTimedOut = errors.New("max-wait elapsed with no change")
func main() { func main() {
// cobra prints the error itself (SilenceErrors stays off); we only need to // cobra prints the error itself (SilenceErrors stays off); we only need to
// turn any command error into a non-zero exit. // turn any command error into a non-zero exit.
if err := newRootCmd().Execute(); err != nil { os.Exit(exitCode(newRootCmd().Execute()))
os.Exit(1) }
// exitCode maps a command result to the process exit status.
func exitCode(err error) int {
switch {
case err == nil:
return 0
case errors.Is(err, errTimedOut):
return exitTimedOut
} }
return 1
} }
// newRootCmd builds the watchpr command tree. It is separated from main so // newRootCmd builds the watchpr command tree. It is separated from main so
// tests can invoke Execute and assert the exit behaviour without spawning a // tests can invoke Execute and assert the exit behaviour without spawning a
// process. // process.
func newRootCmd() *cobra.Command { func newRootCmd() *cobra.Command {
var intervalFlag string var intervalFlag, maxWaitFlag string
var once, jsonMode bool var once, jsonMode bool
root := &cobra.Command{ root := &cobra.Command{
@@ -57,7 +74,9 @@ func newRootCmd() *cobra.Command {
Long: "watchpr polls each PR every --interval and exits (reporting what changed)\n" + Long: "watchpr polls each PR every --interval and exits (reporting what changed)\n" +
"when a PR merges/closes, gets a new non-agent comment, its CI fails, or it\n" + "when a PR merges/closes, gets a new non-agent comment, its CI fails, or it\n" +
"loses mergeability after the baseline. Conditions already true at the\n" + "loses mergeability after the baseline. Conditions already true at the\n" +
"baseline are printed, not alerted on. Refs take owner/repo#N or owner/repo:N.", "baseline are printed, not alerted on. Refs take owner/repo#N or owner/repo:N.\n\n" +
"Exit status: 0 on a change (or --once), 3 when --max-wait elapses with no\n" +
"change (current states are printed), 1 on any error.",
Version: version, Version: version,
Args: cobra.ArbitraryArgs, Args: cobra.ArbitraryArgs,
SilenceUsage: true, SilenceUsage: true,
@@ -69,6 +88,13 @@ func newRootCmd() *cobra.Command {
if err != nil { if err != nil {
return err return err
} }
maxWait, err := parseMaxWait(maxWaitFlag)
if err != nil {
return err
}
if once && maxWait > 0 {
return fmt.Errorf("--max-wait has no effect with --once")
}
refs := make([]agent.PRRef, 0, len(args)) refs := make([]agent.PRRef, 0, len(args))
for _, a := range args { for _, a := range args {
ref, err := agent.ParsePRRef(a) ref, err := agent.ParsePRRef(a)
@@ -81,13 +107,18 @@ func newRootCmd() *cobra.Command {
if once { if once {
return runOnce(c, refs, jsonMode) return runOnce(c, refs, jsonMode)
} }
return runWatch(c, refs, interval, jsonMode) err = runWatch(c, refs, interval, maxWait, jsonMode)
if errors.Is(err, errTimedOut) {
cmd.SilenceErrors = true
}
return err
}, },
} }
root.SetVersionTemplate("{{.Version}}\n") root.SetVersionTemplate("{{.Version}}\n")
f := root.Flags() f := root.Flags()
f.StringVar(&intervalFlag, "interval", "60s", "Polling interval: a duration (30s, 2m, 1h30m) or a bare number of seconds") f.StringVar(&intervalFlag, "interval", "60s", "Polling interval: a duration (30s, 2m, 1h30m) or a bare number of seconds")
f.StringVar(&maxWaitFlag, "max-wait", "0", "Give up after this long with no change and exit 3 (same forms as --interval; 0 = unlimited)")
f.BoolVar(&once, "once", false, "Check once, print current state, and exit") f.BoolVar(&once, "once", false, "Check once, print current state, and exit")
f.BoolVar(&jsonMode, "json", false, "Emit JSON") f.BoolVar(&jsonMode, "json", false, "Emit JSON")
@@ -154,11 +185,17 @@ func runOnce(c *agent.GiteaClient, refs []agent.PRRef, jsonMode bool) error {
// runWatch establishes a baseline then polls until a tracked PR changes // runWatch establishes a baseline then polls until a tracked PR changes
// meaningfully, at which point it reports the change and returns. // meaningfully, at which point it reports the change and returns.
func runWatch(c *agent.GiteaClient, refs []agent.PRRef, interval time.Duration, jsonMode bool) error { func runWatch(c *agent.GiteaClient, refs []agent.PRRef, interval, maxWait time.Duration, jsonMode bool) error {
login := agent.AgentLogin() login := agent.AgentLogin()
ticker := time.NewTicker(interval) ticker := time.NewTicker(interval)
defer ticker.Stop() defer ticker.Stop()
var deadline <-chan time.Time
if maxWait > 0 {
timer := time.NewTimer(maxWait)
defer timer.Stop()
deadline = timer.C
}
onBaseline := func(states []agent.PRState) { onBaseline := func(states []agent.PRState) {
emitBaselines(os.Stderr, states, interval, jsonMode) emitBaselines(os.Stderr, states, interval, jsonMode)
@@ -167,14 +204,65 @@ func runWatch(c *agent.GiteaClient, refs []agent.PRRef, interval time.Duration,
warn(os.Stderr, jsonMode, "polling %s: %v", ref.String(), err) warn(os.Stderr, jsonMode, "polling %s: %v", ref.String(), err)
} }
res, err := agent.Watch(c, refs, login, ticker.C, onBaseline, onError) res, err := agent.Watch(c, refs, login, untilDeadline(ticker.C, deadline), onBaseline, onError)
if err != nil { if err != nil {
return describeFailure(err) return describeFailure(err)
} }
if res.TimedOut {
reportTimeout(os.Stdout, maxWait, res.States, jsonMode)
return errTimedOut
}
report(res.Ref.String(), res.Reason, res.State, jsonMode) report(res.Ref.String(), res.Reason, res.State, jsonMode)
return nil return nil
} }
// parseMaxWait reads --max-wait: 0 means unlimited, anything else parses like
// --interval.
func parseMaxWait(v string) (time.Duration, error) {
if d, err := time.ParseDuration(strings.TrimSpace(v)); err == nil && d == 0 {
return 0, nil
}
return agent.ParseDurationFlag("max-wait", v)
}
// untilDeadline forwards ticks until deadline fires, then closes, which ends
// Watch without waiting out the rest of an interval. A nil deadline never fires.
func untilDeadline(ticks, deadline <-chan time.Time) <-chan time.Time {
out := make(chan time.Time)
go func() {
defer close(out)
for {
select {
case <-deadline:
return
case t := <-ticks:
select {
case out <- t:
case <-deadline:
return
}
}
}
}()
return out
}
// reportTimeout emits each PR's current state when --max-wait elapses.
func reportTimeout(w io.Writer, maxWait time.Duration, states []agent.PRState, jsonMode bool) {
if jsonMode {
_ = json.NewEncoder(w).Encode(struct {
Timeout bool `json:"timeout"`
MaxWait string `json:"max_wait"`
States []agent.PRState `json:"states"`
}{true, maxWait.String(), states})
return
}
_, _ = fmt.Fprintf(w, "timeout: no change within %s\n", maxWait)
for _, st := range states {
_, _ = fmt.Fprintln(w, stateLine(st))
}
}
// baselineRecord is the --json form of baselineLine. Both go to stderr, leaving // baselineRecord is the --json form of baselineLine. Both go to stderr, leaving
// the stdout contract a single result record: a caller automating watchpr is // the stdout contract a single result record: a caller automating watchpr is
// precisely the one who needs to be told the watch started against a PR that is // precisely the one who needs to be told the watch started against a PR that is
+93
View File
@@ -311,3 +311,96 @@ func TestBaselineNamesAnUnknownMergeability(t *testing.T) {
t.Errorf("baselineLine = %q, want the suppression note", got) t.Errorf("baselineLine = %q, want the suppression note", got)
} }
} }
func TestParseMaxWait(t *testing.T) {
for in, want := range map[string]time.Duration{"0": 0, "0s": 0, "30s": 30 * time.Second, "1h55m": 115 * time.Minute, "90": 90 * time.Second} {
got, err := parseMaxWait(in)
if err != nil || got != want {
t.Errorf("parseMaxWait(%q) = %v, %v; want %v", in, got, err, want)
}
}
for _, in := range []string{"soon", "-1m"} {
if _, err := parseMaxWait(in); err == nil || !strings.Contains(err.Error(), "--max-wait") {
t.Errorf("parseMaxWait(%q) err = %v, want a --max-wait error", in, err)
}
}
}
func TestExecuteBadMaxWaitErrors(t *testing.T) {
cmd := newRootCmd()
cmd.SetArgs([]string{"--max-wait", "soon", "unkin/repo#1"})
cmd.SetOut(io.Discard)
cmd.SetErr(io.Discard)
if err := cmd.Execute(); err == nil || errors.Is(err, errTimedOut) {
t.Fatalf("Execute() = %v, want a parse error", err)
}
}
func TestExecuteOnceWithMaxWaitErrors(t *testing.T) {
cmd := newRootCmd()
cmd.SetArgs([]string{"--once", "--max-wait", "5m", "unkin/repo#1"})
cmd.SetOut(io.Discard)
cmd.SetErr(io.Discard)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "--once") {
t.Fatalf("Execute() = %v, want a --once/--max-wait conflict error", err)
}
if got := exitCode(err); got != 1 {
t.Fatalf("exitCode = %d, want 1", got)
}
}
func TestExitCode(t *testing.T) {
for _, tc := range []struct {
err error
want int
}{
{nil, 0},
{errTimedOut, exitTimedOut},
{fmt.Errorf("wrapped: %w", errTimedOut), exitTimedOut},
{errors.New("boom"), 1},
} {
if got := exitCode(tc.err); got != tc.want {
t.Errorf("exitCode(%v) = %d, want %d", tc.err, got, tc.want)
}
}
}
// Ticks pass through until the deadline fires; then the channel closes without
// waiting for another tick.
func TestUntilDeadline(t *testing.T) {
ticks := make(chan time.Time)
deadline := make(chan time.Time)
out := untilDeadline(ticks, deadline)
now := time.Now()
ticks <- now
if got := <-out; !got.Equal(now) {
t.Fatalf("forwarded %v, want %v", got, now)
}
close(deadline)
if _, ok := <-out; ok {
t.Fatal("channel still open after the deadline")
}
}
func TestReportTimeout(t *testing.T) {
st := agent.PRState{Ref: agent.PRRef{Owner: "unkin", Repo: "repo", Number: 3}, State: "open", Mergeable: agent.MergeYes, CIStatus: "pending"}
var buf bytes.Buffer
reportTimeout(&buf, 110*time.Minute, []agent.PRState{st}, false)
lines := strings.Split(strings.TrimSpace(buf.String()), "\n")
if len(lines) != 2 || lines[0] != "timeout: no change within 1h50m0s" || lines[1] != stateLine(st) {
t.Errorf("text output = %q", buf.String())
}
buf.Reset()
reportTimeout(&buf, time.Minute, []agent.PRState{st}, true)
var rec struct {
Timeout bool `json:"timeout"`
States []agent.PRState `json:"states"`
}
if err := json.Unmarshal(buf.Bytes(), &rec); err != nil || !rec.Timeout || len(rec.States) != 1 {
t.Errorf("json output = %q (err %v)", buf.String(), err)
}
}
+12 -6
View File
@@ -124,11 +124,14 @@ func (c *GiteaClient) FetchState(ref PRRef, agentLogin string) (PRState, error)
return FetchState(c, ref, agentLogin) return FetchState(c, ref, agentLogin)
} }
// WatchResult is the change that ended a watch. // WatchResult is the change that ended a watch. When the tick channel closes
// with no change, TimedOut is set and States holds each PR's last seen state.
type WatchResult struct { type WatchResult struct {
Ref PRRef Ref PRRef
Reason string Reason string
State PRState State PRState
TimedOut bool
States []PRState
} }
// terminalState reports whether a PR has reached a final state from which no // terminalState reports whether a PR has reached a final state from which no
@@ -276,6 +279,7 @@ const MaxPollFailures = 20
// look healthy. // look healthy.
// onBaseline, if set, receives every captured baseline once, before the first // onBaseline, if set, receives every captured baseline once, before the first
// tick, so a caller can show what state the watch started from. // tick, so a caller can show what state the watch started from.
// A closed ticks channel ends the watch with a TimedOut result.
func Watch(f StateFetcher, refs []PRRef, agentLogin string, ticks <-chan time.Time, onBaseline func([]PRState), onError func(PRRef, error)) (WatchResult, error) { 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)) watches := make(map[string]*prWatch, len(refs))
baselines := make([]PRState, 0, len(refs)) baselines := make([]PRState, 0, len(refs))
@@ -293,11 +297,12 @@ func Watch(f StateFetcher, refs []PRRef, agentLogin string, ticks <-chan time.Ti
if onBaseline != nil { if onBaseline != nil {
onBaseline(baselines) onBaseline(baselines)
} }
latest := append([]PRState(nil), baselines...)
fails := make(map[string]int, len(refs)) fails := make(map[string]int, len(refs))
// The tick carries the time it fired, which is the clock the conflict // The tick carries the time it fired, which is the clock the conflict
// window is measured on. // window is measured on.
for now := range ticks { for now := range ticks {
for _, ref := range refs { for i, ref := range refs {
key := ref.String() key := ref.String()
cur, err := f.FetchState(ref, agentLogin) cur, err := f.FetchState(ref, agentLogin)
if err != nil { if err != nil {
@@ -315,12 +320,13 @@ func Watch(f StateFetcher, refs []PRRef, agentLogin string, ticks <-chan time.Ti
continue continue
} }
fails[key] = 0 fails[key] = 0
latest[i] = cur
if changed, reason := watches[key].observe(cur, now); changed { if changed, reason := watches[key].observe(cur, now); changed {
return WatchResult{Ref: ref, Reason: reason, State: cur}, nil return WatchResult{Ref: ref, Reason: reason, State: cur}, nil
} }
} }
} }
return WatchResult{}, nil return WatchResult{TimedOut: true, States: latest}, nil
} }
// countNonAgentComments counts comments authored by anyone other than agentLogin. // countNonAgentComments counts comments authored by anyone other than agentLogin.
+20
View File
@@ -1389,3 +1389,23 @@ func TestWatchConflictWindowCountsFromTheZeroTime(t *testing.T) {
t.Fatalf("reason = %q, want the mergeability loss; a run starting at the zero time still counts", res.Reason) t.Fatalf("reason = %q, want the mergeability loss; a run starting at the zero time still counts", res.Reason)
} }
} }
// Closed ticks with no change end the watch as a timeout carrying the latest
// state seen, not the baseline.
func TestWatchTimesOutWithLatestStates(t *testing.T) {
open := base()
pushed := base()
pushed.HeadSHA = "def456"
f := &fakeFetcher{states: []PRState{open, pushed}}
res, err := Watch(f, []PRRef{open.Ref}, "unkin-agent", drainableTicks(2), nil, nil)
if err != nil {
t.Fatalf("Watch: %v", err)
}
if !res.TimedOut {
t.Fatalf("TimedOut = false, want true")
}
if len(res.States) != 1 || res.States[0].HeadSHA != "def456" {
t.Errorf("States = %+v, want the latest polled state", res.States)
}
}