Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 15 additions & 2 deletions pkg/filestore/filestore.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,9 @@ package filestore
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"

"github.com/ethereum/go-ethereum/core/types"
Expand All @@ -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 {
Expand Down
44 changes: 44 additions & 0 deletions pkg/filestore/filestore_internal_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
61 changes: 61 additions & 0 deletions pkg/filestore/filestore_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
}
23 changes: 16 additions & 7 deletions pkg/gzipstore/gzipstore.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package gzipstore

import (
"compress/gzip"
"errors"
"fmt"
"io"
"os"
Expand All @@ -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
Expand Down
31 changes: 31 additions & 0 deletions pkg/gzipstore/gzipstore_internal_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
47 changes: 47 additions & 0 deletions pkg/gzipstore/gzipstore_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading