Luigit
repositories / bugabinga.net

bugabinga.net

personal infrastructure for bugabinga!

owned by admin

services/luci/cmd/luci/watch_test.go

Raw
package 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")
	}
}