From deff51adb8ca6a129c70571f6383364eeb510107 Mon Sep 17 00:00:00 2001 From: dvd233 <111864431+dvd233@users.noreply.github.com> Date: Thu, 3 Sep 2026 19:10:24 +0800 Subject: [PATCH 1/2] fix(output): write reports atomically --- cmd/opencodereview/output_file_test.go | 254 ++++++++++++++++++++++++- cmd/opencodereview/shared.go | 171 ++++++++++++++--- 2 files changed, 392 insertions(+), 33 deletions(-) diff --git a/cmd/opencodereview/output_file_test.go b/cmd/opencodereview/output_file_test.go index 134137e6c..ccee2b626 100644 --- a/cmd/opencodereview/output_file_test.go +++ b/cmd/opencodereview/output_file_test.go @@ -206,6 +206,15 @@ func TestResolveOutputWriter_MissingParent(t *testing.T) { // --- lazyFileWriter --- +func outputTempFiles(t *testing.T, dir string) []string { + t.Helper() + matches, err := filepath.Glob(filepath.Join(dir, ".ocr-out-*")) + if err != nil { + t.Fatalf("glob output temp files: %v", err) + } + return matches +} + // TestLazyFileWriter_NoWriteLeavesExistingFileUntouched pins the core // data-safety contract: a writer that is resolved but never written to (a // failed run, a preview error) must not create or truncate the target file. @@ -247,16 +256,32 @@ func TestLazyFileWriter_NoWriteDoesNotCreateFile(t *testing.T) { } } -func TestLazyFileWriter_WriteCreatesFileAndPrintsHint(t *testing.T) { - path := filepath.Join(t.TempDir(), "out.json") - stderr := captureStderr(t, func() { - w, closeFn, err := resolveOutputWriter(path, "json") +func TestLazyFileWriter_WriteCommitsFileAndPrintsHintOnClose(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "out.json") + var closeFn func() error + stderrBeforeClose := captureStderr(t, func() { + w, close, err := resolveOutputWriter(path, "json") if err != nil { t.Fatalf("resolve: %v", err) } + closeFn = close if _, err := w.Write([]byte(`{"status":"success"}`)); err != nil { t.Fatalf("write: %v", err) } + }) + t.Cleanup(func() { _ = closeFn() }) + if stderrBeforeClose != "" { + t.Fatalf("hint printed before output was committed: %q", stderrBeforeClose) + } + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatalf("target became visible before close, stat err = %v", err) + } + if temps := outputTempFiles(t, dir); len(temps) != 1 { + t.Fatalf("output temp files before close = %v, want one", temps) + } + + stderrAfterClose := captureStderr(t, func() { if err := closeFn(); err != nil { t.Fatalf("close: %v", err) } @@ -268,8 +293,225 @@ func TestLazyFileWriter_WriteCreatesFileAndPrintsHint(t *testing.T) { if string(data) != `{"status":"success"}` { t.Fatalf("file content = %q, want the written bytes", data) } - if !strings.Contains(stderr, "[ocr] Results written to "+path) { - t.Fatalf("expected 'Results written' hint on stderr, got %q", stderr) + if !strings.Contains(stderrAfterClose, "[ocr] Results written to "+path) { + t.Fatalf("expected 'Results written' hint after close, got %q", stderrAfterClose) + } + if temps := outputTempFiles(t, dir); len(temps) != 0 { + t.Fatalf("output temp files after close = %v, want none", temps) + } +} + +func TestLazyFileWriter_ReplacesExistingFileOnlyOnClose(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "out.json") + const oldContent = "previous complete report\n" + const newContent = `{"status":"success"}` + if err := os.WriteFile(path, []byte(oldContent), 0o640); err != nil { + t.Fatalf("write existing: %v", err) + } + before, err := os.Stat(path) + if err != nil { + t.Fatalf("stat existing: %v", err) + } + + w, closeFn, err := resolveOutputWriter(path, "json") + if err != nil { + t.Fatalf("resolve: %v", err) + } + t.Cleanup(func() { _ = closeFn() }) + if _, err := w.Write([]byte(newContent)); err != nil { + t.Fatalf("write: %v", err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read before close: %v", err) + } + if string(data) != oldContent { + t.Fatalf("target changed before close: got %q, want %q", data, oldContent) + } + if temps := outputTempFiles(t, dir); len(temps) != 1 { + t.Fatalf("output temp files before close = %v, want one", temps) + } + + if err := closeFn(); err != nil { + t.Fatalf("close: %v", err) + } + data, err = os.ReadFile(path) + if err != nil { + t.Fatalf("read after close: %v", err) + } + if string(data) != newContent { + t.Fatalf("target after close = %q, want %q", data, newContent) + } + after, err := os.Stat(path) + if err != nil { + t.Fatalf("stat after close: %v", err) + } + if after.Mode().Perm() != before.Mode().Perm() { + t.Fatalf("target mode after close = %v, want %v", after.Mode().Perm(), before.Mode().Perm()) + } + if temps := outputTempFiles(t, dir); len(temps) != 0 { + t.Fatalf("output temp files after close = %v, want none", temps) + } +} + +func TestLazyFileWriter_NewFileUsesCreatePermissions(t *testing.T) { + dir := t.TempDir() + probePath := filepath.Join(dir, "probe") + probe, err := os.Create(probePath) + if err != nil { + t.Fatalf("create mode probe: %v", err) + } + if err := probe.Close(); err != nil { + t.Fatalf("close mode probe: %v", err) + } + probeInfo, err := os.Stat(probePath) + if err != nil { + t.Fatalf("stat mode probe: %v", err) + } + if err := os.Remove(probePath); err != nil { + t.Fatalf("remove mode probe: %v", err) + } + + path := filepath.Join(dir, "out.json") + w, closeFn, err := resolveOutputWriter(path, "json") + if err != nil { + t.Fatalf("resolve: %v", err) + } + t.Cleanup(func() { _ = closeFn() }) + if _, err := w.Write([]byte(`{}`)); err != nil { + t.Fatalf("write: %v", err) + } + if err := closeFn(); err != nil { + t.Fatalf("close: %v", err) + } + info, err := os.Stat(path) + if err != nil { + t.Fatalf("stat output: %v", err) + } + if info.Mode().Perm() != probeInfo.Mode().Perm() { + t.Fatalf("new output mode = %v, want os.Create mode %v", info.Mode().Perm(), probeInfo.Mode().Perm()) + } +} + +func TestLazyFileWriter_RenameFailureLeavesTargetAndCleansTemp(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "out.json") + w, closeFn, err := resolveOutputWriter(path, "json") + if err != nil { + t.Fatalf("resolve: %v", err) + } + t.Cleanup(func() { _ = closeFn() }) + if _, err := w.Write([]byte(`{"status":"success"}`)); err != nil { + t.Fatalf("write: %v", err) + } + if err := os.Mkdir(path, 0o755); err != nil { + t.Fatalf("create conflicting target: %v", err) + } + stderr := captureStderr(t, func() { + if err := closeFn(); err == nil { + t.Fatal("expected atomic rename to fail for a directory target") + } + }) + if strings.Contains(stderr, "Results written") { + t.Fatalf("success hint printed after failed rename: %q", stderr) + } + info, err := os.Stat(path) + if err != nil { + t.Fatalf("stat conflicting target: %v", err) + } + if !info.IsDir() { + t.Fatal("failed rename changed the conflicting target") + } + if temps := outputTempFiles(t, dir); len(temps) != 0 { + t.Fatalf("output temp files after failed rename = %v, want none", temps) + } +} + +func TestLazyFileWriter_WriteFailureLeavesExistingFileUntouched(t *testing.T) { + dir := t.TempDir() + path := filepath.Join(dir, "out.json") + const oldContent = "previous complete report\n" + if err := os.WriteFile(path, []byte(oldContent), 0o644); err != nil { + t.Fatalf("write existing: %v", err) + } + w, closeFn, err := resolveOutputWriter(path, "json") + if err != nil { + t.Fatalf("resolve: %v", err) + } + t.Cleanup(func() { _ = closeFn() }) + if _, err := w.Write([]byte("partial replacement")); err != nil { + t.Fatalf("first write: %v", err) + } + lazy := w.(*lazyFileWriter) + if err := lazy.file.Close(); err != nil { + t.Fatalf("inject close: %v", err) + } + if _, err := w.Write([]byte("must fail")); err == nil { + t.Fatal("expected the injected write failure") + } + if err := closeFn(); err != nil { + t.Fatalf("cleanup after reported write failure: %v", err) + } + data, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read existing: %v", err) + } + if string(data) != oldContent { + t.Fatalf("existing file changed after write failure: got %q, want %q", data, oldContent) + } + if temps := outputTempFiles(t, dir); len(temps) != 0 { + t.Fatalf("output temp files after write failure = %v, want none", temps) + } +} + +func TestLazyFileWriter_PreservesSymlinkTarget(t *testing.T) { + dir := t.TempDir() + target := filepath.Join(dir, "target.json") + link := filepath.Join(dir, "out.json") + const oldContent = "previous complete report\n" + const newContent = `{"status":"success"}` + if err := os.WriteFile(target, []byte(oldContent), 0o644); err != nil { + t.Fatalf("write target: %v", err) + } + if err := os.Symlink(filepath.Base(target), link); err != nil { + t.Skipf("symlinks are unavailable: %v", err) + } + + w, closeFn, err := resolveOutputWriter(link, "json") + if err != nil { + t.Fatalf("resolve: %v", err) + } + t.Cleanup(func() { _ = closeFn() }) + if _, err := w.Write([]byte(newContent)); err != nil { + t.Fatalf("write: %v", err) + } + data, err := os.ReadFile(target) + if err != nil { + t.Fatalf("read target before close: %v", err) + } + if string(data) != oldContent { + t.Fatalf("symlink target changed before close: got %q, want %q", data, oldContent) + } + if err := closeFn(); err != nil { + t.Fatalf("close: %v", err) + } + info, err := os.Lstat(link) + if err != nil { + t.Fatalf("lstat link: %v", err) + } + if info.Mode()&os.ModeSymlink == 0 { + t.Fatal("output symlink was replaced instead of preserving its target") + } + data, err = os.ReadFile(target) + if err != nil { + t.Fatalf("read target after close: %v", err) + } + if string(data) != newContent { + t.Fatalf("symlink target after close = %q, want %q", data, newContent) + } + if temps := outputTempFiles(t, dir); len(temps) != 0 { + t.Fatalf("output temp files after close = %v, want none", temps) } } diff --git a/cmd/opencodereview/shared.go b/cmd/opencodereview/shared.go index 790452dab..55995b4e0 100644 --- a/cmd/opencodereview/shared.go +++ b/cmd/opencodereview/shared.go @@ -5,6 +5,8 @@ package main import ( "context" + "crypto/rand" + "errors" "fmt" "io" "net/url" @@ -519,31 +521,118 @@ func (w *stripAnsiWriter) Write(p []byte) (int, error) { return len(p), nil } -// lazyFileWriter defers os.Create until the first Write so a run that never -// produces output (LLM failure, preview error, interruption) leaves an -// existing target file untouched instead of truncating it to zero bytes. The -// "Results written" hint is printed to stderr only after the first successful -// Write, so agents never see a path hint for a file that stayed empty or was -// never persisted. +const outputTempPrefix = ".ocr-out-" + +// createOutputTemp creates a collision-resistant file in the target directory. +// Unlike os.CreateTemp, the 0666 mode preserves os.Create's umask-sensitive +// permissions for new output files. +func createOutputTemp(dir string) (*os.File, error) { + for range 100 { + path := filepath.Join(dir, outputTempPrefix+rand.Text()) + f, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE|os.O_EXCL, 0o666) + if err == nil { + return f, nil + } + if !os.IsExist(err) { + return nil, err + } + } + return nil, fmt.Errorf("too many temporary output filename collisions in %s", dir) +} + +// resolveOutputCommitPath preserves os.Create's behavior for an existing +// symlink: write through the link to its target instead of replacing the link +// itself during the final rename. +func resolveOutputCommitPath(path string) (string, error) { + info, err := os.Lstat(path) + if err != nil { + if os.IsNotExist(err) { + return path, nil + } + return "", fmt.Errorf("inspect output path %s: %w", path, err) + } + if info.Mode()&os.ModeSymlink == 0 { + return path, nil + } + resolved, err := filepath.EvalSymlinks(path) + if err != nil { + return "", fmt.Errorf("resolve output symlink %s: %w", path, err) + } + return resolved, nil +} + +func cleanupOutputTemp(file *os.File, path string) error { + var cleanupErr error + if file != nil { + if err := file.Close(); err != nil && !errors.Is(err, os.ErrClosed) { + cleanupErr = errors.Join(cleanupErr, fmt.Errorf("close temporary output file: %w", err)) + } + } + if path != "" { + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + cleanupErr = errors.Join(cleanupErr, fmt.Errorf("remove temporary output file %s: %w", path, err)) + } + } + return cleanupErr +} + +// lazyFileWriter defers creating a same-directory temporary file until the +// first Write. Close syncs and atomically renames a complete result into place, +// so a run that never writes or encounters a write failure leaves the previous +// target untouched. Replacing the directory entry intentionally gives readers +// a consistent old-or-new snapshot; open descriptors and hard links keep +// referring to the old inode. type lazyFileWriter struct { - path string - strip bool // strip ANSI when the target format is text - once sync.Once - file *os.File - stripper *stripAnsiWriter - err error // os.Create error - writeErr error // first error from a Write - hinted bool // hint already printed after a successful Write + path string + strip bool // strip ANSI when the target format is text + once sync.Once + file *os.File // temporary file committed by Close + tempPath string + commitPath string + stripper *stripAnsiWriter + err error // temporary-file creation or setup error + writeErr error // first error from a Write + closed bool + closeErr error } func (w *lazyFileWriter) Write(p []byte) (int, error) { + if w.closed { + return 0, os.ErrClosed + } w.once.Do(func() { - f, err := os.Create(w.path) + commitPath, err := resolveOutputCommitPath(w.path) + if err != nil { + w.err = err + return + } + + var existingMode os.FileMode + hasExistingMode := false + if info, statErr := os.Stat(commitPath); statErr == nil { + existingMode = info.Mode() & (os.ModePerm | os.ModeSetuid | os.ModeSetgid | os.ModeSticky) + hasExistingMode = true + } else if !os.IsNotExist(statErr) { + w.err = fmt.Errorf("inspect output file %s: %w", w.path, statErr) + return + } + + f, err := createOutputTemp(filepath.Dir(commitPath)) if err != nil { - w.err = fmt.Errorf("create output file %s: %w", w.path, err) + w.err = fmt.Errorf("create temporary output file for %s: %w", w.path, err) return } + if hasExistingMode { + if err := f.Chmod(existingMode); err != nil { + cleanupErr := cleanupOutputTemp(f, f.Name()) + w.err = errors.Join(fmt.Errorf("preserve output permissions for %s: %w", w.path, err), cleanupErr) + return + } + } + w.file = f + w.tempPath = f.Name() + w.commitPath = commitPath if w.strip { w.stripper = &stripAnsiWriter{dst: f} } @@ -561,10 +650,6 @@ func (w *lazyFileWriter) Write(p []byte) (int, error) { if err != nil && w.writeErr == nil { w.writeErr = err } - if err == nil && !w.hinted { - w.hinted = true - fmt.Fprintf(os.Stderr, "[ocr] Results written to %s\n", w.path) - } return n, err } @@ -589,22 +674,54 @@ func writeOutError(out io.Writer) error { return nil } -// Close closes the underlying file. It is a no-op when the file was never -// created (no output produced), so failure paths cannot leave a fresh empty -// file behind. +// Close commits a complete output file atomically. It is a no-op when no Write +// occurred, and it discards the temporary file when a Write already failed. func (w *lazyFileWriter) Close() error { + if w.closed { + return w.closeErr + } + w.closed = true if w.file == nil { return nil } - return w.file.Close() + if w.writeErr != nil { + w.closeErr = cleanupOutputTemp(w.file, w.tempPath) + w.file = nil + w.tempPath = "" + return w.closeErr + } + if err := w.file.Sync(); err != nil { + cleanupErr := cleanupOutputTemp(w.file, w.tempPath) + w.file = nil + w.tempPath = "" + w.closeErr = errors.Join(fmt.Errorf("sync output file %s: %w", w.path, err), cleanupErr) + return w.closeErr + } + if err := w.file.Close(); err != nil { + cleanupErr := cleanupOutputTemp(w.file, w.tempPath) + w.file = nil + w.tempPath = "" + w.closeErr = errors.Join(fmt.Errorf("close output file %s: %w", w.path, err), cleanupErr) + return w.closeErr + } + w.file = nil + if err := os.Rename(w.tempPath, w.commitPath); err != nil { + cleanupErr := cleanupOutputTemp(nil, w.tempPath) + w.tempPath = "" + w.closeErr = errors.Join(fmt.Errorf("replace output file %s: %w", w.path, err), cleanupErr) + return w.closeErr + } + w.tempPath = "" + fmt.Fprintf(os.Stderr, "[ocr] Results written to %s\n", w.path) + return nil } // resolveOutputWriter resolves the --output target into a writer plus a // cleanup function. // - "" or "-" → os.Stdout with a no-op cleanup (colors preserved, no hint) -// - otherwise → a lazyFileWriter over os.Create(path), deferred until the -// first Write; text format wraps the file in stripAnsiWriter so ANSI -// colors never reach the result file. +// - otherwise → a lazyFileWriter over a same-directory temporary file, +// deferred until the first Write and atomically committed by Close; text +// format wraps the file in stripAnsiWriter so ANSI colors never reach it. // // Fail-fast checks (directory target, missing parent) run here without // creating or truncating anything; deeper errors (permissions, disk) surface From 7ae6442c8b0e34d88804836ba605fb83060a8a8e Mon Sep 17 00:00:00 2001 From: dvd233 <111864431+dvd233@users.noreply.github.com> Date: Mon, 7 Sep 2026 19:41:49 +0800 Subject: [PATCH 2/2] fix(output): commit reports before MCP shutdown --- cmd/opencodereview/review_cmd.go | 37 ++++++++++---- .../review_output_order_test.go | 51 +++++++++++++++++++ 2 files changed, 79 insertions(+), 9 deletions(-) create mode 100644 cmd/opencodereview/review_output_order_test.go diff --git a/cmd/opencodereview/review_cmd.go b/cmd/opencodereview/review_cmd.go index 69f9f538b..9d0ba8b52 100644 --- a/cmd/opencodereview/review_cmd.go +++ b/cmd/opencodereview/review_cmd.go @@ -115,9 +115,20 @@ func executeReviewContext(ctx context.Context, opts reviewOptions) (retErr error if err != nil { return err } - defer func() { + closeOutPending := true + finishOutput := func() error { + if !closeOutPending { + return nil + } + closeOutPending = false if cerr := closeOut(); cerr != nil { - retErr = errors.Join(retErr, fmt.Errorf("close output file: %w", cerr)) + return fmt.Errorf("close output file: %w", cerr) + } + return nil + } + defer func() { + if cerr := finishOutput(); cerr != nil { + retErr = errors.Join(retErr, cerr) } }() @@ -195,13 +206,7 @@ func executeReviewContext(ctx context.Context, opts reviewOptions) (retErr error tools := buildToolRegistry(rt.Collector, fileReader) mcpClients := initMCPClients(ctx, rt.AppCfg, tools, cc.RepoDir, Version) - defer func() { - for _, mc := range mcpClients { - if err := mc.Close(); err != nil { - fmt.Fprintf(os.Stderr, "[ocr] WARNING: failed to close MCP server %q: %v\n", mc.Name(), err) - } - } - }() + defer closeReviewMCPClients(mcpClients) mcpToolDefs := mcp.CollectToolDefs(mcpClients, tools) rt.PlanToolDefs = append(rt.PlanToolDefs, mcpToolDefs...) @@ -291,6 +296,10 @@ func executeReviewContext(ctx context.Context, opts reviewOptions) (retErr error emitErr = emitRunResult(runCtx, ag, comments, startTime, opts.outputFormat, opts.audience, q, llmIdentity, out, retryReport) if emitErr != nil { emitErr = fmt.Errorf("emit review result: %w", emitErr) + } else { + // Commit the report before potentially slow MCP shutdown. The deferred + // close remains a fallback for every earlier return and emit failure. + emitErr = finishOutput() } } if resultErr != nil { @@ -577,6 +586,16 @@ func initMCPClients(ctx context.Context, cfg *Config, tools *tool.Registry, repo return clients } +// closeReviewMCPClients is a variable so tests can observe the shutdown +// boundary without starting an intentionally unresponsive subprocess. +var closeReviewMCPClients = func(clients []*mcp.Client) { + for _, mc := range clients { + if err := mc.Close(); err != nil { + fmt.Fprintf(os.Stderr, "[ocr] WARNING: failed to close MCP server %q: %v\n", mc.Name(), err) + } + } +} + func buildToolRegistry(collector *tool.CommentCollector, fr *tool.FileReader) *tool.Registry { reg := tool.NewRegistry() reg.Register(tool.NewFileRead(fr)) diff --git a/cmd/opencodereview/review_output_order_test.go b/cmd/opencodereview/review_output_order_test.go new file mode 100644 index 000000000..2ed3893e3 --- /dev/null +++ b/cmd/opencodereview/review_output_order_test.go @@ -0,0 +1,51 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright 2026 alibaba/open-code-review Contributors + +package main + +import ( + "encoding/json" + "os" + "path/filepath" + "testing" + + "github.com/alibaba/open-code-review/internal/mcp" +) + +func TestReviewOutputCommittedBeforeMCPShutdown(t *testing.T) { + repoDir := retryTestRepo(t) + startFakeLLM(t, newFakeLLM()) + outputPath := filepath.Join(t.TempDir(), "review.json") + + originalClose := closeReviewMCPClients + t.Cleanup(func() { closeReviewMCPClients = originalClose }) + shutdownObserved := false + closeReviewMCPClients = func(clients []*mcp.Client) { + shutdownObserved = true + data, err := os.ReadFile(outputPath) + if err != nil { + t.Errorf("read output at MCP shutdown boundary: %v", err) + } else { + var report jsonOutput + if err := json.Unmarshal(data, &report); err != nil { + t.Errorf("output at MCP shutdown boundary is incomplete JSON: %v", err) + } + } + originalClose(clients) + } + + err := runReview([]string{ + "--repo", repoDir, + "--from", "HEAD~1", + "--to", "HEAD", + "--format", "json", + "--audience", "agent", + "--output", outputPath, + }) + if err != nil { + t.Fatalf("review must succeed: %v", err) + } + if !shutdownObserved { + t.Fatal("MCP shutdown boundary was not observed") + } +}