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
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -11,4 +11,4 @@ lint:
golangci-lint run --timeout 5m ./...

bench: *.go
go test -run='^$$' -bench=BenchmarkScan -benchmem -benchtime=2ms ./...
go test -run='^$$' -bench=BenchmarkScan -benchmem -benchtime=5s ./...
9 changes: 6 additions & 3 deletions linebuf.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,18 +57,19 @@ func New(opts ...Option) *Scanner {
return s
}

var errIterationStopped = errors.New("iteration stopped")

// Iter returns a sequence of lines from the provided reader. Any error encountered during iteration can be retrieved using IterErr.
func (s *Scanner) Iter(r io.Reader) iter.Seq[[]byte] {
stopErr := errors.New("iteration stopped")
s.iterErr = nil
return func(yield func([]byte) bool) {
err := s.Scan(r, func(data []byte) error {
if !yield(data) {
return stopErr
return errIterationStopped
}
return nil
})
if err != nil && err != stopErr { //nolint:errorlint,staticcheck
if err != nil && err != errIterationStopped { //nolint:errorlint,staticcheck
s.iterErr = err
}
}
Expand Down Expand Up @@ -126,6 +127,8 @@ func callCB(cb CB, line []byte) error {

// Scan reads from the provided reader and invokes the callback for each line. It returns an error if any occurs during scanning.
func (s *Scanner) Scan(r io.Reader, cb CB) error {
// reset the offset before starting the scan
s.offset = 0
for {
err := s.scanInternal(r, cb)
if err != nil && err == io.EOF { //nolint:staticcheck,errorlint
Expand Down
100 changes: 61 additions & 39 deletions linebuf_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,38 +4,16 @@ import (
"bufio"
"fmt"
"io"
"math"
"math/rand/v2"
"os"
"path/filepath"
"strings"
"testing"
"time"

"github.com/stretchr/testify/require"
)

func radixInput(n int, distribution string) []float64 {
r := rand.New(rand.NewPCG(1, 2))
points := make([]float64, n)
for i := range points {
switch distribution {
case "duplicates":
points[i] = float64(i%100) + 0.5
case "random":
points[i] = r.Float64() * 10000
case "wide":
points[i] = math.Float64frombits(r.Uint64() & 0x7fefffffffffffff)
case "sorted":
points[i] = float64(i)
case "equal":
points[i] = 42
case "response_time":
points[i] = float64(r.IntN(500)) / 1000
}
}
return points
}

func TestScanBufferProcessChunk(t *testing.T) {
lines := []string{}
cb := func(data []byte) error {
Expand Down Expand Up @@ -206,6 +184,26 @@ func TestScanBufferScanFileLongLine(t *testing.T) {
require.Equal(t, []string{strings.TrimSuffix(longLine, "\n")}, lines)
}

func TestResetScannerOffset(t *testing.T) {
sb := New(WithStartBufSize(32))
lines := []string{}
cb := func(data []byte) error {
lines = append(lines, string(data))
return nil
}
r := strings.NewReader("a\nb\nc")
err := sb.Scan(r, cb)
require.NoError(t, err)
require.Equal(t, []string{"a", "b", "c"}, lines)

r = strings.NewReader("d\ne\nf")
lines = []string{}
err = sb.Scan(r, cb)
require.NoError(t, err)
require.Equal(t, []string{"d", "e", "f"}, lines)

}

func TestCallCB(t *testing.T) {
tests := []struct {
name string
Expand Down Expand Up @@ -283,48 +281,71 @@ func TestIterErr(t *testing.T) {
require.Equal(t, testErr, s.IterErr())
}

func generateBenchmarkFile(b testing.TB, w io.Writer, numLines int) error {
b.Helper()
r := rand.New(rand.NewPCG(1, 2))
for i := range numLines {
line := fmt.Sprintf(`{"time": "%s", "status": "%d", "reqtime": "%f", "host": "%s", "req": "%s", "method": "%s", "size": "%d", "ua": "%s"}`,
time.Now().Format(time.RFC3339),
200+i%5,
float64(r.IntN(500))/1000,
"10.20.30.40",
"GET /example/path HTTP/1.1",
"GET",
941,
"Mozilla/5.0 (Linux; Android 4.4.2; SO-01F) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/73.0.3683.90 Mobile Safari/537.36",
)
_, err := w.Write([]byte(line + "\n"))
if err != nil {
return err
}
}

return nil
}

func testFileBuilder(b testing.TB, count int) (*os.File, error) {
dir := b.TempDir()
filePath := filepath.Join(dir, "testfile.txt")
file, err := os.Create(filePath)
if err != nil {
b.Fatal(err)
}

for _, f := range radixInput(count, "random") {
_, err := fmt.Fprintf(file, "%.3f\n", f)
if err != nil {
b.Fatal(err)
}
err = generateBenchmarkFile(b, file, count)
if err != nil {
b.Fatal(err)
}
return file, nil
}

var benchmarkStartBufSize = 8192
var benchmarkFileLines = 10000

func BenchmarkScan_bufio(b *testing.B) {
file, err := testFileBuilder(b, 10000)
file, err := testFileBuilder(b, benchmarkFileLines)
if err != nil {
b.Fatal(err)
}
defer file.Close()
b.ResetTimer()
b.ReportAllocs()
buf := make([]byte, 0, benchmarkStartBufSize)
for b.Loop() {
_, _ = file.Seek(0, io.SeekStart)
scanner := bufio.NewScanner(file)
scanner.Buffer(buf, 64*1024)
total := 0
for scanner.Scan() {
b := scanner.Bytes()
total += len(b)
}
if err := scanner.Err(); err != nil {
b.Fatal(err)
}
require.NoError(b, scanner.Err())
}
}

func BenchmarkScan_linebuf(b *testing.B) {

file, err := testFileBuilder(b, 10000)
file, err := testFileBuilder(b, benchmarkFileLines)
if err != nil {
b.Fatal(err)
}
Expand All @@ -336,27 +357,28 @@ func BenchmarkScan_linebuf(b *testing.B) {
total += len(data)
return nil
}
s := New(WithStartBufSize(benchmarkStartBufSize), WithMaxBufSize(64*1024))
for b.Loop() {
_, _ = file.Seek(0, io.SeekStart)
total = 0
err := Scan(file, cb, WithStartBufSize(4096), WithMaxBufSize(64*1024))
err := s.Scan(file, cb)
require.NoError(b, err)
}
}

func BenchmarkScan_iter(b *testing.B) {
file, err := testFileBuilder(b, 10000)
file, err := testFileBuilder(b, benchmarkFileLines)
if err != nil {
b.Fatal(err)
}
defer file.Close()
b.ResetTimer()
b.ReportAllocs()
total := 0
// reusable scanner
s := New(WithStartBufSize(benchmarkStartBufSize), WithMaxBufSize(64*1024))
for b.Loop() {
_, _ = file.Seek(0, io.SeekStart)
total = 0
s := New(WithStartBufSize(4096), WithMaxBufSize(64*1024))
total := 0
for res := range s.Iter(file) {
total += len(res)
}
Expand Down
Loading