package history import ( "bufio" "encoding/json" "errors" "fmt" "io" "os" "path/filepath" "sort" "strings" "sync" "time" "github.com/klauspost/compress/zstd" ) type Store struct { DataDir string Now func() time.Time } type zstdReader interface { io.Reader Close() } type zstdWriter interface { io.Writer Close() error } type appendFile interface { io.Writer Close() error } var newReader = func(r io.Reader) (zstdReader, error) { return zstd.NewReader(r) } var newWriter = func(w io.Writer) (zstdWriter, error) { return zstd.NewWriter(w) } var marshal = json.Marshal var openAppend = func(path string) (appendFile, error) { return os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644) } var appendMu sync.Mutex type Event struct { RunID string `json:"run_id"` ChildID string `json:"child_id,omitempty"` Repo string `json:"repo"` Job string `json:"job"` Rev string `json:"rev"` Ref string `json:"ref,omitempty"` Trigger string `json:"trigger,omitempty"` Status string `json:"status"` Subject string `json:"subject,omitempty"` Detail string `json:"detail,omitempty"` Matrix map[string]string `json:"matrix,omitempty"` StartedAt time.Time `json:"started_at,omitzero"` Time time.Time `json:"time"` } func ReadAll(dataDir string) ([]Event, error) { dir := filepath.Join(dataDir, "events") entries, err := os.ReadDir(dir) if errors.Is(err, os.ErrNotExist) { return nil, nil } if err != nil { return nil, err } var paths []string for _, entry := range entries { if !entry.IsDir() && strings.HasSuffix(entry.Name(), ".jsonl.zst") { paths = append(paths, filepath.Join(dir, entry.Name())) } } sort.Strings(paths) var events []Event for _, path := range paths { file, err := os.Open(path) if err != nil { return nil, err } decoder, err := newReader(file) if err != nil { file.Close() return nil, err } scanner := bufio.NewScanner(decoder) for scanner.Scan() { var event Event if err := json.Unmarshal(scanner.Bytes(), &event); err != nil { decoder.Close() file.Close() return nil, err } events = append(events, event) } scanErr := scanner.Err() decoder.Close() file.Close() if scanErr != nil && scanErr != io.EOF { return nil, scanErr } } return events, nil } func (s Store) Append(event Event) (string, error) { event.Repo = bounded(event.Repo, 4096) event.Job = bounded(event.Job, 4096) event.Rev = bounded(event.Rev, 4096) event.Ref = bounded(event.Ref, 4096) event.Trigger = bounded(event.Trigger, 256) event.Status = bounded(event.Status, 256) event.Detail = bounded(event.Detail, 8192) if event.Matrix != nil { matrix := make(map[string]string, len(event.Matrix)) for key, value := range event.Matrix { matrix[bounded(key, 1024)] = bounded(value, 4096) } event.Matrix = matrix } now := time.Now() if s.Now != nil { now = s.Now() } if event.Time.IsZero() { event.Time = now } event.Time = event.Time.UTC() if !event.StartedAt.IsZero() { event.StartedAt = event.StartedAt.UTC() } data, err := marshal(event) if err != nil { return "", err } if len(data) > 60<<10 { event.Matrix = nil data, err = marshal(event) if err != nil { return "", err } } if len(data) > 60<<10 { return "", fmt.Errorf("history event exceeds 60 KiB") } path := filepath.Join(s.DataDir, "events", now.Format("2006-01")+".jsonl.zst") dir := filepath.Dir(path) if err := os.MkdirAll(dir, 0o755); err != nil { return "", err } appendMu.Lock() defer appendMu.Unlock() temp, err := os.CreateTemp(dir, ".history-*.tmp") if err != nil { return "", err } tmp := temp.Name() if err := temp.Close(); err != nil { os.Remove(tmp) return "", err } defer os.Remove(tmp) file, err := openAppend(tmp) if err != nil { return "", err } if previous, err := os.Open(path); err == nil { _, copyErr := io.Copy(file, previous) closeErr := previous.Close() if copyErr != nil || closeErr != nil { file.Close() return "", errors.Join(copyErr, closeErr) } } else if !errors.Is(err, os.ErrNotExist) { file.Close() return "", err } encoder, err := newWriter(file) if err != nil { file.Close() return "", err } data = append(data, '\n') if _, err := encoder.Write(data); err != nil { encoder.Close() file.Close() return "", err } if err := encoder.Close(); err != nil { file.Close() return "", err } if syncer, ok := file.(interface{ Sync() error }); ok { if err := syncer.Sync(); err != nil { file.Close() return "", err } } if err := file.Close(); err != nil { return "", err } if err := os.Rename(tmp, path); err != nil { return "", err } directory, err := os.Open(dir) if err != nil { return "", err } defer directory.Close() if err := directory.Sync(); err != nil { return "", err } return path, nil } func bounded(value string, max int) string { if len(value) <= max { return value } return value[:max] + "..." } func (s Store) Contains(runID, childID string) (bool, error) { events, err := ReadAll(s.DataDir) if err != nil { return false, err } for _, event := range events { if event.RunID == runID && event.ChildID == childID { return true, nil } } return false, nil } func (s Store) AppendOnce(event Event) (bool, error) { found, err := s.Contains(event.RunID, event.ChildID) if err != nil || found { return false, err } _, err = s.Append(event) return err == nil, err }