Luigit
repositories / bugabinga.net

bugabinga.net

personal infrastructure for bugabinga!

owned by admin

services/luci/internal/inbox/inbox.go

Raw
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
}