repositories / bugabinga.net
bugabinga.net
personal infrastructure for bugabinga!
owned by admin
services/luci/cmd/luci/watch_test.go
Rawpackage main
import (
"bytes"
"context"
"encoding/json"
"errors"
"io"
"os"
"path/filepath"
"strings"
"testing"
"time"
"bugabinga.net/luci/internal/config"
"bugabinga.net/luci/internal/history"
"bugabinga.net/luci/internal/inbox"
"bugabinga.net/luci/internal/logs"
"bugabinga.net/luci/internal/runstate"
"bugabinga.net/luci/internal/trigger"
)
func TestWatchTerminalJSONUsesCanonicalValuesAndCounts(t *testing.T) {
root := t.TempDir()
setRequiredEnv(t, root)
for _, dir := range []string{filepath.Join(root, "inbox"), filepath.Join(root, "data", "active")} {
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
}
now := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
store := history.Store{DataDir: filepath.Join(root, "data"), Now: func() time.Time { return now }}
for _, event := range []history.Event{
{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Rev: "abcdef012345", Status: "success", StartedAt: now.Add(-time.Minute)},
{RunID: "42", ChildID: "002", Repo: "repo", Job: "later", Rev: "abcdef012345", Status: "skipped"},
{RunID: "42", ChildID: "run", Repo: "repo", Job: "push", Rev: "abcdef012345", Status: "success", StartedAt: now.Add(-time.Minute)},
} {
if _, err := store.Append(event); err != nil {
t.Fatal(err)
}
}
var output, stderr bytes.Buffer
code := watchRun(context.Background(), config.Config{DataDir: filepath.Join(root, "data"), InboxDir: filepath.Join(root, "inbox")}, "42", watchOptions{json: true}, &output, &stderr, make(chan time.Time))
if code != 0 || stderr.Len() != 0 {
t.Fatalf("code=%d output=%q stderr=%q", code, output.String(), stderr.String())
}
lines := strings.Split(strings.TrimSpace(output.String()), "\n")
if len(lines) != 4 {
t.Fatalf("lines=%q", output.String())
}
for _, line := range lines {
var event watchEvent
if err := json.Unmarshal([]byte(line), &event); err != nil {
t.Fatalf("invalid NDJSON %q: %v", line, err)
}
if event.At == "" || event.Type == "" {
t.Fatalf("missing stable fields: %#v", event)
}
}
if !strings.Contains(output.String(), `"run_id":"42"`) || !strings.Contains(output.String(), `"revision":"abcdef012345"`) || !strings.Contains(output.String(), `"unsuccessful":1`) {
t.Fatalf("canonical values/count absent: %q", output.String())
}
}
func TestWatchTransitionsQueuedRunningTerminal(t *testing.T) {
root := t.TempDir()
for _, dir := range []string{filepath.Join(root, "inbox"), filepath.Join(root, "data", "active")} {
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
}
ref := "refs/heads/trunk"
if err := inbox.WriteID(filepath.Join(root, "inbox"), "42", trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Ref: &ref}); err != nil {
t.Fatal(err)
}
cfg := config.Config{DataDir: filepath.Join(root, "data"), InboxDir: filepath.Join(root, "inbox")}
ticks := make(chan time.Time, 3)
writer := &signalWatchWriter{ready: make(chan struct{}, 8)}
done := make(chan int, 1)
go func() { done <- watchRun(context.Background(), cfg, "42", watchOptions{}, writer, io.Discard, ticks) }()
<-writer.ready
if err := (runstate.Store{DataDir: cfg.DataDir}).Write(runstate.Active{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Rev: "abc", Ref: ref, Trigger: "manual", StartedAt: time.Now()}); err != nil {
t.Fatal(err)
}
ticks <- time.Now()
<-writer.ready
<-writer.ready
store := history.Store{DataDir: cfg.DataDir}
for _, event := range []history.Event{{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Rev: "abc", Status: "success"}, {RunID: "42", ChildID: "run", Repo: "repo", Job: "push", Rev: "abc", Status: "success"}} {
if _, err := store.Append(event); err != nil {
t.Fatal(err)
}
}
if err := (runstate.Store{DataDir: cfg.DataDir}).Remove("42"); err != nil {
t.Fatal(err)
}
ticks <- time.Now()
if code := <-done; code != 0 {
t.Fatalf("code=%d output=%q", code, writer.String())
}
for _, want := range []string{"SNAPSHOT QUEUED babab-babab-babab-babop", "PARENT RUNNING babab-babab-babab-babop", "CHILD 001 RUNNING repo / check", "PARENT SUCCESS babab-babab-babab-babop", "CHILD 001 SUCCESS repo / check", "SUMMARY SUCCESS children=1 success=1 failed=0 unsuccessful=0"} {
if !strings.Contains(writer.String(), want) {
t.Fatalf("missing %q in %q", want, writer.String())
}
}
if !strings.Contains(writer.String(), "revision: abc") || !strings.Contains(writer.String(), "started:") {
t.Fatalf("human snapshot lacks evidence: %q", writer.String())
}
}
type signalWatchWriter struct {
bytes.Buffer
ready chan struct{}
}
func (w *signalWatchWriter) Write(data []byte) (int, error) {
n, err := w.Buffer.Write(data)
w.ready <- struct{}{}
return n, err
}
func TestWatchLogsRemainIncrementalAcrossSanitizerAndUTF8Boundaries(t *testing.T) {
reader := &chunkWatchReader{chunks: [][]byte{{'a', 0x1b, '['}, {'3', '1', 'm', 0xc3}, {0xa9}}}
stream := newWatchLogStream(reader, false)
streams := map[string]*watchLogStream{"001": stream}
children := map[string]watchRecord{"001": {RunID: "42", ChildID: "001", Status: "running"}}
var output bytes.Buffer
for range 3 {
if _, err := emitWatchLogs("", "42", children, streams, watchOptions{json: true}, &output); err != nil {
t.Fatal(err)
}
}
if got := output.String(); !strings.Contains(got, `"text":"a"`) || !strings.Contains(got, `"text":"é"`) || strings.Contains(got, "\\ufffd") {
t.Fatalf("split ANSI/UTF-8 corrupted: %q", got)
}
children["001"] = watchRecord{RunID: "42", ChildID: "001", Status: "success"}
if _, err := emitWatchLogs("", "42", children, streams, watchOptions{}, io.Discard); err != nil {
t.Fatal(err)
}
closeWatchLogStreams(streams)
if !reader.closed || reader.closeCalls != 1 || len(streams) != 1 || !streams["001"].eof || !streams["001"].closed {
t.Fatalf("terminal stream not finalized once: closed=%v closes=%d streams=%d", reader.closed, reader.closeCalls, len(streams))
}
}
type chunkWatchReader struct {
chunks [][]byte
closed bool
closeCalls int
}
func (r *chunkWatchReader) Read(data []byte) (int, error) {
if len(r.chunks) == 0 {
return 0, io.EOF
}
chunk := r.chunks[0]
r.chunks = r.chunks[1:]
return copy(data, chunk), nil
}
func (r *chunkWatchReader) Close() error { r.closed = true; r.closeCalls++; return nil }
func TestWatchLogsFollowLiveFDThroughFinalizationAndClose(t *testing.T) {
root := t.TempDir()
store := logs.Store{DataDir: root}
file, err := store.CreateLive("42-001")
if err != nil {
t.Fatal(err)
}
if _, err := io.WriteString(file, "first "); err != nil {
t.Fatal(err)
}
if err := file.Close(); err != nil {
t.Fatal(err)
}
children := map[string]watchRecord{"001": {RunID: "42", ChildID: "001", Status: "running"}}
streams := map[string]*watchLogStream{}
var output bytes.Buffer
if _, err := emitWatchLogs(root, "42", children, streams, watchOptions{}, &output); err != nil {
t.Fatal(err)
}
file, err = os.OpenFile(filepath.Join(root, "live", "42-001.log"), os.O_APPEND|os.O_WRONLY, 0)
if err != nil {
t.Fatal(err)
}
if _, err := io.WriteString(file, "tail"); err != nil {
t.Fatal(err)
}
if err := file.Close(); err != nil {
t.Fatal(err)
}
if _, err := store.Finalize("42-001"); err != nil {
t.Fatal(err)
}
children["001"] = watchRecord{RunID: "42", ChildID: "001", Status: "success"}
parent := makeWatchView("42", []history.Event{{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Status: "success"}}, nil, nil)
if parent.terminal || parent.parent.Status != "running" {
t.Fatalf("child completion inferred terminal parent: %#v", parent)
}
for range 4 {
if _, err := emitWatchLogs(root, "42", children, streams, watchOptions{}, &output); err != nil {
t.Fatal(err)
}
}
parent = makeWatchView("42", []history.Event{{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Status: "success"}, {RunID: "42", ChildID: "run", Repo: "repo", Job: "push", Status: "success"}}, nil, nil)
if !parent.terminal || parent.parent.Status != "success" {
t.Fatalf("explicit terminal parent not honored: %#v", parent)
}
if _, err := emitWatchLogs(root, "42", children, streams, watchOptions{}, &output); err != nil {
t.Fatal(err)
}
if got := output.String(); strings.Count(got, "first ") != 1 || strings.Count(got, "tail") != 1 || len(streams) != 1 || !streams["001"].eof || !streams["001"].closed {
t.Fatalf("tail lost, duplicated, or completion marker missing: %q streams=%d", got, len(streams))
}
}
func TestParseWatchRejectsInvalidFlagsAndDuplicateSelector(t *testing.T) {
for _, args := range [][]string{{"one", "two"}, {"--unknown"}, {"--timeout"}, {"--timeout", "0s"}, {"--timeout", "bad"}} {
if _, err := parseWatch(args); err == nil {
t.Fatalf("accepted %q", args)
}
}
}
func TestWatchExitCodesAndTerminalStatuses(t *testing.T) {
for _, test := range []struct {
name, status string
want int
}{{"success", "success", 0}, {"failed", "failed", 1}, {"skipped", "skipped", 1}, {"cancelled", "cancelled", 1}} {
t.Run(test.name, func(t *testing.T) {
view := watchView{known: true, terminal: true, parent: watchRecord{RunID: "42", Status: test.status}, children: map[string]watchRecord{}}
if got := watchRunFromView(context.Background(), config.Config{}, "42", watchOptions{}, io.Discard, io.Discard, make(chan time.Time), view); got != test.want {
t.Fatalf("code=%d want=%d", got, test.want)
}
})
}
cancelled, cancel := context.WithCancel(context.Background())
cancel()
view := watchView{known: true, parent: watchRecord{RunID: "42", Status: "running"}, children: map[string]watchRecord{}}
if got := watchRunFromView(cancelled, config.Config{}, "42", watchOptions{}, io.Discard, io.Discard, make(chan time.Time), view); got != 130 {
t.Fatalf("cancel code=%d", got)
}
timedOut, cancel := context.WithTimeout(context.Background(), 0)
defer cancel()
if got := watchRunFromView(timedOut, config.Config{}, "42", watchOptions{}, io.Discard, io.Discard, make(chan time.Time), view); got != 124 {
t.Fatalf("timeout code=%d", got)
}
dataFile := filepath.Join(t.TempDir(), "file")
if err := os.WriteFile(dataFile, nil, 0o644); err != nil {
t.Fatal(err)
}
if got := watchRun(context.Background(), config.Config{DataDir: dataFile, InboxDir: ""}, "42", watchOptions{}, io.Discard, io.Discard, make(chan time.Time)); got != 2 {
t.Fatalf("operational code=%d", got)
}
}
func TestWatchInitialSnapshotStaysPinnedAcrossClaimGapAndNewerRun(t *testing.T) {
root := t.TempDir()
cfg := config.Config{DataDir: filepath.Join(root, "data"), InboxDir: filepath.Join(root, "inbox")}
for _, dir := range []string{cfg.DataDir, cfg.InboxDir, filepath.Join(cfg.DataDir, "active")} {
if err := os.MkdirAll(dir, 0o755); err != nil {
t.Fatal(err)
}
}
initial := watchView{known: true, parent: watchRecord{RunID: "42", Status: "queued", Repo: "repo", Job: "check"}, children: map[string]watchRecord{}}
ticks := make(chan time.Time, 2)
var output bytes.Buffer
done := make(chan int, 1)
go func() {
done <- watchRunFromView(context.Background(), cfg, "42", watchOptions{}, &output, io.Discard, ticks, initial)
}()
store := history.Store{DataDir: cfg.DataDir}
if _, err := store.Append(history.Event{RunID: "99", ChildID: "run", Repo: "repo", Job: "new", Status: "success"}); err != nil {
t.Fatal(err)
}
if _, err := store.Append(history.Event{RunID: "42", ChildID: "run", Repo: "repo", Job: "check", Status: "success"}); err != nil {
t.Fatal(err)
}
ticks <- time.Now()
if code := <-done; code != 0 {
t.Fatalf("code=%d", code)
}
if got := output.String(); strings.Contains(got, "99") || !strings.Contains(got, "SNAPSHOT QUEUED") || !strings.Contains(got, "PARENT SUCCESS") {
t.Fatalf("pinned snapshot lost: %q", got)
}
}
func TestWatchViewSynthesizesRunningParentAcrossHistoryGap(t *testing.T) {
child := history.Event{RunID: "42", ChildID: "001", Repo: "repo", Job: "check", Rev: "abcdef", Ref: "refs/heads/main", Trigger: "push", Status: "success"}
gap := makeWatchView("42", []history.Event{child}, nil, nil)
if !gap.known || gap.terminal || gap.parent.RunID != "42" || gap.parent.ChildID != "run" || gap.parent.Status != "running" || gap.parent.Repo != "repo" || gap.parent.Revision != "abcdef" {
t.Fatalf("invalid synthesized parent: %#v", gap)
}
next := makeWatchView("42", []history.Event{child}, nil, []runstate.Active{{RunID: "42", ChildID: "002", Repo: "repo", Job: "test", Rev: "abcdef", Ref: "refs/heads/main", Trigger: "push"}})
if next.parent.RunID != "42" || next.parent.Status != "running" || next.parent.Repo != "repo" || next.parent.Revision != "abcdef" {
t.Fatalf("parent changed across next child: %#v", next.parent)
}
terminal := makeWatchView("42", []history.Event{child, {RunID: "42", ChildID: "run", Repo: "repo", Job: "push", Rev: "abcdef", Ref: "refs/heads/main", Trigger: "push", Status: "success"}}, nil, nil)
if !terminal.terminal || terminal.parent.RunID != "42" || terminal.parent.Status != "success" || terminal.parent.Repo != "repo" || terminal.parent.Revision != "abcdef" {
t.Fatalf("explicit terminal parent lost metadata: %#v", terminal)
}
}
func TestWatchViewKeepsTerminalChildOverStaleActive(t *testing.T) {
view := makeWatchView("42", []history.Event{{RunID: "42", ChildID: "001", Status: "success"}, {RunID: "42", ChildID: "run", Status: "success"}}, nil, []runstate.Active{{RunID: "42", ChildID: "001"}})
if view.children["001"].Status != "success" {
t.Fatalf("stale active overwrote terminal child: %#v", view.children["001"])
}
}
type failingWatchWriter struct{}
func (failingWatchWriter) Write([]byte) (int, error) { return 0, errors.New("broken pipe") }
func TestWatchPropagatesOutputFailure(t *testing.T) {
if err := emitWatch(failingWatchWriter{}, watchOptions{}, "snapshot", &watchRecord{RunID: "42", Status: "success"}, nil, nil); err == nil {
t.Fatal("output failure ignored")
}
}