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 | trusted proxy; from anyone else the header is ignored. Empty, the | 81 | trusted proxy; from anyone else the header is ignored. Empty, the |
| 82 | default, is right when gitbayd terminates TLS itself. | 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 | ** [web] | 88 | ** [web] |
| 85 | - =mode= — =view_only= (default) | =accounts=. In view_only the mutating | 89 | - =mode= — =view_only= (default) | =accounts=. In view_only the mutating |
| 86 | web routes are never registered; in accounts, browser sessions are | 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 | push; a default-branch push registers or updates them. Tag jobs run on | 38 | push; a default-branch push registers or updates them. Tag jobs run on |
| 39 | a matching tag push and nothing else. | 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 | * The table | 52 | * The table |
| 42 | 53 | ||
| 43 | One =ci.yml= with five jobs, each isolating one rule: | 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 | | build list filters (ref, status, job) | yes | yes | yes | | 210 | | build list filters (ref, status, job) | yes | yes | yes | |
| 211 | | build show (one build) | yes | yes | yes | | 211 | | build show (one build) | yes | yes | yes | |
| 212 | | build log | yes | yes | yes | | 212 | | build log | yes | yes | yes | |
| 213 | | build log follow (until it ends) | yes | yes | no | | ||
| 213 | | build jobs | yes | yes | yes | | 214 | | build jobs | yes | yes | yes | |
| 214 | | build trigger | yes | yes | yes | | 215 | | build trigger | yes | yes | yes | |
| 215 | | build cancel | yes | yes | yes | | 216 | | build cancel | yes | yes | yes | |
cmd/gitbay/main.go +1 −1
| @@ -48,7 +48,7 @@ func newRoot() *cobra.Command { | |||
| 48 | group("build", "CI builds", | 48 | group("build", "CI builds", |
| 49 | pass("list", "recent builds: <owner/name>", passOpts{server: []string{"build", "list"}, needsRepo: true}), | 49 | pass("list", "recent builds: <owner/name>", passOpts{server: []string{"build", "list"}, needsRepo: true}), |
| 50 | pass("show", "one build: <owner/name> <n>", passOpts{server: []string{"build", "show"}, needsRepo: true}), | 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 | pass("jobs", "list the jobs a trigger can name", passOpts{server: []string{"build", "jobs"}, needsRepo: true}), | 52 | pass("jobs", "list the jobs a trigger can name", passOpts{server: []string{"build", "jobs"}, needsRepo: true}), |
| 53 | pass("trigger", "queue a job now: <job>", passOpts{server: []string{"build", "trigger"}, needsRepo: true}), | 53 | pass("trigger", "queue a job now: <job>", passOpts{server: []string{"build", "trigger"}, needsRepo: true}), |
| 54 | pass("cancel", "withdraw a queued build: <n>", passOpts{server: []string{"build", "cancel"}, needsRepo: true}), | 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 | fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n") | 93 | fmt.Fprintf(os.Stderr, "gitbay control plane: interactive shells are not available.\nTry: ssh <host> help\n") |
| 94 | os.Exit(protocol.ExitUsage) | 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 | st.Close() | 97 | st.Close() |
| 98 | os.Exit(code) | 98 | os.Exit(code) |
| 99 | return nil | 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 | build1 := inst.base() + "/alice/app/builds/1" | 46 | build1 := inst.base() + "/alice/app/builds/1" |
| 47 | 47 | ||
| 48 | // Build 1 is queued: the control is on the page, for someone with | 48 | // Build 1 is queued: the control is on the page, for someone with |
| 49 | // write access. | 49 | // write access. follow=0: a live build's page streams until it ends. |
| 50 | _, body := browserGet(t, alice, build1) | 50 | _, body := browserGet(t, alice, build1+"?follow=0") |
| 51 | if !strings.Contains(body, `action="/alice/app/builds/1/cancel"`) { | 51 | if !strings.Contains(body, `action="/alice/app/builds/1/cancel"`) { |
| 52 | t.Fatalf("no cancel control on a queued build:\n%s", body) | 52 | t.Fatalf("no cancel control on a queued build:\n%s", body) |
| 53 | } | 53 | } |
| 54 | // A reader gets no control. | 54 | // A reader gets no control. |
| 55 | _, anon := browserGet(t, newBrowser(t), build1) | 55 | _, anon := browserGet(t, newBrowser(t), build1+"?follow=0") |
| 56 | if strings.Contains(anon, "/cancel") { | 56 | if strings.Contains(anon, "/cancel") { |
| 57 | t.Fatal("anonymous visitor sees the cancel control") | 57 | t.Fatal("anonymous visitor sees the cancel control") |
| 58 | } | 58 | } |
| @@ -78,7 +78,7 @@ func TestBuildCancelWeb(t *testing.T) { | |||
| 78 | t.Fatalf("runner claim: %s", out) | 78 | t.Fatalf("runner claim: %s", out) |
| 79 | } | 79 | } |
| 80 | build2 := inst.base() + "/alice/app/builds/2" | 80 | build2 := inst.base() + "/alice/app/builds/2" |
| 81 | _, body = browserGet(t, alice, build2) | 81 | _, body = browserGet(t, alice, build2+"?follow=0") |
| 82 | if !strings.Contains(body, "running") { | 82 | if !strings.Contains(body, "running") { |
| 83 | t.Fatalf("build 2 not running:\n%s", body) | 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 | return priv | 160 | return priv |
| 161 | } | 161 | } |
| 162 | 162 | ||
| 163 | // ssh runs the real OpenSSH client against the instance with the given key. | 163 | // sshCmd is the ssh invocation ssh runs, for a test that reads the output |
| 164 | func (i *instance) ssh(t *testing.T, key string, stdin string, args ...string) (string, string, int) { | 164 | // as it arrives. |
| 165 | t.Helper() | 165 | func (i *instance) sshCmd(key string, args ...string) *exec.Cmd { |
| 166 | base := []string{ | 166 | base := []string{ |
| 167 | "-p", fmt.Sprint(i.port), | 167 | "-p", fmt.Sprint(i.port), |
| 168 | "-i", key, | 168 | "-i", key, |
| @@ -172,7 +172,13 @@ func (i *instance) ssh(t *testing.T, key string, stdin string, args ...string) ( | |||
| 172 | "-o", "BatchMode=yes", | 172 | "-o", "BatchMode=yes", |
| 173 | "git@127.0.0.1", | 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 | if stdin != "" { | 182 | if stdin != "" { |
| 177 | cmd.Stdin = strings.NewReader(stdin) | 183 | cmd.Stdin = strings.NewReader(stdin) |
| 178 | } | 184 | } |
internal/control/build.go +10 −3
| @@ -27,8 +27,8 @@ func init() { | |||
| 27 | Summary: "show one build", | 27 | Summary: "show one build", |
| 28 | Usage: "build show <owner/name> <n>", ReadOnly: true, Run: runBuildShow}) | 28 | Usage: "build show <owner/name> <n>", ReadOnly: true, Run: runBuildShow}) |
| 29 | register(Command{Path: []string{"build", "log"}, | 29 | register(Command{Path: []string{"build", "log"}, |
| 30 | Summary: "print a build's log", | 30 | Summary: "print a build's log, or follow it until the build ends", |
| 31 | Usage: "build log <owner/name> <n>", ReadOnly: true, Run: runBuildLog}) | 31 | Usage: "build log <owner/name> <n> [--follow]", ReadOnly: true, Run: runBuildLog}) |
| 32 | 32 | ||
| 33 | register(Command{Path: []string{"build", "jobs"}, | 33 | register(Command{Path: []string{"build", "jobs"}, |
| 34 | Summary: "list the jobs a trigger can name", | 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 | func runBuildLog(c *Ctx, args []string) int { | 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 | if code >= 0 { | 204 | if code >= 0 { |
| 201 | return code | 205 | return code |
| 202 | } | 206 | } |
| 207 | if f.Has("--follow") { | ||
| 208 | return followBuildLog(c, b) | ||
| 209 | } | ||
| 203 | log, err := c.Store.BuildLog(b.ID) | 210 | log, err := c.Store.BuildLog(b.ID) |
| 204 | if err != nil { | 211 | if err != nil { |
| 205 | return c.fail(protocol.ExitFailure, "%v", err) | 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 | // newQueueTestRepo returns a store with one public repo (default branch | 43 | // newQueueTestRepo returns a store with one public repo (default branch |
| 44 | // "main", matching the schema default) and the uid to queue builds as. | 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 | func newQueueTestRepo(t *testing.T) (*store.Store, store.Repo, int64) { | 49 | func newQueueTestRepo(t *testing.T) (*store.Store, store.Repo, int64) { |
| 46 | t.Helper() | 50 | t.Helper() |
| 47 | st, err := store.Open(":memory:") | 51 | st, err := store.Open(filepath.Join(t.TempDir(), "gitbay.db")) |
| 48 | if err != nil { | 52 | if err != nil { |
| 49 | t.Fatal(err) | 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 | // Cmd is the command being run, set by Dispatch, so a usage error can | 40 | // Cmd is the command being run, set by Dispatch, so a usage error can |
| 41 | // print the registered usage rather than a copy of it. | 41 | // print the registered usage rather than a copy of it. |
| 42 | Cmd Command | 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 | // usage reports a bad invocation with the command's registered usage, | 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 | JSON: true, | 73 | JSON: true, |
| 74 | ViaAPI: true, | 74 | ViaAPI: true, |
| 75 | ReadOnly: scope == "read", | 75 | ReadOnly: scope == "read", |
| 76 | Done: r.Context().Done(), | ||
| 76 | } | 77 | } |
| 77 | code := control.Dispatch(ctx, req.Argv) | 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 | JSON: true, | 62 | JSON: true, |
| 63 | ViaAPI: true, | 63 | ViaAPI: true, |
| 64 | ReadOnly: true, | 64 | ReadOnly: true, |
| 65 | Done: r.Context().Done(), | ||
| 65 | } | 66 | } |
| 66 | code := control.Dispatch(ctx, argv) | 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 | func TestBuildPageRendersCommandOutput(t *testing.T) { | 60 | func TestBuildPageRendersCommandOutput(t *testing.T) { |
| 61 | var sb strings.Builder | 61 | var sb strings.Builder |
| 62 | err := web.Render(&sb, "build.html", struct { | 62 | err := web.Render(&sb, "build.html", buildView{ |
| 63 | repoPage | 63 | repoPage: testRepoPage(), |
| 64 | Build control.BuildOut | 64 | Build: control.BuildOut{ |
| 65 | Log string | ||
| 66 | CanWrite bool | ||
| 67 | Notice string | ||
| 68 | }{ | ||
| 69 | testRepoPage(), | ||
| 70 | control.BuildOut{ | ||
| 71 | Number: 60, Job: "build", Status: "success", | 65 | Number: 60, Job: "build", Status: "success", |
| 72 | SHA: "ff6271a9d4570cd46f169091637a9d2e40ad5c2b", Ref: "cli-coverage", | 66 | SHA: "ff6271a9d4570cd46f169091637a9d2e40ad5c2b", Ref: "cli-coverage", |
| 73 | CreatedAt: "2026-08-28T04:42:54Z", FinishedAt: "2026-08-28T04:43:06Z", | 67 | CreatedAt: "2026-08-28T04:42:54Z", FinishedAt: "2026-08-28T04:43:06Z", |
| 74 | }, | 68 | }, |
| 75 | "step 1 ok", | 69 | Log: "step 1 ok", |
| 76 | true, "", | 70 | CanWrite: true, |
| 77 | }) | 71 | }) |
| 78 | if err != nil { | 72 | if err != nil { |
| 79 | t.Fatalf("render: %v", err) | 73 | t.Fatalf("render: %v", err) |
internal/httpd/builds.go +102 −8
| @@ -1,12 +1,20 @@ | |||
| 1 | package httpd | 1 | package httpd |
| 2 | 2 | ||
| 3 | import ( | 3 | import ( |
| 4 | "bytes" | ||
| 5 | "fmt" | ||
| 6 | "html/template" | ||
| 7 | "io" | ||
| 4 | "net/http" | 8 | "net/http" |
| 5 | "net/url" | 9 | "net/url" |
| 6 | "slices" | 10 | "slices" |
| 7 | "strconv" | 11 | "strconv" |
| 12 | "strings" | ||
| 8 | 13 | ||
| 9 | "gitbay.org/gitbay/internal/control" | 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 | // buildFilter is the builds page's GET filter: branch, status and job, | 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 | s.notFound(w, r) | 303 | s.notFound(w, r) |
| 296 | return | 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 { | 315 | type buildView struct { |
| 301 | repoPage | 316 | repoPage |
| 302 | Build control.BuildOut | 317 | Build control.BuildOut |
| 303 | Log string | 318 | Log string |
| 304 | CanWrite bool | 319 | Live bool |
| 305 | Notice string | 320 | CanWrite bool |
| 306 | }{p, b, log, s.canWriteRepo(r, p.Repo), s.takeFlash(w, r)}) | 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 | g.gz.Close() | 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 | "io" | 6 | "io" |
| 7 | "net/http" | 7 | "net/http" |
| 8 | "net/http/httptest" | 8 | "net/http/httptest" |
| 9 | "strings" | ||
| 9 | "testing" | 10 | "testing" |
| 10 | 11 | ||
| 11 | "gitbay.org/gitbay/internal/config" | 12 | "gitbay.org/gitbay/internal/config" |
| @@ -79,3 +80,31 @@ func TestBinaryResponsesPassThrough(t *testing.T) { | |||
| 79 | t.Fatalf("git transport touched: %q %q", rec.Header().Get("Content-Encoding"), rec.Body.String()) | 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 | import ( | 3 | import ( |
| 4 | "bytes" | 4 | "bytes" |
| 5 | "encoding/json" | 5 | "encoding/json" |
| 6 | "io" | ||
| 6 | "net/http" | 7 | "net/http" |
| 7 | "strings" | 8 | "strings" |
| 8 | 9 | ||
| @@ -50,6 +51,27 @@ func (s *Server) runControlCode(u store.User, argv []string) (out string, msg st | |||
| 50 | return stdout.String(), m, code | 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 | // done finishes a form action by exit code: back to the page on success, | 75 | // done finishes a form action by exit code: back to the page on success, |
| 54 | // the 404 page when the thing does not exist, and back to the page with | 76 | // the 404 page when the thing does not exist, and back to the page with |
| 55 | // the message for anything else. A refusal is feedback on the page a | 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 | continue | 238 | continue |
| 239 | } | 239 | } |
| 240 | req.Reply(true, nil) | 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 | sendExit(ch, code) | 252 | sendExit(ch, code) |
| 243 | return | 253 | return |
| 244 | case "shell": | 254 | case "shell": |
| @@ -260,7 +270,7 @@ func sendExit(ch ssh.Channel, code int) { | |||
| 260 | ch.SendRequest("exit-status", false, ssh.Marshal(&msg)) | 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 | ext := sconn.Permissions.Extensions | 274 | ext := sconn.Permissions.Extensions |
| 265 | if blob := ext["anon-key"]; blob != "" { | 275 | if blob := ext["anon-key"]; blob != "" { |
| 266 | return s.runAnonymous(ch, blob, cmdline) | 276 | return s.runAnonymous(ch, blob, cmdline) |
| @@ -273,7 +283,7 @@ func (s *Server) runExec(sconn *ssh.ServerConn, ch ssh.Channel, cmdline string) | |||
| 273 | return protocol.ExitDenied | 283 | return protocol.ExitDenied |
| 274 | } | 284 | } |
| 275 | _ = s.st.TouchSSHKey(keyID) | 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 | // runAnonymous handles a session from an unregistered key: the register | 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 | // single dispatch path shared by the embedded listener and the system-sshd | 314 | // single dispatch path shared by the embedded listener and the system-sshd |
| 305 | // forced command (gitbayd shell). | 315 | // forced command (gitbayd shell). |
| 306 | func Exec(cfg config.Config, st *store.Store, user store.User, scope, source, cmdline string, | 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 | if user.Disabled { | 318 | if user.Disabled { |
| 309 | fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") | 319 | fmt.Fprintln(stderr, "this account is disabled; contact the instance admin") |
| 310 | return protocol.ExitDenied | 320 | return protocol.ExitDenied |
| @@ -341,6 +351,7 @@ func Exec(cfg config.Config, st *store.Store, user store.User, scope, source, cm | |||
| 341 | Stdin: stdin, | 351 | Stdin: stdin, |
| 342 | Stdout: stdout, | 352 | Stdout: stdout, |
| 343 | Stderr: stderr, | 353 | Stderr: stderr, |
| 354 | Done: done, | ||
| 344 | } | 355 | } |
| 345 | return control.Dispatch(ctx, argv) | 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 | return err | 233 | return err |
| 234 | } | 234 | } |
| 235 | if n, _ := res.RowsAffected(); n > 0 { | 235 | if n, _ := res.RowsAffected(); n > 0 { |
| 236 | s.wakeBuild(id) | ||
| 236 | return nil | 237 | return nil |
| 237 | } | 238 | } |
| 238 | // Over the cap. The bounds match exactly once: appending the notice puts | 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 | UPDATE builds SET log = log || ? | 242 | UPDATE builds SET log = log || ? |
| 242 | WHERE id = ? AND length(log) >= ? AND length(log) < ?`, | 243 | WHERE id = ? AND length(log) >= ? AND length(log) < ?`, |
| 243 | truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice)) | 244 | truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice)) |
| 245 | if err == nil { | ||
| 246 | s.wakeBuild(id) | ||
| 247 | } | ||
| 244 | return err | 248 | return err |
| 245 | } | 249 | } |
| 246 | 250 | ||
| @@ -255,6 +259,7 @@ func (s *Store) FinishBuild(id int64, status string) error { | |||
| 255 | if n, _ := res.RowsAffected(); n == 0 { | 259 | if n, _ := res.RowsAffected(); n == 0 { |
| 256 | return ErrNotFound | 260 | return ErrNotFound |
| 257 | } | 261 | } |
| 262 | s.wakeBuild(id) | ||
| 258 | return nil | 263 | return nil |
| 259 | } | 264 | } |
| 260 | 265 | ||
| @@ -331,6 +336,48 @@ func (s *Store) BuildLog(id int64) ([]byte, error) { | |||
| 331 | return log, err | 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 | // LatestBuild returns the newest build for a repo, optionally narrowed to | 381 | // LatestBuild returns the newest build for a repo, optionally narrowed to |
| 335 | // one job. It is what a status badge reports. | 382 | // one job. It is what a status badge reports. |
| 336 | func (s *Store) LatestBuild(repoID int64, job string) (Build, error) { | 383 | func (s *Store) LatestBuild(repoID int64, job string) (Build, error) { |
| @@ -399,6 +446,7 @@ func (s *Store) CancelBuild(id int64) error { | |||
| 399 | if n, _ := res.RowsAffected(); n == 0 { | 446 | if n, _ := res.RowsAffected(); n == 0 { |
| 400 | return ErrNotFound | 447 | return ErrNotFound |
| 401 | } | 448 | } |
| 449 | s.wakeBuild(id) | ||
| 402 | return nil | 450 | return nil |
| 403 | } | 451 | } |
| 404 | 452 | ||
internal/store/builds_test.go +101
| @@ -426,3 +426,104 @@ func TestClaimBuildSkipsUntrustedUnlessAsked(t *testing.T) { | |||
| 426 | t.Fatalf("untrusted claim: err=%v ok=%v number=%d, want %d", err, ok, b.Number, forkBuild) | 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 | "sort" | 12 | "sort" |
| 13 | "strconv" | 13 | "strconv" |
| 14 | "strings" | 14 | "strings" |
| 15 | "sync" | ||
| 15 | 16 | ||
| 16 | "modernc.org/sqlite" | 17 | "modernc.org/sqlite" |
| 17 | ) | 18 | ) |
| @@ -21,6 +22,11 @@ var migrationFS embed.FS | |||
| 21 | 22 | ||
| 22 | type Store struct { | 23 | type Store struct { |
| 23 | DB *sql.DB | 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 | // Open opens (creating if needed) the database at path with WAL mode and | 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 | {{end}} | 12 | {{end}} |
| 13 | </div> | 13 | </div> |
| 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> | 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 | {{end}} | 18 | {{end}} |