diff --git a/README.md b/README.md index d8a0d84a2..7fa64e98d 100644 --- a/README.md +++ b/README.md @@ -331,6 +331,13 @@ 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. +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 ### 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..3d50a0b3c 100644 --- a/internal/cron/append_run_test.go +++ b/internal/cron/append_run_test.go @@ -1,10 +1,61 @@ 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) { + 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 +75,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 +83,256 @@ 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, Error: strings.Repeat("x", 1024*1024)} + if err := store.AppendRun(job.ID, rec); err == nil { + t.Fatal("expected oversized new record error") + } + got, err := os.ReadFile(path) + if err != nil || string(got) != data { + t.Fatalf("failed append changed history: %v", err) + } + } +} + +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]) + } + } + } +} + +// 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) + } +} + +// 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) { + 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"}) + 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..36d6f656f 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 @@ -56,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 { @@ -309,10 +317,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 +326,98 @@ 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. + if err := f.Close(); err != nil { + return err + } + path := filepath.Join(dir, "runs.jsonl") + // 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) || !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. + // 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) + } + 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() } -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 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) { return nil, nil @@ -352,14 +426,48 @@ 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 + oversized := false + 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' { + if !oversized { + decode() + } + oversized = false + } else if !oversized { + if len(line) >= maxRunBytes-1 { + line = line[:0] + oversized = true + continue + } + 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) } } 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) + } + } +}