build log --follow, streamed to the build page !462
27 files changed, +2118 −38
Layout: unified · split
.gitbay/wiki/Admin.org +4
| @@ -81,6 +81,10 @@ rate limiting, to the last =X-Forwarded-For= hop that is not itself a | ||
| 81 | 81 | trusted proxy; from anyone else the header is ignored. Empty, the |
| 82 | 82 | default, is right when gitbayd terminates TLS itself. |
| 83 | 83 | |
| 84 | A reverse proxy in front of gitbay must not buffer responses, or the | |
| 85 | build page's live log arrives only when the build ends; | |
| 86 | =X-Accel-Buffering: no= covers nginx. | |
| 87 | ||
| 84 | 88 | ** [web] |
| 85 | 89 | - =mode= — =view_only= (default) | =accounts=. In view_only the mutating |
| 86 | 90 | web routes are never registered; in accounts, browser sessions are |
.gitbay/wiki/CI.org +11
| @@ -38,6 +38,17 @@ Scheduled jobs run on their cron against the default branch, never on | ||
| 38 | 38 | push; a default-branch push registers or updates them. Tag jobs run on |
| 39 | 39 | a matching tag push and nothing else. |
| 40 | 40 | |
| 41 | A running build is followed with =build log <owner/name> <n> --follow=: | |
| 42 | the stored log, then output as the runner sends it, then the outcome | |
| 43 | as =build <n> <status>= on stderr once the build ends. The exit code is | |
| 44 | 0 whatever the outcome. A follow of a build still queued after ten | |
| 45 | minutes ends on its own and says so on stderr, since nothing reaps a | |
| 46 | queued build. The build page does the same without JavaScript while a | |
| 47 | build is queued or running; =?follow=0= renders it once. An account | |
| 48 | holds at most eight follows open, and signed-out viewers share one | |
| 49 | account's eight. Over the JSON API the command answers when the build | |
| 50 | ends, with the whole log. | |
| 51 | ||
| 41 | 52 | * The table |
| 42 | 53 | |
| 43 | 54 | One =ci.yml= with five jobs, each isolating one rule: |
.gitbay/wiki/Parity.org +1
| @@ -210,6 +210,7 @@ rather than the one the web page shows. | ||
| 210 | 210 | | build list filters (ref, status, job) | yes | yes | yes | |
| 211 | 211 | | build show (one build) | yes | yes | yes | |
| 212 | 212 | | build log | yes | yes | yes | |
| 213 | | build log follow (until it ends) | yes | yes | no | | |
| 213 | 214 | | build jobs | yes | yes | yes | |
| 214 | 215 | | build trigger | yes | yes | yes | |
| 215 | 216 | | build cancel | yes | yes | yes | |
cmd/gitbay/main.go +1 −1
| @@ -48,7 +48,7 @@ func newRoot() *cobra.Command { | ||
| 48 | 48 | group("build", "CI builds", |
| 49 | 49 | pass("list", "recent builds: <owner/name>", passOpts{server: []string{"build", "list"}, needsRepo: true}), |
| 50 | 50 | pass("show", "one build: <owner/name> <n>", passOpts{server: []string{"build", "show"}, needsRepo: true}), |
| 51 | pass("log", "a build's log: <owner/name> <n>", passOpts{server: []string{"build", "log"}, needsRepo: true}), | |
| 51 | pass("log", "a build's log: <owner/name> <n> [--follow]", passOpts{server: []string{"build", "log"}, needsRepo: true}), | |
| 52 | 52 | pass("jobs", "list the jobs a trigger can name", passOpts{server: []string{"build", "jobs"}, needsRepo: true}), |
| 53 | 53 | pass("trigger", "queue a job now: <job>", passOpts{server: []string{"build", "trigger"}, needsRepo: true}), |
| 54 | 54 | pass("cancel", "withdraw a queued build: <n>", passOpts{server: []string{"build", "cancel"}, needsRepo: true}), |
cmd/gitbayd/system.go +1 −1
| @@ -93,7 +93,7 @@ func shellCmd() *cobra.Command { | ||
| 93 | 93 | fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n") |
| 94 | 94 | os.Exit(protocol.ExitUsage) |
| 95 | 95 | } |
| 96 | code := sshd.Exec(cfg, st, user, key.Scope, key.Fingerprint, cmdline, os.Stdin, os.Stdout, os.Stderr) | |
| 96 | code := sshd.Exec(cfg, st, user, key.Scope, key.Fingerprint, cmdline, os.Stdin, os.Stdout, os.Stderr, nil) | |
| 97 | 97 | st.Close() |
| 98 | 98 | os.Exit(code) |
| 99 | 99 | return nil |
docs/plans/2026-09-23-build-log-follow.md added +1100
| @@ -0,0 +1,1100 @@ | ||
| 1 | # build log --follow Implementation Plan | |
| 2 | ||
| 3 | > **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. | |
| 4 | ||
| 5 | **Goal:** `build log <owner/name> <n> --follow` streams a build's log until the build has an outcome, and the web build page streams the same command without JavaScript. | |
| 6 | ||
| 7 | **Architecture:** The store wakes waiters when a build's row changes and reads the log from a byte offset. The control command loops read → write → wait. The web handler renders `build.html` with a marker where the log goes, writes the part before it, dispatches the command with a writer that HTML-escapes and flushes, then writes the rest. | |
| 8 | ||
| 9 | **Tech Stack:** Go, SQLite (modernc), `golang.org/x/crypto/ssh`, `html/template`, `net/http`. | |
| 10 | ||
| 11 | **Spec:** `docs/specs/2026-09-23-build-log-follow-design.md` | |
| 12 | ||
| 13 | ## Global Constraints | |
| 14 | ||
| 15 | - Work in `/Users/cmc/git/krz/gitbay-follow` (branch `build-log-follow`). Never touch `/Users/cmc/git/krz/gitbay`. | |
| 16 | - Every commit is signed (the repo signs by config; do not pass `--no-gpg-sign`). Messages reference `Ref #250`; the last commit says `Closes #250`. No attribution lines of any kind. | |
| 17 | - Write files with the editor tool, not heredocs. | |
| 18 | - Locally run: `go build ./...`, `go vet` on touched packages, unit tests of touched packages, and only the one e2e test being written. CI runs the rest. | |
| 19 | - Comments: plain, factual, match the surrounding density. No before/after narration. | |
| 20 | - Follow cap: 8 per account. Fallback re-read: 2 seconds. Settle after an outcome: 1 second. | |
| 21 | - The outcome line on stderr is exactly `build <n> <status>`; exit 0 whatever the outcome. | |
| 22 | - The web outcome line is exactly `<p class="notice" role="status">build finished: <status></p>`. | |
| 23 | ||
| 24 | --- | |
| 25 | ||
| 26 | ### Task 1: Store — wake followers and read from an offset | |
| 27 | ||
| 28 | **Files:** | |
| 29 | - Modify: `internal/store/store.go` (the `Store` struct, ~line 22) | |
| 30 | - Modify: `internal/store/builds.go` (`AppendBuildLog` ~228, `FinishBuild` ~248, `CancelBuild` ~393; new functions after `BuildLog` ~325) | |
| 31 | - Test: `internal/store/builds_test.go` | |
| 32 | ||
| 33 | **Interfaces:** | |
| 34 | - Produces: `func (s *Store) BuildLogWait(id int64) <-chan struct{}`; `func (s *Store) BuildLogFrom(id, offset int64) (status string, chunk []byte, err error)`; unexported `func (s *Store) wakeBuild(id int64)`. | |
| 35 | ||
| 36 | - [ ] **Step 1: Write the failing tests** — append to `internal/store/builds_test.go`: | |
| 37 | ||
| 38 | ```go | |
| 39 | // A follower's channel closes on each kind of change to its build, and | |
| 40 | // only its build. | |
| 41 | func TestBuildLogWaitWakes(t *testing.T) { | |
| 42 | s := open(t) | |
| 43 | if err := s.MigrateUp(); err != nil { | |
| 44 | t.Fatal(err) | |
| 45 | } | |
| 46 | uid, err := s.CreateUser("cmc", true) | |
| 47 | if err != nil { | |
| 48 | t.Fatal(err) | |
| 49 | } | |
| 50 | if _, err := s.CreateRepo("user", uid, "orgo", "public"); err != nil { | |
| 51 | t.Fatal(err) | |
| 52 | } | |
| 53 | newBuild := func() int64 { | |
| 54 | t.Helper() | |
| 55 | id, err := s.CreateBuild(1, "test", "abc123", "main", `["true"]`, "", "", true) | |
| 56 | if err != nil { | |
| 57 | t.Fatal(err) | |
| 58 | } | |
| 59 | return id | |
| 60 | } | |
| 61 | closed := func(ch <-chan struct{}) bool { | |
| 62 | select { | |
| 63 | case <-ch: | |
| 64 | return true | |
| 65 | default: | |
| 66 | return false | |
| 67 | } | |
| 68 | } | |
| 69 | ||
| 70 | a, b := newBuild(), newBuild() | |
| 71 | wa, wb := s.BuildLogWait(a), s.BuildLogWait(b) | |
| 72 | if err := s.AppendBuildLog(a, []byte("x")); err != nil { | |
| 73 | t.Fatal(err) | |
| 74 | } | |
| 75 | if !closed(wa) { | |
| 76 | t.Error("append did not wake its build") | |
| 77 | } | |
| 78 | if closed(wb) { | |
| 79 | t.Error("append woke another build") | |
| 80 | } | |
| 81 | ||
| 82 | mustClaim(t, s, 1) // claims a, the oldest | |
| 83 | wa = s.BuildLogWait(a) | |
| 84 | if err := s.FinishBuild(a, "success"); err != nil { | |
| 85 | t.Fatal(err) | |
| 86 | } | |
| 87 | if !closed(wa) { | |
| 88 | t.Error("finish did not wake") | |
| 89 | } | |
| 90 | ||
| 91 | if err := s.CancelBuild(b); err != nil { | |
| 92 | t.Fatal(err) | |
| 93 | } | |
| 94 | if !closed(wb) { | |
| 95 | t.Error("cancel did not wake") | |
| 96 | } | |
| 97 | } | |
| 98 | ||
| 99 | // Offsets are bytes, not characters: || stores the log as text, and a | |
| 100 | // multibyte character must not shift where the next read starts. | |
| 101 | func TestBuildLogFrom(t *testing.T) { | |
| 102 | s := open(t) | |
| 103 | if err := s.MigrateUp(); err != nil { | |
| 104 | t.Fatal(err) | |
| 105 | } | |
| 106 | uid, err := s.CreateUser("cmc", true) | |
| 107 | if err != nil { | |
| 108 | t.Fatal(err) | |
| 109 | } | |
| 110 | if _, err := s.CreateRepo("user", uid, "orgo", "public"); err != nil { | |
| 111 | t.Fatal(err) | |
| 112 | } | |
| 113 | id, err := s.CreateBuild(1, "test", "abc123", "main", `["true"]`, "", "", true) | |
| 114 | if err != nil { | |
| 115 | t.Fatal(err) | |
| 116 | } | |
| 117 | first := "héllo — ok\n" | |
| 118 | for _, c := range []string{first, "wörld\n"} { | |
| 119 | if err := s.AppendBuildLog(id, []byte(c)); err != nil { | |
| 120 | t.Fatal(err) | |
| 121 | } | |
| 122 | } | |
| 123 | status, all, err := s.BuildLogFrom(id, 0) | |
| 124 | if err != nil || status != "pending" || string(all) != first+"wörld\n" { | |
| 125 | t.Fatalf("from 0: %q %q %v", status, all, err) | |
| 126 | } | |
| 127 | _, rest, err := s.BuildLogFrom(id, int64(len(first))) | |
| 128 | if err != nil || string(rest) != "wörld\n" { | |
| 129 | t.Fatalf("from %d: %q %v", len(first), rest, err) | |
| 130 | } | |
| 131 | _, none, err := s.BuildLogFrom(id, int64(len(all))) | |
| 132 | if err != nil || len(none) != 0 { | |
| 133 | t.Fatalf("from the end: %q %v", none, err) | |
| 134 | } | |
| 135 | if _, _, err := s.BuildLogFrom(9999, 0); err != ErrNotFound { | |
| 136 | t.Fatalf("missing build: %v", err) | |
| 137 | } | |
| 138 | } | |
| 139 | ``` | |
| 140 | ||
| 141 | - [ ] **Step 2: Run to see them fail** | |
| 142 | ||
| 143 | Run: `go test ./internal/store/ -run 'TestBuildLogWaitWakes|TestBuildLogFrom' -count=1` | |
| 144 | Expected: build failure, `s.BuildLogWait undefined`. | |
| 145 | ||
| 146 | - [ ] **Step 3: Implement** | |
| 147 | ||
| 148 | In `internal/store/store.go`, add `"sync"` to the imports and the fields: | |
| 149 | ||
| 150 | ```go | |
| 151 | type Store struct { | |
| 152 | DB *sql.DB | |
| 153 | ||
| 154 | // logWait holds one channel per build someone is following, closed | |
| 155 | // by the next change to that build's row (BuildLogWait). | |
| 156 | logMu sync.Mutex | |
| 157 | logWait map[int64]chan struct{} | |
| 158 | } | |
| 159 | ``` | |
| 160 | ||
| 161 | In `internal/store/builds.go`, after `BuildLog`: | |
| 162 | ||
| 163 | ```go | |
| 164 | // BuildLogWait returns a channel closed by the next append to, finish of | |
| 165 | // or cancel of the build. Take it before reading, so a change between the | |
| 166 | // read and the wait still wakes the reader. Only this process's writes | |
| 167 | // wake it. | |
| 168 | func (s *Store) BuildLogWait(id int64) <-chan struct{} { | |
| 169 | s.logMu.Lock() | |
| 170 | defer s.logMu.Unlock() | |
| 171 | if s.logWait == nil { | |
| 172 | s.logWait = map[int64]chan struct{}{} | |
| 173 | } | |
| 174 | ch, ok := s.logWait[id] | |
| 175 | if !ok { | |
| 176 | ch = make(chan struct{}) | |
| 177 | s.logWait[id] = ch | |
| 178 | } | |
| 179 | return ch | |
| 180 | } | |
| 181 | ||
| 182 | func (s *Store) wakeBuild(id int64) { | |
| 183 | s.logMu.Lock() | |
| 184 | defer s.logMu.Unlock() | |
| 185 | if ch, ok := s.logWait[id]; ok { | |
| 186 | close(ch) | |
| 187 | delete(s.logWait, id) | |
| 188 | } | |
| 189 | } | |
| 190 | ||
| 191 | // BuildLogFrom returns the build's status and its log past offset bytes, | |
| 192 | // read together so a terminal status comes with every byte before it. | |
| 193 | // The cast matters: || stores the log as text, and substr on text counts | |
| 194 | // characters. | |
| 195 | func (s *Store) BuildLogFrom(id, offset int64) (string, []byte, error) { | |
| 196 | var status string | |
| 197 | var chunk []byte | |
| 198 | err := s.DB.QueryRow(`SELECT status, substr(CAST(log AS BLOB), ?) FROM builds WHERE id = ?`, | |
| 199 | offset+1, id).Scan(&status, &chunk) | |
| 200 | if errors.Is(err, sql.ErrNoRows) { | |
| 201 | return "", nil, ErrNotFound | |
| 202 | } | |
| 203 | return status, chunk, err | |
| 204 | } | |
| 205 | ``` | |
| 206 | ||
| 207 | Wake after each successful write. In `AppendBuildLog`, replace the tail so both paths wake: | |
| 208 | ||
| 209 | ```go | |
| 210 | if n, _ := res.RowsAffected(); n > 0 { | |
| 211 | s.wakeBuild(id) | |
| 212 | return nil | |
| 213 | } | |
| 214 | // Over the cap. The bounds match exactly once: appending the notice puts | |
| 215 | // the log past the upper bound, so later chunks fall through silently. | |
| 216 | _, err = s.DB.Exec(` | |
| 217 | UPDATE builds SET log = log || ? | |
| 218 | WHERE id = ? AND length(log) >= ? AND length(log) < ?`, | |
| 219 | truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice)) | |
| 220 | if err == nil { | |
| 221 | s.wakeBuild(id) | |
| 222 | } | |
| 223 | return err | |
| 224 | ``` | |
| 225 | ||
| 226 | In `FinishBuild`, before the final `return nil`: `s.wakeBuild(id)`. In `CancelBuild`, the same, before its final `return nil`. | |
| 227 | ||
| 228 | - [ ] **Step 4: Run the tests and vet** | |
| 229 | ||
| 230 | Run: `go test ./internal/store/ -count=1 && go vet ./internal/store/` | |
| 231 | Expected: `ok`, vet silent (no copylocks: `Store` is only ever `&Store{...}` in `Open`). | |
| 232 | ||
| 233 | - [ ] **Step 5: Commit** | |
| 234 | ||
| 235 | ```bash | |
| 236 | git add internal/store/store.go internal/store/builds.go internal/store/builds_test.go | |
| 237 | git commit -m "store: wake build log followers, read the log from an offset | |
| 238 | ||
| 239 | Ref #250" | |
| 240 | ``` | |
| 241 | ||
| 242 | --- | |
| 243 | ||
| 244 | ### Task 2: Control — `build log --follow` | |
| 245 | ||
| 246 | **Files:** | |
| 247 | - Modify: `internal/control/control.go` (`Ctx`, ~line 21) | |
| 248 | - Modify: `internal/control/build.go` (registration ~29, `runBuildLog` ~198) | |
| 249 | - Create: `internal/control/buildfollow.go` | |
| 250 | - Create: `internal/control/buildfollow_test.go` | |
| 251 | - Modify: `cmd/gitbay/main.go:51` (help string) | |
| 252 | ||
| 253 | **Interfaces:** | |
| 254 | - Consumes: `Store.BuildLogWait`, `Store.BuildLogFrom` (Task 1). | |
| 255 | - Produces: `Ctx.Done <-chan struct{}`; the command `build log <owner/name> <n> [--follow]`; package vars `followSettle`, `followPoll` (tests shorten them); `maxFollows = 8`. | |
| 256 | ||
| 257 | - [ ] **Step 1: Write the failing tests** — `internal/control/buildfollow_test.go`: | |
| 258 | ||
| 259 | ```go | |
| 260 | package control | |
| 261 | ||
| 262 | import ( | |
| 263 | "bytes" | |
| 264 | "strings" | |
| 265 | "testing" | |
| 266 | "time" | |
| 267 | ||
| 268 | "gitbay.org/gitbay/internal/protocol" | |
| 269 | "gitbay.org/gitbay/internal/store" | |
| 270 | ) | |
| 271 | ||
| 272 | // follow starts build log --follow on build 1 of repo and returns the | |
| 273 | // buffers and a channel carrying the exit code. | |
| 274 | func follow(t *testing.T, st *store.Store, uid int64, repo store.Repo, done <-chan struct{}) (*bytes.Buffer, *bytes.Buffer, chan int) { | |
| 275 | t.Helper() | |
| 276 | u, err := st.UserByID(uid) | |
| 277 | if err != nil { | |
| 278 | t.Fatal(err) | |
| 279 | } | |
| 280 | var out, errOut bytes.Buffer | |
| 281 | c := &Ctx{User: u, Scope: "full", Store: st, Stdin: strings.NewReader(""), | |
| 282 | Stdout: &out, Stderr: &errOut, Done: done} | |
| 283 | res := make(chan int, 1) | |
| 284 | go func() { res <- Dispatch(c, []string{"build", "log", repo.Path(), "1", "--follow"}) }() | |
| 285 | return &out, &errOut, res | |
| 286 | } | |
| 287 | ||
| 288 | func waitExit(t *testing.T, res chan int) int { | |
| 289 | t.Helper() | |
| 290 | select { | |
| 291 | case code := <-res: | |
| 292 | return code | |
| 293 | case <-time.After(10 * time.Second): | |
| 294 | t.Fatal("follow did not end") | |
| 295 | return -1 | |
| 296 | } | |
| 297 | } | |
| 298 | ||
| 299 | func shortFollowTimers(t *testing.T) { | |
| 300 | settle, poll := followSettle, followPoll | |
| 301 | followSettle, followPoll = 200*time.Millisecond, 50*time.Millisecond | |
| 302 | t.Cleanup(func() { followSettle, followPoll = settle, poll }) | |
| 303 | } | |
| 304 | ||
| 305 | // The follow prints the stored log, then what arrives, and ends with the | |
| 306 | // outcome on stderr once the build finishes. | |
| 307 | func TestBuildLogFollow(t *testing.T) { | |
| 308 | shortFollowTimers(t) | |
| 309 | st, repo, uid := newQueueTestRepo(t) | |
| 310 | id, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true) | |
| 311 | if err != nil { | |
| 312 | t.Fatal(err) | |
| 313 | } | |
| 314 | st.AppendBuildLog(id, []byte("queued\n")) | |
| 315 | out, errOut, res := follow(t, st, uid, repo, nil) | |
| 316 | ||
| 317 | if _, ok, err := st.ClaimBuild([]int64{repo.ID}, false); err != nil || !ok { | |
| 318 | t.Fatalf("claim: %v %v", ok, err) | |
| 319 | } | |
| 320 | st.AppendBuildLog(id, []byte("step one\n")) | |
| 321 | st.AppendBuildLog(id, []byte("step two\n")) | |
| 322 | if err := st.FinishBuild(id, "success"); err != nil { | |
| 323 | t.Fatal(err) | |
| 324 | } | |
| 325 | if code := waitExit(t, res); code != protocol.ExitOK { | |
| 326 | t.Fatalf("exit %d: %s", code, errOut) | |
| 327 | } | |
| 328 | if got := out.String(); got != "queued\nstep one\nstep two\n" { | |
| 329 | t.Errorf("stdout %q", got) | |
| 330 | } | |
| 331 | if got := strings.TrimSpace(errOut.String()); got != "build 1 success" { | |
| 332 | t.Errorf("stderr %q", got) | |
| 333 | } | |
| 334 | } | |
| 335 | ||
| 336 | // A cancel ends the follow, and the line the cancel appends after the | |
| 337 | // status change still arrives. | |
| 338 | func TestBuildLogFollowCancel(t *testing.T) { | |
| 339 | shortFollowTimers(t) | |
| 340 | st, repo, uid := newQueueTestRepo(t) | |
| 341 | id, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true) | |
| 342 | if err != nil { | |
| 343 | t.Fatal(err) | |
| 344 | } | |
| 345 | out, errOut, res := follow(t, st, uid, repo, nil) | |
| 346 | if err := st.CancelBuild(id); err != nil { | |
| 347 | t.Fatal(err) | |
| 348 | } | |
| 349 | st.AppendBuildLog(id, []byte("cancelled by alice before a runner claimed it\n")) | |
| 350 | if code := waitExit(t, res); code != protocol.ExitOK { | |
| 351 | t.Fatalf("exit %d: %s", code, errOut) | |
| 352 | } | |
| 353 | if !strings.Contains(out.String(), "cancelled by alice") { | |
| 354 | t.Errorf("the cancel line did not arrive: %q", out) | |
| 355 | } | |
| 356 | if got := strings.TrimSpace(errOut.String()); got != "build 1 cancelled" { | |
| 357 | t.Errorf("stderr %q", got) | |
| 358 | } | |
| 359 | } | |
| 360 | ||
| 361 | // Closing Done ends a follow of a build that is still running. | |
| 362 | func TestBuildLogFollowDone(t *testing.T) { | |
| 363 | shortFollowTimers(t) | |
| 364 | st, repo, uid := newQueueTestRepo(t) | |
| 365 | if _, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true); err != nil { | |
| 366 | t.Fatal(err) | |
| 367 | } | |
| 368 | done := make(chan struct{}) | |
| 369 | _, _, res := follow(t, st, uid, repo, done) | |
| 370 | close(done) | |
| 371 | if code := waitExit(t, res); code != protocol.ExitFailure { | |
| 372 | t.Fatalf("exit %d, want %d", code, protocol.ExitFailure) | |
| 373 | } | |
| 374 | } | |
| 375 | ||
| 376 | // An account holding maxFollows is refused another. | |
| 377 | func TestBuildLogFollowCap(t *testing.T) { | |
| 378 | st, repo, uid := newQueueTestRepo(t) | |
| 379 | if _, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true); err != nil { | |
| 380 | t.Fatal(err) | |
| 381 | } | |
| 382 | followMu.Lock() | |
| 383 | follows[uid] = maxFollows | |
| 384 | followMu.Unlock() | |
| 385 | t.Cleanup(func() { | |
| 386 | followMu.Lock() | |
| 387 | delete(follows, uid) | |
| 388 | followMu.Unlock() | |
| 389 | }) | |
| 390 | _, errOut, res := follow(t, st, uid, repo, nil) | |
| 391 | if code := waitExit(t, res); code != protocol.ExitDenied { | |
| 392 | t.Fatalf("exit %d, want %d", code, protocol.ExitDenied) | |
| 393 | } | |
| 394 | if !strings.Contains(errOut.String(), "8 follows are already open") { | |
| 395 | t.Errorf("stderr %q", errOut) | |
| 396 | } | |
| 397 | } | |
| 398 | ``` | |
| 399 | ||
| 400 | - [ ] **Step 2: Run to see them fail** | |
| 401 | ||
| 402 | Run: `go test ./internal/control/ -run 'TestBuildLogFollow' -count=1` | |
| 403 | Expected: build failure, `unknown field Done` / `undefined: followSettle`. | |
| 404 | ||
| 405 | - [ ] **Step 3: Implement** | |
| 406 | ||
| 407 | `internal/control/control.go`, in `Ctx` after `Cmd`: | |
| 408 | ||
| 409 | ```go | |
| 410 | // Done, when the surface has one, closes when nobody is reading any | |
| 411 | // more: the SSH channel closed or the HTTP request ended. A command | |
| 412 | // that runs until something happens (build log --follow) stops on it. | |
| 413 | Done <-chan struct{} | |
| 414 | ``` | |
| 415 | ||
| 416 | `internal/control/build.go` registration: | |
| 417 | ||
| 418 | ```go | |
| 419 | register(Command{Path: []string{"build", "log"}, | |
| 420 | Summary: "print a build's log, or follow it until the build ends", | |
| 421 | Usage: "build log <owner/name> <n> [--follow]", ReadOnly: true, Run: runBuildLog}) | |
| 422 | ``` | |
| 423 | ||
| 424 | `runBuildLog`: | |
| 425 | ||
| 426 | ```go | |
| 427 | func runBuildLog(c *Ctx, args []string) int { | |
| 428 | f, err := parseFlags(args, flagSpec{Bools: []string{"--follow"}, MaxPos: 2, Usage: c.Cmd.Usage}) | |
| 429 | if err != nil { | |
| 430 | return c.fail(protocol.ExitUsage, "%v", err) | |
| 431 | } | |
| 432 | _, b, code := buildRef(c, f.Pos) | |
| 433 | if code >= 0 { | |
| 434 | return code | |
| 435 | } | |
| 436 | if f.Has("--follow") { | |
| 437 | return followBuildLog(c, b) | |
| 438 | } | |
| 439 | log, err := c.Store.BuildLog(b.ID) | |
| 440 | if err != nil { | |
| 441 | return c.fail(protocol.ExitFailure, "%v", err) | |
| 442 | } | |
| 443 | c.Stdout.Write(log) | |
| 444 | return protocol.ExitOK | |
| 445 | } | |
| 446 | ``` | |
| 447 | ||
| 448 | Check how `runBuildList` (build.go ~125) reports a `parseFlags` error and match it exactly if it differs from the above. | |
| 449 | ||
| 450 | `internal/control/buildfollow.go`: | |
| 451 | ||
| 452 | ```go | |
| 453 | package control | |
| 454 | ||
| 455 | import ( | |
| 456 | "fmt" | |
| 457 | "sync" | |
| 458 | "time" | |
| 459 | ||
| 460 | "gitbay.org/gitbay/internal/protocol" | |
| 461 | "gitbay.org/gitbay/internal/store" | |
| 462 | ) | |
| 463 | ||
| 464 | // maxFollows is how many build log follows one account holds open at | |
| 465 | // once. Signed-out web viewers are account 0 and share it. | |
| 466 | const maxFollows = 8 | |
| 467 | ||
| 468 | var ( | |
| 469 | // followPoll bounds a wait with no wake. A write from another process | |
| 470 | // (gitbayd admin, or any session under gitbayd shell) wakes nobody; | |
| 471 | // this is how its bytes still arrive. | |
| 472 | followPoll = 2 * time.Second | |
| 473 | // followSettle is how long a follow keeps reading after the build has | |
| 474 | // an outcome: a cancel appends its line after the status changes, and | |
| 475 | // a cancelled runner's stream runs on until its next check. | |
| 476 | followSettle = time.Second | |
| 477 | ) | |
| 478 | ||
| 479 | var ( | |
| 480 | followMu sync.Mutex | |
| 481 | follows = map[int64]int{} | |
| 482 | ) | |
| 483 | ||
| 484 | func takeFollow(uid int64) bool { | |
| 485 | followMu.Lock() | |
| 486 | defer followMu.Unlock() | |
| 487 | if follows[uid] >= maxFollows { | |
| 488 | return false | |
| 489 | } | |
| 490 | follows[uid]++ | |
| 491 | return true | |
| 492 | } | |
| 493 | ||
| 494 | func dropFollow(uid int64) { | |
| 495 | followMu.Lock() | |
| 496 | defer followMu.Unlock() | |
| 497 | if follows[uid]--; follows[uid] <= 0 { | |
| 498 | delete(follows, uid) | |
| 499 | } | |
| 500 | } | |
| 501 | ||
| 502 | // followBuildLog writes the build's log as it grows and returns once the | |
| 503 | // build has an outcome and its last bytes are written. The outcome goes | |
| 504 | // to stderr, so stdout is the log byte for byte. | |
| 505 | func followBuildLog(c *Ctx, b store.Build) int { | |
| 506 | if !takeFollow(c.User.ID) { | |
| 507 | return c.fail(protocol.ExitDenied, "%d follows are already open for this account; close one and retry", maxFollows) | |
| 508 | } | |
| 509 | defer dropFollow(c.User.ID) | |
| 510 | ||
| 511 | var off int64 | |
| 512 | var settleBy time.Time | |
| 513 | for { | |
| 514 | wake := c.Store.BuildLogWait(b.ID) | |
| 515 | status, chunk, err := c.Store.BuildLogFrom(b.ID, off) | |
| 516 | if err != nil { | |
| 517 | return c.fail(protocol.ExitFailure, "%v", err) | |
| 518 | } | |
| 519 | if len(chunk) > 0 { | |
| 520 | if _, err := c.Stdout.Write(chunk); err != nil { | |
| 521 | return protocol.ExitFailure | |
| 522 | } | |
| 523 | off += int64(len(chunk)) | |
| 524 | } | |
| 525 | wait := followPoll | |
| 526 | if status != "pending" && status != "running" { | |
| 527 | if settleBy.IsZero() { | |
| 528 | settleBy = time.Now().Add(followSettle) | |
| 529 | } | |
| 530 | left := time.Until(settleBy) | |
| 531 | if left <= 0 && len(chunk) == 0 { | |
| 532 | fmt.Fprintf(c.Stderr, "build %d %s\n", b.Number, status) | |
| 533 | return protocol.ExitOK | |
| 534 | } | |
| 535 | wait = min(wait, max(left, 0)) | |
| 536 | } | |
| 537 | t := time.NewTimer(wait) | |
| 538 | select { | |
| 539 | case <-wake: | |
| 540 | case <-t.C: | |
| 541 | case <-c.Done: | |
| 542 | t.Stop() | |
| 543 | return protocol.ExitFailure | |
| 544 | } | |
| 545 | t.Stop() | |
| 546 | } | |
| 547 | } | |
| 548 | ``` | |
| 549 | ||
| 550 | Check the loop against the spec before moving on: once the status is terminal it keeps reading until `followSettle` has passed *and* a read came back empty, then prints the outcome. A deadline, not a `time.After` channel: a timer channel delivers once, and a second check of it would block. A nil `c.Done` never fires in the select, which is what a surface without one wants. `min`/`max` are Go 1.21 builtins; check `go.mod`'s go line is at least 1.21. | |
| 551 | ||
| 552 | `cmd/gitbay/main.go:51`: | |
| 553 | ||
| 554 | ```go | |
| 555 | pass("log", "a build's log: <owner/name> <n> [--follow]", passOpts{server: []string{"build", "log"}, needsRepo: true}), | |
| 556 | ``` | |
| 557 | ||
| 558 | - [ ] **Step 4: Run the tests, the race detector, and vet** | |
| 559 | ||
| 560 | Run: `go test ./internal/control/ -run 'TestBuildLog' -count=1 -race && go vet ./internal/control/ ./cmd/gitbay/ && go test ./cmd/gitbay/ -count=1` | |
| 561 | Expected: `ok` for each. | |
| 562 | ||
| 563 | Then run the whole control package once, since `build log` is covered elsewhere too: `go test ./internal/control/ -count=1`. | |
| 564 | ||
| 565 | - [ ] **Step 5: Commit** | |
| 566 | ||
| 567 | ```bash | |
| 568 | git add internal/control/control.go internal/control/build.go internal/control/buildfollow.go internal/control/buildfollow_test.go cmd/gitbay/main.go | |
| 569 | git commit -m "build log --follow: stream a build's log until it ends | |
| 570 | ||
| 571 | Ref #250" | |
| 572 | ``` | |
| 573 | ||
| 574 | --- | |
| 575 | ||
| 576 | ### Task 3: Surfaces pass Done — SSH channel close, HTTP request end | |
| 577 | ||
| 578 | **Files:** | |
| 579 | - Modify: `internal/sshd/sshd.go` (`handleSession` ~230, `runExec` ~270, `Exec` ~305, the `control.Ctx` at ~335) | |
| 580 | - Modify: `cmd/gitbayd/system.go:96` | |
| 581 | - Modify: `internal/httpd/api.go` (~64), `internal/httpd/apiread.go` (~53) | |
| 582 | ||
| 583 | **Interfaces:** | |
| 584 | - Consumes: `Ctx.Done` (Task 2). | |
| 585 | - Produces: `sshd.Exec(cfg, st, user, scope, source, cmdline string, stdin io.Reader, stdout, stderr io.Writer, done <-chan struct{}) int`. | |
| 586 | ||
| 587 | - [ ] **Step 1: SSH.** In `handleSession`'s `"exec"` case, replace the two lines after `req.Reply(true, nil)`: | |
| 588 | ||
| 589 | ```go | |
| 590 | req.Reply(true, nil) | |
| 591 | // x/crypto closes reqs when the client closes the channel. That | |
| 592 | // is how a follow learns nobody is reading: the CLI's shared | |
| 593 | // connection outlives a Ctrl-C, the channel does not. | |
| 594 | done := make(chan struct{}) | |
| 595 | go func() { | |
| 596 | for r := range reqs { | |
| 597 | r.Reply(false, nil) | |
| 598 | } | |
| 599 | close(done) | |
| 600 | }() | |
| 601 | code := s.runExec(sconn, ch, payload.Command, done) | |
| 602 | sendExit(ch, code) | |
| 603 | return | |
| 604 | ``` | |
| 605 | ||
| 606 | `runExec` gains `done <-chan struct{}` as its last parameter and passes it to `Exec`. `Exec` gains `done <-chan struct{}` as its last parameter and sets `Done: done` in the `control.Ctx` it builds. `cmd/gitbayd/system.go:96` passes `nil` (that process ends with its session). | |
| 607 | ||
| 608 | Find every other caller: `grep -rn 'sshd.Exec(\|\.runExec(' --include='*.go' .` and update them, test files included. | |
| 609 | ||
| 610 | - [ ] **Step 2: API.** In `internal/httpd/api.go` and `internal/httpd/apiread.go`, add `Done: r.Context().Done(),` to the `control.Ctx` literal. | |
| 611 | ||
| 612 | - [ ] **Step 3: Build, vet, test** | |
| 613 | ||
| 614 | Run: `go build ./... && go vet ./internal/sshd/ ./internal/httpd/ ./cmd/gitbayd/ && go test ./internal/sshd/ ./internal/httpd/ -count=1` | |
| 615 | Expected: `ok`. | |
| 616 | ||
| 617 | - [ ] **Step 4: Commit** | |
| 618 | ||
| 619 | ```bash | |
| 620 | git add internal/sshd/sshd.go cmd/gitbayd/system.go internal/httpd/api.go internal/httpd/apiread.go | |
| 621 | git commit -m "sshd, api: end a command when its reader goes away | |
| 622 | ||
| 623 | Ref #250" | |
| 624 | ``` | |
| 625 | ||
| 626 | --- | |
| 627 | ||
| 628 | ### Task 4: gzipWriter passes a flush through | |
| 629 | ||
| 630 | **Files:** | |
| 631 | - Modify: `internal/httpd/compress.go` | |
| 632 | - Test: `internal/httpd/compress_test.go` | |
| 633 | ||
| 634 | **Interfaces:** | |
| 635 | - Produces: `(*gzipWriter).Flush()`, `(*gzipWriter).Unwrap() http.ResponseWriter`. | |
| 636 | ||
| 637 | - [ ] **Step 1: Write the failing test** — append to `internal/httpd/compress_test.go` (add missing imports: `bytes`, `compress/gzip`, `io`, `net/http`, `net/http/httptest`, `strings`): | |
| 638 | ||
| 639 | ```go | |
| 640 | // A flush mid-response reaches the connection with what was written so | |
| 641 | // far decodable, which is what lets a page stream through gzip. | |
| 642 | func TestGzipWriterFlushes(t *testing.T) { | |
| 643 | rec := httptest.NewRecorder() | |
| 644 | h := compressed(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { | |
| 645 | w.Header().Set("Content-Type", "text/html; charset=utf-8") | |
| 646 | io.WriteString(w, "<p>first</p>") | |
| 647 | if err := http.NewResponseController(w).Flush(); err != nil { | |
| 648 | t.Fatalf("flush: %v", err) | |
| 649 | } | |
| 650 | if !rec.Flushed { | |
| 651 | t.Fatal("the flush did not reach the connection") | |
| 652 | } | |
| 653 | zr, err := gzip.NewReader(bytes.NewReader(rec.Body.Bytes())) | |
| 654 | if err != nil { | |
| 655 | t.Fatalf("gzip header: %v", err) | |
| 656 | } | |
| 657 | got, _ := io.ReadAll(zr) // no trailer yet: ends in ErrUnexpectedEOF | |
| 658 | if !strings.Contains(string(got), "<p>first</p>") { | |
| 659 | t.Fatalf("flushed body decodes to %q", got) | |
| 660 | } | |
| 661 | io.WriteString(w, "<p>second</p>") | |
| 662 | })) | |
| 663 | req := httptest.NewRequest("GET", "/", nil) | |
| 664 | req.Header.Set("Accept-Encoding", "gzip") | |
| 665 | h.ServeHTTP(rec, req) | |
| 666 | } | |
| 667 | ``` | |
| 668 | ||
| 669 | - [ ] **Step 2: Run to see it fail** | |
| 670 | ||
| 671 | Run: `go test ./internal/httpd/ -run TestGzipWriterFlushes -count=1` | |
| 672 | Expected: FAIL, `flush: feature not supported`. | |
| 673 | ||
| 674 | - [ ] **Step 3: Implement** — in `internal/httpd/compress.go`, after `Close`: | |
| 675 | ||
| 676 | ```go | |
| 677 | // Flush sends what the gzip stream holds, then flushes the connection, so | |
| 678 | // a streamed page reaches the browser as it is written. | |
| 679 | func (g *gzipWriter) Flush() { | |
| 680 | if !g.decided { | |
| 681 | g.decide(http.StatusOK) | |
| 682 | } | |
| 683 | if g.gz != nil { | |
| 684 | g.gz.Flush() | |
| 685 | } | |
| 686 | http.NewResponseController(g.ResponseWriter).Flush() | |
| 687 | } | |
| 688 | ||
| 689 | func (g *gzipWriter) Unwrap() http.ResponseWriter { return g.ResponseWriter } | |
| 690 | ``` | |
| 691 | ||
| 692 | - [ ] **Step 4: Run** | |
| 693 | ||
| 694 | Run: `go test ./internal/httpd/ -count=1 && go vet ./internal/httpd/` | |
| 695 | Expected: `ok`. | |
| 696 | ||
| 697 | - [ ] **Step 5: Commit** | |
| 698 | ||
| 699 | ```bash | |
| 700 | git add internal/httpd/compress.go internal/httpd/compress_test.go | |
| 701 | git commit -m "httpd: gzipWriter passes a flush through | |
| 702 | ||
| 703 | Ref #250" | |
| 704 | ``` | |
| 705 | ||
| 706 | --- | |
| 707 | ||
| 708 | ### Task 5: The build page streams a live build | |
| 709 | ||
| 710 | **Files:** | |
| 711 | - Modify: `internal/httpd/builds.go` (`build`, ~280) | |
| 712 | - Modify: `internal/httpd/control.go` (new `runControlStream` after `runControlCode`) | |
| 713 | - Modify: `internal/web/templates/build.html` | |
| 714 | ||
| 715 | **Interfaces:** | |
| 716 | - Consumes: `build log --follow` (Task 2), `gzipWriter.Flush` (Task 4), `Ctx.Done`. | |
| 717 | - Produces: `func (s *Server) runControlStream(u store.User, argv []string, out io.Writer, done <-chan struct{}) (msg string, code int)`; `type buildView`; `const liveLogMarker`. | |
| 718 | ||
| 719 | - [ ] **Step 1: Template.** Replace the last content line of `build.html`: | |
| 720 | ||
| 721 | ``` | |
| 722 | {{if .Log}}<pre class="code buildlog" tabindex="0">{{.Log}}</pre>{{else}}<p class="empty-note">no log yet</p>{{end}} | |
| 723 | ``` | |
| 724 | ||
| 725 | with: | |
| 726 | ||
| 727 | ``` | |
| 728 | {{if .Live}}<p class="meta">Live: the log streams here until the build ends. If it stops without a “build finished” line, reload to pick it up again. <a href="?follow=0">Show it without updates</a></p> | |
| 729 | <pre class="code buildlog" tabindex="0">{{.Log}}</pre> | |
| 730 | {{else if .Log}}<pre class="code buildlog" tabindex="0">{{.Log}}</pre>{{else}}<p class="empty-note">no log yet</p>{{end}} | |
| 731 | ``` | |
| 732 | ||
| 733 | - [ ] **Step 2: `runControlStream`** in `internal/httpd/control.go` after `runControlCode` (add `io` to imports): | |
| 734 | ||
| 735 | ```go | |
| 736 | // runControlStream runs a command whose output is written as it is | |
| 737 | // produced: stdout goes to out, and done ends the command when the | |
| 738 | // request does. msg is stderr. | |
| 739 | func (s *Server) runControlStream(u store.User, argv []string, out io.Writer, done <-chan struct{}) (msg string, code int) { | |
| 740 | var stderr bytes.Buffer | |
| 741 | ctx := &control.Ctx{ | |
| 742 | User: u, | |
| 743 | Source: "web", | |
| 744 | Scope: "full", | |
| 745 | Store: s.st, | |
| 746 | Cfg: s.cfg, | |
| 747 | Stdin: strings.NewReader(""), | |
| 748 | Stdout: out, | |
| 749 | Stderr: &stderr, | |
| 750 | ViaAPI: true, | |
| 751 | Done: done, | |
| 752 | } | |
| 753 | code = control.Dispatch(ctx, argv) | |
| 754 | return strings.TrimSpace(stderr.String()), code | |
| 755 | } | |
| 756 | ``` | |
| 757 | ||
| 758 | - [ ] **Step 3: Handler.** In `internal/httpd/builds.go`, replace `build` from the `log, _, _ := s.runControl(...)` line to the end with: | |
| 759 | ||
| 760 | ```go | |
| 761 | v := buildView{repoPage: p, Build: b, CanWrite: s.canWriteRepo(r, p.Repo), Notice: s.takeFlash(w, r)} | |
| 762 | if (b.Status == "pending" || b.Status == "running") && r.URL.Query().Get("follow") != "0" { | |
| 763 | s.streamBuild(w, r, v, viewer, n) | |
| 764 | return | |
| 765 | } | |
| 766 | v.Log, _, _ = s.runControl(viewer, []string{"build", "log", p.Repo.Path(), n}) | |
| 767 | s.render(w, "build.html", v) | |
| 768 | } | |
| 769 | ||
| 770 | type buildView struct { | |
| 771 | repoPage | |
| 772 | Build control.BuildOut | |
| 773 | Log string | |
| 774 | Live bool | |
| 775 | CanWrite bool | |
| 776 | Notice string | |
| 777 | } | |
| 778 | ||
| 779 | // liveLogMarker stands in for the log when build.html is rendered for a | |
| 780 | // live build; streamBuild splits the page there and streams the log into | |
| 781 | // the gap. Git refs, paths and job names cannot hold the control byte. | |
| 782 | const liveLogMarker = "\x1elive-log\x1e" | |
| 783 | ||
| 784 | // streamBuild writes the build page with the log following the build: | |
| 785 | // the page up to the log, then build log --follow escaped and flushed as | |
| 786 | // it arrives, then the outcome and the rest of the page. | |
| 787 | func (s *Server) streamBuild(w http.ResponseWriter, r *http.Request, v buildView, viewer store.User, n string) { | |
| 788 | v.Live, v.Log = true, liveLogMarker | |
| 789 | var buf bytes.Buffer | |
| 790 | if err := web.Render(&buf, "build.html", v); err != nil { | |
| 791 | http.Error(w, "template error: "+err.Error(), http.StatusInternalServerError) | |
| 792 | return | |
| 793 | } | |
| 794 | head, tail, ok := strings.Cut(buf.String(), liveLogMarker) | |
| 795 | if !ok || !strings.HasPrefix(tail, "</pre>") { | |
| 796 | http.Error(w, "template error: build.html has no live log slot", http.StatusInternalServerError) | |
| 797 | return | |
| 798 | } | |
| 799 | tail = strings.TrimPrefix(tail, "</pre>") | |
| 800 | ||
| 801 | h := w.Header() | |
| 802 | h.Set("Content-Type", "text/html; charset=utf-8") | |
| 803 | h.Set("Cache-Control", "no-store") | |
| 804 | h.Set("X-Accel-Buffering", "no") | |
| 805 | rc := http.NewResponseController(w) | |
| 806 | io.WriteString(w, head) | |
| 807 | rc.Flush() | |
| 808 | ||
| 809 | path := v.Repo.Path() | |
| 810 | msg, code := s.runControlStream(viewer, []string{"build", "log", path, n, "--follow"}, | |
| 811 | htmlStream{w: w, rc: rc}, r.Context().Done()) | |
| 812 | if code == protocol.ExitDenied { | |
| 813 | // The follow cap: the stored log once, and why it is not live. | |
| 814 | log, _, _ := s.runControl(viewer, []string{"build", "log", path, n}) | |
| 815 | template.HTMLEscape(w, []byte(log)) | |
| 816 | } | |
| 817 | io.WriteString(w, "</pre>") | |
| 818 | switch { | |
| 819 | case code == protocol.ExitOK: | |
| 820 | var b control.BuildOut | |
| 821 | if _, ok := s.runControlInto(viewer, []string{"build", "show", path, n}, &b); ok { | |
| 822 | fmt.Fprintf(w, `<p class="notice" role="status">build finished: %s</p>`, template.HTMLEscapeString(b.Status)) | |
| 823 | } | |
| 824 | case code == protocol.ExitDenied: | |
| 825 | fmt.Fprintf(w, `<p class="error" role="alert">%s</p>`, template.HTMLEscapeString(msg)) | |
| 826 | } | |
| 827 | io.WriteString(w, tail) | |
| 828 | } | |
| 829 | ||
| 830 | // htmlStream escapes each chunk of a streamed log into the page and | |
| 831 | // flushes it, so the browser draws it as it arrives. | |
| 832 | type htmlStream struct { | |
| 833 | w io.Writer | |
| 834 | rc *http.ResponseController | |
| 835 | } | |
| 836 | ||
| 837 | func (h htmlStream) Write(p []byte) (int, error) { | |
| 838 | template.HTMLEscape(h.w, p) | |
| 839 | if err := h.rc.Flush(); err != nil { | |
| 840 | return 0, err | |
| 841 | } | |
| 842 | return len(p), nil | |
| 843 | } | |
| 844 | ``` | |
| 845 | ||
| 846 | Add the imports `builds.go` now needs (`bytes`, `fmt`, `html/template`, `io`, `strings`, `gitbay.org/gitbay/internal/protocol`, `gitbay.org/gitbay/internal/web`) — only those not already present. If `builds.go` already imports `text/template` or another `template`, alias accordingly. | |
| 847 | ||
| 848 | - [ ] **Step 4: Build, vet, unit tests** | |
| 849 | ||
| 850 | Run: `go build ./... && go vet ./internal/httpd/ ./internal/web/ && go test ./internal/httpd/ ./internal/web/ -count=1` | |
| 851 | Expected: `ok`. (`TestMainWidthClass` needs nothing: no new template.) | |
| 852 | ||
| 853 | - [ ] **Step 5: Commit** | |
| 854 | ||
| 855 | ```bash | |
| 856 | git add internal/httpd/builds.go internal/httpd/control.go internal/web/templates/build.html | |
| 857 | git commit -m "web: the build page streams a live build's log | |
| 858 | ||
| 859 | Ref #250" | |
| 860 | ``` | |
| 861 | ||
| 862 | --- | |
| 863 | ||
| 864 | ### Task 6: e2e — follow over ssh and on the page | |
| 865 | ||
| 866 | **Files:** | |
| 867 | - Modify: `e2e/ssh_test.go` (extract `sshCmd` from `ssh`, ~line 164) | |
| 868 | - Create: `e2e/buildfollow_test.go` | |
| 869 | ||
| 870 | **Interfaces:** | |
| 871 | - Consumes: everything above; e2e helpers `startInstance`, `newKey`, `admin`, `ssh`, `gitEnv`, `sshURL`, `mustGit`, `httpPort`. | |
| 872 | - Produces: `func (i *instance) sshCmd(key string, args ...string) *exec.Cmd`. | |
| 873 | ||
| 874 | - [ ] **Step 1: Extract `sshCmd`.** In `e2e/ssh_test.go`, split `ssh`: | |
| 875 | ||
| 876 | ```go | |
| 877 | // sshCmd is the ssh invocation ssh runs, for a test that reads the output | |
| 878 | // as it arrives. | |
| 879 | func (i *instance) sshCmd(key string, args ...string) *exec.Cmd { | |
| 880 | base := []string{ | |
| 881 | "-p", fmt.Sprint(i.port), | |
| 882 | "-i", key, | |
| 883 | "-o", "IdentitiesOnly=yes", | |
| 884 | "-o", "StrictHostKeyChecking=no", | |
| 885 | "-o", "UserKnownHostsFile=" + filepath.Join(i.sshDir, "known_hosts"), | |
| 886 | "-o", "BatchMode=yes", | |
| 887 | "git@127.0.0.1", | |
| 888 | } | |
| 889 | return exec.Command("ssh", append(base, args...)...) | |
| 890 | } | |
| 891 | ``` | |
| 892 | ||
| 893 | and have `ssh` start with `cmd := i.sshCmd(key, args...)` in place of building `base` itself. | |
| 894 | ||
| 895 | - [ ] **Step 2: Write the test** — `e2e/buildfollow_test.go`: | |
| 896 | ||
| 897 | ```go | |
| 898 | package e2e | |
| 899 | ||
| 900 | import ( | |
| 901 | "encoding/json" | |
| 902 | "fmt" | |
| 903 | "io" | |
| 904 | "net/http" | |
| 905 | "os" | |
| 906 | "path/filepath" | |
| 907 | "strings" | |
| 908 | "testing" | |
| 909 | "time" | |
| 910 | ) | |
| 911 | ||
| 912 | // streamReader collects what r delivers, so a test can wait for text to | |
| 913 | // arrive while the writer is still going. | |
| 914 | type streamReader struct { | |
| 915 | ch chan []byte | |
| 916 | buf strings.Builder | |
| 917 | } | |
| 918 | ||
| 919 | func newStreamReader(r io.Reader) *streamReader { | |
| 920 | s := &streamReader{ch: make(chan []byte, 16)} | |
| 921 | go func() { | |
| 922 | b := make([]byte, 4096) | |
| 923 | for { | |
| 924 | n, err := r.Read(b) | |
| 925 | if n > 0 { | |
| 926 | s.ch <- append([]byte(nil), b[:n]...) | |
| 927 | } | |
| 928 | if err != nil { | |
| 929 | close(s.ch) | |
| 930 | return | |
| 931 | } | |
| 932 | } | |
| 933 | }() | |
| 934 | return s | |
| 935 | } | |
| 936 | ||
| 937 | func (s *streamReader) waitFor(t *testing.T, want string) string { | |
| 938 | t.Helper() | |
| 939 | deadline := time.After(20 * time.Second) | |
| 940 | for !strings.Contains(s.buf.String(), want) { | |
| 941 | select { | |
| 942 | case b, ok := <-s.ch: | |
| 943 | if !ok { | |
| 944 | t.Fatalf("stream ended before %q:\n%s", want, s.buf.String()) | |
| 945 | } | |
| 946 | s.buf.Write(b) | |
| 947 | case <-deadline: | |
| 948 | t.Fatalf("no %q after 20s:\n%s", want, s.buf.String()) | |
| 949 | } | |
| 950 | } | |
| 951 | return s.buf.String() | |
| 952 | } | |
| 953 | ||
| 954 | // A running build is followed over ssh and on its page: output the runner | |
| 955 | // sends arrives while the build runs, and both end with the outcome. | |
| 956 | func TestBuildLogFollow(t *testing.T) { | |
| 957 | t.Parallel() | |
| 958 | inst := startInstance(t) | |
| 959 | aliceKey := inst.newKey(t, "alice") | |
| 960 | runnerKey := inst.newKey(t, "ci") | |
| 961 | inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub") | |
| 962 | inst.admin(t, "admin", "user", "create", "ci", "--key", runnerKey+".pub", "--admin") | |
| 963 | if _, _, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/app"); code != 0 { | |
| 964 | t.Fatal("repo create failed") | |
| 965 | } | |
| 966 | work := t.TempDir() | |
| 967 | env := inst.gitEnv(aliceKey) | |
| 968 | mustGit(t, work, env, "clone", inst.sshURL("alice/app"), "w") | |
| 969 | dir := filepath.Join(work, "w") | |
| 970 | os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755) | |
| 971 | os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte("jobs:\n unit:\n steps:\n - echo fine\n"), 0o644) | |
| 972 | mustGit(t, dir, env, "checkout", "-q", "-b", "main") | |
| 973 | mustGit(t, dir, env, "add", ".") | |
| 974 | mustGit(t, dir, env, "commit", "-q", "-m", "ci") | |
| 975 | mustGit(t, dir, env, "push", "-q", "origin", "main") | |
| 976 | ||
| 977 | // Claim build 1 by hand, so the test decides when output arrives. | |
| 978 | out, errOut, code := inst.ssh(t, runnerKey, "", "runner", "next", "--json") | |
| 979 | if code != 0 { | |
| 980 | t.Fatalf("runner next: %s", errOut) | |
| 981 | } | |
| 982 | var claim struct { | |
| 983 | Data struct { | |
| 984 | ID int64 `json:"id"` | |
| 985 | } `json:"data"` | |
| 986 | } | |
| 987 | if err := json.Unmarshal([]byte(out), &claim); err != nil || claim.Data.ID == 0 { | |
| 988 | t.Fatalf("runner next output %q: %v", out, err) | |
| 989 | } | |
| 990 | id := fmt.Sprint(claim.Data.ID) | |
| 991 | ||
| 992 | cmd := inst.sshCmd(aliceKey, "build", "log", "alice/app", "1", "--follow") | |
| 993 | stdout, err := cmd.StdoutPipe() | |
| 994 | if err != nil { | |
| 995 | t.Fatal(err) | |
| 996 | } | |
| 997 | var stderr strings.Builder | |
| 998 | cmd.Stderr = &stderr | |
| 999 | if err := cmd.Start(); err != nil { | |
| 1000 | t.Fatal(err) | |
| 1001 | } | |
| 1002 | follow := newStreamReader(stdout) | |
| 1003 | ||
| 1004 | page, err := http.Get(fmt.Sprintf("http://127.0.0.1:%d/alice/app/builds/1", inst.httpPort)) | |
| 1005 | if err != nil { | |
| 1006 | t.Fatal(err) | |
| 1007 | } | |
| 1008 | defer page.Body.Close() | |
| 1009 | web := newStreamReader(page.Body) | |
| 1010 | web.waitFor(t, "Live: the log streams here") | |
| 1011 | ||
| 1012 | // A static render while the build runs returns at once. | |
| 1013 | static := &http.Client{Timeout: 10 * time.Second} | |
| 1014 | resp, err := static.Get(fmt.Sprintf("http://127.0.0.1:%d/alice/app/builds/1?follow=0", inst.httpPort)) | |
| 1015 | if err != nil { | |
| 1016 | t.Fatalf("?follow=0 did not return: %v", err) | |
| 1017 | } | |
| 1018 | body, _ := io.ReadAll(resp.Body) | |
| 1019 | resp.Body.Close() | |
| 1020 | if strings.Contains(string(body), "Live:") { | |
| 1021 | t.Fatalf("?follow=0 rendered the live page:\n%s", body) | |
| 1022 | } | |
| 1023 | ||
| 1024 | if _, errOut, code := inst.ssh(t, runnerKey, "hello from the runner <b>\n", "runner", "log", id); code != 0 { | |
| 1025 | t.Fatalf("runner log: %s", errOut) | |
| 1026 | } | |
| 1027 | follow.waitFor(t, "hello from the runner <b>\n") | |
| 1028 | web.waitFor(t, "hello from the runner <b>") | |
| 1029 | ||
| 1030 | if _, errOut, code := inst.ssh(t, runnerKey, "", "runner", "done", id, "success"); code != 0 { | |
| 1031 | t.Fatalf("runner done: %s", errOut) | |
| 1032 | } | |
| 1033 | web.waitFor(t, `<p class="notice" role="status">build finished: success</p>`) | |
| 1034 | web.waitFor(t, "</html>") | |
| 1035 | if err := cmd.Wait(); err != nil { | |
| 1036 | t.Fatalf("follow exited: %v\n%s", err, stderr.String()) | |
| 1037 | } | |
| 1038 | if got := strings.TrimSpace(stderr.String()); got != "build 1 success" { | |
| 1039 | t.Errorf("follow stderr %q", got) | |
| 1040 | } | |
| 1041 | } | |
| 1042 | ``` | |
| 1043 | ||
| 1044 | - [ ] **Step 3: Run it** | |
| 1045 | ||
| 1046 | Run: `go test ./e2e/ -run 'TestBuildLogFollow$' -count=1 -v 2>&1 | tail -20` | |
| 1047 | Expected: `--- PASS: TestBuildLogFollow`. If `</html>` is not how the layout ends, use the last line `layout.html` renders. | |
| 1048 | ||
| 1049 | - [ ] **Step 4: Check the neighbours.** Run `go vet ./e2e/` and the tests that use `inst.ssh` heavily and build pages: `go test ./e2e/ -run 'TestBuildCancel$|TestControlPlaneOverBareSSH$' -count=1`. | |
| 1050 | ||
| 1051 | - [ ] **Step 5: Commit** | |
| 1052 | ||
| 1053 | ```bash | |
| 1054 | git add e2e/ssh_test.go e2e/buildfollow_test.go | |
| 1055 | git commit -m "e2e: follow a running build over ssh and on its page | |
| 1056 | ||
| 1057 | Ref #250" | |
| 1058 | ``` | |
| 1059 | ||
| 1060 | --- | |
| 1061 | ||
| 1062 | ### Task 7: Docs | |
| 1063 | ||
| 1064 | **Files:** | |
| 1065 | - Modify: `.gitbay/wiki/Parity.org` (build rows, ~line 212) | |
| 1066 | - Modify: `.gitbay/wiki/CI.org` | |
| 1067 | ||
| 1068 | - [ ] **Step 1: Parity.** After `| build log | yes | yes | yes |` add (columns are cli, web, ios): | |
| 1069 | ||
| 1070 | ``` | |
| 1071 | | build log follow (until it ends) | yes | yes | no | | |
| 1072 | ``` | |
| 1073 | ||
| 1074 | - [ ] **Step 2: CI.org.** After the paragraph that begins "Scheduled jobs run on their cron", add: | |
| 1075 | ||
| 1076 | ``` | |
| 1077 | A running build is followed with =build log <owner/name> <n> --follow=: | |
| 1078 | the stored log, then output as the runner sends it, then the outcome | |
| 1079 | as =build <n> <status>= on stderr once the build ends. The exit code is | |
| 1080 | 0 whatever the outcome. The build page does the same without | |
| 1081 | JavaScript while a build is queued or running; =?follow=0= renders it | |
| 1082 | once. An account holds at most eight follows open, and signed-out | |
| 1083 | viewers share one account's eight. Over the JSON API the command | |
| 1084 | answers when the build ends, with the whole log. | |
| 1085 | ``` | |
| 1086 | ||
| 1087 | - [ ] **Step 3: Commit** | |
| 1088 | ||
| 1089 | ```bash | |
| 1090 | git add .gitbay/wiki/Parity.org .gitbay/wiki/CI.org | |
| 1091 | git commit -m "wiki: build log --follow | |
| 1092 | ||
| 1093 | Closes #250" | |
| 1094 | ``` | |
| 1095 | ||
| 1096 | --- | |
| 1097 | ||
| 1098 | ## Finish | |
| 1099 | ||
| 1100 | Push `build-log-follow`, open the MR with `gitbay mr create --source build-log-follow --target main --title "build log --follow, streamed to the build page" --file - < <body file>`, wait for CI, then `gitbay mr merge <n> --strategy ff` and delete the branch in both places. | |
docs/specs/2026-09-23-build-log-follow-design.md added +150
| @@ -0,0 +1,150 @@ | ||
| 1 | # Following a running build | |
| 2 | ||
| 3 | Closes #250. `build log <owner/name> <n> --follow` streams a build's log | |
| 4 | until the build reaches an outcome, and the web build page streams the | |
| 5 | same command without JavaScript. | |
| 6 | ||
| 7 | ## Problem | |
| 8 | ||
| 9 | No surface can follow a running build. `build log` prints what is stored | |
| 10 | and exits; the build page renders the log once. Watching a build means | |
| 11 | re-running the command or reloading the page. | |
| 12 | ||
| 13 | ## Decision | |
| 14 | ||
| 15 | The capability is a flag on the existing control command. The web page | |
| 16 | dispatches that command with a writer that escapes and flushes each | |
| 17 | chunk, so the CLI, stock ssh and the page share one implementation. | |
| 18 | ||
| 19 | Rejected: a `Refresh` header on the build page. It re-renders the page | |
| 20 | and re-reads the whole log per reload per viewer, interrupts selection | |
| 21 | and screen readers, needs an off switch for WCAG 2.2.1, and gives only | |
| 22 | the web a way to follow. | |
| 23 | ||
| 24 | ## Store: waking followers | |
| 25 | ||
| 26 | `Store` gains an in-memory waiter table keyed by build id: | |
| 27 | ||
| 28 | - `BuildLogWait(id int64) <-chan struct{}` returns a channel closed by | |
| 29 | the next change to that build. | |
| 30 | - `AppendBuildLog`, `FinishBuild` and `CancelBuild` close and drop the | |
| 31 | build's channel after their write succeeds (`wakeBuild(id)`). | |
| 32 | - `BuildLogFrom(id, offset int64) (status string, chunk []byte, err | |
| 33 | error)` reads `status` and `substr(log, offset+1)` in one query, so a | |
| 34 | follower reads each byte once. | |
| 35 | ||
| 36 | gitbayd is one process and its SSH and HTTP servers share one | |
| 37 | `*store.Store`, so an in-memory table reaches every follower. Writers in | |
| 38 | another process do not wake anyone: `gitbayd admin` subcommands, and | |
| 39 | every session under the system-sshd forced command (`gitbayd shell`), | |
| 40 | where each session is its own process. The follow loop also re-reads | |
| 41 | every 2 seconds, which bounds that case. | |
| 42 | ||
| 43 | ## Command | |
| 44 | ||
| 45 | `build log <owner/name> <n> [--follow]`, still `ReadOnly`. | |
| 46 | ||
| 47 | Without `--follow`, unchanged. | |
| 48 | ||
| 49 | With `--follow`: | |
| 50 | ||
| 51 | 1. Write the stored log. | |
| 52 | 2. Loop: take a wait channel, read from the offset, write any new bytes. | |
| 53 | If the status is no longer `pending` or `running` and the read | |
| 54 | returned nothing new, stop. Otherwise wait on the channel, the | |
| 55 | 2-second timer, or `Ctx.Done`. | |
| 56 | 3. Write `build <n> <status>` to stderr and exit 0, whatever the | |
| 57 | outcome. Stdout stays the log, byte for byte. | |
| 58 | ||
| 59 | The wait channel is taken before the read, so a change between the read | |
| 60 | and the wait still wakes the loop. | |
| 61 | ||
| 62 | `Ctx` gains `Done <-chan struct{}`, nil when the surface has none. The | |
| 63 | embedded sshd closes it when the session's channel closes (the CLI's | |
| 64 | shared connection outlives a Ctrl-C, the channel does not); `gitbayd | |
| 65 | shell` still passes nil, but its process does not end with the | |
| 66 | session: OpenSSH closes the child's pipes and sends no signal to a | |
| 67 | session with no pty, so a follow there ends at its next write, at the | |
| 68 | build's outcome, or at the queued limit below; the per-account cap is | |
| 69 | per process in that mode. httpd sets it from `r.Context()` on the web | |
| 70 | and both API endpoints. A write error also ends the loop. On `Done` | |
| 71 | the command returns `protocol.ExitFailure` with no message; nobody is | |
| 72 | reading. | |
| 73 | ||
| 74 | At most 8 follows per account run at once (a counter in `control`, | |
| 75 | decremented on return). The ninth exits 4: "8 follows are already open | |
| 76 | for this account; close one and retry". Signed-out web viewers are | |
| 77 | account 0 and share the 8; the ninth gets the stored log once with the | |
| 78 | refusal under it. | |
| 79 | ||
| 80 | Nothing reaps a queued build (`ReapStaleBuilds` only reaps `running` | |
| 81 | builds), and a running one is already bounded by the reaper's | |
| 82 | deadline, so a follow of a build that stays `pending` ends on its own | |
| 83 | after `followQueued` (10 minutes), writing to stderr `build <n> is | |
| 84 | still queued; nothing claimed it in 10m0s. Follow again once a runner | |
| 85 | has.` and exiting `protocol.ExitFailure`. The clock runs only while the | |
| 86 | follow has seen the build `pending`; once it sees `running` or a | |
| 87 | terminal status the limit no longer applies. | |
| 88 | ||
| 89 | The CLI's `pass("log", …)` help in `cmd/gitbay/main.go` names | |
| 90 | `--follow`. | |
| 91 | ||
| 92 | ## Web | |
| 93 | ||
| 94 | `GET /{owner}/{repo}/builds/{n}` streams when the build is `pending` or | |
| 95 | `running`, the query has no `follow=0`, and the method is `GET`. A HEAD | |
| 96 | request (the route also matches it) renders once, like `?follow=0`. | |
| 97 | ||
| 98 | Streaming: | |
| 99 | ||
| 100 | 1. Render `build.html` into a buffer with `Log` set to a marker and | |
| 101 | `Live` true, and split the output at the marker. | |
| 102 | 2. Write the head, flush. | |
| 103 | 3. Dispatch `build log <repo> <n> --follow` with `Stdout` an escaping | |
| 104 | writer (`template.HTMLEscape` per chunk, then flush through | |
| 105 | `http.ResponseController`) and `Done` from the request context. | |
| 106 | 4. If the request context is done (the client left), write nothing | |
| 107 | more. Otherwise write the `</pre>` that begins the tail, then, by | |
| 108 | the command's exit code: `ExitOK` reads the build and writes | |
| 109 | `<p class="notice" role="status">build finished: <status></p>`; | |
| 110 | `ExitDenied` (the follow cap) writes the stored log once above the | |
| 111 | `</pre>` and an error paragraph with the refusal, except a | |
| 112 | signed-out viewer (`viewer.ID == 0`) gets "Too many signed-out | |
| 113 | viewers are watching live builds. This is the log so far; reload to | |
| 114 | try again, or sign in." instead of the command's account-scoped | |
| 115 | wording; `ExitFailure` with a message (the queued limit) writes it | |
| 116 | as a `<p class="notice" role="status">`. Then the rest of the tail. | |
| 117 | ||
| 118 | `build.html`, when `Live`, puts a line above the log: the log streams | |
| 119 | until the build ends; a stream that stops with no "build finished" line | |
| 120 | resumes on reload; and a link to `?follow=0`, the same page rendered | |
| 121 | once, as the way to stop the updates (WCAG 2.2.2). The stored log | |
| 122 | renders inside the stream from the first write, so there is no separate | |
| 123 | "no log yet" state while live. | |
| 124 | ||
| 125 | `gzipWriter` gains `Flush()` (flush the gzip stream, then the underlying | |
| 126 | writer) and `Unwrap()`. The HTTP server has no `WriteTimeout`, so a long | |
| 127 | silent step does not end the response. The handler sets | |
| 128 | `X-Accel-Buffering: no` for a proxy in front of the instance. | |
| 129 | ||
| 130 | ## API | |
| 131 | ||
| 132 | `/api/v1/cmd` and `/api/v1/read` buffer stdout, so there `--follow` | |
| 133 | returns the whole log when the build ends. Streaming a JSON response is | |
| 134 | out of scope. | |
| 135 | ||
| 136 | ## Tests | |
| 137 | ||
| 138 | - `internal/store`: `BuildLogWait` closes on append, finish and cancel; | |
| 139 | `BuildLogFrom` returns bytes past the offset and the status. | |
| 140 | - `internal/control`: follow a running build while another goroutine | |
| 141 | appends and finishes; stdout is the full log, stderr ends with the | |
| 142 | outcome, exit 0. Cancel ends a follow. A closed `Done` ends a follow. | |
| 143 | The ninth concurrent follow exits 4. | |
| 144 | - `internal/httpd`: `gzipWriter` passes a flush through. | |
| 145 | - `e2e`: ssh `build log --follow` on a running build while the runner | |
| 146 | key appends with `runner log` and reports with `runner done`; the web | |
| 147 | page of a running build arrives complete with "build finished: | |
| 148 | success"; `?follow=0` renders without the live line. | |
| 149 | ||
| 150 | Docs: the CI wiki page and Parity get the flag. | |
e2e/buildcancelweb_test.go +4 −4
| @@ -46,13 +46,13 @@ func TestBuildCancelWeb(t *testing.T) { | ||
| 46 | 46 | build1 := inst.base() + "/alice/app/builds/1" |
| 47 | 47 | |
| 48 | 48 | // Build 1 is queued: the control is on the page, for someone with |
| 49 | // write access. | |
| 50 | _, body := browserGet(t, alice, build1) | |
| 49 | // write access. follow=0: a live build's page streams until it ends. | |
| 50 | _, body := browserGet(t, alice, build1+"?follow=0") | |
| 51 | 51 | if !strings.Contains(body, `action="/alice/app/builds/1/cancel"`) { |
| 52 | 52 | t.Fatalf("no cancel control on a queued build:\n%s", body) |
| 53 | 53 | } |
| 54 | 54 | // A reader gets no control. |
| 55 | _, anon := browserGet(t, newBrowser(t), build1) | |
| 55 | _, anon := browserGet(t, newBrowser(t), build1+"?follow=0") | |
| 56 | 56 | if strings.Contains(anon, "/cancel") { |
| 57 | 57 | t.Fatal("anonymous visitor sees the cancel control") |
| 58 | 58 | } |
| @@ -78,7 +78,7 @@ func TestBuildCancelWeb(t *testing.T) { | ||
| 78 | 78 | t.Fatalf("runner claim: %s", out) |
| 79 | 79 | } |
| 80 | 80 | build2 := inst.base() + "/alice/app/builds/2" |
| 81 | _, body = browserGet(t, alice, build2) | |
| 81 | _, body = browserGet(t, alice, build2+"?follow=0") | |
| 82 | 82 | if !strings.Contains(body, "running") { |
| 83 | 83 | t.Fatalf("build 2 not running:\n%s", body) |
| 84 | 84 | } |
e2e/buildfollow_test.go added +153
| @@ -0,0 +1,153 @@ | ||
| 1 | package e2e | |
| 2 | ||
| 3 | import ( | |
| 4 | "encoding/json" | |
| 5 | "fmt" | |
| 6 | "io" | |
| 7 | "net/http" | |
| 8 | "os" | |
| 9 | "path/filepath" | |
| 10 | "strings" | |
| 11 | "testing" | |
| 12 | "time" | |
| 13 | ) | |
| 14 | ||
| 15 | // streamReader collects what r delivers, so a test can wait for text to | |
| 16 | // arrive while the writer is still going. | |
| 17 | type streamReader struct { | |
| 18 | ch chan []byte | |
| 19 | buf strings.Builder | |
| 20 | } | |
| 21 | ||
| 22 | func newStreamReader(r io.Reader) *streamReader { | |
| 23 | s := &streamReader{ch: make(chan []byte, 16)} | |
| 24 | go func() { | |
| 25 | b := make([]byte, 4096) | |
| 26 | for { | |
| 27 | n, err := r.Read(b) | |
| 28 | if n > 0 { | |
| 29 | s.ch <- append([]byte(nil), b[:n]...) | |
| 30 | } | |
| 31 | if err != nil { | |
| 32 | close(s.ch) | |
| 33 | return | |
| 34 | } | |
| 35 | } | |
| 36 | }() | |
| 37 | return s | |
| 38 | } | |
| 39 | ||
| 40 | func (s *streamReader) waitFor(t *testing.T, want string) string { | |
| 41 | t.Helper() | |
| 42 | deadline := time.After(20 * time.Second) | |
| 43 | for !strings.Contains(s.buf.String(), want) { | |
| 44 | select { | |
| 45 | case b, ok := <-s.ch: | |
| 46 | if !ok { | |
| 47 | t.Fatalf("stream ended before %q:\n%s", want, s.buf.String()) | |
| 48 | } | |
| 49 | s.buf.Write(b) | |
| 50 | case <-deadline: | |
| 51 | t.Fatalf("no %q after 20s:\n%s", want, s.buf.String()) | |
| 52 | } | |
| 53 | } | |
| 54 | return s.buf.String() | |
| 55 | } | |
| 56 | ||
| 57 | // A running build is followed over ssh and on its page: output the runner | |
| 58 | // sends arrives while the build runs, and both end with the outcome. | |
| 59 | func TestBuildLogFollow(t *testing.T) { | |
| 60 | t.Parallel() | |
| 61 | inst := startInstance(t) | |
| 62 | aliceKey := inst.newKey(t, "alice") | |
| 63 | runnerKey := inst.newKey(t, "ci") | |
| 64 | inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub") | |
| 65 | inst.admin(t, "admin", "user", "create", "ci", "--key", runnerKey+".pub", "--admin") | |
| 66 | if _, _, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/app"); code != 0 { | |
| 67 | t.Fatal("repo create failed") | |
| 68 | } | |
| 69 | work := t.TempDir() | |
| 70 | env := inst.gitEnv(aliceKey) | |
| 71 | mustGit(t, work, env, "clone", inst.sshURL("alice/app"), "w") | |
| 72 | dir := filepath.Join(work, "w") | |
| 73 | os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755) | |
| 74 | os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte("jobs:\n unit:\n steps:\n - echo fine\n"), 0o644) | |
| 75 | mustGit(t, dir, env, "checkout", "-q", "-b", "main") | |
| 76 | mustGit(t, dir, env, "add", ".") | |
| 77 | mustGit(t, dir, env, "commit", "-q", "-m", "ci") | |
| 78 | mustGit(t, dir, env, "push", "-q", "origin", "main") | |
| 79 | ||
| 80 | // Claim build 1 by hand, so the test decides when output arrives. | |
| 81 | out, errOut, code := inst.ssh(t, runnerKey, "", "runner", "next", "--json") | |
| 82 | if code != 0 { | |
| 83 | t.Fatalf("runner next: %s", errOut) | |
| 84 | } | |
| 85 | var claim struct { | |
| 86 | Data struct { | |
| 87 | ID int64 `json:"id"` | |
| 88 | } `json:"data"` | |
| 89 | } | |
| 90 | if err := json.Unmarshal([]byte(out), &claim); err != nil || claim.Data.ID == 0 { | |
| 91 | t.Fatalf("runner next output %q: %v", out, err) | |
| 92 | } | |
| 93 | id := fmt.Sprint(claim.Data.ID) | |
| 94 | ||
| 95 | cmd := inst.sshCmd(aliceKey, "build", "log", "alice/app", "1", "--follow") | |
| 96 | stdout, err := cmd.StdoutPipe() | |
| 97 | if err != nil { | |
| 98 | t.Fatal(err) | |
| 99 | } | |
| 100 | var stderr strings.Builder | |
| 101 | cmd.Stderr = &stderr | |
| 102 | if err := cmd.Start(); err != nil { | |
| 103 | t.Fatal(err) | |
| 104 | } | |
| 105 | // If the test fails before the final cmd.Wait below, kill the follow | |
| 106 | // instead of leaving it running. Once that Wait has run, | |
| 107 | // cmd.ProcessState is set and this is a no-op. | |
| 108 | defer func() { | |
| 109 | if cmd.ProcessState == nil { | |
| 110 | cmd.Process.Kill() | |
| 111 | cmd.Wait() | |
| 112 | } | |
| 113 | }() | |
| 114 | follow := newStreamReader(stdout) | |
| 115 | ||
| 116 | page, err := http.Get(fmt.Sprintf("http://127.0.0.1:%d/alice/app/builds/1", inst.httpPort)) | |
| 117 | if err != nil { | |
| 118 | t.Fatal(err) | |
| 119 | } | |
| 120 | defer page.Body.Close() | |
| 121 | web := newStreamReader(page.Body) | |
| 122 | web.waitFor(t, "Live: the log streams here") | |
| 123 | ||
| 124 | // A static render while the build runs returns at once. | |
| 125 | static := &http.Client{Timeout: 10 * time.Second} | |
| 126 | resp, err := static.Get(fmt.Sprintf("http://127.0.0.1:%d/alice/app/builds/1?follow=0", inst.httpPort)) | |
| 127 | if err != nil { | |
| 128 | t.Fatalf("?follow=0 did not return: %v", err) | |
| 129 | } | |
| 130 | body, _ := io.ReadAll(resp.Body) | |
| 131 | resp.Body.Close() | |
| 132 | if strings.Contains(string(body), "Live:") { | |
| 133 | t.Fatalf("?follow=0 rendered the live page:\n%s", body) | |
| 134 | } | |
| 135 | ||
| 136 | if _, errOut, code := inst.ssh(t, runnerKey, "hello from the runner <b>\n", "runner", "log", id); code != 0 { | |
| 137 | t.Fatalf("runner log: %s", errOut) | |
| 138 | } | |
| 139 | follow.waitFor(t, "hello from the runner <b>\n") | |
| 140 | web.waitFor(t, "hello from the runner <b>") | |
| 141 | ||
| 142 | if _, errOut, code := inst.ssh(t, runnerKey, "", "runner", "done", id, "success"); code != 0 { | |
| 143 | t.Fatalf("runner done: %s", errOut) | |
| 144 | } | |
| 145 | web.waitFor(t, `<p class="notice" role="status">build finished: success</p>`) | |
| 146 | web.waitFor(t, "</html>") | |
| 147 | if err := cmd.Wait(); err != nil { | |
| 148 | t.Fatalf("follow exited: %v\n%s", err, stderr.String()) | |
| 149 | } | |
| 150 | if got := strings.TrimSpace(stderr.String()); got != "build 1 success" { | |
| 151 | t.Errorf("follow stderr %q", got) | |
| 152 | } | |
| 153 | } | |
e2e/ssh_test.go +10 −4
| @@ -160,9 +160,9 @@ func (i *instance) newKey(t *testing.T, name string) string { | ||
| 160 | 160 | return priv |
| 161 | 161 | } |
| 162 | 162 | |
| 163 | // ssh runs the real OpenSSH client against the instance with the given key. | |
| 164 | func (i *instance) ssh(t *testing.T, key string, stdin string, args ...string) (string, string, int) { | |
| 165 | t.Helper() | |
| 163 | // sshCmd is the ssh invocation ssh runs, for a test that reads the output | |
| 164 | // as it arrives. | |
| 165 | func (i *instance) sshCmd(key string, args ...string) *exec.Cmd { | |
| 166 | 166 | base := []string{ |
| 167 | 167 | "-p", fmt.Sprint(i.port), |
| 168 | 168 | "-i", key, |
| @@ -172,7 +172,13 @@ func (i *instance) ssh(t *testing.T, key string, stdin string, args ...string) ( | ||
| 172 | 172 | "-o", "BatchMode=yes", |
| 173 | 173 | "git@127.0.0.1", |
| 174 | 174 | } |
| 175 | cmd := exec.Command("ssh", append(base, args...)...) | |
| 175 | return exec.Command("ssh", append(base, args...)...) | |
| 176 | } | |
| 177 | ||
| 178 | // ssh runs the real OpenSSH client against the instance with the given key. | |
| 179 | func (i *instance) ssh(t *testing.T, key string, stdin string, args ...string) (string, string, int) { | |
| 180 | t.Helper() | |
| 181 | cmd := i.sshCmd(key, args...) | |
| 176 | 182 | if stdin != "" { |
| 177 | 183 | cmd.Stdin = strings.NewReader(stdin) |
| 178 | 184 | } |
internal/control/build.go +10 −3
| @@ -27,8 +27,8 @@ func init() { | ||
| 27 | 27 | Summary: "show one build", |
| 28 | 28 | Usage: "build show <owner/name> <n>", ReadOnly: true, Run: runBuildShow}) |
| 29 | 29 | register(Command{Path: []string{"build", "log"}, |
| 30 | Summary: "print a build's log", | |
| 31 | Usage: "build log <owner/name> <n>", ReadOnly: true, Run: runBuildLog}) | |
| 30 | Summary: "print a build's log, or follow it until the build ends", | |
| 31 | Usage: "build log <owner/name> <n> [--follow]", ReadOnly: true, Run: runBuildLog}) | |
| 32 | 32 | |
| 33 | 33 | register(Command{Path: []string{"build", "jobs"}, |
| 34 | 34 | Summary: "list the jobs a trigger can name", |
| @@ -196,10 +196,17 @@ func runBuildShow(c *Ctx, args []string) int { | ||
| 196 | 196 | } |
| 197 | 197 | |
| 198 | 198 | func runBuildLog(c *Ctx, args []string) int { |
| 199 | _, b, code := buildRef(c, args) | |
| 199 | f, err := parseFlags(args, flagSpec{Bools: []string{"--follow"}, MaxPos: 2, Usage: c.Cmd.Usage}) | |
| 200 | if err != nil { | |
| 201 | return c.fail(protocol.ExitUsage, "%v", err) | |
| 202 | } | |
| 203 | _, b, code := buildRef(c, f.Pos) | |
| 200 | 204 | if code >= 0 { |
| 201 | 205 | return code |
| 202 | 206 | } |
| 207 | if f.Has("--follow") { | |
| 208 | return followBuildLog(c, b) | |
| 209 | } | |
| 203 | 210 | log, err := c.Store.BuildLog(b.ID) |
| 204 | 211 | if err != nil { |
| 205 | 212 | return c.fail(protocol.ExitFailure, "%v", err) |
internal/control/build_test.go +5 −1
| @@ -42,9 +42,13 @@ func gitRunner(t *testing.T) func(dir string, args ...string) string { | ||
| 42 | 42 | |
| 43 | 43 | // newQueueTestRepo returns a store with one public repo (default branch |
| 44 | 44 | // "main", matching the schema default) and the uid to queue builds as. |
| 45 | // | |
| 46 | // A real file, not ":memory:": ":memory:" gives each connection its own | |
| 47 | // database, so a goroutine querying while another writes (build log | |
| 48 | // --follow) sees an empty schema (see internal/store/contention_test.go). | |
| 45 | 49 | func newQueueTestRepo(t *testing.T) (*store.Store, store.Repo, int64) { |
| 46 | 50 | t.Helper() |
| 47 | st, err := store.Open(":memory:") | |
| 51 | st, err := store.Open(filepath.Join(t.TempDir(), "gitbay.db")) | |
| 48 | 52 | if err != nil { |
| 49 | 53 | t.Fatal(err) |
| 50 | 54 | } |
internal/control/buildfollow.go added +117
| @@ -0,0 +1,117 @@ | ||
| 1 | package control | |
| 2 | ||
| 3 | import ( | |
| 4 | "fmt" | |
| 5 | "sync" | |
| 6 | "time" | |
| 7 | ||
| 8 | "gitbay.org/gitbay/internal/protocol" | |
| 9 | "gitbay.org/gitbay/internal/store" | |
| 10 | ) | |
| 11 | ||
| 12 | // maxFollows is how many build log follows one account holds open at | |
| 13 | // once. Signed-out web viewers are account 0 and share it. | |
| 14 | const maxFollows = 8 | |
| 15 | ||
| 16 | var ( | |
| 17 | // followPoll bounds a wait with no wake. A write from another process | |
| 18 | // (gitbayd admin, or any session under gitbayd shell) wakes nobody; | |
| 19 | // this is how its bytes still arrive. | |
| 20 | followPoll = 2 * time.Second | |
| 21 | // followSettle is how long a follow keeps reading after the build has | |
| 22 | // an outcome: a cancel appends its line after the status changes, and | |
| 23 | // a cancelled runner's stream runs on until its next check. | |
| 24 | followSettle = time.Second | |
| 25 | // followQueued bounds how long a follow waits on a build that stays | |
| 26 | // pending: nothing reaps a queued build (ReapStaleBuilds only reaps | |
| 27 | // running builds), and a running one is already bounded by the | |
| 28 | // reaper's deadline, so a follow needs its own limit for the queued | |
| 29 | // case or it never ends. | |
| 30 | followQueued = 10 * time.Minute | |
| 31 | ) | |
| 32 | ||
| 33 | var ( | |
| 34 | followMu sync.Mutex | |
| 35 | follows = map[int64]int{} | |
| 36 | ) | |
| 37 | ||
| 38 | func takeFollow(uid int64) bool { | |
| 39 | followMu.Lock() | |
| 40 | defer followMu.Unlock() | |
| 41 | if follows[uid] >= maxFollows { | |
| 42 | return false | |
| 43 | } | |
| 44 | follows[uid]++ | |
| 45 | return true | |
| 46 | } | |
| 47 | ||
| 48 | func dropFollow(uid int64) { | |
| 49 | followMu.Lock() | |
| 50 | defer followMu.Unlock() | |
| 51 | if follows[uid]--; follows[uid] <= 0 { | |
| 52 | delete(follows, uid) | |
| 53 | } | |
| 54 | } | |
| 55 | ||
| 56 | // followBuildLog writes the build's log as it grows and returns once the | |
| 57 | // build has an outcome and its last bytes are written. The outcome goes | |
| 58 | // to stderr, so stdout is the log byte for byte. | |
| 59 | func followBuildLog(c *Ctx, b store.Build) int { | |
| 60 | if !takeFollow(c.User.ID) { | |
| 61 | return c.fail(protocol.ExitDenied, "%d follows are already open for this account; close one and retry", maxFollows) | |
| 62 | } | |
| 63 | defer dropFollow(c.User.ID) | |
| 64 | ||
| 65 | var off int64 | |
| 66 | var settleBy time.Time | |
| 67 | var queuedSince time.Time | |
| 68 | for { | |
| 69 | wake := c.Store.BuildLogWait(b.ID) | |
| 70 | status, chunk, err := c.Store.BuildLogFrom(b.ID, off) | |
| 71 | if err != nil { | |
| 72 | return c.fail(protocol.ExitFailure, "%v", err) | |
| 73 | } | |
| 74 | if len(chunk) > 0 { | |
| 75 | if _, err := c.Stdout.Write(chunk); err != nil { | |
| 76 | return protocol.ExitFailure | |
| 77 | } | |
| 78 | off += int64(len(chunk)) | |
| 79 | } | |
| 80 | if status == "pending" { | |
| 81 | if queuedSince.IsZero() { | |
| 82 | queuedSince = time.Now() | |
| 83 | } | |
| 84 | } else { | |
| 85 | queuedSince = time.Time{} | |
| 86 | } | |
| 87 | wait := followPoll | |
| 88 | if status == "pending" { | |
| 89 | left := queuedSince.Add(followQueued).Sub(time.Now()) | |
| 90 | if left <= 0 { | |
| 91 | fmt.Fprintf(c.Stderr, "build %d is still queued; nothing claimed it in %v. Follow again once a runner has.\n", b.Number, followQueued) | |
| 92 | return protocol.ExitFailure | |
| 93 | } | |
| 94 | wait = min(wait, left) | |
| 95 | } | |
| 96 | if status != "pending" && status != "running" { | |
| 97 | if settleBy.IsZero() { | |
| 98 | settleBy = time.Now().Add(followSettle) | |
| 99 | } | |
| 100 | left := time.Until(settleBy) | |
| 101 | if left <= 0 && len(chunk) == 0 { | |
| 102 | fmt.Fprintf(c.Stderr, "build %d %s\n", b.Number, status) | |
| 103 | return protocol.ExitOK | |
| 104 | } | |
| 105 | wait = min(wait, max(left, 0)) | |
| 106 | } | |
| 107 | t := time.NewTimer(wait) | |
| 108 | select { | |
| 109 | case <-wake: | |
| 110 | case <-t.C: | |
| 111 | case <-c.Done: | |
| 112 | t.Stop() | |
| 113 | return protocol.ExitFailure | |
| 114 | } | |
| 115 | t.Stop() | |
| 116 | } | |
| 117 | } | |
internal/control/buildfollow_test.go added +190
| @@ -0,0 +1,190 @@ | ||
| 1 | package control | |
| 2 | ||
| 3 | import ( | |
| 4 | "bytes" | |
| 5 | "strings" | |
| 6 | "sync" | |
| 7 | "testing" | |
| 8 | "time" | |
| 9 | ||
| 10 | "gitbay.org/gitbay/internal/protocol" | |
| 11 | "gitbay.org/gitbay/internal/store" | |
| 12 | ) | |
| 13 | ||
| 14 | // syncBuffer is a bytes.Buffer guarded by a mutex, safe for a test to poll | |
| 15 | // while the follow goroutine is still writing to it. | |
| 16 | type syncBuffer struct { | |
| 17 | mu sync.Mutex | |
| 18 | buf bytes.Buffer | |
| 19 | } | |
| 20 | ||
| 21 | func (b *syncBuffer) Write(p []byte) (int, error) { | |
| 22 | b.mu.Lock() | |
| 23 | defer b.mu.Unlock() | |
| 24 | return b.buf.Write(p) | |
| 25 | } | |
| 26 | ||
| 27 | func (b *syncBuffer) String() string { | |
| 28 | b.mu.Lock() | |
| 29 | defer b.mu.Unlock() | |
| 30 | return b.buf.String() | |
| 31 | } | |
| 32 | ||
| 33 | // follow starts build log --follow on build 1 of repo and returns the | |
| 34 | // buffers and a channel carrying the exit code. | |
| 35 | func follow(t *testing.T, st *store.Store, uid int64, repo store.Repo, done <-chan struct{}) (*syncBuffer, *syncBuffer, chan int) { | |
| 36 | t.Helper() | |
| 37 | u, err := st.UserByID(uid) | |
| 38 | if err != nil { | |
| 39 | t.Fatal(err) | |
| 40 | } | |
| 41 | var out, errOut syncBuffer | |
| 42 | c := &Ctx{User: u, Scope: "full", Store: st, Stdin: strings.NewReader(""), | |
| 43 | Stdout: &out, Stderr: &errOut, Done: done} | |
| 44 | res := make(chan int, 1) | |
| 45 | go func() { res <- Dispatch(c, []string{"build", "log", repo.Path(), "1", "--follow"}) }() | |
| 46 | return &out, &errOut, res | |
| 47 | } | |
| 48 | ||
| 49 | func waitExit(t *testing.T, res chan int) int { | |
| 50 | t.Helper() | |
| 51 | select { | |
| 52 | case code := <-res: | |
| 53 | return code | |
| 54 | case <-time.After(10 * time.Second): | |
| 55 | t.Fatal("follow did not end") | |
| 56 | return -1 | |
| 57 | } | |
| 58 | } | |
| 59 | ||
| 60 | func shortFollowTimers(t *testing.T) { | |
| 61 | settle, poll, queued := followSettle, followPoll, followQueued | |
| 62 | followSettle, followPoll = 200*time.Millisecond, 50*time.Millisecond | |
| 63 | t.Cleanup(func() { followSettle, followPoll, followQueued = settle, poll, queued }) | |
| 64 | } | |
| 65 | ||
| 66 | // The follow prints the stored log, then what arrives, and ends with the | |
| 67 | // outcome on stderr once the build finishes. | |
| 68 | func TestBuildLogFollow(t *testing.T) { | |
| 69 | shortFollowTimers(t) | |
| 70 | st, repo, uid := newQueueTestRepo(t) | |
| 71 | id, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true) | |
| 72 | if err != nil { | |
| 73 | t.Fatal(err) | |
| 74 | } | |
| 75 | st.AppendBuildLog(id, []byte("queued\n")) | |
| 76 | out, errOut, res := follow(t, st, uid, repo, nil) | |
| 77 | ||
| 78 | if _, ok, err := st.ClaimBuild([]int64{repo.ID}, false); err != nil || !ok { | |
| 79 | t.Fatalf("claim: %v %v", ok, err) | |
| 80 | } | |
| 81 | st.AppendBuildLog(id, []byte("step one\n")) | |
| 82 | st.AppendBuildLog(id, []byte("step two\n")) | |
| 83 | if err := st.FinishBuild(id, "success"); err != nil { | |
| 84 | t.Fatal(err) | |
| 85 | } | |
| 86 | if code := waitExit(t, res); code != protocol.ExitOK { | |
| 87 | t.Fatalf("exit %d: %s", code, errOut) | |
| 88 | } | |
| 89 | if got := out.String(); got != "queued\nstep one\nstep two\n" { | |
| 90 | t.Errorf("stdout %q", got) | |
| 91 | } | |
| 92 | if got := strings.TrimSpace(errOut.String()); got != "build 1 success" { | |
| 93 | t.Errorf("stderr %q", got) | |
| 94 | } | |
| 95 | } | |
| 96 | ||
| 97 | // A cancel ends the follow, and the line the cancel appends after the | |
| 98 | // status change still arrives. | |
| 99 | func TestBuildLogFollowCancel(t *testing.T) { | |
| 100 | shortFollowTimers(t) | |
| 101 | st, repo, uid := newQueueTestRepo(t) | |
| 102 | id, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true) | |
| 103 | if err != nil { | |
| 104 | t.Fatal(err) | |
| 105 | } | |
| 106 | out, errOut, res := follow(t, st, uid, repo, nil) | |
| 107 | if err := st.CancelBuild(id); err != nil { | |
| 108 | t.Fatal(err) | |
| 109 | } | |
| 110 | st.AppendBuildLog(id, []byte("cancelled by alice before a runner claimed it\n")) | |
| 111 | if code := waitExit(t, res); code != protocol.ExitOK { | |
| 112 | t.Fatalf("exit %d: %s", code, errOut) | |
| 113 | } | |
| 114 | if !strings.Contains(out.String(), "cancelled by alice") { | |
| 115 | t.Errorf("the cancel line did not arrive: %q", out) | |
| 116 | } | |
| 117 | if got := strings.TrimSpace(errOut.String()); got != "build 1 cancelled" { | |
| 118 | t.Errorf("stderr %q", got) | |
| 119 | } | |
| 120 | } | |
| 121 | ||
| 122 | // Closing Done ends a follow of a build that is still running, even while | |
| 123 | // it is blocked waiting for the next change: the build gets a line, the | |
| 124 | // follower is confirmed to have read it (so it is back in its wait), then | |
| 125 | // Done closes. | |
| 126 | func TestBuildLogFollowDone(t *testing.T) { | |
| 127 | shortFollowTimers(t) | |
| 128 | st, repo, uid := newQueueTestRepo(t) | |
| 129 | id, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true) | |
| 130 | if err != nil { | |
| 131 | t.Fatal(err) | |
| 132 | } | |
| 133 | st.AppendBuildLog(id, []byte("step one\n")) | |
| 134 | done := make(chan struct{}) | |
| 135 | out, _, res := follow(t, st, uid, repo, done) | |
| 136 | ||
| 137 | deadline := time.After(2 * time.Second) | |
| 138 | for !strings.Contains(out.String(), "step one") { | |
| 139 | select { | |
| 140 | case <-deadline: | |
| 141 | t.Fatal("follow never read the appended line") | |
| 142 | case <-time.After(10 * time.Millisecond): | |
| 143 | } | |
| 144 | } | |
| 145 | close(done) | |
| 146 | if code := waitExit(t, res); code != protocol.ExitFailure { | |
| 147 | t.Fatalf("exit %d, want %d", code, protocol.ExitFailure) | |
| 148 | } | |
| 149 | } | |
| 150 | ||
| 151 | // A build that stays pending ends its own follow: nothing reaps a queued | |
| 152 | // build, so the follow must give up on its own. | |
| 153 | func TestBuildLogFollowQueued(t *testing.T) { | |
| 154 | shortFollowTimers(t) | |
| 155 | followQueued = 150 * time.Millisecond | |
| 156 | st, repo, uid := newQueueTestRepo(t) | |
| 157 | if _, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true); err != nil { | |
| 158 | t.Fatal(err) | |
| 159 | } | |
| 160 | _, errOut, res := follow(t, st, uid, repo, nil) | |
| 161 | if code := waitExit(t, res); code != protocol.ExitFailure { | |
| 162 | t.Fatalf("exit %d, want %d: %s", code, protocol.ExitFailure, errOut) | |
| 163 | } | |
| 164 | if !strings.Contains(errOut.String(), "still queued") { | |
| 165 | t.Errorf("stderr %q", errOut) | |
| 166 | } | |
| 167 | } | |
| 168 | ||
| 169 | // An account holding maxFollows is refused another. | |
| 170 | func TestBuildLogFollowCap(t *testing.T) { | |
| 171 | st, repo, uid := newQueueTestRepo(t) | |
| 172 | if _, err := st.CreateBuild(repo.ID, "unit", "abc", "main", `["true"]`, "", "", true); err != nil { | |
| 173 | t.Fatal(err) | |
| 174 | } | |
| 175 | followMu.Lock() | |
| 176 | follows[uid] = maxFollows | |
| 177 | followMu.Unlock() | |
| 178 | t.Cleanup(func() { | |
| 179 | followMu.Lock() | |
| 180 | delete(follows, uid) | |
| 181 | followMu.Unlock() | |
| 182 | }) | |
| 183 | _, errOut, res := follow(t, st, uid, repo, nil) | |
| 184 | if code := waitExit(t, res); code != protocol.ExitDenied { | |
| 185 | t.Fatalf("exit %d, want %d", code, protocol.ExitDenied) | |
| 186 | } | |
| 187 | if !strings.Contains(errOut.String(), "8 follows are already open") { | |
| 188 | t.Errorf("stderr %q", errOut) | |
| 189 | } | |
| 190 | } | |
internal/control/control.go +4
| @@ -40,6 +40,10 @@ type Ctx struct { | ||
| 40 | 40 | // Cmd is the command being run, set by Dispatch, so a usage error can |
| 41 | 41 | // print the registered usage rather than a copy of it. |
| 42 | 42 | Cmd Command |
| 43 | // Done, when the surface has one, closes when nobody is reading any | |
| 44 | // more: the SSH channel closed or the HTTP request ended. A command | |
| 45 | // that runs until something happens (build log --follow) stops on it. | |
| 46 | Done <-chan struct{} | |
| 43 | 47 | } |
| 44 | 48 | |
| 45 | 49 | // usage reports a bad invocation with the command's registered usage, |
internal/httpd/api.go +1
| @@ -73,6 +73,7 @@ func (s *Server) apiCmd(w http.ResponseWriter, r *http.Request) { | ||
| 73 | 73 | JSON: true, |
| 74 | 74 | ViaAPI: true, |
| 75 | 75 | ReadOnly: scope == "read", |
| 76 | Done: r.Context().Done(), | |
| 76 | 77 | } |
| 77 | 78 | code := control.Dispatch(ctx, req.Argv) |
| 78 | 79 | |
internal/httpd/apiread.go +1
| @@ -62,6 +62,7 @@ func (s *Server) apiRead(w http.ResponseWriter, r *http.Request) { | ||
| 62 | 62 | JSON: true, |
| 63 | 63 | ViaAPI: true, |
| 64 | 64 | ReadOnly: true, |
| 65 | Done: r.Context().Done(), | |
| 65 | 66 | } |
| 66 | 67 | code := control.Dispatch(ctx, argv) |
| 67 | 68 | |
internal/httpd/buildpages_test.go +5 −11
| @@ -59,21 +59,15 @@ func TestBuildsPageRendersCommandOutput(t *testing.T) { | ||
| 59 | 59 | |
| 60 | 60 | func TestBuildPageRendersCommandOutput(t *testing.T) { |
| 61 | 61 | var sb strings.Builder |
| 62 | err := web.Render(&sb, "build.html", struct { | |
| 63 | repoPage | |
| 64 | Build control.BuildOut | |
| 65 | Log string | |
| 66 | CanWrite bool | |
| 67 | Notice string | |
| 68 | }{ | |
| 69 | testRepoPage(), | |
| 70 | control.BuildOut{ | |
| 62 | err := web.Render(&sb, "build.html", buildView{ | |
| 63 | repoPage: testRepoPage(), | |
| 64 | Build: control.BuildOut{ | |
| 71 | 65 | Number: 60, Job: "build", Status: "success", |
| 72 | 66 | SHA: "ff6271a9d4570cd46f169091637a9d2e40ad5c2b", Ref: "cli-coverage", |
| 73 | 67 | CreatedAt: "2026-08-28T04:42:54Z", FinishedAt: "2026-08-28T04:43:06Z", |
| 74 | 68 | }, |
| 75 | "step 1 ok", | |
| 76 | true, "", | |
| 69 | Log: "step 1 ok", | |
| 70 | CanWrite: true, | |
| 77 | 71 | }) |
| 78 | 72 | if err != nil { |
| 79 | 73 | t.Fatalf("render: %v", err) |
internal/httpd/builds.go +102 −8
| @@ -1,12 +1,20 @@ | ||
| 1 | 1 | package httpd |
| 2 | 2 | |
| 3 | 3 | import ( |
| 4 | "bytes" | |
| 5 | "fmt" | |
| 6 | "html/template" | |
| 7 | "io" | |
| 4 | 8 | "net/http" |
| 5 | 9 | "net/url" |
| 6 | 10 | "slices" |
| 7 | 11 | "strconv" |
| 12 | "strings" | |
| 8 | 13 | |
| 9 | 14 | "gitbay.org/gitbay/internal/control" |
| 15 | "gitbay.org/gitbay/internal/protocol" | |
| 16 | "gitbay.org/gitbay/internal/store" | |
| 17 | "gitbay.org/gitbay/internal/web" | |
| 10 | 18 | ) |
| 11 | 19 | |
| 12 | 20 | // buildFilter is the builds page's GET filter: branch, status and job, |
| @@ -295,13 +303,99 @@ func (s *Server) build(w http.ResponseWriter, r *http.Request) { | ||
| 295 | 303 | s.notFound(w, r) |
| 296 | 304 | return |
| 297 | 305 | } |
| 298 | log, _, _ := s.runControl(viewer, []string{"build", "log", p.Repo.Path(), n}) | |
| 306 | v := buildView{repoPage: p, Build: b, CanWrite: s.canWriteRepo(r, p.Repo), Notice: s.takeFlash(w, r)} | |
| 307 | if (b.Status == "pending" || b.Status == "running") && r.URL.Query().Get("follow") != "0" && r.Method == http.MethodGet { | |
| 308 | s.streamBuild(w, r, v, viewer, n) | |
| 309 | return | |
| 310 | } | |
| 311 | v.Log, _, _ = s.runControl(viewer, []string{"build", "log", p.Repo.Path(), n}) | |
| 312 | s.render(w, "build.html", v) | |
| 313 | } | |
| 299 | 314 | |
| 300 | s.render(w, "build.html", struct { | |
| 301 | repoPage | |
| 302 | Build control.BuildOut | |
| 303 | Log string | |
| 304 | CanWrite bool | |
| 305 | Notice string | |
| 306 | }{p, b, log, s.canWriteRepo(r, p.Repo), s.takeFlash(w, r)}) | |
| 315 | type buildView struct { | |
| 316 | repoPage | |
| 317 | Build control.BuildOut | |
| 318 | Log string | |
| 319 | Live bool | |
| 320 | CanWrite bool | |
| 321 | Notice string | |
| 322 | } | |
| 323 | ||
| 324 | // liveLogMarker stands in for the log when build.html is rendered for a | |
| 325 | // live build; streamBuild splits the page there and streams the log into | |
| 326 | // the gap. Git refs, paths and job names cannot hold the control byte. | |
| 327 | const liveLogMarker = "\x1elive-log\x1e" | |
| 328 | ||
| 329 | // streamBuild writes the build page with the log following the build: | |
| 330 | // the page up to the log, then build log --follow escaped and flushed as | |
| 331 | // it arrives, then the outcome and the rest of the page. | |
| 332 | func (s *Server) streamBuild(w http.ResponseWriter, r *http.Request, v buildView, viewer store.User, n string) { | |
| 333 | v.Live, v.Log = true, liveLogMarker | |
| 334 | var buf bytes.Buffer | |
| 335 | if err := web.Render(&buf, "build.html", v); err != nil { | |
| 336 | http.Error(w, "template error: "+err.Error(), http.StatusInternalServerError) | |
| 337 | return | |
| 338 | } | |
| 339 | head, tail, ok := strings.Cut(buf.String(), liveLogMarker) | |
| 340 | if !ok || !strings.HasPrefix(tail, "</pre>") { | |
| 341 | http.Error(w, "template error: build.html has no live log slot", http.StatusInternalServerError) | |
| 342 | return | |
| 343 | } | |
| 344 | tail = strings.TrimPrefix(tail, "</pre>") | |
| 345 | ||
| 346 | h := w.Header() | |
| 347 | h.Set("Content-Type", "text/html; charset=utf-8") | |
| 348 | h.Set("Cache-Control", "no-store") | |
| 349 | h.Set("X-Accel-Buffering", "no") | |
| 350 | rc := http.NewResponseController(w) | |
| 351 | io.WriteString(w, head) | |
| 352 | rc.Flush() | |
| 353 | ||
| 354 | path := v.Repo.Path() | |
| 355 | msg, code := s.runControlStream(viewer, []string{"build", "log", path, n, "--follow"}, | |
| 356 | htmlStream{w: w, rc: rc}, r.Context().Done()) | |
| 357 | if r.Context().Err() != nil { | |
| 358 | // The client left; nothing more to write. | |
| 359 | return | |
| 360 | } | |
| 361 | if code == protocol.ExitDenied { | |
| 362 | // The follow cap: the stored log once, and why it is not live. | |
| 363 | log, _, _ := s.runControl(viewer, []string{"build", "log", path, n}) | |
| 364 | template.HTMLEscape(w, []byte(log)) | |
| 365 | } | |
| 366 | io.WriteString(w, "</pre>") | |
| 367 | switch { | |
| 368 | case code == protocol.ExitOK: | |
| 369 | var b control.BuildOut | |
| 370 | if _, ok := s.runControlInto(viewer, []string{"build", "show", path, n}, &b); ok { | |
| 371 | fmt.Fprintf(w, `<p class="notice" role="status">build finished: %s</p>`, template.HTMLEscapeString(b.Status)) | |
| 372 | } | |
| 373 | case code == protocol.ExitDenied: | |
| 374 | if viewer.ID == 0 { | |
| 375 | msg = "Too many signed-out viewers are watching live builds. This is the log so far; reload to try again, or sign in." | |
| 376 | } | |
| 377 | fmt.Fprintf(w, `<p class="error" role="alert">%s</p>`, template.HTMLEscapeString(msg)) | |
| 378 | case code == protocol.ExitFailure && msg != "": | |
| 379 | fmt.Fprintf(w, `<p class="notice" role="status">%s</p>`, template.HTMLEscapeString(msg)) | |
| 380 | } | |
| 381 | io.WriteString(w, tail) | |
| 382 | } | |
| 383 | ||
| 384 | // htmlStream escapes each chunk of a streamed log into the page and | |
| 385 | // flushes it, so the browser draws it as it arrives. | |
| 386 | type htmlStream struct { | |
| 387 | w io.Writer | |
| 388 | rc *http.ResponseController | |
| 389 | } | |
| 390 | ||
| 391 | func (h htmlStream) Write(p []byte) (int, error) { | |
| 392 | var buf bytes.Buffer | |
| 393 | template.HTMLEscape(&buf, p) | |
| 394 | if _, err := h.w.Write(buf.Bytes()); err != nil { | |
| 395 | return 0, err | |
| 396 | } | |
| 397 | if err := h.rc.Flush(); err != nil { | |
| 398 | return 0, err | |
| 399 | } | |
| 400 | return len(p), nil | |
| 307 | 401 | } |
internal/httpd/compress.go +24
| @@ -83,3 +83,27 @@ func (g *gzipWriter) Close() { | ||
| 83 | 83 | g.gz.Close() |
| 84 | 84 | } |
| 85 | 85 | } |
| 86 | ||
| 87 | // Flush sends what the gzip stream holds, then flushes the connection, so | |
| 88 | // a streamed page reaches the browser as it is written. | |
| 89 | func (g *gzipWriter) Flush() { | |
| 90 | g.FlushError() | |
| 91 | } | |
| 92 | ||
| 93 | // FlushError is Flush with the error a failed flush produces. | |
| 94 | // http.ResponseController.Flush prefers this over Flush when both are | |
| 95 | // implemented, so a write that fails partway through a stream is reported | |
| 96 | // instead of silently dropped. | |
| 97 | func (g *gzipWriter) FlushError() error { | |
| 98 | if !g.decided { | |
| 99 | g.decide(http.StatusOK) | |
| 100 | } | |
| 101 | if g.gz != nil { | |
| 102 | if err := g.gz.Flush(); err != nil { | |
| 103 | return err | |
| 104 | } | |
| 105 | } | |
| 106 | return http.NewResponseController(g.ResponseWriter).Flush() | |
| 107 | } | |
| 108 | ||
| 109 | func (g *gzipWriter) Unwrap() http.ResponseWriter { return g.ResponseWriter } | |
internal/httpd/compress_test.go +29
| @@ -6,6 +6,7 @@ import ( | ||
| 6 | 6 | "io" |
| 7 | 7 | "net/http" |
| 8 | 8 | "net/http/httptest" |
| 9 | "strings" | |
| 9 | 10 | "testing" |
| 10 | 11 | |
| 11 | 12 | "gitbay.org/gitbay/internal/config" |
| @@ -79,3 +80,31 @@ func TestBinaryResponsesPassThrough(t *testing.T) { | ||
| 79 | 80 | t.Fatalf("git transport touched: %q %q", rec.Header().Get("Content-Encoding"), rec.Body.String()) |
| 80 | 81 | } |
| 81 | 82 | } |
| 83 | ||
| 84 | // A flush mid-response reaches the connection with what was written so | |
| 85 | // far decodable, which is what lets a page stream through gzip. | |
| 86 | func TestGzipWriterFlushes(t *testing.T) { | |
| 87 | rec := httptest.NewRecorder() | |
| 88 | h := compressed(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { | |
| 89 | w.Header().Set("Content-Type", "text/html; charset=utf-8") | |
| 90 | io.WriteString(w, "<p>first</p>") | |
| 91 | if err := http.NewResponseController(w).Flush(); err != nil { | |
| 92 | t.Fatalf("flush: %v", err) | |
| 93 | } | |
| 94 | if !rec.Flushed { | |
| 95 | t.Fatal("the flush did not reach the connection") | |
| 96 | } | |
| 97 | zr, err := gzip.NewReader(bytes.NewReader(rec.Body.Bytes())) | |
| 98 | if err != nil { | |
| 99 | t.Fatalf("gzip header: %v", err) | |
| 100 | } | |
| 101 | got, _ := io.ReadAll(zr) // no trailer yet: ends in ErrUnexpectedEOF | |
| 102 | if !strings.Contains(string(got), "<p>first</p>") { | |
| 103 | t.Fatalf("flushed body decodes to %q", got) | |
| 104 | } | |
| 105 | io.WriteString(w, "<p>second</p>") | |
| 106 | })) | |
| 107 | req := httptest.NewRequest("GET", "/", nil) | |
| 108 | req.Header.Set("Accept-Encoding", "gzip") | |
| 109 | h.ServeHTTP(rec, req) | |
| 110 | } | |
internal/httpd/control.go +22
| @@ -3,6 +3,7 @@ package httpd | ||
| 3 | 3 | import ( |
| 4 | 4 | "bytes" |
| 5 | 5 | "encoding/json" |
| 6 | "io" | |
| 6 | 7 | "net/http" |
| 7 | 8 | "strings" |
| 8 | 9 | |
| @@ -50,6 +51,27 @@ func (s *Server) runControlCode(u store.User, argv []string) (out string, msg st | ||
| 50 | 51 | return stdout.String(), m, code |
| 51 | 52 | } |
| 52 | 53 | |
| 54 | // runControlStream runs a command whose output is written as it is | |
| 55 | // produced: stdout goes to out, and done ends the command when the | |
| 56 | // request does. msg is stderr. | |
| 57 | func (s *Server) runControlStream(u store.User, argv []string, out io.Writer, done <-chan struct{}) (msg string, code int) { | |
| 58 | var stderr bytes.Buffer | |
| 59 | ctx := &control.Ctx{ | |
| 60 | User: u, | |
| 61 | Source: "web", | |
| 62 | Scope: "full", | |
| 63 | Store: s.st, | |
| 64 | Cfg: s.cfg, | |
| 65 | Stdin: strings.NewReader(""), | |
| 66 | Stdout: out, | |
| 67 | Stderr: &stderr, | |
| 68 | ViaAPI: true, | |
| 69 | Done: done, | |
| 70 | } | |
| 71 | code = control.Dispatch(ctx, argv) | |
| 72 | return strings.TrimSpace(stderr.String()), code | |
| 73 | } | |
| 74 | ||
| 53 | 75 | // done finishes a form action by exit code: back to the page on success, |
| 54 | 76 | // the 404 page when the thing does not exist, and back to the page with |
| 55 | 77 | // the message for anything else. A refusal is feedback on the page a |
internal/sshd/sshd.go +15 −4
| @@ -238,7 +238,17 @@ func (s *Server) handleSession(sconn *ssh.ServerConn, ch ssh.Channel, reqs <-cha | ||
| 238 | 238 | continue |
| 239 | 239 | } |
| 240 | 240 | req.Reply(true, nil) |
| 241 | code := s.runExec(sconn, ch, payload.Command) | |
| 241 | // x/crypto closes reqs when the client closes the channel. That | |
| 242 | // is how a follow learns nobody is reading: the CLI's shared | |
| 243 | // connection outlives a Ctrl-C, the channel does not. | |
| 244 | done := make(chan struct{}) | |
| 245 | go func() { | |
| 246 | for r := range reqs { | |
| 247 | r.Reply(false, nil) | |
| 248 | } | |
| 249 | close(done) | |
| 250 | }() | |
| 251 | code := s.runExec(sconn, ch, payload.Command, done) | |
| 242 | 252 | sendExit(ch, code) |
| 243 | 253 | return |
| 244 | 254 | case "shell": |
| @@ -260,7 +270,7 @@ func sendExit(ch ssh.Channel, code int) { | ||
| 260 | 270 | ch.SendRequest("exit-status", false, ssh.Marshal(&msg)) |
| 261 | 271 | } |
| 262 | 272 | |
| 263 | func (s *Server) runExec(sconn *ssh.ServerConn, ch ssh.Channel, cmdline string) int { | |
| 273 | func (s *Server) runExec(sconn *ssh.ServerConn, ch ssh.Channel, cmdline string, done <-chan struct{}) int { | |
| 264 | 274 | ext := sconn.Permissions.Extensions |
| 265 | 275 | if blob := ext["anon-key"]; blob != "" { |
| 266 | 276 | return s.runAnonymous(ch, blob, cmdline) |
| @@ -273,7 +283,7 @@ func (s *Server) runExec(sconn *ssh.ServerConn, ch ssh.Channel, cmdline string) | ||
| 273 | 283 | return protocol.ExitDenied |
| 274 | 284 | } |
| 275 | 285 | _ = s.st.TouchSSHKey(keyID) |
| 276 | return Exec(s.cfg, s.st, user, ext["scope"], ext["key-fp"], cmdline, ch, ch, ch.Stderr()) | |
| 286 | return Exec(s.cfg, s.st, user, ext["scope"], ext["key-fp"], cmdline, ch, ch, ch.Stderr(), done) | |
| 277 | 287 | } |
| 278 | 288 | |
| 279 | 289 | // runAnonymous handles a session from an unregistered key: the register |
| @@ -304,7 +314,7 @@ func (s *Server) runAnonymous(ch ssh.Channel, keyB64, cmdline string) int { | ||
| 304 | 314 | // single dispatch path shared by the embedded listener and the system-sshd |
| 305 | 315 | // forced command (gitbayd shell). |
| 306 | 316 | func Exec(cfg config.Config, st *store.Store, user store.User, scope, source, cmdline string, |
| 307 | stdin io.Reader, stdout, stderr io.Writer) int { | |
| 317 | stdin io.Reader, stdout, stderr io.Writer, done <-chan struct{}) int { | |
| 308 | 318 | if user.Disabled { |
| 309 | 319 | fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") |
| 310 | 320 | return protocol.ExitDenied |
| @@ -341,6 +351,7 @@ func Exec(cfg config.Config, st *store.Store, user store.User, scope, source, cm | ||
| 341 | 351 | Stdin: stdin, |
| 342 | 352 | Stdout: stdout, |
| 343 | 353 | Stderr: stderr, |
| 354 | Done: done, | |
| 344 | 355 | } |
| 345 | 356 | return control.Dispatch(ctx, argv) |
| 346 | 357 | } |
internal/store/builds.go +48
| @@ -233,6 +233,7 @@ func (s *Store) AppendBuildLog(id int64, chunk []byte) error { | ||
| 233 | 233 | return err |
| 234 | 234 | } |
| 235 | 235 | if n, _ := res.RowsAffected(); n > 0 { |
| 236 | s.wakeBuild(id) | |
| 236 | 237 | return nil |
| 237 | 238 | } |
| 238 | 239 | // Over the cap. The bounds match exactly once: appending the notice puts |
| @@ -241,6 +242,9 @@ func (s *Store) AppendBuildLog(id int64, chunk []byte) error { | ||
| 241 | 242 | UPDATE builds SET log = log || ? |
| 242 | 243 | WHERE id = ? AND length(log) >= ? AND length(log) < ?`, |
| 243 | 244 | truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice)) |
| 245 | if err == nil { | |
| 246 | s.wakeBuild(id) | |
| 247 | } | |
| 244 | 248 | return err |
| 245 | 249 | } |
| 246 | 250 | |
| @@ -255,6 +259,7 @@ func (s *Store) FinishBuild(id int64, status string) error { | ||
| 255 | 259 | if n, _ := res.RowsAffected(); n == 0 { |
| 256 | 260 | return ErrNotFound |
| 257 | 261 | } |
| 262 | s.wakeBuild(id) | |
| 258 | 263 | return nil |
| 259 | 264 | } |
| 260 | 265 | |
| @@ -331,6 +336,48 @@ func (s *Store) BuildLog(id int64) ([]byte, error) { | ||
| 331 | 336 | return log, err |
| 332 | 337 | } |
| 333 | 338 | |
| 339 | // BuildLogWait returns a channel closed by the next append to, finish of | |
| 340 | // or cancel of the build. Take it before reading, so a change between the | |
| 341 | // read and the wait still wakes the reader. Only this process's writes | |
| 342 | // wake it. | |
| 343 | func (s *Store) BuildLogWait(id int64) <-chan struct{} { | |
| 344 | s.logMu.Lock() | |
| 345 | defer s.logMu.Unlock() | |
| 346 | if s.logWait == nil { | |
| 347 | s.logWait = map[int64]chan struct{}{} | |
| 348 | } | |
| 349 | ch, ok := s.logWait[id] | |
| 350 | if !ok { | |
| 351 | ch = make(chan struct{}) | |
| 352 | s.logWait[id] = ch | |
| 353 | } | |
| 354 | return ch | |
| 355 | } | |
| 356 | ||
| 357 | func (s *Store) wakeBuild(id int64) { | |
| 358 | s.logMu.Lock() | |
| 359 | defer s.logMu.Unlock() | |
| 360 | if ch, ok := s.logWait[id]; ok { | |
| 361 | close(ch) | |
| 362 | delete(s.logWait, id) | |
| 363 | } | |
| 364 | } | |
| 365 | ||
| 366 | // BuildLogFrom returns the build's status and its log past offset bytes, | |
| 367 | // read together so a terminal status comes with every byte before it. | |
| 368 | // The cast matters: || stores the log as text, and substr on text counts | |
| 369 | // characters. | |
| 370 | func (s *Store) BuildLogFrom(id, offset int64) (string, []byte, error) { | |
| 371 | var status string | |
| 372 | var chunk []byte | |
| 373 | err := s.DB.QueryRow(`SELECT status, substr(CAST(log AS BLOB), ?) FROM builds WHERE id = ?`, | |
| 374 | offset+1, id).Scan(&status, &chunk) | |
| 375 | if errors.Is(err, sql.ErrNoRows) { | |
| 376 | return "", nil, ErrNotFound | |
| 377 | } | |
| 378 | return status, chunk, err | |
| 379 | } | |
| 380 | ||
| 334 | 381 | // LatestBuild returns the newest build for a repo, optionally narrowed to |
| 335 | 382 | // one job. It is what a status badge reports. |
| 336 | 383 | func (s *Store) LatestBuild(repoID int64, job string) (Build, error) { |
| @@ -399,6 +446,7 @@ func (s *Store) CancelBuild(id int64) error { | ||
| 399 | 446 | if n, _ := res.RowsAffected(); n == 0 { |
| 400 | 447 | return ErrNotFound |
| 401 | 448 | } |
| 449 | s.wakeBuild(id) | |
| 402 | 450 | return nil |
| 403 | 451 | } |
| 404 | 452 | |
internal/store/builds_test.go +101
| @@ -426,3 +426,104 @@ func TestClaimBuildSkipsUntrustedUnlessAsked(t *testing.T) { | ||
| 426 | 426 | t.Fatalf("untrusted claim: err=%v ok=%v number=%d, want %d", err, ok, b.Number, forkBuild) |
| 427 | 427 | } |
| 428 | 428 | } |
| 429 | ||
| 430 | // A follower's channel closes on each kind of change to its build, and | |
| 431 | // only its build. | |
| 432 | func TestBuildLogWaitWakes(t *testing.T) { | |
| 433 | s := open(t) | |
| 434 | if err := s.MigrateUp(); err != nil { | |
| 435 | t.Fatal(err) | |
| 436 | } | |
| 437 | uid, err := s.CreateUser("cmc", true) | |
| 438 | if err != nil { | |
| 439 | t.Fatal(err) | |
| 440 | } | |
| 441 | if _, err := s.CreateRepo("user", uid, "orgo", "public"); err != nil { | |
| 442 | t.Fatal(err) | |
| 443 | } | |
| 444 | newBuild := func() int64 { | |
| 445 | t.Helper() | |
| 446 | id, err := s.CreateBuild(1, "test", "abc123", "main", `["true"]`, "", "", true) | |
| 447 | if err != nil { | |
| 448 | t.Fatal(err) | |
| 449 | } | |
| 450 | return id | |
| 451 | } | |
| 452 | closed := func(ch <-chan struct{}) bool { | |
| 453 | select { | |
| 454 | case <-ch: | |
| 455 | return true | |
| 456 | default: | |
| 457 | return false | |
| 458 | } | |
| 459 | } | |
| 460 | ||
| 461 | a, b := newBuild(), newBuild() | |
| 462 | wa, wb := s.BuildLogWait(a), s.BuildLogWait(b) | |
| 463 | if err := s.AppendBuildLog(a, []byte("x")); err != nil { | |
| 464 | t.Fatal(err) | |
| 465 | } | |
| 466 | if !closed(wa) { | |
| 467 | t.Error("append did not wake its build") | |
| 468 | } | |
| 469 | if closed(wb) { | |
| 470 | t.Error("append woke another build") | |
| 471 | } | |
| 472 | ||
| 473 | mustClaim(t, s, 1) // claims a, the oldest | |
| 474 | wa = s.BuildLogWait(a) | |
| 475 | if err := s.FinishBuild(a, "success"); err != nil { | |
| 476 | t.Fatal(err) | |
| 477 | } | |
| 478 | if !closed(wa) { | |
| 479 | t.Error("finish did not wake") | |
| 480 | } | |
| 481 | ||
| 482 | if err := s.CancelBuild(b); err != nil { | |
| 483 | t.Fatal(err) | |
| 484 | } | |
| 485 | if !closed(wb) { | |
| 486 | t.Error("cancel did not wake") | |
| 487 | } | |
| 488 | } | |
| 489 | ||
| 490 | // Offsets are bytes, not characters: || stores the log as text, and a | |
| 491 | // multibyte character must not shift where the next read starts. | |
| 492 | func TestBuildLogFrom(t *testing.T) { | |
| 493 | s := open(t) | |
| 494 | if err := s.MigrateUp(); err != nil { | |
| 495 | t.Fatal(err) | |
| 496 | } | |
| 497 | uid, err := s.CreateUser("cmc", true) | |
| 498 | if err != nil { | |
| 499 | t.Fatal(err) | |
| 500 | } | |
| 501 | if _, err := s.CreateRepo("user", uid, "orgo", "public"); err != nil { | |
| 502 | t.Fatal(err) | |
| 503 | } | |
| 504 | id, err := s.CreateBuild(1, "test", "abc123", "main", `["true"]`, "", "", true) | |
| 505 | if err != nil { | |
| 506 | t.Fatal(err) | |
| 507 | } | |
| 508 | first := "héllo — ok\n" | |
| 509 | for _, c := range []string{first, "wörld\n"} { | |
| 510 | if err := s.AppendBuildLog(id, []byte(c)); err != nil { | |
| 511 | t.Fatal(err) | |
| 512 | } | |
| 513 | } | |
| 514 | status, all, err := s.BuildLogFrom(id, 0) | |
| 515 | if err != nil || status != "pending" || string(all) != first+"wörld\n" { | |
| 516 | t.Fatalf("from 0: %q %q %v", status, all, err) | |
| 517 | } | |
| 518 | _, rest, err := s.BuildLogFrom(id, int64(len(first))) | |
| 519 | if err != nil || string(rest) != "wörld\n" { | |
| 520 | t.Fatalf("from %d: %q %v", len(first), rest, err) | |
| 521 | } | |
| 522 | _, none, err := s.BuildLogFrom(id, int64(len(all))) | |
| 523 | if err != nil || len(none) != 0 { | |
| 524 | t.Fatalf("from the end: %q %v", none, err) | |
| 525 | } | |
| 526 | if _, _, err := s.BuildLogFrom(9999, 0); err != ErrNotFound { | |
| 527 | t.Fatalf("missing build: %v", err) | |
| 528 | } | |
| 529 | } | |
internal/store/store.go +6
| @@ -12,6 +12,7 @@ import ( | ||
| 12 | 12 | "sort" |
| 13 | 13 | "strconv" |
| 14 | 14 | "strings" |
| 15 | "sync" | |
| 15 | 16 | |
| 16 | 17 | "modernc.org/sqlite" |
| 17 | 18 | ) |
| @@ -21,6 +22,11 @@ var migrationFS embed.FS | ||
| 21 | 22 | |
| 22 | 23 | type Store struct { |
| 23 | 24 | DB *sql.DB |
| 25 | ||
| 26 | // logWait holds one channel per build someone is following, closed | |
| 27 | // by the next change to that build's row (BuildLogWait). | |
| 28 | logMu sync.Mutex | |
| 29 | logWait map[int64]chan struct{} | |
| 24 | 30 | } |
| 25 | 31 | |
| 26 | 32 | // Open opens (creating if needed) the database at path with WAL mode and |
internal/web/templates/build.html +3 −1
| @@ -12,5 +12,7 @@ | ||
| 12 | 12 | {{end}} |
| 13 | 13 | </div> |
| 14 | 14 | <p class="meta">{{.Build.Job}} on {{.Build.Ref}} · <code><a href="/{{.Repo.OwnerName}}/{{.Repo.Name}}/commit/{{.Build.SHA}}">{{printf "%.10s" .Build.SHA}}</a></code> · queued {{when .Build.CreatedAt}}{{if .Build.FinishedAt}} · finished {{when .Build.FinishedAt}}{{end}}</p> |
| 15 | {{if .Log}}<pre class="code buildlog" tabindex="0">{{.Log}}</pre>{{else}}<p class="empty-note">no log yet</p>{{end}} | |
| 15 | {{if .Live}}<p class="meta">Live: the log streams here until the build ends. If it stops without a “build finished” line, reload to pick it up again. <a href="?follow=0">Show it without updates</a></p> | |
| 16 | <pre class="code buildlog" tabindex="0">{{.Log}}</pre> | |
| 17 | {{else if .Log}}<pre class="code buildlog" tabindex="0">{{.Log}}</pre>{{else}}<p class="empty-note">no log yet</p>{{end}} | |
| 16 | 18 | {{end}} |