diff --git a/CHANGELOG.md b/CHANGELOG.md index 690eeb2..d142585 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ## v0.14.9 - Unreleased +- Recover truncated scheduler history tails after interrupted writes while preserving valid final records without a newline and reporting write or cleanup failures. Thanks @SebTardif. - Update SQLite to v1.58.0 with its required libc v1.75.6 runtime, refresh x/crypto and go-runewidth, and prefer Go 1.27.1 while retaining the Go 1.27.0 minimum. - Refresh the pinned TruffleHog secret-scanning action to v3.97.4. Thanks @dependabot. diff --git a/README.md b/README.md index 4f33859..2f914e5 100644 --- a/README.md +++ b/README.md @@ -94,6 +94,8 @@ See the [package guide](docs/packages.md) for the complete inventory and [Go pac `crawlctl` discovers installed crawl apps through their machine-readable metadata, runs configured refresh jobs under a single-process lock, and records JSONL run history. +If a write is interrupted, history reads ignore a truncated final JSON value and the next run repairs that tail before appending. Valid final records without a trailing newline are retained; complete corrupt records still report an error. + | Command | Purpose | | --- | --- | | `init` | Discover crawl apps and write a controller config | diff --git a/scheduler/run.go b/scheduler/run.go index 1b06f0f..cdf5bac 100644 --- a/scheduler/run.go +++ b/scheduler/run.go @@ -2,6 +2,7 @@ package scheduler import ( "bufio" + "bytes" "context" "crypto/rand" "encoding/hex" @@ -279,13 +280,89 @@ func appendHistory(path string, record RunRecord) error { if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil { return err } - file, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600) + // The scheduler lock serializes writes. Avoid O_APPEND: Windows append + // handles lack the write-data access required for recovery truncation. + file, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0o600) if err != nil { return err } - defer file.Close() - enc := json.NewEncoder(file) - return enc.Encode(record) + return appendHistoryFile(file, record) +} + +type historyFile interface { + io.Writer + io.ReaderAt + io.Seeker + Stat() (os.FileInfo, error) + Truncate(int64) error + Close() error +} + +func appendHistoryFile(file historyFile, record RunRecord) (err error) { + defer func() { err = errors.Join(err, file.Close()) }() + info, err := file.Stat() + if err != nil { + return err + } + start := info.Size() + tailStart, tail, err := readHistoryTail(file, start) + if err != nil { + return err + } + if len(tail) > 0 { + var previous RunRecord + if err := json.Unmarshal(tail, &previous); err != nil { + if !incompleteHistoryRecord(tail) { + return err + } + if err := file.Truncate(tailStart); err != nil { + return err + } + start = tailStart + tail = nil + } + } + var buf bytes.Buffer + if len(tail) > 0 { + // A valid EOF record may lack only its separator; retain every byte. + buf.WriteByte('\n') + } + if err := json.NewEncoder(&buf).Encode(record); err != nil { + return err + } + if _, err := file.Seek(start, io.SeekStart); err != nil { + return err + } + n, err := file.Write(buf.Bytes()) + if err == nil && n != buf.Len() { + err = io.ErrShortWrite + } + if err != nil { + return errors.Join(err, file.Truncate(start)) + } + return nil +} + +func readHistoryTail(file io.ReaderAt, end int64) (int64, []byte, error) { + var tail []byte + for end > 0 { + start := max(int64(0), end-4096) + block := make([]byte, end-start) + if _, err := file.ReadAt(block, start); err != nil { + return 0, nil, err + } + if n := bytes.LastIndexByte(block, '\n'); n >= 0 { + return start + int64(n) + 1, append(block[n+1:], tail...), nil + } + tail = append(block, tail...) + end = start + } + return 0, tail, nil +} + +func incompleteHistoryRecord(data []byte) bool { + var value json.RawMessage + return errors.Is(json.NewDecoder(bytes.NewReader(data)).Decode(&value), io.ErrUnexpectedEOF) } func ReadHistory(path string) ([]RunRecord, error) { @@ -299,9 +376,20 @@ func ReadHistory(path string) ([]RunRecord, error) { defer file.Close() var records []RunRecord scanner := bufio.NewScanner(file) + terminated := false + scanner.Split(func(data []byte, atEOF bool) (int, []byte, error) { + advance, token, err := bufio.ScanLines(data, atEOF) + if advance > 0 { + terminated = data[advance-1] == '\n' + } + return advance, token, err + }) for scanner.Scan() { var record RunRecord if err := json.Unmarshal(scanner.Bytes(), &record); err != nil { + if !terminated && incompleteHistoryRecord(scanner.Bytes()) { + break + } return nil, err } records = append(records, record) diff --git a/scheduler/scheduler_test.go b/scheduler/scheduler_test.go index 6cacc98..3198335 100644 --- a/scheduler/scheduler_test.go +++ b/scheduler/scheduler_test.go @@ -1,7 +1,11 @@ package scheduler import ( + "bytes" "context" + "encoding/json" + "errors" + "io" "os" "path/filepath" "runtime" @@ -255,3 +259,265 @@ func TestDefaultPathsCustomConfigKeepsStateNearby(t *testing.T) { t.Fatalf("history = %s, want state next to config", paths.History) } } + +func TestReadHistoryIgnoresTruncatedLastLine(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "runs.jsonl") + complete := `{"id":"1","job":"ok","command":["true"],"started_at":"2026-01-01T00:00:00Z","finished_at":"2026-01-01T00:00:01Z","duration_ms":1,"exit_code":0,"status":"success","log_path":"ok.log"}` + "\n" + if err := os.WriteFile(path, []byte(complete+`{"id":"2","job":"ok"`), 0o600); err != nil { + t.Fatal(err) + } + history, err := ReadHistory(path) + if err != nil { + t.Fatalf("history: %v", err) + } + if len(history) != 1 || history[0].ID != "1" { + t.Fatalf("history = %#v", history) + } +} + +func TestRunDoesNotRefuseOnTruncatedHistoryLine(t *testing.T) { + if runtime.GOOS == "windows" { + t.Skip("shell command path differs on windows") + } + dir := t.TempDir() + paths := Paths{LogDir: filepath.Join(dir, "logs"), StateDir: filepath.Join(dir, "state"), LockPath: filepath.Join(dir, "state", "lock"), History: filepath.Join(dir, "state", "runs.jsonl")} + if err := os.MkdirAll(paths.StateDir, 0o755); err != nil { + t.Fatal(err) + } + complete := `{"id":"1","job":"ok","command":["true"],"started_at":"2026-01-01T00:00:00Z","finished_at":"2026-01-01T00:00:01Z","duration_ms":1,"exit_code":0,"status":"success","log_path":"ok.log"}` + "\n" + if err := os.WriteFile(paths.History, []byte(complete+`{"id":"2","job":"ok"`), 0o600); err != nil { + t.Fatal(err) + } + cfg := DefaultConfig() + cfg.Jobs["ok"] = Job{Enabled: true, Command: []string{"sh", "-c", "echo ok"}} + records, err := Run(context.Background(), RunOptions{Config: cfg, Paths: paths, Names: []string{"ok"}}) + if err != nil { + t.Fatalf("run: %v", err) + } + if len(records) != 1 || records[0].Status != "success" { + t.Fatalf("records = %#v", records) + } + history, err := ReadHistory(paths.History) + if err != nil { + t.Fatalf("history: %v", err) + } + if len(history) != 2 || history[0].ID != "1" || history[1].ID == "" || history[1].ID == "1" { + t.Fatalf("history = %#v", history) + } + data, err := os.ReadFile(paths.History) + if err != nil { + t.Fatal(err) + } + if len(data) == 0 || data[len(data)-1] != '\n' { + t.Fatalf("history file = %q, want complete JSONL", data) + } +} + +func TestAppendHistoryWritesCompleteJSONLLine(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "runs.jsonl") + record := RunRecord{ + ID: "rec1", + Job: "ok", + Command: []string{"echo", "ok"}, + Status: "success", + StartedAt: "2026-08-29T00:00:00Z", + FinishedAt: "2026-08-29T00:00:01Z", + DurationMs: 1000, + LogPath: "/tmp/ok.log", + } + if err := appendHistory(path, record); err != nil { + t.Fatalf("appendHistory: %v", err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read: %v", err) + } + if len(data) == 0 || data[len(data)-1] != '\n' { + t.Fatalf("history = %q, want one newline-terminated JSONL line", data) + } + if bytes.Count(data, []byte{'\n'}) != 1 { + t.Fatalf("history = %q, want exactly one line", data) + } + history, err := ReadHistory(path) + if err != nil { + t.Fatalf("ReadHistory: %v", err) + } + if len(history) != 1 || history[0].ID != record.ID || history[0].Job != record.Job { + t.Fatalf("history = %#v", history) + } +} + +func TestHistoryRetainsValidRecordAtEOF(t *testing.T) { + for _, prefix := range []string{"", "{\"id\":\"first\"}\n"} { + t.Run(prefix, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "runs.jsonl") + original := prefix + `{"id":"last","job":"ok","error":"` + strings.Repeat("x", 5000) + `"}` + "\r" + if err := os.WriteFile(path, []byte(original), 0o600); err != nil { + t.Fatal(err) + } + history, err := ReadHistory(path) + if err != nil || len(history) == 0 || history[len(history)-1].ID != "last" { + t.Fatalf("valid EOF record lost: history=%v err=%v", history, err) + } + if err := appendHistory(path, RunRecord{ID: "next"}); err != nil { + t.Fatal(err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + if !bytes.HasPrefix(data, []byte(original+"\n")) { + t.Fatalf("append changed existing bytes: %q", data) + } + after, err := ReadHistory(path) + if err != nil || len(after) != len(history)+1 || after[len(after)-1].ID != "next" { + t.Fatalf("append history=%v err=%v", after, err) + } + }) + } +} + +func TestHistoryRejectsCorruptFinalRecord(t *testing.T) { + for _, tail := range []string{`{"id":!}`, `{"duration_ms":"wrong type"}`, `{"id":"one"}{"id":`, "{bad json}\n"} { + t.Run(tail, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "runs.jsonl") + original := []byte("{\"id\":\"keep\"}\n" + tail) + if err := os.WriteFile(path, original, 0o600); err != nil { + t.Fatal(err) + } + if _, err := ReadHistory(path); err == nil { + t.Fatal("corruption must remain visible") + } + if !bytes.HasSuffix(original, []byte{'\n'}) { + if err := appendHistory(path, RunRecord{ID: "next"}); err == nil { + t.Fatal("append must reject corrupt tail") + } + } + data, err := os.ReadFile(path) + if err != nil || !bytes.Equal(data, original) { + t.Fatalf("corrupt history changed: %q err=%v", data, err) + } + }) + } +} + +func TestHistoryRecoversEveryPartialRecordPrefix(t *testing.T) { + // Every interrupted byte boundary of a normal encoded record is recoverable. + line, err := json.Marshal(RunRecord{ID: "partial", Job: "café", Status: "success"}) + if err != nil { + t.Fatal(err) + } + path := filepath.Join(t.TempDir(), "runs.jsonl") + prefix := []byte("{\"id\":\"keep\"}\n") + for cut := 1; cut < len(line); cut++ { + data := append(bytes.Clone(prefix), line[:cut]...) + if err := os.WriteFile(path, data, 0o600); err != nil { + t.Fatal(err) + } + history, err := ReadHistory(path) + if err != nil || len(history) != 1 || history[0].ID != "keep" { + t.Fatalf("cut %d: history=%v err=%v", cut, history, err) + } + if err := appendHistory(path, RunRecord{ID: "next"}); err != nil { + t.Fatalf("cut %d: append: %v", cut, err) + } + history, err = ReadHistory(path) + if err != nil || len(history) != 2 || history[0].ID != "keep" || history[1].ID != "next" { + t.Fatalf("cut %d: recovered history=%v err=%v", cut, history, err) + } + } +} + +type failingHistoryFile struct { + *os.File + failWrite bool + writeErr error + truncateErr error + closeErr error +} + +func (f failingHistoryFile) Write(p []byte) (int, error) { + if !f.failWrite { + return f.File.Write(p) + } + n, err := f.File.Write(p[:len(p)/2]) + return n, errors.Join(err, f.writeErr) +} + +func (f failingHistoryFile) Truncate(size int64) error { + if f.truncateErr != nil { + return f.truncateErr + } + return f.File.Truncate(size) +} + +func (f failingHistoryFile) Close() error { + return errors.Join(f.File.Close(), f.closeErr) +} + +func TestHistoryWriteFailurePreservesPriorRecords(t *testing.T) { + writeErr := errors.New("injected disk write failure") + truncateErr := errors.New("injected truncate failure") + for _, tc := range []struct { + name string + writeErr error + truncateErr error + }{ + {name: "short write"}, + {name: "write error", writeErr: writeErr}, + {name: "failed rollback", writeErr: writeErr, truncateErr: truncateErr}, + } { + t.Run(tc.name, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "runs.jsonl") + original := []byte(`{"id":"keep"}`) + if err := os.WriteFile(path, original, 0o600); err != nil { + t.Fatal(err) + } + file, err := os.OpenFile(path, os.O_RDWR, 0o600) + if err != nil { + t.Fatal(err) + } + err = appendHistoryFile(failingHistoryFile{File: file, failWrite: true, writeErr: tc.writeErr, truncateErr: tc.truncateErr}, RunRecord{ID: "failed"}) + wantErr := tc.writeErr + if wantErr == nil { + wantErr = io.ErrShortWrite + } + if !errors.Is(err, wantErr) || (tc.truncateErr != nil && !errors.Is(err, truncateErr)) { + t.Fatalf("append error=%v, want write and cleanup failures", err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + if tc.truncateErr == nil && !bytes.Equal(data, original) { + t.Fatalf("rollback changed prior bytes: %q", data) + } + history, err := ReadHistory(path) + if err != nil || len(history) != 1 || history[0].ID != "keep" { + t.Fatalf("partial write hid prior record: %v, %v", history, err) + } + if err := appendHistory(path, RunRecord{ID: "next"}); err != nil { + t.Fatal(err) + } + history, err = ReadHistory(path) + if err != nil || len(history) != 2 || history[0].ID != "keep" || history[1].ID != "next" { + t.Fatalf("recovery history=%v err=%v", history, err) + } + }) + } +} + +func TestHistoryReturnsCloseFailure(t *testing.T) { + path := filepath.Join(t.TempDir(), "runs.jsonl") + file, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + t.Fatal(err) + } + closeErr := errors.New("injected close failure") + err = appendHistoryFile(failingHistoryFile{File: file, closeErr: closeErr}, RunRecord{ID: "complete"}) + if !errors.Is(err, closeErr) { + t.Fatalf("close error lost: %v", err) + } +}