package inbox import ( "encoding/json" "errors" "fmt" "log" "os" "path/filepath" "sort" "strings" "sync" "time" "bugabinga.net/luci/internal/trigger" ) type Envelope struct { ID string Path string Event trigger.Event } type Pending struct { ID string Phase string Event trigger.Event } var writeFile = os.WriteFile var rename = os.Rename var remove = os.Remove var readFile = os.ReadFile var stateMu sync.Mutex func Write(dir string, event trigger.Event) (string, error) { if err := event.Validate(); err != nil { return "", err } if err := os.MkdirAll(dir, 0o755); err != nil { return "", err } id := fmt.Sprintf("%d", time.Now().UnixNano()) final := filepath.Join(dir, id+".json") tmp := final + ".tmp" data, err := json.Marshal(event) if err != nil { return "", err } data = append(data, '\n') if err := writeFile(tmp, data, 0o644); err != nil { return "", err } if err := rename(tmp, final); err != nil { _ = remove(tmp) return "", err } return id, nil } func WriteID(dir, id string, event trigger.Event) error { if err := cleanID(id); err != nil { return err } if err := event.Validate(); err != nil { return err } if err := os.MkdirAll(dir, 0o755); err != nil { return err } data, err := json.Marshal(event) if err != nil { return err } data = append(data, '\n') file, err := os.CreateTemp(dir, "."+id+"-*.tmp") if err != nil { return err } tmp := file.Name() defer os.Remove(tmp) if _, err := file.Write(data); err != nil { file.Close() return err } if err := file.Sync(); err != nil { file.Close() return err } if err := file.Close(); err != nil { return err } if err := os.Link(tmp, filepath.Join(dir, id+".json")); err != nil { return err } return syncDir(dir) } func ClaimOne(dir string) (*Envelope, error) { stateMu.Lock() defer stateMu.Unlock() paths, err := filepath.Glob(filepath.Join(dir, "*.json")) if err != nil { return nil, err } if len(paths) == 0 { return nil, nil } sort.Strings(paths) selected := 0 for i, candidate := range paths { data, err := os.ReadFile(candidate) if err != nil { selected = i break } event, err := trigger.Decode(data) if err != nil || event.EventKind() != "schedule" { selected = i break } } path := paths[selected] + ".processing" if err := rename(paths[selected], path); err != nil { return nil, err } id := strings.TrimSuffix(filepath.Base(path), ".json.processing") envelope := Envelope{ID: id, Path: path} data, err := readFile(path) if err != nil { return nil, err } event, err := trigger.Decode(data) if err != nil { _, failErr := Fail(envelope) return nil, errors.Join(err, failErr) } envelope.Event = event return &envelope, nil } func List(dir string) ([]Pending, []error, error) { stateMu.Lock() defer stateMu.Unlock() entries, err := os.ReadDir(dir) if err != nil { return nil, nil, err } done := map[string]bool{} for _, entry := range entries { if strings.HasSuffix(entry.Name(), ".done") { done[strings.TrimSuffix(entry.Name(), ".done")] = true } } var pending []Pending var warnings []error for _, entry := range entries { name := entry.Name() phase, suffix := "", "" switch { case strings.HasSuffix(name, ".json.processing"): phase, suffix = "preparing", ".json.processing" case strings.HasSuffix(name, ".json"): phase, suffix = "queued", ".json" default: continue } id := strings.TrimSuffix(name, suffix) if done[id] { continue } data, err := readFile(filepath.Join(dir, name)) if err != nil { if errors.Is(err, os.ErrNotExist) { continue } warnings = append(warnings, fmt.Errorf("%s: %w", name, err)) continue } event, err := trigger.Decode(data) if err != nil { warnings = append(warnings, fmt.Errorf("%s: %w", name, err)) continue } pending = append(pending, Pending{ID: id, Phase: phase, Event: event}) } sort.Slice(pending, func(i, j int) bool { return pending[i].ID < pending[j].ID }) return pending, warnings, nil } func Processing(dir string) ([]Envelope, error) { paths, err := filepath.Glob(filepath.Join(dir, "*.json.processing")) if err != nil { return nil, err } sort.Strings(paths) out := make([]Envelope, 0, len(paths)) for _, path := range paths { envelope, err := readEnvelope(path) if err != nil { if _, failErr := Fail(envelope); failErr != nil { return nil, errors.Join(err, failErr) } log.Printf("luci inbox warning: %s: %v", path, err) continue } out = append(out, envelope) } return out, nil } func DoneIDs(dir string) ([]string, error) { paths, err := filepath.Glob(filepath.Join(dir, "*.done")) if err != nil { return nil, err } sort.Strings(paths) out := make([]string, 0, len(paths)) for _, path := range paths { out = append(out, strings.TrimSuffix(filepath.Base(path), ".done")) } return out, nil } func MarkDone(dir, id string) error { if err := cleanID(id); err != nil { return err } return writeFile(filepath.Join(dir, id+".done"), nil, 0o644) } func AckDone(dir, id string) error { if err := cleanID(id); err != nil { return err } if err := remove(filepath.Join(dir, id+".json.processing")); err != nil && !errors.Is(err, os.ErrNotExist) { return err } return remove(filepath.Join(dir, id+".done")) } func Ack(envelope Envelope) error { return remove(envelope.Path) } func Fail(envelope Envelope) (string, error) { failed := filepath.Join(filepath.Dir(envelope.Path), envelope.ID+".failed") return failed, rename(envelope.Path, failed) } func readEnvelope(path string) (Envelope, error) { id := strings.TrimSuffix(filepath.Base(path), ".json.processing") envelope := Envelope{ID: id, Path: path} data, err := readFile(path) if err != nil { return envelope, err } event, err := trigger.Decode(data) if err != nil { return envelope, err } envelope.Event = event return envelope, nil } func Sync(dir string) error { return syncDir(dir) } func syncDir(path string) error { dir, err := os.Open(path) if err != nil { return err } defer dir.Close() return dir.Sync() } func cleanID(id string) error { if id == "" || filepath.Base(id) != id { return fmt.Errorf("unsafe inbox id %q", id) } return nil }