repositories / bugabinga.net
bugabinga.net
personal infrastructure for bugabinga!
owned by admin
services/luci/internal/inbox/inbox.go
Rawpackage 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
}