From 064d8b33bc1be73470fa297ac5760bbc1717f3b7 Mon Sep 17 00:00:00 2001 From: Amp Date: Sat, 26 Sep 2026 08:37:54 +0000 Subject: [PATCH 1/4] fix(cron): bound run history and read recent records from the tail Retain the newest 1,000 outcomes under the job lock and publish complete history snapshots. Add bounded reverse reads and regression coverage for legacy logs, failures, and concurrent stores. Fixes Gitlawb/zero#1086 Amp-Thread-ID: https://ampcode.com/threads/T-01a0dcd6-ad63-74d8-8396-45573c57dcec Co-authored-by: Pierre Bruno --- README.md | 5 + internal/cli/cron_run_test.go | 4 +- internal/cron/append_run_test.go | 156 ++++++++++++++++++++++++++++++- internal/cron/store.go | 98 ++++++++++++++++--- internal/cron/store_test.go | 4 +- 5 files changed, 247 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index d8a0d84a2..c52288339 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/internal/cli/cron_run_test.go b/internal/cli/cron_run_test.go index 90c53a1e7..8603e6ef8 100644 --- a/internal/cli/cron_run_test.go +++ b/internal/cli/cron_run_test.go @@ -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)) } @@ -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) } diff --git a/internal/cron/append_run_test.go b/internal/cron/append_run_test.go index 18f8b0061..98da23f41 100644 --- a/internal/cron/append_run_test.go +++ b/internal/cron/append_run_test.go @@ -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). @@ -24,7 +70,7 @@ 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) } @@ -32,3 +78,111 @@ func TestAppendRunDoesNotResurrectRemovedJob(t *testing.T) { 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 + } +} diff --git a/internal/cron/store.go b/internal/cron/store.go index d23e79145..45c7ae00b 100644 --- a/internal/cron/store.go +++ b/internal/cron/store.go @@ -7,6 +7,7 @@ import ( "fmt" "os" "path/filepath" + "slices" "strings" "time" @@ -16,6 +17,10 @@ import ( 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 @@ -309,10 +314,8 @@ func (s *Store) AppendRun(id string, rec RunRecord) error { 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 @@ -320,30 +323,67 @@ func (s *Store) AppendRun(id string, rec RunRecord) error { 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) { 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 @@ -352,14 +392,42 @@ func (s *Store) Runs(id string) ([]RunRecord, error) { 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 } diff --git a/internal/cron/store_test.go b/internal/cron/store_test.go index 9eb4d5789..213334e78 100644 --- a/internal/cron/store_test.go +++ b/internal/cron/store_test.go @@ -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) } @@ -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) } } From 37778a7286b6b1a6b2f670b0392f19734b4d720f Mon Sep 17 00:00:00 2001 From: Amp Date: Mon, 28 Sep 2026 14:27:45 +0000 Subject: [PATCH 2/4] fix(cron): skip oversized legacy history records Discard entire oversized lines during bounded tail reads so compaction can record new outcomes without weakening atomic publication. Cover exact size boundaries and oversized first, middle, and unterminated final records. Amp-Thread-ID: https://ampcode.com/threads/T-01a0e861-ebee-7764-be1d-e67517cd914c Co-authored-by: Pierre Bruno --- README.md | 2 ++ internal/cron/append_run_test.go | 60 +++++++++++++++++++++++++++++--- internal/cron/store.go | 15 +++++--- 3 files changed, 68 insertions(+), 9 deletions(-) diff --git a/README.md b/README.md index c52288339..7fa64e98d 100644 --- a/README.md +++ b/README.md @@ -335,6 +335,8 @@ 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. +Malformed records and legacy lines of 1 MiB or more are skipped when reading or +compacting history; new records of that size are rejected. ## Extending Zero diff --git a/internal/cron/append_run_test.go b/internal/cron/append_run_test.go index 98da23f41..7ef2f1474 100644 --- a/internal/cron/append_run_test.go +++ b/internal/cron/append_run_test.go @@ -130,12 +130,9 @@ func TestAppendRunFailurePreservesHistory(t *testing.T) { 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) - } + rec := RunRecord{ExitCode: 8, Error: strings.Repeat("x", 1024*1024)} if err := store.AppendRun(job.ID, rec); err == nil { - t.Fatal("expected oversized new or existing record error") + t.Fatal("expected oversized new record error") } got, err := os.ReadFile(path) if err != nil || string(got) != data { @@ -144,6 +141,59 @@ func TestAppendRunFailurePreservesHistory(t *testing.T) { } } +func TestAppendRunWithOversizedHistory(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") + // An oversized unterminated legacy record must not block future outcomes. + if err := os.WriteFile(path, []byte(strings.Repeat("x", maxRunBytes)), 0o600); err != nil { + t.Fatal(err) + } + if err := store.AppendRun(job.ID, RunRecord{ExitCode: 42}); err != nil { + t.Fatal(err) + } + runs, err := store.Runs(job.ID, 1) + if err != nil || len(runs) != 1 || runs[0].ExitCode != 42 { + t.Fatalf("new outcome lost: %+v, %v", runs, err) + } +} + +func TestRunsSkipOversizedLines(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 _, size := range []int{maxRunBytes - 1, maxRunBytes, maxRunBytes + 4096} { + // Valid JSON plus trailing whitespace catches accidental decoding of a + // prefix after dropping only the oversized suffix. The exact boundary + // also distinguishes an oversized line from the largest allowed line. + line := `{"exitCode":99}` + line += strings.Repeat(" ", size-len(line)) + data := line + "\n{\"exitCode\":7}\n" + line + "\n{\"exitCode\":42}\n" + line + if err := os.WriteFile(path, []byte(data), 0o600); err != nil { + t.Fatal(err) + } + runs, err := store.Runs(job.ID, 10) + want := []int{7, 42} + if size < maxRunBytes { + want = []int{99, 7, 99, 42, 99} + } + if err != nil || len(runs) != len(want) { + t.Fatalf("size %d: got %+v, err %v; want %v", size, runs, err, want) + } + for i, rec := range runs { + if rec.ExitCode != want[i] { + t.Fatalf("size %d: record %d = %d, want %d", size, i, rec.ExitCode, want[i]) + } + } + } +} + func TestRunHistoryConcurrentStores(t *testing.T) { store := newTestStore(t) job, err := store.Add(Job{Prompt: "x"}) diff --git a/internal/cron/store.go b/internal/cron/store.go index 45c7ae00b..1f15c8bdd 100644 --- a/internal/cron/store.go +++ b/internal/cron/store.go @@ -382,7 +382,8 @@ func (s *Store) Runs(id string, limit int) ([]RunRecord, error) { } // 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. +// line. Malformed and oversized legacy lines are skipped so they cannot prevent +// recording new outcomes during compaction. 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) { @@ -398,6 +399,7 @@ func (s *Store) readRuns(id string, limit int) ([]RunRecord, error) { } var runs []RunRecord var line []byte + oversized := false decode := func() { slices.Reverse(line) var rec RunRecord @@ -415,10 +417,15 @@ func (s *Store) readRuns(id string, limit int) ([]RunRecord, error) { } for i := n - 1; i >= 0 && len(runs) < limit; i-- { if block[i] == '\n' { - decode() - } else { + if !oversized { + decode() + } + oversized = false + } else if !oversized { if len(line) >= maxRunBytes-1 { - return nil, fmt.Errorf("cron run record exceeds %d bytes", maxRunBytes-1) + line = line[:0] + oversized = true + continue } line = append(line, block[i]) } From d081bf0492982b24ada005cbd36fb1fdc83a6586 Mon Sep 17 00:00:00 2001 From: PierrunoYT Date: Tue, 29 Sep 2026 10:45:12 +0200 Subject: [PATCH 3/4] fix(cron): append the outcome when the history log cannot be replaced On Windows the atomic replace of runs.jsonl fails while another process holds the log open, which dropped every outcome for as long as it stayed open. Keep the replace, and only when it fails append the record on a fresh line. The leading newline terminates any unterminated tail; the blank line it may leave is already skipped by readRuns. Zero's own readers take the job lock, so they never observe a partial append. Co-Authored-By: Claude Opus 5.5 --- internal/cron/append_run_test.go | 47 ++++++++++++++++++++++++++++++++ internal/cron/store.go | 28 ++++++++++++++++++- 2 files changed, 74 insertions(+), 1 deletion(-) diff --git a/internal/cron/append_run_test.go b/internal/cron/append_run_test.go index 7ef2f1474..1ae09c0b4 100644 --- a/internal/cron/append_run_test.go +++ b/internal/cron/append_run_test.go @@ -194,6 +194,53 @@ func TestRunsSkipOversizedLines(t *testing.T) { } } +// An outside reader holding runs.jsonl open makes the Windows replace fail; +// the outcome must still be recorded rather than dropped. +func TestAppendRunWhileHistoryHeldOpen(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") + if err := os.WriteFile(path, []byte("{\"exitCode\":7}\n"), 0o600); err != nil { + t.Fatal(err) + } + held, err := os.Open(path) + if err != nil { + t.Fatal(err) + } + defer held.Close() + if err := store.AppendRun(job.ID, RunRecord{ExitCode: 42}); err != nil { + t.Fatalf("AppendRun with history held open: %v", err) + } + runs, err := store.Runs(job.ID, 10) + if err != nil || len(runs) != 2 || runs[0].ExitCode != 7 || runs[1].ExitCode != 42 { + t.Fatalf("outcome lost while history held open: %+v, %v", runs, err) + } +} + +// The fallback append must start a fresh line even when the log ends in an +// unterminated record, so neither record is corrupted. +func TestAppendRunLineAfterUnterminatedTail(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") + if err := os.WriteFile(path, []byte("{\"exitCode\":7}"), 0o600); err != nil { + t.Fatal(err) + } + if err := appendRunLine(path, []byte("{\"exitCode\":42}")); err != nil { + t.Fatal(err) + } + runs, err := store.Runs(job.ID, 10) + if err != nil || len(runs) != 2 || runs[0].ExitCode != 7 || runs[1].ExitCode != 42 { + t.Fatalf("fallback append corrupted history: %+v, %v", runs, err) + } +} + func TestRunHistoryConcurrentStores(t *testing.T) { store := newTestStore(t) job, err := store.Add(Job{Prompt: "x"}) diff --git a/internal/cron/store.go b/internal/cron/store.go index 1f15c8bdd..578dbe8a3 100644 --- a/internal/cron/store.go +++ b/internal/cron/store.go @@ -360,7 +360,33 @@ func (s *Store) AppendRun(id string, rec RunRecord) error { if err := f.Close(); err != nil { return err } - return fsutil.RenameWithRetry(f.Name(), filepath.Join(dir, "runs.jsonl"), nil) + path := filepath.Join(dir, "runs.jsonl") + err = fsutil.RenameWithRetry(f.Name(), path, nil) + var committed *fsutil.CommittedReplacementCleanupError + if err == nil || errors.As(err, &committed) { + return err + } + // Windows refuses to replace a log another process holds open. Append the + // record instead so the outcome survives; compaction resumes on a later run. + // The leading newline ends any unterminated tail, and readRuns skips the + // blank line it may leave. Zero's own readers hold the job lock. + if appendErr := appendRunLine(path, line); appendErr != nil { + return errors.Join(err, appendErr) + } + return nil +} + +func appendRunLine(path string, line []byte) error { + f, err := os.OpenFile(path, os.O_WRONLY|os.O_APPEND, 0) + if err != nil { + return err + } + record := append(append([]byte{'\n'}, line...), '\n') + if _, err := f.Write(record); err != nil { + _ = f.Close() + return err + } + return f.Close() } // Runs returns the newest limit valid records in append order (oldest first). From 4f4d4b982224f8a9d57ac7f3f6e7ebdcc39bcea8 Mon Sep 17 00:00:00 2001 From: PierrunoYT Date: Tue, 29 Sep 2026 11:00:37 +0200 Subject: [PATCH 4/4] fix(cron): preserve the history DACL and append only on sharing violations Compacting runs.jsonl renamed a temporary file over it, which on Windows replaced a DACL applied to the log with the job directory's inherited one. Publish through ReplaceWithRetry so the log keeps its descriptor. Fall back to appending only when the replace fails with a Windows sharing or lock violation, the error a reader holding the log open produces. Any other failure, such as a persistent access denial, is returned so a log that can never be replaced cannot grow past the retention limit. fsutil.IsSharingOrLockViolation exposes the existing classifier, gated to Windows because errno 32 and 33 mean EPIPE and EDOM elsewhere. Co-Authored-By: Claude Opus 5.5 --- internal/cron/append_run_test.go | 53 +++++++++++++++++++++ internal/cron/store.go | 15 ++++-- internal/cron/store_windows_test.go | 73 +++++++++++++++++++++++++++++ internal/fsutil/rename.go | 8 ++++ internal/fsutil/rename_test.go | 20 ++++++++ 5 files changed, 165 insertions(+), 4 deletions(-) create mode 100644 internal/cron/store_windows_test.go diff --git a/internal/cron/append_run_test.go b/internal/cron/append_run_test.go index 1ae09c0b4..3d50a0b3c 100644 --- a/internal/cron/append_run_test.go +++ b/internal/cron/append_run_test.go @@ -3,11 +3,16 @@ package cron import ( "bytes" "encoding/json" + "errors" "os" "path/filepath" + "runtime" "strings" "sync" + "syscall" "testing" + + "github.com/Gitlawb/zero/internal/fsutil" ) func TestAppendRunRetainsNewestThousand(t *testing.T) { @@ -220,6 +225,54 @@ func TestAppendRunWhileHistoryHeldOpen(t *testing.T) { } } +// Only a sharing or lock violation falls back to appending. Any other replace +// failure, including a persistent access denial, and a committed replacement +// whose backup cleanup failed must leave the log untouched, so history can +// neither grow past retention nor record an outcome twice. +func TestAppendRunFallbackOnlyForSharingViolation(t *testing.T) { + const seed = "{\"exitCode\":7}\n" + for _, tc := range []struct { + name string + err error + append bool + }{ + {"sharing violation", syscall.Errno(32), runtime.GOOS == "windows"}, + {"access denied", os.ErrPermission, false}, + {"other", errors.New("boom"), false}, + {"committed", &fsutil.CommittedReplacementCleanupError{Cause: syscall.Errno(32)}, false}, + } { + t.Run(tc.name, func(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") + if err := os.WriteFile(path, []byte(seed), 0o600); err != nil { + t.Fatal(err) + } + store.replace = func(string, string) error { return tc.err } + err = store.AppendRun(job.ID, RunRecord{ExitCode: 42}) + data, readErr := os.ReadFile(path) + if readErr != nil { + t.Fatal(readErr) + } + if tc.append { + if err != nil || !strings.HasPrefix(string(data), seed) || !strings.Contains(string(data), "\"exitCode\":42") { + t.Fatalf("sharing violation dropped the outcome: err %v, log %q", err, data) + } + return + } + if !errors.Is(err, tc.err) { + t.Fatalf("AppendRun error = %v, want %v", err, tc.err) + } + if string(data) != seed { + t.Fatalf("non-transient replace failure changed the log: %q", data) + } + }) + } +} + // The fallback append must start a fresh line even when the log ends in an // unterminated record, so neither record is corrupted. func TestAppendRunLineAfterUnterminatedTail(t *testing.T) { diff --git a/internal/cron/store.go b/internal/cron/store.go index 578dbe8a3..36d6f656f 100644 --- a/internal/cron/store.go +++ b/internal/cron/store.go @@ -61,6 +61,9 @@ type StoreOptions struct { type Store struct { root string now func() time.Time + // replace overrides the run-history replacement primitive in tests; nil + // uses the platform default. + replace func(src, dst string) error } func NewStore(opts StoreOptions) *Store { @@ -361,15 +364,19 @@ func (s *Store) AppendRun(id string, rec RunRecord) error { return err } path := filepath.Join(dir, "runs.jsonl") - err = fsutil.RenameWithRetry(f.Name(), path, nil) + // ReplaceWithRetry keeps a DACL applied to runs.jsonl itself; a rename would + // publish the temporary file's inherited directory DACL instead. + err = fsutil.ReplaceWithRetry(f.Name(), path, s.replace) var committed *fsutil.CommittedReplacementCleanupError - if err == nil || errors.As(err, &committed) { + if err == nil || errors.As(err, &committed) || !fsutil.IsSharingOrLockViolation(err) { return err } // Windows refuses to replace a log another process holds open. Append the // record instead so the outcome survives; compaction resumes on a later run. - // The leading newline ends any unterminated tail, and readRuns skips the - // blank line it may leave. Zero's own readers hold the job lock. + // Other failures are returned so a log that can never be replaced cannot + // grow past the retention limit. The leading newline ends any unterminated + // tail, and readRuns skips the blank line it may leave. Zero's own readers + // hold the job lock. if appendErr := appendRunLine(path, line); appendErr != nil { return errors.Join(err, appendErr) } diff --git a/internal/cron/store_windows_test.go b/internal/cron/store_windows_test.go new file mode 100644 index 000000000..cc93c2458 --- /dev/null +++ b/internal/cron/store_windows_test.go @@ -0,0 +1,73 @@ +//go:build windows + +package cron + +import ( + "os" + "path/filepath" + "strings" + "testing" + + "golang.org/x/sys/windows" +) + +// Compacting runs.jsonl publishes a temporary file carrying the job directory's +// inherited DACL. A rename would replace a restrictive DACL applied to the log +// itself, so AppendRun must keep the log's own descriptor. +func TestAppendRunPreservesHistoryDACL(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") + if err := os.WriteFile(path, []byte("{\"exitCode\":7}\n"), 0o600); err != nil { + t.Fatal(err) + } + // A protected DACL granting only the owner: distinct from whatever the + // temporary file inherits from the directory. + restricted, err := windows.SecurityDescriptorFromString("D:P(A;;FA;;;OW)") + if err != nil { + t.Skipf("cannot build a test security descriptor: %v", err) + } + dacl, _, err := restricted.DACL() + if err != nil { + t.Skipf("cannot read the test DACL: %v", err) + } + if err := windows.SetNamedSecurityInfo( + path, + windows.SE_FILE_OBJECT, + windows.DACL_SECURITY_INFORMATION|windows.PROTECTED_DACL_SECURITY_INFORMATION, + nil, nil, dacl, nil, + ); err != nil { + t.Skipf("cannot apply a restrictive DACL on this filesystem: %v", err) + } + want := describeDACL(t, path) + if inherited := describeDACL(t, store.jobDir(job.ID)); strings.Contains(inherited, want) { + t.Skip("the job directory already carries the same DACL; this filesystem cannot show the difference") + } + + if err := store.AppendRun(job.ID, RunRecord{ExitCode: 42}); err != nil { + t.Fatal(err) + } + if got := describeDACL(t, path); got != want { + t.Fatalf("DACL after AppendRun = %q, want the log's own %q", got, want) + } + runs, err := store.Runs(job.ID, 10) + if err != nil || len(runs) != 2 || runs[1].ExitCode != 42 { + t.Fatalf("outcome not published by replacement: %+v, %v", runs, err) + } +} + +func describeDACL(t *testing.T, path string) string { + t.Helper() + sd, err := windows.GetNamedSecurityInfo(path, windows.SE_FILE_OBJECT, windows.DACL_SECURITY_INFORMATION) + if err != nil { + t.Skipf("cannot read the security descriptor of %s: %v", path, err) + } + text := sd.String() + if index := strings.Index(text, "D:"); index >= 0 { + return text[index:] + } + return text +} diff --git a/internal/fsutil/rename.go b/internal/fsutil/rename.go index 4f74e8544..a432f012f 100644 --- a/internal/fsutil/rename.go +++ b/internal/fsutil/rename.go @@ -75,6 +75,14 @@ func RenameWithRetry(src, dst string, rename func(src, dst string) error) error return err } +// IsSharingOrLockViolation reports whether err is a Windows sharing or lock +// violation: another open handle on the file denies the requested access. It +// is always false on other platforms, where the same errno values mean +// unrelated errors. +func IsSharingOrLockViolation(err error) bool { + return runtime.GOOS == "windows" && isWindowsSharingOrLockViolation(err) +} + func isWindowsSharingOrLockViolation(err error) bool { var errno syscall.Errno if errors.As(err, &errno) { diff --git a/internal/fsutil/rename_test.go b/internal/fsutil/rename_test.go index 8a64852d2..efa10a0ee 100644 --- a/internal/fsutil/rename_test.go +++ b/internal/fsutil/rename_test.go @@ -2,6 +2,7 @@ package fsutil import ( "errors" + "fmt" "runtime" "syscall" "testing" @@ -45,3 +46,22 @@ func TestRenameWithRetryNonRetryableError(t *testing.T) { t.Errorf("expected only 1 attempt for non-retryable error, got %d", attempts) } } + +// Errno 32 and 33 are EPIPE and EDOM on Unix, so only Windows may classify them. +func TestIsSharingOrLockViolation(t *testing.T) { + onWindows := runtime.GOOS == "windows" + for _, tc := range []struct { + err error + want bool + }{ + {syscall.Errno(32), onWindows}, + {fmt.Errorf("wrapped: %w", syscall.Errno(33)), onWindows}, + {syscall.Errno(5), false}, // ERROR_ACCESS_DENIED + {errors.New("other"), false}, + {nil, false}, + } { + if got := IsSharingOrLockViolation(tc.err); got != tc.want { + t.Errorf("IsSharingOrLockViolation(%v) = %v, want %v", tc.err, got, tc.want) + } + } +}