diff --git a/pkg/filestore/filestore.go b/pkg/filestore/filestore.go index 9617158..b07134d 100644 --- a/pkg/filestore/filestore.go +++ b/pkg/filestore/filestore.go @@ -3,7 +3,9 @@ package filestore import ( "context" "encoding/json" + "errors" "fmt" + "io" "os" "github.com/ethereum/go-ethereum/core/types" @@ -15,9 +17,20 @@ func SaveLogsAsync(ctx context.Context, logChan <-chan types.Log, filePath strin if err != nil { return fmt.Errorf("error creating file: %w", err) } - defer file.Close() + return saveLogs(ctx, logChan, file) +} + +func saveLogs(ctx context.Context, logChan <-chan types.Log, w io.WriteCloser) (err error) { + // OS-buffered writes can surface a failure only when the file is flushed + // at close time, so a swallowed Close error would report an incomplete + // export as success. + defer func() { + if cerr := w.Close(); cerr != nil { + err = errors.Join(err, fmt.Errorf("error closing file: %w", cerr)) + } + }() - encoder := json.NewEncoder(file) + encoder := json.NewEncoder(w) for { select { diff --git a/pkg/filestore/filestore_internal_test.go b/pkg/filestore/filestore_internal_test.go new file mode 100644 index 0000000..3cd0653 --- /dev/null +++ b/pkg/filestore/filestore_internal_test.go @@ -0,0 +1,44 @@ +package filestore + +import ( + "context" + "errors" + "testing" + + "github.com/ethereum/go-ethereum/core/types" +) + +// closeFailWriter accepts all writes but fails on Close, mimicking a file +// whose OS-buffered data cannot be flushed (disk full, quota, network fs). +type closeFailWriter struct{ closeErr error } + +func (w *closeFailWriter) Write(p []byte) (int, error) { return len(p), nil } +func (w *closeFailWriter) Close() error { return w.closeErr } + +func TestSaveLogsReportsCloseError(t *testing.T) { + closeErr := errors.New("flush to disk failed") + + logChan := make(chan types.Log, 1) + logChan <- types.Log{BlockNumber: 1} + close(logChan) + + err := saveLogs(context.Background(), logChan, &closeFailWriter{closeErr: closeErr}) + if !errors.Is(err, closeErr) { + t.Fatalf("got %v, want close error %v", err, closeErr) + } +} + +func TestSaveLogsKeepsContextErrorOnCloseFailure(t *testing.T) { + closeErr := errors.New("flush to disk failed") + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + err := saveLogs(ctx, make(chan types.Log), &closeFailWriter{closeErr: closeErr}) + if !errors.Is(err, context.Canceled) { + t.Fatalf("got %v, want context.Canceled", err) + } + if !errors.Is(err, closeErr) { + t.Fatalf("got %v, want close error %v", err, closeErr) + } +} diff --git a/pkg/filestore/filestore_test.go b/pkg/filestore/filestore_test.go new file mode 100644 index 0000000..8add67e --- /dev/null +++ b/pkg/filestore/filestore_test.go @@ -0,0 +1,61 @@ +package filestore_test + +import ( + "bufio" + "context" + "encoding/json" + "os" + "path/filepath" + "testing" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethersphere/batch-export/pkg/filestore" +) + +func TestSaveLogsAsyncWritesNDJSON(t *testing.T) { + path := filepath.Join(t.TempDir(), "export.ndjson") + + logChan := make(chan types.Log, 2) + for i := uint64(1); i <= 2; i++ { + logChan <- types.Log{ + Address: common.HexToAddress("0x000000000000000000000000000000000000bEEF"), + Topics: []common.Hash{common.HexToHash("0x11")}, + Data: []byte{0xde, 0xad}, + BlockNumber: i, + } + } + close(logChan) + + if err := filestore.SaveLogsAsync(context.Background(), logChan, path); err != nil { + t.Fatalf("SaveLogsAsync: %v", err) + } + + file, err := os.Open(path) + if err != nil { + t.Fatalf("open output: %v", err) + } + defer file.Close() + + var got []types.Log + scanner := bufio.NewScanner(file) + for scanner.Scan() { + var l types.Log + if err := json.Unmarshal(scanner.Bytes(), &l); err != nil { + t.Fatalf("decode line %d: %v", len(got)+1, err) + } + got = append(got, l) + } + if err := scanner.Err(); err != nil { + t.Fatalf("scan output: %v", err) + } + + if len(got) != 2 { + t.Fatalf("got %d logs, want 2", len(got)) + } + for i, l := range got { + if l.BlockNumber != uint64(i+1) { + t.Errorf("log %d: blockNumber got %d want %d", i, l.BlockNumber, i+1) + } + } +} diff --git a/pkg/gzipstore/gzipstore.go b/pkg/gzipstore/gzipstore.go index 6763ad0..2b079ed 100644 --- a/pkg/gzipstore/gzipstore.go +++ b/pkg/gzipstore/gzipstore.go @@ -2,6 +2,7 @@ package gzipstore import ( "compress/gzip" + "errors" "fmt" "io" "os" @@ -23,14 +24,22 @@ func CompressFile(inputFilePath string, outputFilePath string) error { } defer outputFile.Close() - // create a new gzip writer that writes to the output file - gzipWriter := gzip.NewWriter(outputFile) - defer gzipWriter.Close() + return compress(outputFile, inputFile) +} - // copy the contents from the input file to the gzip writer - _, err = io.Copy(gzipWriter, inputFile) - if err != nil { - return fmt.Errorf("failed to write compressed data to '%s': %w", outputFilePath, err) +func compress(dst io.Writer, src io.Reader) (err error) { + gzipWriter := gzip.NewWriter(dst) + // Close flushes the remaining compressed bytes and the gzip footer; a + // swallowed error here would report a truncated archive as success. + defer func() { + if cerr := gzipWriter.Close(); cerr != nil { + err = errors.Join(err, fmt.Errorf("failed to finalize gzip stream: %w", cerr)) + } + }() + + // copy the contents from the source to the gzip writer + if _, err := io.Copy(gzipWriter, src); err != nil { + return fmt.Errorf("failed to write compressed data: %w", err) } return nil diff --git a/pkg/gzipstore/gzipstore_internal_test.go b/pkg/gzipstore/gzipstore_internal_test.go new file mode 100644 index 0000000..5da04ff --- /dev/null +++ b/pkg/gzipstore/gzipstore_internal_test.go @@ -0,0 +1,31 @@ +package gzipstore + +import ( + "errors" + "strings" + "testing" +) + +// headerOnlyWriter accepts the first write (the gzip header) and fails every +// write after it, so the compressed payload flush at Close is what fails. +type headerOnlyWriter struct { + writes int + err error +} + +func (w *headerOnlyWriter) Write(p []byte) (int, error) { + w.writes++ + if w.writes > 1 { + return 0, w.err + } + return len(p), nil +} + +func TestCompressReportsFlushError(t *testing.T) { + writeErr := errors.New("disk full") + + err := compress(&headerOnlyWriter{err: writeErr}, strings.NewReader("payload")) + if !errors.Is(err, writeErr) { + t.Fatalf("got %v, want flush error %v", err, writeErr) + } +} diff --git a/pkg/gzipstore/gzipstore_test.go b/pkg/gzipstore/gzipstore_test.go new file mode 100644 index 0000000..61a173a --- /dev/null +++ b/pkg/gzipstore/gzipstore_test.go @@ -0,0 +1,47 @@ +package gzipstore_test + +import ( + "bytes" + "compress/gzip" + "io" + "os" + "path/filepath" + "testing" + + "github.com/ethersphere/batch-export/pkg/gzipstore" +) + +func TestCompressFileRoundTrip(t *testing.T) { + dir := t.TempDir() + inputPath := filepath.Join(dir, "in.ndjson") + outputPath := filepath.Join(dir, "out.gzip") + + content := []byte("{\"blockNumber\":\"0x1\"}\n{\"blockNumber\":\"0x2\"}\n") + if err := os.WriteFile(inputPath, content, 0o644); err != nil { + t.Fatalf("write input: %v", err) + } + + if err := gzipstore.CompressFile(inputPath, outputPath); err != nil { + t.Fatalf("CompressFile: %v", err) + } + + file, err := os.Open(outputPath) + if err != nil { + t.Fatalf("open output: %v", err) + } + defer file.Close() + + gzipReader, err := gzip.NewReader(file) + if err != nil { + t.Fatalf("gzip reader: %v", err) + } + defer gzipReader.Close() + + got, err := io.ReadAll(gzipReader) + if err != nil { + t.Fatalf("decompress: %v", err) + } + if !bytes.Equal(got, content) { + t.Errorf("round trip mismatch: got %q want %q", got, content) + } +}