package inbox import ( "errors" "os" "path/filepath" "strings" "testing" "bugabinga.net/luci/internal/trigger" ) type badJSONEvent struct{} func (badJSONEvent) Validate() error { return nil } func (badJSONEvent) EventKind() string { return "bad" } func (badJSONEvent) MarshalJSON() ([]byte, error) { return nil, errors.New("bad json") } func TestWriteCreatesFinalJSONOnly(t *testing.T) { dir := t.TempDir() ref := "main" id, err := Write(dir, trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Ref: &ref, RequestedBy: "oli"}) if err != nil { t.Fatalf("Write() error = %v", err) } if _, err := os.Stat(filepath.Join(dir, id+".json")); err != nil { t.Fatalf("final event missing: %v", err) } if matches, _ := filepath.Glob(filepath.Join(dir, "*.tmp")); len(matches) != 0 { t.Fatalf("tmp files remain: %#v", matches) } } func TestWriteIDRejectsDuplicateWithoutChangingAcceptedEvent(t *testing.T) { dir := t.TempDir() ref := "main" first := trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "first", Ref: &ref} second := trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "second", Ref: &ref} if err := WriteID(dir, "42", first); err != nil { t.Fatal(err) } if err := WriteID(dir, "42", second); !errors.Is(err, os.ErrExist) { t.Fatalf("duplicate err=%v", err) } data, err := os.ReadFile(filepath.Join(dir, "42.json")) if err != nil || !strings.Contains(string(data), `"job":"first"`) { t.Fatalf("data=%s err=%v", data, err) } } func TestWriteRejectsInvalidEvent(t *testing.T) { if _, err := Write(t.TempDir(), trigger.PushEvent{Kind: "push"}); err == nil { t.Fatal("invalid event written") } } func TestWriteInjectedFilesystemFailures(t *testing.T) { ref := "main" valid := trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Ref: &ref} boom := errors.New("boom") oldWrite := writeFile writeFile = func(string, []byte, os.FileMode) error { return boom } if _, err := Write(t.TempDir(), valid); !errors.Is(err, boom) { t.Fatalf("write err = %v", err) } writeFile = oldWrite oldRename := rename oldRemove := remove removed := false rename = func(string, string) error { return boom } remove = func(string) error { removed = true; return nil } if _, err := Write(t.TempDir(), valid); !errors.Is(err, boom) || !removed { t.Fatalf("rename err = %v removed=%v", err, removed) } rename = oldRename remove = oldRemove } func TestWriteRejectsFilesystemAndMarshalFailures(t *testing.T) { ref := "main" valid := trigger.ManualEvent{Kind: "manual", Repo: "repo", Job: "check", Ref: &ref} root := t.TempDir() file := filepath.Join(root, "file") if err := os.WriteFile(file, []byte("x"), 0o644); err != nil { t.Fatal(err) } if _, err := Write(filepath.Join(file, "inbox"), valid); err == nil { t.Fatal("Write accepted file parent") } if _, err := Write(t.TempDir(), badJSONEvent{}); err == nil { t.Fatal("Write accepted marshal failure") } } func TestClaimOnePrioritizesPushAndManualOverSchedule(t *testing.T) { dir := t.TempDir() scheduled := `{"kind":"schedule","repo":"repo","job":"nightly","rev":"abc","ref":"refs/heads/main","minute":"2026-03-04T03:00Z"}` push := `{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}` if err := os.WriteFile(filepath.Join(dir, "1.json"), []byte(scheduled), 0o644); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "2.json"), []byte(push), 0o644); err != nil { t.Fatal(err) } envelope, err := ClaimOne(dir) if err != nil || envelope.ID != "2" || envelope.Event.EventKind() != "push" { t.Fatalf("envelope=%#v err=%v", envelope, err) } } func TestClaimOneKeepsEventUntilAck(t *testing.T) { dir := t.TempDir() writePush := func(name string) { if err := os.WriteFile(filepath.Join(dir, name+".json"), []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } } writePush("b") writePush("a") envelope, err := ClaimOne(dir) if err != nil { t.Fatalf("ClaimOne() error = %v", err) } if envelope == nil || envelope.ID != "a" { t.Fatalf("envelope = %#v", envelope) } if _, err := os.Stat(envelope.Path); err != nil { t.Fatalf("claimed event missing: %v", err) } if err := Ack(*envelope); err != nil { t.Fatalf("Ack() error = %v", err) } if _, err := os.Stat(envelope.Path); !os.IsNotExist(err) { t.Fatalf("acked event remains: %v", err) } } func TestFailDeadLettersClaimedEvent(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "a.json"), []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } envelope, err := ClaimOne(dir) if err != nil { t.Fatal(err) } failed, err := Fail(*envelope) if err != nil { t.Fatalf("Fail() error = %v", err) } if failed != filepath.Join(dir, "a.failed") { t.Fatalf("failed path = %q", failed) } if _, err := os.Stat(failed); err != nil { t.Fatalf("failed event missing: %v", err) } } func TestClaimOneIgnoresProcessingEvents(t *testing.T) { dir := t.TempDir() path := filepath.Join(dir, "a.json.processing") if err := os.WriteFile(path, []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } envelope, err := ClaimOne(dir) if err != nil || envelope != nil { t.Fatalf("envelope = %#v err=%v", envelope, err) } } func TestClaimOneRetainsEventAfterReadFailure(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "a.json"), []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } oldRead := readFile readFile = func(string) ([]byte, error) { return nil, errors.New("boom") } if _, err := ClaimOne(dir); err == nil { t.Fatal("read failure hidden") } readFile = oldRead if _, err := os.Stat(filepath.Join(dir, "a.json.processing")); err != nil { t.Fatalf("event lost after read failure: %v", err) } } func TestClaimOneDeadLettersInvalidEvent(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "bad.json"), []byte(`{"kind":`), 0o644); err != nil { t.Fatal(err) } if _, err := ClaimOne(dir); err == nil { t.Fatal("invalid JSON accepted") } if _, err := os.Stat(filepath.Join(dir, "bad.failed")); err != nil { t.Fatalf("invalid event not dead-lettered: %v", err) } } func TestClaimOneMissingDirIsEmpty(t *testing.T) { envelope, err := ClaimOne(filepath.Join(t.TempDir(), "missing")) if err != nil || envelope != nil { t.Fatalf("envelope=%#v err=%v", envelope, err) } } func TestClaimOneRejectsBadGlobReadAndJSON(t *testing.T) { if _, err := ClaimOne("["); err == nil { t.Fatal("bad glob accepted") } dir := t.TempDir() if err := os.Mkdir(filepath.Join(dir, "bad.json"), 0o755); err != nil { t.Fatal(err) } if _, err := ClaimOne(dir); err == nil { t.Fatal("directory event accepted") } dir = t.TempDir() if err := os.WriteFile(filepath.Join(dir, "bad.json"), []byte(`{"kind":`), 0o644); err != nil { t.Fatal(err) } if _, err := ClaimOne(dir); err == nil { t.Fatal("invalid JSON accepted") } } func TestClaimAckAndFailPropagateFilesystemErrors(t *testing.T) { boom := errors.New("boom") dir := t.TempDir() path := filepath.Join(dir, "a.json") if err := os.WriteFile(path, []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } oldRename := rename rename = func(string, string) error { return boom } if _, err := ClaimOne(dir); !errors.Is(err, boom) { t.Fatalf("claim err = %v", err) } rename = oldRename envelope, err := ClaimOne(dir) if err != nil { t.Fatal(err) } oldRemove := remove remove = func(string) error { return boom } if err := Ack(*envelope); !errors.Is(err, boom) { t.Fatalf("ack err = %v", err) } remove = oldRemove rename = func(string, string) error { return boom } if _, err := Fail(*envelope); !errors.Is(err, boom) { t.Fatalf("fail err = %v", err) } rename = oldRename } func TestListClassifiesPendingAndExcludesDone(t *testing.T) { dir := t.TempDir() data := []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`) for _, name := range []string{"1.json", "2.json.processing"} { if err := os.WriteFile(filepath.Join(dir, name), data, 0o644); err != nil { t.Fatal(err) } } if err := os.WriteFile(filepath.Join(dir, "2.done"), nil, 0o644); err != nil { t.Fatal(err) } pending, warnings, err := List(dir) if err != nil || len(warnings) != 0 || len(pending) != 1 || pending[0].ID != "1" || pending[0].Phase != "queued" { t.Fatalf("pending=%#v warnings=%v err=%v", pending, warnings, err) } ids, err := DoneIDs(dir) if err != nil || len(ids) != 1 || ids[0] != "2" { t.Fatalf("ids=%v err=%v", ids, err) } } func TestDoneAcknowledgesProcessingBeforeMarker(t *testing.T) { dir := t.TempDir() processing := filepath.Join(dir, "1.json.processing") if err := os.WriteFile(processing, []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } if err := MarkDone(dir, "1"); err != nil { t.Fatal(err) } if err := AckDone(dir, "1"); err != nil { t.Fatal(err) } for _, path := range []string{processing, filepath.Join(dir, "1.done")} { if _, err := os.Stat(path); !os.IsNotExist(err) { t.Fatalf("remains %s: %v", path, err) } } } func TestProcessingListsClaimedEvents(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "1.json.processing"), []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } got, err := Processing(dir) if err != nil || len(got) != 1 || got[0].ID != "1" { t.Fatalf("got=%#v err=%v", got, err) } } func TestProcessingDeadLettersMalformedAndContinues(t *testing.T) { dir := t.TempDir() if err := os.WriteFile(filepath.Join(dir, "bad.json.processing"), []byte("{"), 0o644); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "good.json.processing"), []byte(`{"kind":"push","repo":"repo","old":"0","new":"abc","ref":"refs/heads/main"}`), 0o644); err != nil { t.Fatal(err) } got, err := Processing(dir) if err != nil || len(got) != 1 || got[0].ID != "good" { t.Fatalf("got=%#v err=%v", got, err) } if _, err := os.Stat(filepath.Join(dir, "bad.failed")); err != nil { t.Fatalf("bad event not dead-lettered: %v", err) } }