runs/log/wait/wait-quiet/release against the Gitea API. Log fetch tries API (token) -> public web route -> ssh+zstd file storage (hex-bucket path verified against the real Pi5 layout). wait-quiet encodes the "serialize heavy builds" rule with private repos as warnings, not failures. Injectable fetchers make the wait logic testable. Spec: doc/tool-parity.md 4.1. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011KikHkfCiC3yELbsMN8fT9
281 lines
7.7 KiB
Go
281 lines
7.7 KiB
Go
// Package gitea wraps the slice of Gitea's REST API that agents need
|
|
// daily: Actions runs, job logs and releases — plus the wait logic
|
|
// that encodes the homelab's "serialize heavy builds" rule.
|
|
package gitea
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"net/url"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// Task is one Actions job as returned by /actions/tasks. Gitea calls
|
|
// these tasks; each run can hold several.
|
|
type Task struct {
|
|
ID int64 `json:"id"`
|
|
Name string `json:"name"`
|
|
HeadBranch string `json:"head_branch"`
|
|
RunNumber int64 `json:"run_number"`
|
|
Status string `json:"status"`
|
|
DisplayTitle string `json:"display_title"`
|
|
WorkflowID string `json:"workflow_id"`
|
|
URL string `json:"url"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
UpdatedAt time.Time `json:"updated_at"`
|
|
RunStartedAt time.Time `json:"run_started_at"`
|
|
}
|
|
|
|
// Active reports whether the task still occupies (or will occupy) the
|
|
// runner.
|
|
func (t Task) Active() bool {
|
|
switch t.Status {
|
|
case "running", "waiting", "blocked":
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
// OK reports whether the task ended well (skipped counts as ok).
|
|
func (t Task) OK() bool {
|
|
return t.Status == "success" || t.Status == "skipped"
|
|
}
|
|
|
|
func (t Task) Duration() time.Duration {
|
|
if t.RunStartedAt.IsZero() || t.UpdatedAt.IsZero() || t.Active() {
|
|
return 0
|
|
}
|
|
d := t.UpdatedAt.Sub(t.RunStartedAt)
|
|
if d < 0 {
|
|
return 0
|
|
}
|
|
return d
|
|
}
|
|
|
|
// Release is a Gitea release with its downloadable assets.
|
|
type Release struct {
|
|
TagName string `json:"tag_name"`
|
|
Name string `json:"name"`
|
|
PublishedAt time.Time `json:"published_at"`
|
|
Assets []struct {
|
|
Name string `json:"name"`
|
|
Size int64 `json:"size"`
|
|
DownloadURL string `json:"browser_download_url"`
|
|
} `json:"assets"`
|
|
}
|
|
|
|
type Client struct {
|
|
BaseURL string // e.g. https://gitea.brasse-pc.eu
|
|
Owner string
|
|
Token string // optional; required for private repos
|
|
HTTP *http.Client
|
|
}
|
|
|
|
func New(baseURL, owner, token string) *Client {
|
|
return &Client{
|
|
BaseURL: strings.TrimRight(baseURL, "/"),
|
|
Owner: owner,
|
|
Token: token,
|
|
HTTP: &http.Client{Timeout: 60 * time.Second},
|
|
}
|
|
}
|
|
|
|
func (c *Client) get(path string) (*http.Response, error) {
|
|
req, err := http.NewRequest("GET", c.BaseURL+path, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if c.Token != "" {
|
|
req.Header.Set("Authorization", "token "+c.Token)
|
|
}
|
|
resp, err := c.HTTP.Do(req)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if resp.StatusCode == 404 {
|
|
resp.Body.Close()
|
|
return nil, fmt.Errorf("not found (private repo without token in the config?)")
|
|
}
|
|
if resp.StatusCode >= 300 {
|
|
body, _ := io.ReadAll(io.LimitReader(resp.Body, 2048))
|
|
resp.Body.Close()
|
|
return nil, fmt.Errorf("gitea answered %s: %s", resp.Status, strings.TrimSpace(string(body)))
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
// Tasks lists a repo's Actions jobs, newest first.
|
|
func (c *Client) Tasks(repo string, limit int) ([]Task, error) {
|
|
resp, err := c.get(fmt.Sprintf("/api/v1/repos/%s/%s/actions/tasks?limit=%d",
|
|
url.PathEscape(c.Owner), url.PathEscape(repo), limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer resp.Body.Close()
|
|
var out struct {
|
|
WorkflowRuns []Task `json:"workflow_runs"`
|
|
}
|
|
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
|
return nil, err
|
|
}
|
|
sort.Slice(out.WorkflowRuns, func(i, j int) bool {
|
|
return out.WorkflowRuns[i].ID > out.WorkflowRuns[j].ID
|
|
})
|
|
if limit > 0 && len(out.WorkflowRuns) > limit {
|
|
out.WorkflowRuns = out.WorkflowRuns[:limit]
|
|
}
|
|
return out.WorkflowRuns, nil
|
|
}
|
|
|
|
// LogAPI fetches a job log through the token-authenticated API.
|
|
func (c *Client) LogAPI(repo string, taskID int64) (string, error) {
|
|
resp, err := c.get(fmt.Sprintf("/api/v1/repos/%s/%s/actions/jobs/%d/logs",
|
|
url.PathEscape(c.Owner), url.PathEscape(repo), taskID))
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer resp.Body.Close()
|
|
data, err := io.ReadAll(io.LimitReader(resp.Body, 64<<20))
|
|
return string(data), err
|
|
}
|
|
|
|
// LogWeb fetches a job log through the web route, which works
|
|
// anonymously for public repos (job = index within the run, not id).
|
|
func (c *Client) LogWeb(repo string, runNumber int64, jobIndex int) (string, error) {
|
|
resp, err := c.get(fmt.Sprintf("/%s/%s/actions/runs/%d/jobs/%d/logs",
|
|
url.PathEscape(c.Owner), url.PathEscape(repo), runNumber, jobIndex))
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer resp.Body.Close()
|
|
data, err := io.ReadAll(io.LimitReader(resp.Body, 64<<20))
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if strings.HasPrefix(strings.TrimSpace(string(data)), "<") {
|
|
return "", fmt.Errorf("got HTML instead of a log (login page? private repo needs a token)")
|
|
}
|
|
return string(data), nil
|
|
}
|
|
|
|
// ZstPath is where Gitea's file storage keeps a task's compressed log
|
|
// on the server: <dir>/<owner>/<repo>/<hex(taskID%256)>/<taskID>.log.zst.
|
|
// Used by the ssh fallback for logs the HTTP routes no longer serve.
|
|
func ZstPath(dir, owner, repo string, taskID int64) string {
|
|
return fmt.Sprintf("%s/%s/%s/%02x/%d.log.zst", dir, owner, repo, taskID%256, taskID)
|
|
}
|
|
|
|
// Releases lists a repo's releases, newest first.
|
|
func (c *Client) Releases(repo string, limit int) ([]Release, error) {
|
|
resp, err := c.get(fmt.Sprintf("/api/v1/repos/%s/%s/releases?limit=%d",
|
|
url.PathEscape(c.Owner), url.PathEscape(repo), limit))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer resp.Body.Close()
|
|
var out []Release
|
|
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
|
return nil, err
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// --- wait logic (injectable fetch so it is testable) ---
|
|
|
|
type TaskFetcher func(repo string) ([]Task, error)
|
|
|
|
type WaitResult struct {
|
|
Done bool // latest run finished within the timeout
|
|
AllOK bool
|
|
Tasks []Task // the latest run's tasks (or last seen)
|
|
Attempts int
|
|
}
|
|
|
|
// WaitRun polls until every task of repo's newest run has finished.
|
|
func WaitRun(repo string, fetch TaskFetcher, interval, timeout time.Duration, log func(string)) (WaitResult, error) {
|
|
start := time.Now()
|
|
res := WaitResult{}
|
|
for {
|
|
res.Attempts++
|
|
tasks, err := fetch(repo)
|
|
if err != nil {
|
|
return res, err
|
|
}
|
|
if len(tasks) == 0 {
|
|
return res, fmt.Errorf("repo has no Actions runs")
|
|
}
|
|
latest := tasks[0].RunNumber
|
|
var runTasks []Task
|
|
active := false
|
|
allOK := true
|
|
for _, t := range tasks {
|
|
if t.RunNumber != latest {
|
|
continue
|
|
}
|
|
runTasks = append(runTasks, t)
|
|
if t.Active() {
|
|
active = true
|
|
} else if !t.OK() {
|
|
allOK = false
|
|
}
|
|
}
|
|
res.Tasks = runTasks
|
|
if !active {
|
|
res.Done, res.AllOK = true, allOK
|
|
return res, nil
|
|
}
|
|
if log != nil {
|
|
log(fmt.Sprintf("run #%d still active (%d jobs), waiting...", latest, len(runTasks)))
|
|
}
|
|
if time.Since(start)+interval > timeout {
|
|
return res, nil
|
|
}
|
|
time.Sleep(interval)
|
|
}
|
|
}
|
|
|
|
// CountActive counts active jobs across the given repos.
|
|
func CountActive(repos []string, fetch TaskFetcher) (int, []string, error) {
|
|
count := 0
|
|
var busy []string
|
|
for _, r := range repos {
|
|
tasks, err := fetch(r)
|
|
if err != nil {
|
|
return 0, nil, fmt.Errorf("%s: %w", r, err)
|
|
}
|
|
for _, t := range tasks {
|
|
if t.Active() {
|
|
count++
|
|
busy = append(busy, fmt.Sprintf("%s#%d(%s)", r, t.RunNumber, t.Status))
|
|
}
|
|
}
|
|
}
|
|
return count, busy, nil
|
|
}
|
|
|
|
// WaitQuiet blocks until at most maxActive jobs run across repos —
|
|
// the "serialize pushes, >2 heavy builds take the Pi5 down" rule.
|
|
func WaitQuiet(repos []string, maxActive int, fetch TaskFetcher, interval, timeout time.Duration, log func(string)) (bool, error) {
|
|
start := time.Now()
|
|
for {
|
|
n, busy, err := CountActive(repos, fetch)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if n <= maxActive {
|
|
return true, nil
|
|
}
|
|
if log != nil {
|
|
log(fmt.Sprintf("%d active builds (max %d): %s", n, maxActive, strings.Join(busy, " ")))
|
|
}
|
|
if time.Since(start)+interval > timeout {
|
|
return false, nil
|
|
}
|
|
time.Sleep(interval)
|
|
}
|
|
}
|