Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -331,6 +331,11 @@ zero update --check check for newer releases
zero upgrade download, verify, and install the latest release
```

Cron keeps the newest 1,000 run outcomes per job in `runs.jsonl`. Each new run
atomically replaces the history with the retained tail; existing larger histories
are trimmed on their next run. Archive the file before that run if you need older
outcomes. This retention does not delete session directories or reset fire counts.

## Extending Zero

### Project and personal instructions
Expand Down
4 changes: 2 additions & 2 deletions internal/cli/cron_run_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ func TestCronRunOnceFiresDueJobs(t *testing.T) {
if d.FireCount != 1 || !d.NextRunAt.After(now) {
t.Fatalf("due job not advanced: %+v", d)
}
runs, _ := store.Runs(due.ID)
runs, _ := store.Runs(due.ID, 0)
if len(runs) != 1 {
t.Fatalf("expected 1 run record, got %d", len(runs))
}
Expand Down Expand Up @@ -146,7 +146,7 @@ func TestCronRunPausesUnadvanceableJob(t *testing.T) {
if d.Status != cron.StatusPaused {
t.Fatalf("unadvanceable job must be paused, got status=%q", d.Status)
}
runs, _ := store.Runs(job.ID)
runs, _ := store.Runs(job.ID, 0)
if len(runs) != 1 || runs[0].Error == "" {
t.Fatalf("expected one run record with an error, got %+v", runs)
}
Expand Down
156 changes: 155 additions & 1 deletion internal/cron/append_run_test.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,56 @@
package cron

import (
"bytes"
"encoding/json"
"os"
"path/filepath"
"strings"
"sync"
"testing"
)

func TestAppendRunRetainsNewestThousand(t *testing.T) {
store := newTestStore(t)
job, err := store.Add(Job{Expr: "* * * * *", Prompt: "x"})
if err != nil {
t.Fatal(err)
}
// Seed an oversized legacy log, then cross the boundary again after compaction.
var history bytes.Buffer
for i := 0; i < 1207; i++ {
if err := json.NewEncoder(&history).Encode(RunRecord{ExitCode: i}); err != nil {
t.Fatal(err)
}
}
path := filepath.Join(store.jobDir(job.ID), "runs.jsonl")
if err := os.WriteFile(path, history.Bytes(), 0o600); err != nil {
t.Fatal(err)
}
for _, next := range []int{1207, 1208} {
if err := store.AppendRun(job.ID, RunRecord{ExitCode: next}); err != nil {
t.Fatal(err)
}
data, err := os.ReadFile(path)
if err != nil {
t.Fatal(err)
}
lines := bytes.Split(bytes.TrimSpace(data), []byte("\n"))
if len(lines) != 1000 {
t.Fatalf("retained %d records, want 1000", len(lines))
}
for i, line := range lines {
var rec RunRecord
if err := json.Unmarshal(line, &rec); err != nil {
t.Fatal(err)
}
if rec.ExitCode != next-999+i {
t.Fatalf("record %d = %d, want %d", i, rec.ExitCode, next-999+i)
}
}
}
}

// AppendRun on a job removed mid-run must NOT resurrect its directory (which the
// old unconditional MkdirAll did, leaving an orphaned runs.jsonl with no
// metadata.json).
Expand All @@ -24,11 +70,119 @@ func TestAppendRunDoesNotResurrectRemovedJob(t *testing.T) {
if _, err := os.Stat(store.jobDir(job.ID)); !os.IsNotExist(err) {
t.Fatalf("AppendRun resurrected the removed job directory (stat err = %v)", err)
}
runs, err := store.Runs(job.ID)
runs, err := store.Runs(job.ID, 0)
if err != nil {
t.Fatalf("Runs: %v", err)
}
if len(runs) != 0 {
t.Fatalf("expected no runs for a removed job, got %d", len(runs))
}
}

func TestRunsTail(t *testing.T) {
store := newTestStore(t)
job, err := store.Add(Job{Prompt: "x"})
if err != nil {
t.Fatal(err)
}
var history bytes.Buffer
// A forward scan would fail on this old oversized line. A recent-window
// reader must stop before reaching it, rather than merely cap its output.
history.WriteString(strings.Repeat("x", 1024*1024+1) + "\n")
for i := 0; i < 1103; i++ {
if err := json.NewEncoder(&history).Encode(RunRecord{ExitCode: i, SessionTitle: strings.Repeat("界", 1500)}); err != nil {
t.Fatal(err)
}
}
history.WriteString("bad json\n\n{\"exitCode\":1103}") // no final newline
path := filepath.Join(store.jobDir(job.ID), "runs.jsonl")
if err := os.WriteFile(path, history.Bytes(), 0o600); err != nil {
t.Fatal(err)
}
for _, limit := range []int{1, 7, 999, 1000, 1001, 0, -1} {
runs, err := store.Runs(job.ID, limit)
want := limit
if want <= 0 || want > 1000 {
want = 1000
}
if err != nil || len(runs) != want {
t.Fatalf("limit %d: got %d records, err %v", limit, len(runs), err)
}
for i, run := range runs {
if run.ExitCode != 1104-want+i {
t.Fatalf("limit %d, record %d = %d", limit, i, run.ExitCode)
}
if run.ExitCode != 1103 && run.SessionTitle != strings.Repeat("界", 1500) {
t.Fatal("record spanning read blocks was corrupted")
}
}
}
}

func TestAppendRunFailurePreservesHistory(t *testing.T) {
store := newTestStore(t)
job, err := store.Add(Job{Prompt: "x"})
if err != nil {
t.Fatal(err)
}
path := filepath.Join(store.jobDir(job.ID), "runs.jsonl")
for _, data := range []string{"{\"exitCode\":7}\n", strings.Repeat("x", 1024*1024)} {
if err := os.WriteFile(path, []byte(data), 0o600); err != nil {
t.Fatal(err)
}
rec := RunRecord{ExitCode: 8}
if strings.HasPrefix(data, "{") {
rec.Error = strings.Repeat("x", 1024*1024)
}
if err := store.AppendRun(job.ID, rec); err == nil {
t.Fatal("expected oversized new or existing record error")
}
got, err := os.ReadFile(path)
if err != nil || string(got) != data {
t.Fatalf("failed append changed history: %v", err)
}
}
}

func TestRunHistoryConcurrentStores(t *testing.T) {
store := newTestStore(t)
job, err := store.Add(Job{Prompt: "x"})
if err != nil {
t.Fatal(err)
}
var history bytes.Buffer
for i := 0; i < 1000; i++ {
if err := json.NewEncoder(&history).Encode(RunRecord{ExitCode: -1}); err != nil {
t.Fatal(err)
}
}
if err := os.WriteFile(filepath.Join(store.jobDir(job.ID), "runs.jsonl"), history.Bytes(), 0o600); err != nil {
t.Fatal(err)
}
var wg sync.WaitGroup
for i := 0; i < 12; i++ {
wg.Go(func() {
other := NewStore(StoreOptions{RootDir: store.root})
if err := other.AppendRun(job.ID, RunRecord{ExitCode: i}); err != nil {
t.Error(err)
return
}
runs, err := other.Runs(job.ID, 1000)
if err != nil || len(runs) != 1000 {
t.Errorf("inconsistent history: %d records, err %v", len(runs), err)
}
})
}
wg.Wait()
runs, err := store.Runs(job.ID, 12)
if err != nil || len(runs) != 12 {
t.Fatalf("final tail: %d records, err %v", len(runs), err)
}
seen := make(map[int]bool)
for _, run := range runs {
if run.ExitCode < 0 || run.ExitCode >= 12 || seen[run.ExitCode] {
t.Fatalf("lost or duplicate concurrent append: %+v", runs)
}
seen[run.ExitCode] = true
}
}
98 changes: 83 additions & 15 deletions internal/cron/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
"fmt"
"os"
"path/filepath"
"slices"
"strings"
"time"

Expand All @@ -16,6 +17,10 @@
const (
StatusActive = "active"
StatusPaused = "paused"

// MaxRunHistory is the number of recent outcomes retained per job.
MaxRunHistory = 1000
maxRunBytes = 1024 * 1024
)

// ErrJobNotFound is returned (wrapped) by Get when a job's metadata file is
Expand Down Expand Up @@ -309,41 +314,76 @@
return err
}
defer unlock()
// Bail if the job was removed (e.g. mid-run) — otherwise the MkdirAll below
// would resurrect a deleted job's directory with an orphaned runs.jsonl and no
// metadata.json (which Runs would then still return). Under the per-job lock
// this stat-then-write is race-free against a concurrent Remove.
// Do not resurrect a job removed mid-run. The lock serializes this check
// and publication against Remove and other history writers/readers.
if _, err := os.Stat(filepath.Join(s.jobDir(id), "metadata.json")); err != nil {
if errors.Is(err, os.ErrNotExist) {
return nil
}
return err
}
dir := s.jobDir(id)
if err := os.MkdirAll(dir, 0o700); err != nil {
line, err := json.Marshal(rec)
if err != nil {
return err
}
f, err := os.OpenFile(filepath.Join(dir, "runs.jsonl"), os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o600)
if len(line) >= maxRunBytes {
return fmt.Errorf("cron run record exceeds %d bytes", maxRunBytes-1)
}
runs, err := s.readRuns(id, MaxRunHistory-1)
if err != nil {
return err
}
line, err := json.Marshal(rec)
// Publish a complete snapshot even below the retention boundary so readers
// never see a partial append. Legacy oversized logs compact on their next run.
f, err := os.CreateTemp(dir, "runs-*.tmp")
if err != nil {
return err
}
defer os.Remove(f.Name())
w := bufio.NewWriter(f)
for _, run := range runs {
if err := json.NewEncoder(w).Encode(run); err != nil {
_ = f.Close()
return err
}
}
if _, err := w.Write(append(line, '\n')); err != nil {
_ = f.Close()
return err
}
if _, err := f.Write(append(line, '\n')); err != nil {
if err := w.Flush(); err != nil {
_ = f.Close()
return err
}
// Surface a buffered-write failure that only materializes on Close.
return f.Close()
if err := f.Close(); err != nil {
return err
}
return fsutil.RenameWithRetry(f.Name(), filepath.Join(dir, "runs.jsonl"), nil)
}

func (s *Store) Runs(id string) ([]RunRecord, error) {
// Runs returns the newest limit valid records in append order (oldest first).
// Non-positive limits and limits above MaxRunHistory use MaxRunHistory. It reads
// backwards from the tail, including for legacy logs not yet compacted.
func (s *Store) Runs(id string, limit int) ([]RunRecord, error) {

Check failure on line 369 in internal/cron/store.go

View workflow job for this annotation

GitHub Actions / Code Quality & Lint

unreachable func: Store.Runs
if !validID(id) {
return nil, fmt.Errorf("invalid cron job id %q", id)
}
if limit <= 0 || limit > MaxRunHistory {
limit = MaxRunHistory
}
unlock, err := s.lockJob(id)
if err != nil {
return nil, err
}
defer unlock()
return s.readRuns(id, limit)
}

// readRuns requires the job lock. Memory is bounded by limit records plus one
// line; malformed JSON lines are skipped, as in the original forward reader.
func (s *Store) readRuns(id string, limit int) ([]RunRecord, error) {
f, err := os.Open(filepath.Join(s.jobDir(id), "runs.jsonl"))
if errors.Is(err, os.ErrNotExist) {
return nil, nil
Expand All @@ -352,14 +392,42 @@
return nil, err
}
defer f.Close()
info, err := f.Stat()
if err != nil {
return nil, err
}
var runs []RunRecord
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for scanner.Scan() {
var line []byte
decode := func() {
slices.Reverse(line)
var rec RunRecord
if json.Unmarshal(scanner.Bytes(), &rec) == nil {
if json.Unmarshal(line, &rec) == nil {
runs = append(runs, rec)
}
line = line[:0]
}
var block [4096]byte
for end := info.Size(); end > 0 && len(runs) < limit; {
start := max(int64(0), end-int64(len(block)))
n, err := f.ReadAt(block[:end-start], start)
if err != nil {
return nil, err
}
for i := n - 1; i >= 0 && len(runs) < limit; i-- {
if block[i] == '\n' {
decode()
} else {
if len(line) >= maxRunBytes-1 {
return nil, fmt.Errorf("cron run record exceeds %d bytes", maxRunBytes-1)
}
line = append(line, block[i])
}
}
end = start
}
if len(runs) < limit && len(line) > 0 {
decode()
}
return runs, scanner.Err()
slices.Reverse(runs)
return runs, nil
}
4 changes: 2 additions & 2 deletions internal/cron/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ func TestStoreAppendRun(t *testing.T) {
t.Fatalf("AppendRun: %v", err)
}
}
runs, err := s.Runs(job.ID)
runs, err := s.Runs(job.ID, 0)
if err != nil || len(runs) != 3 || runs[2].ExitCode != 2 {
t.Fatalf("Runs=%v err=%v", runs, err)
}
Expand Down Expand Up @@ -280,7 +280,7 @@ func TestStoreRejectsUnsafeID(t *testing.T) {
if _, err := s.Get(id); err == nil {
t.Fatalf("Get(%q) must be rejected", id)
}
if _, err := s.Runs(id); err == nil {
if _, err := s.Runs(id, 0); err == nil {
t.Fatalf("Runs(%q) must be rejected", id)
}
}
Expand Down
Loading