From 85d044f36ed32162c1de90db583c1ea17856a9b7 Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:40:12 +0530 Subject: [PATCH 01/10] refactor(prob): centralize input traversal and selection --- cmd/bf/main.go | 40 ++----------- internal/prob/flags.go | 1 + internal/prob/stream.go | 111 +++++++++++++++++++++++------------ internal/prob/stream_test.go | 95 ++++++++++++++++++++++++++++-- 4 files changed, 172 insertions(+), 75 deletions(-) diff --git a/cmd/bf/main.go b/cmd/bf/main.go index ce76cfc..2815f42 100644 --- a/cmd/bf/main.go +++ b/cmd/bf/main.go @@ -2,7 +2,6 @@ package main import ( "bufio" - "bytes" "encoding/json" "flag" "fmt" @@ -88,8 +87,6 @@ func build(args []string, in io.Reader, out, errOut io.Writer) error { noSizeLimit := addSizeLimitFlag(fs) var input prob.InputOptions prob.AddInputFlags(fs, &input) - var fields prob.FieldOptions - prob.AddFieldFlags(fs, &fields) fs.Usage = func() { fmt.Fprint(fs.Output(), `Usage: bf build --expected-items --false-positive-rate

[--delimiter value --field n] [--no-size-limit] [file...] > filter.bf @@ -104,7 +101,7 @@ Options: if err := fs.Parse(args); err != nil { return err } - if err := fields.Validate(); err != nil { + if err := input.Fields.Validate(); err != nil { return fmt.Errorf("usage: %w", err) } @@ -112,7 +109,7 @@ Options: if err != nil { return err } - if err := eachSelectedInput(fs.Args(), in, input, fields, func(item, _ []byte) error { + if err := prob.EachSelectedInputFrom(fs.Args(), in, input, func(item, _ []byte) error { f.Add(item) return nil }); err != nil { @@ -130,8 +127,6 @@ func test(args []string, in io.Reader, out io.Writer) error { noSizeLimit := addSizeLimitFlag(fs) var input prob.InputOptions prob.AddInputFlags(fs, &input) - var fields prob.FieldOptions - prob.AddFieldFlags(fs, &fields) fs.Usage = func() { fmt.Fprint(fs.Output(), `Usage: bf test [--invert] [--delimiter value --field n] [--no-size-limit] [file...] @@ -147,7 +142,7 @@ Options: if err := fs.Parse(args); err != nil { return err } - if err := fields.Validate(); err != nil { + if err := input.Fields.Validate(); err != nil { return fmt.Errorf("usage: %w", err) } if fs.NArg() < 1 { @@ -167,7 +162,7 @@ Options: return err } buffered := bufio.NewWriter(out) - err = eachSelectedInput(paths, in, input, fields, func(item, output []byte) error { + err = prob.EachSelectedInputFrom(paths, in, input, func(item, output []byte) error { present := f.Test(item) if present != *invert { return writeItem(buffered, output, input.NUL) @@ -185,8 +180,6 @@ func dedupe(args []string, in io.Reader, out, errOut io.Writer) error { noSizeLimit := addSizeLimitFlag(fs) var input prob.InputOptions prob.AddInputFlags(fs, &input) - var fields prob.FieldOptions - prob.AddFieldFlags(fs, &fields) fs.Usage = func() { fmt.Fprint(fs.Output(), `Usage: bf dedupe --expected-items --false-positive-rate

[--delimiter value --field n] [--no-size-limit] [file...] @@ -201,7 +194,7 @@ Options: if err := fs.Parse(args); err != nil { return err } - if err := fields.Validate(); err != nil { + if err := input.Fields.Validate(); err != nil { return fmt.Errorf("usage: %w", err) } f, err := bf.NewWithLimit(*expected, *rate, filterSizeLimit(*noSizeLimit)) @@ -209,7 +202,7 @@ Options: return err } buffered := bufio.NewWriter(out) - err = eachSelectedInput(fs.Args(), in, input, fields, func(item, output []byte) error { + err = prob.EachSelectedInputFrom(fs.Args(), in, input, func(item, output []byte) error { if f.Test(item) { return nil } @@ -358,27 +351,6 @@ func countStdin(paths []string) int { return count } -func eachSelectedInput(paths []string, in io.Reader, input prob.InputOptions, fields prob.FieldOptions, fn func(item, output []byte) error) error { - rawInput := prob.InputOptions{NUL: input.NUL} - return prob.EachInputFrom(paths, in, rawInput, func(record []byte) error { - item, err := fields.Select(record) - if err != nil { - return err - } - if input.Trim { - item = bytes.TrimSpace(item) - } - if input.IgnoreEmpty && len(item) == 0 { - return nil - } - output := item - if fields.Enabled() { - output = record - } - return fn(item, output) - }) -} - func writeItem(out *bufio.Writer, item []byte, nul bool) error { if _, err := out.Write(item); err != nil { return err diff --git a/internal/prob/flags.go b/internal/prob/flags.go index 5f29280..eaa7fe4 100644 --- a/internal/prob/flags.go +++ b/internal/prob/flags.go @@ -7,6 +7,7 @@ func AddInputFlags(fs *flag.FlagSet, opts *InputOptions) { fs.BoolVar(&opts.NUL, "nul", false, "read NUL-delimited input") fs.BoolVar(&opts.Trim, "trim", false, "trim surrounding whitespace") fs.BoolVar(&opts.IgnoreEmpty, "ignore-empty", false, "ignore empty input items") + AddFieldFlags(fs, &opts.Fields) } func AddFieldFlags(fs *flag.FlagSet, opts *FieldOptions) { diff --git a/internal/prob/stream.go b/internal/prob/stream.go index a0232c1..724c470 100644 --- a/internal/prob/stream.go +++ b/internal/prob/stream.go @@ -12,6 +12,7 @@ type InputOptions struct { NUL bool Trim bool IgnoreEmpty bool + Fields FieldOptions } func EachInput(paths []string, opts InputOptions, fn func([]byte) error) error { @@ -20,10 +21,51 @@ func EachInput(paths []string, opts InputOptions, fn func([]byte) error) error { // EachInputFrom reads stdin for an empty path list or an explicit "-" path. func EachInputFrom(paths []string, stdin io.Reader, opts InputOptions, fn func([]byte) error) error { - if len(paths) == 0 { - return eachReader("", stdin, opts, fn) + return EachSelectedInputFrom(paths, stdin, opts, func(item, _ []byte) error { + return fn(item) + }) +} + +// EachSelectedInputFrom supplies the selected value and the output record. +// With field selection, output retains the complete record without its delimiter. +func EachSelectedInputFrom(paths []string, stdin io.Reader, opts InputOptions, fn func(item, output []byte) error) error { + if err := opts.Fields.Validate(); err != nil { + return err } + return EachRecordFrom(paths, stdin, opts.NUL, func(record []byte) error { + delim := byte('\n') + if opts.NUL { + delim = 0 + } + if record[len(record)-1] == delim { + record = record[:len(record)-1] + if !opts.NUL && len(record) > 0 && record[len(record)-1] == '\r' { + record = record[:len(record)-1] + } + } + item, err := opts.Fields.Select(record) + if err != nil { + return err + } + if opts.Trim { + item = bytes.TrimSpace(item) + } + if opts.IgnoreEmpty && len(item) == 0 { + return nil + } + output := item + if opts.Fields.Enabled() { + output = record + } + return fn(item, output) + }) +} +// EachFile visits files in order, using stdin for no paths or one "-" path. +func EachFile(paths []string, stdin io.Reader, fn func(string, io.Reader) error) error { + if len(paths) == 0 { + return fn("", stdin) + } stdinUsed := false for _, path := range paths { if path == "-" { @@ -31,7 +73,11 @@ func EachInputFrom(paths []string, stdin io.Reader, opts InputOptions, fn func([ return fmt.Errorf("stdin may be read only once") } stdinUsed = true - if err := eachReader("", stdin, opts, fn); err != nil { + } + } + for _, path := range paths { + if path == "-" { + if err := fn("", stdin); err != nil { return err } continue @@ -40,7 +86,7 @@ func EachInputFrom(paths []string, stdin io.Reader, opts InputOptions, fn func([ if err != nil { return fmt.Errorf("open %s: %w", path, err) } - err = eachReader(path, f, opts, fn) + err = fn(path, f) closeErr := f.Close() if err != nil { return err @@ -52,45 +98,38 @@ func EachInputFrom(paths []string, stdin io.Reader, opts InputOptions, fn func([ return nil } -func eachReader(name string, r io.Reader, opts InputOptions, fn func([]byte) error) error { +// EachRecordFrom preserves record bytes, including any LF or NUL delimiter. +// Callback slices are valid only until the callback returns. +func EachRecordFrom(paths []string, stdin io.Reader, nul bool, fn func([]byte) error) error { delim := byte('\n') - if opts.NUL { + if nul { delim = 0 } - - br := bufio.NewReader(r) - var continued []byte - for { - item, err := br.ReadSlice(delim) - if err == bufio.ErrBufferFull { - continued = append(continued, item...) - continue - } - if len(continued) != 0 { - item = append(continued, item...) - continued = nil - } - if len(item) > 0 { - if item[len(item)-1] == delim { - item = item[:len(item)-1] - if !opts.NUL && len(item) > 0 && item[len(item)-1] == '\r' { - item = item[:len(item)-1] - } + return EachFile(paths, stdin, func(name string, r io.Reader) error { + br := bufio.NewReader(r) + var continued []byte + for { + record, err := br.ReadSlice(delim) + if err == bufio.ErrBufferFull { + continued = append(continued, record...) + continue } - if opts.Trim { - item = bytes.TrimSpace(item) + if len(continued) != 0 { + continued = append(continued, record...) + record = continued } - if !opts.IgnoreEmpty || len(item) > 0 { - if err := fn(item); err != nil { + if len(record) > 0 { + if err := fn(record); err != nil { return err } } + continued = continued[:0] + if err == io.EOF { + return nil + } + if err != nil { + return fmt.Errorf("read %s: %w", name, err) + } } - if err == io.EOF { - return nil - } - if err != nil { - return fmt.Errorf("read %s: %w", name, err) - } - } + }) } diff --git a/internal/prob/stream_test.go b/internal/prob/stream_test.go index af67282..da936f5 100644 --- a/internal/prob/stream_test.go +++ b/internal/prob/stream_test.go @@ -2,6 +2,10 @@ package prob import ( "bytes" + "errors" + "io" + "os" + "path/filepath" "reflect" "strings" "testing" @@ -12,7 +16,7 @@ func TestEachReaderCanTrimAndIgnoreEmptyLines(t *testing.T) { var got []string spacedLines := strings.Join([]string{" a ", "", "b"}, "\n") + "\n" - err := eachReader("test", strings.NewReader(spacedLines), InputOptions{Trim: true, IgnoreEmpty: true}, func(item []byte) error { + err := EachInputFrom(nil, strings.NewReader(spacedLines), InputOptions{Trim: true, IgnoreEmpty: true}, func(item []byte) error { got = append(got, string(item)) return nil }) @@ -29,7 +33,7 @@ func TestEachReaderNULDelimited(t *testing.T) { var got []string input := string([]byte{'a', 0, 'b', 0}) - err := eachReader("test", strings.NewReader(input), InputOptions{NUL: true}, func(item []byte) error { + err := EachInputFrom(nil, strings.NewReader(input), InputOptions{NUL: true}, func(item []byte) error { got = append(got, string(item)) return nil }) @@ -45,7 +49,7 @@ func TestEachReaderRemovesCRLFDelimiter(t *testing.T) { t.Parallel() var got []string - err := eachReader("test", strings.NewReader("a\r\nb\r\n"), InputOptions{}, func(item []byte) error { + err := EachInputFrom(nil, strings.NewReader("a\r\nb\r\n"), InputOptions{}, func(item []byte) error { got = append(got, string(item)) return nil }) @@ -81,7 +85,7 @@ func TestEachReaderAcceptsLargeItems(t *testing.T) { want := strings.Repeat("x", 8<<10) var got string - if err := eachReader("test", strings.NewReader(want+"\n"), InputOptions{}, func(item []byte) error { + if err := EachInputFrom(nil, strings.NewReader(want+"\n"), InputOptions{}, func(item []byte) error { got = string(item) return nil }); err != nil { @@ -97,7 +101,88 @@ func BenchmarkEachReader(b *testing.B) { b.SetBytes(int64(len(input))) b.ReportAllocs() for b.Loop() { - if err := eachReader("bench", bytes.NewReader(input), InputOptions{}, func([]byte) error { return nil }); err != nil { + if err := EachInputFrom(nil, bytes.NewReader(input), InputOptions{}, func([]byte) error { return nil }); err != nil { + b.Fatal(err) + } + } +} + +func TestSharedInputPreservesFileBoundariesAndSelection(t *testing.T) { + t.Parallel() + path := filepath.Join(t.TempDir(), "records") + if err := os.WriteFile(path, []byte("first:: a ::payload"), 0600); err != nil { + t.Fatal(err) + } + opts := InputOptions{Trim: true, IgnoreEmpty: true, Fields: FieldOptions{Delimiter: "::", Field: 2}} + var values, records []string + err := EachSelectedInputFrom([]string{path, "-"}, strings.NewReader("second:: b ::payload\r\nskip:: ::payload\n"), opts, func(item, output []byte) error { + values = append(values, string(item)) + records = append(records, string(output)) + return nil + }) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(values, []string{"a", "b"}) || !reflect.DeepEqual(records, []string{"first:: a ::payload", "second:: b ::payload"}) { + t.Fatalf("values = %q, records = %q", values, records) + } + if err := EachInputFrom(nil, strings.NewReader("missing\n"), opts, func([]byte) error { return nil }); err == nil { + t.Fatal("accepted a missing field") + } +} + +func TestRawInputPreservesLongBinaryRecords(t *testing.T) { + t.Parallel() + for _, nul := range []bool{false, true} { + delim := "\n" + if nul { + delim = "\x00" + } + want := []string{delim, strings.Repeat("x", 16<<10) + "\r" + delim, strings.Repeat("y", 16<<10) + delim, "last\r"} + var got []string + err := EachRecordFrom(nil, strings.NewReader(strings.Join(want, "")), nul, func(record []byte) error { + got = append(got, string(record)) + return nil + }) + if err != nil || !reflect.DeepEqual(got, want) { + t.Fatalf("NUL=%v: record preservation failed: %v", nul, err) + } + } +} + +func TestSharedInputPropagatesErrorsAndRejectsRepeatedStdinBeforeReading(t *testing.T) { + t.Parallel() + sentinel := errors.New("input failure") + for _, run := range []func() error{ + func() error { + return EachInputFrom(nil, strings.NewReader("value\n"), InputOptions{}, func([]byte) error { return sentinel }) + }, + func() error { + return EachRecordFrom(nil, io.MultiReader(strings.NewReader("value\n"), failingReader{sentinel}), false, func([]byte) error { return nil }) + }, + } { + if err := run(); !errors.Is(err, sentinel) { + t.Fatalf("error = %v, want %v", err, sentinel) + } + } + if err := EachFile([]string{"-", "-"}, strings.NewReader("value"), func(string, io.Reader) error { + t.Fatal("consumed input before rejecting repeated stdin") + return nil + }); err == nil { + t.Fatal("accepted repeated stdin") + } +} + +type failingReader struct{ err error } + +func (r failingReader) Read([]byte) (int, error) { return 0, r.err } + +func BenchmarkLongRecords(b *testing.B) { + input := bytes.Repeat(append(bytes.Repeat([]byte("x"), 16<<10), '\n'), 1000) + b.SetBytes(int64(len(input))) + b.ReportAllocs() + for b.Loop() { + if err := EachInputFrom(nil, bytes.NewReader(input), InputOptions{}, func([]byte) error { return nil }); err != nil { b.Fatal(err) } } From 04139587007dc64c2112a70fd2ed93da861d671c Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:40:33 +0530 Subject: [PATCH 02/10] feat(prob): select delimited fields in hll and heavy --- cmd/heavy/main.go | 6 ++++ cmd/hll/main.go | 34 +++++++++++------- cmd/hll/main_test.go | 68 ++++++++++++++++++++++++++++++++++++ internal/heavy/heavy.go | 4 +-- internal/heavy/heavy_test.go | 34 ++++++++++++++++++ 5 files changed, 131 insertions(+), 15 deletions(-) create mode 100644 cmd/hll/main_test.go diff --git a/cmd/heavy/main.go b/cmd/heavy/main.go index a2ce8d3..b5538d6 100644 --- a/cmd/heavy/main.go +++ b/cmd/heavy/main.go @@ -30,6 +30,7 @@ func main() { fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s [options] [file...] Find frequent values in newline-delimited input. +Use --delimiter and --field to rank one field per record. Default mode uses bounded-memory approximate counts. Use --exact only when all distinct values fit in memory. Examples: @@ -64,6 +65,11 @@ Options: os.Exit(2) } + if err := input.Fields.Validate(); err != nil { + fmt.Fprintf(os.Stderr, "heavy: %v\n", err) + os.Exit(2) + } + results, err := heavy.Run(flag.Args(), heavy.Config{ Top: *top, Capacity: *capacity, diff --git a/cmd/hll/main.go b/cmd/hll/main.go index 91cb3d0..a42566f 100644 --- a/cmd/hll/main.go +++ b/cmd/hll/main.go @@ -4,6 +4,7 @@ import ( "encoding/json" "flag" "fmt" + "io" "os" "strings" @@ -28,9 +29,9 @@ func main() { var err error switch os.Args[1] { case "count": - err = count(os.Args[2:]) + err = count(os.Args[2:], os.Stdin, os.Stdout) case "build": - err = build(os.Args[2:]) + err = build(os.Args[2:], os.Stdin, os.Stdout) case "estimate": err = estimate(os.Args[2:]) case "merge": @@ -75,7 +76,7 @@ Run "hll -h" for command-specific flags. `) } -func count(args []string) error { +func count(args []string, in io.Reader, out io.Writer) error { fs := flag.NewFlagSet("hll count", flag.ExitOnError) precision := fs.Uint("precision", uint(hll.DefaultP), "HLL precision, 4..20") jsonOut := fs.Bool("json", false, "write JSON output") @@ -85,7 +86,8 @@ func count(args []string) error { fmt.Fprint(fs.Output(), `Usage: hll count [options] [file...] Read newline-delimited values from files or stdin and print an approximate unique count. -Empty lines and surrounding whitespace are significant unless input flags change that. +Use --delimiter and --field to count one field per record. +Empty values and surrounding whitespace are significant unless input flags change that. Example: awk '{print $1}' access.log | hll count --ignore-empty @@ -98,6 +100,9 @@ Options: return err } + if err := input.Fields.Validate(); err != nil { + return fmt.Errorf("usage: %w", err) + } p, err := hll.Precision(*precision) if err != nil { return err @@ -106,16 +111,16 @@ Options: if err != nil { return err } - if err := prob.EachInput(fs.Args(), input, func(item []byte) error { + if err := prob.EachInputFrom(fs.Args(), in, input, func(item []byte) error { s.Add(item) return nil }); err != nil { return err } - return writeEstimate(os.Stdout, s, *jsonOut) + return writeEstimate(out, s, *jsonOut) } -func build(args []string) error { +func build(args []string, in io.Reader, out io.Writer) error { fs := flag.NewFlagSet("hll build", flag.ExitOnError) precision := fs.Uint("precision", uint(hll.DefaultP), "HLL precision, 4..20") var input prob.InputOptions @@ -124,6 +129,7 @@ func build(args []string) error { fmt.Fprint(fs.Output(), `Usage: hll build [options] [file...] > file.hll Read values from files or stdin and write a binary HyperLogLog sketch to stdout. +Use --delimiter and --field to hash one field per record. Redirect stdout to save the sketch. Options: @@ -134,6 +140,9 @@ Options: return err } + if err := input.Fields.Validate(); err != nil { + return fmt.Errorf("usage: %w", err) + } p, err := hll.Precision(*precision) if err != nil { return err @@ -142,13 +151,13 @@ Options: if err != nil { return err } - if err := prob.EachInput(fs.Args(), input, func(item []byte) error { + if err := prob.EachInputFrom(fs.Args(), in, input, func(item []byte) error { s.Add(item) return nil }); err != nil { return err } - return hll.Write(os.Stdout, s) + return hll.Write(out, s) } func estimate(args []string) error { @@ -271,7 +280,7 @@ func countStdin(paths []string) int { return count } -func writeEstimate(out *os.File, s *hll.Sketch, jsonOut bool) error { +func writeEstimate(out io.Writer, s *hll.Sketch, jsonOut bool) error { m := s.Metadata() if jsonOut { return json.NewEncoder(out).Encode(struct { @@ -279,7 +288,6 @@ func writeEstimate(out *os.File, s *hll.Sketch, jsonOut bool) error { RelativeError float64 `json:"relative_error"` }{m.ApproxUnique, m.RelativeError}) } - fmt.Fprintf(out, "approx_unique=%d\n", m.ApproxUnique) - fmt.Fprintf(out, "relative_error=%.2f%%\n", m.RelativeError*100) - return nil + _, err := fmt.Fprintf(out, "approx_unique=%d\nrelative_error=%.2f%%\n", m.ApproxUnique, m.RelativeError*100) + return err } diff --git a/cmd/hll/main_test.go b/cmd/hll/main_test.go new file mode 100644 index 0000000..75ee280 --- /dev/null +++ b/cmd/hll/main_test.go @@ -0,0 +1,68 @@ +package main + +import ( + "bytes" + "errors" + "io" + "strings" + "testing" + + "github.com/jo-cube/toolbox/internal/hll" +) + +func TestCountAndBuildSelectTheSameValues(t *testing.T) { + t.Parallel() + for _, nul := range []bool{false, true} { + args := []string{"--delimiter", "::", "--field", "2", "--trim", "--ignore-empty"} + input := "1:: a ::x\n2::b::y\n3::a::z\n4:: ::empty\n5::c" + if nul { + args = append(args, "-0") + input = strings.ReplaceAll(input, "\n", "\x00") + } + var direct, artifact, saved bytes.Buffer + if err := count(args, strings.NewReader(input), &direct); err != nil { + t.Fatal(err) + } + if err := build(args, strings.NewReader(input), &artifact); err != nil { + t.Fatal(err) + } + sketch, err := hll.Read(&artifact) + if err != nil { + t.Fatal(err) + } + if err := writeEstimate(&saved, sketch, false); err != nil { + t.Fatal(err) + } + if direct.String() != saved.String() || !strings.Contains(direct.String(), "approx_unique=3\n") { + t.Fatalf("direct = %q, saved = %q", direct.String(), saved.String()) + } + } +} + +func TestInputErrorsDoNotProduceASketchOrCount(t *testing.T) { + t.Parallel() + for _, run := range []func([]string, io.Reader, io.Writer) error{count, build} { + for _, args := range [][]string{{"--field", "2"}, {"-d", "::", "-f", "2"}} { + var out bytes.Buffer + err := run(args, strings.NewReader("missing-field\n"), &out) + if err == nil || out.Len() != 0 { + t.Fatalf("args = %q, error = %v, output bytes = %d", args, err, out.Len()) + } + if len(args) == 2 && !strings.HasPrefix(err.Error(), "usage:") { + t.Fatalf("invalid field options must be a usage error: %v", err) + } + } + } +} + +func TestCountPropagatesOutputFailure(t *testing.T) { + t.Parallel() + sentinel := errors.New("output failed") + if err := count(nil, strings.NewReader("a\n"), failingWriter{sentinel}); !errors.Is(err, sentinel) { + t.Fatalf("error = %v, want %v", err, sentinel) + } +} + +type failingWriter struct{ err error } + +func (w failingWriter) Write([]byte) (int, error) { return 0, w.err } diff --git a/internal/heavy/heavy.go b/internal/heavy/heavy.go index f2f25f6..8f1b2f3 100644 --- a/internal/heavy/heavy.go +++ b/internal/heavy/heavy.go @@ -92,12 +92,12 @@ func approximate(paths []string, cfg Config) ([]Result, error) { tracked := map[string]*trackedItem{} var items minItems if err := prob.EachInput(paths, cfg.Input, func(item []byte) error { - key := string(item) - if existing, ok := tracked[key]; ok { + if existing, ok := tracked[string(item)]; ok { existing.count++ heap.Fix(&items, existing.index) return nil } + key := string(item) if len(tracked) < cfg.Capacity { entry := &trackedItem{item: key, count: 1} tracked[key] = entry diff --git a/internal/heavy/heavy_test.go b/internal/heavy/heavy_test.go index 565ceb8..2057bf1 100644 --- a/internal/heavy/heavy_test.go +++ b/internal/heavy/heavy_test.go @@ -6,6 +6,8 @@ import ( "path/filepath" "strings" "testing" + + "github.com/jo-cube/toolbox/internal/prob" ) const repeatedLetters = `b @@ -85,3 +87,35 @@ func writeInput(t *testing.T, content string) string { } return path } + +func TestFieldRankingWorksInBothModes(t *testing.T) { + t.Parallel() + path := writeInput(t, "1:: a ::x\x002::b::y\x003:: a ::z\x004:: ::empty\x005::b\x006::a") + for _, exact := range []bool{false, true} { + got, err := Run([]string{path}, Config{ + Top: 2, Exact: exact, + Input: prob.InputOptions{NUL: true, Trim: true, IgnoreEmpty: true, Fields: prob.FieldOptions{Delimiter: "::", Field: 2}}, + }) + if err != nil { + t.Fatal(err) + } + if len(got) != 2 || got[0].Item != "a" || got[0].CountEstimate != 3 || got[0].CountLowerBound != 3 || got[1].Item != "b" || got[1].CountEstimate != 2 { + t.Fatalf("exact=%v: results = %#v", exact, got) + } + } +} + +func BenchmarkRepeatedValues(b *testing.B) { + path := filepath.Join(b.TempDir(), "input") + input := strings.Repeat("dominant-value\n", 10000) + if err := os.WriteFile(path, []byte(input), 0600); err != nil { + b.Fatal(err) + } + b.SetBytes(int64(len(input))) + b.ReportAllocs() + for b.Loop() { + if _, err := Run([]string{path}, Config{Top: 1}); err != nil { + b.Fatal(err) + } + } +} From f22b6744078f7218af3212664334df6c058b097d Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:40:48 +0530 Subject: [PATCH 03/10] feat(sample): add complementary rate sampling --- cmd/sample/main.go | 5 ++- internal/sample/sample.go | 74 +++++----------------------------- internal/sample/sample_test.go | 53 ++++++++++++++++++++++++ 3 files changed, 67 insertions(+), 65 deletions(-) diff --git a/cmd/sample/main.go b/cmd/sample/main.go index cec74fd..8216359 100644 --- a/cmd/sample/main.go +++ b/cmd/sample/main.go @@ -18,6 +18,7 @@ func main() { rate := flag.Float64("rate", 0, "sample probability, 0..1") count := flag.Int("count", 0, "reservoir sample size") stable := flag.Bool("stable", false, "use deterministic hash sampling with --rate") + invert := flag.Bool("invert", false, "emit records excluded by the rate sample") seed := flag.Int64("seed", 0, "random or stable hash seed") nul := flag.Bool("nul", false, "read and write NUL-delimited records") flag.BoolVar(nul, "0", false, "read and write NUL-delimited records") @@ -26,7 +27,7 @@ func main() { flag.Usage = func() { name := filepath.Base(os.Args[0]) - fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s (--rate

[--stable] | --count ) [file...] + fmt.Fprintf(flag.CommandLine.Output(), `Usage: %s (--rate

[--stable] [--invert] | --count ) [file...] Emit a subset of input records while preserving emitted records exactly. Set exactly one of --rate or --count. @@ -40,6 +41,7 @@ Examples: Notes: --rate samples each record independently unless --stable is set. --stable hashes the full record or a selected field without its trailing delimiter. + --invert emits the complementary rate sample; use --stable or the same --seed to repeat a split. -0 and --nul preserve NUL-delimited records instead of newline-delimited records. --count uses reservoir sampling and writes selected records after reading input. @@ -64,6 +66,7 @@ Options: Count: *count, CountSet: flagWasSet("count"), Stable: *stable, + Invert: *invert, Seed: *seed, SeedSet: flagWasSet("seed"), NUL: *nul, diff --git a/internal/sample/sample.go b/internal/sample/sample.go index db75daa..7b0344d 100644 --- a/internal/sample/sample.go +++ b/internal/sample/sample.go @@ -19,6 +19,7 @@ type Config struct { Count int CountSet bool Stable bool + Invert bool Seed int64 SeedSet bool NUL bool @@ -40,6 +41,9 @@ func Validate(cfg Config) error { if hasCount && cfg.Count <= 0 { return fmt.Errorf("count must be a positive integer") } + if cfg.Invert && hasCount { + return fmt.Errorf("--invert can only be used with --rate") + } if cfg.Stable && hasCount { return fmt.Errorf("--stable can only be used with --rate") } @@ -77,8 +81,8 @@ func RunFrom(paths []string, cfg Config, out io.Writer, stdin io.Reader) error { func rateRandom(paths []string, cfg Config, out io.Writer, stdin io.Reader) error { rng := rand.New(rand.NewSource(seed(cfg))) - return eachRaw(paths, stdin, delimiter(cfg), func(record []byte) error { - if rng.Float64() < cfg.Rate { + return prob.EachRecordFrom(paths, stdin, cfg.NUL, func(record []byte) error { + if (rng.Float64() < cfg.Rate) != cfg.Invert { _, err := out.Write(record) return err } @@ -89,7 +93,7 @@ func rateRandom(paths []string, cfg Config, out io.Writer, stdin io.Reader) erro func rateStable(paths []string, cfg Config, out io.Writer, stdin io.Reader) error { threshold := uint64(cfg.Rate * float64(math.MaxUint64)) delim := delimiter(cfg) - return eachRaw(paths, stdin, delim, func(record []byte) error { + return prob.EachRecordFrom(paths, stdin, cfg.NUL, func(record []byte) error { key := record if len(key) > 0 && key[len(key)-1] == delim { key = key[:len(key)-1] @@ -101,7 +105,8 @@ func rateStable(paths []string, cfg Config, out io.Writer, stdin io.Reader) erro if err != nil { return err } - if cfg.Rate >= 1 || prob.Hash64(key, uint64(cfg.Seed)) < threshold { + selected := cfg.Rate >= 1 || prob.Hash64(key, uint64(cfg.Seed)) < threshold + if selected != cfg.Invert { _, err := out.Write(record) return err } @@ -118,7 +123,7 @@ func reservoir(paths []string, cfg Config, out io.Writer, stdin io.Reader) error var items []selected var seen int64 - if err := eachRaw(paths, stdin, delimiter(cfg), func(record []byte) error { + if err := prob.EachRecordFrom(paths, stdin, cfg.NUL, func(record []byte) error { seen++ if len(items) < cfg.Count { items = append(items, selected{record: append([]byte(nil), record...), order: seen}) @@ -156,62 +161,3 @@ func delimiter(cfg Config) byte { } return '\n' } - -func eachRaw(paths []string, stdin io.Reader, delim byte, fn func([]byte) error) error { - if len(paths) == 0 { - return eachRawReader("", stdin, delim, fn) - } - stdinUsed := false - for _, path := range paths { - if path == "-" { - if stdinUsed { - return fmt.Errorf("stdin may be read only once") - } - stdinUsed = true - if err := eachRawReader("", stdin, delim, fn); err != nil { - return err - } - continue - } - f, err := os.Open(path) - if err != nil { - return fmt.Errorf("open %s: %w", path, err) - } - err = eachRawReader(path, f, delim, fn) - closeErr := f.Close() - if err != nil { - return err - } - if closeErr != nil { - return fmt.Errorf("close %s: %w", path, closeErr) - } - } - return nil -} - -func eachRawReader(name string, r io.Reader, delim byte, fn func([]byte) error) error { - br := bufio.NewReader(r) - var continued []byte - for { - record, err := br.ReadSlice(delim) - if err == bufio.ErrBufferFull { - continued = append(continued, record...) - continue - } - if len(continued) != 0 { - record = append(continued, record...) - continued = nil - } - if len(record) > 0 { - if err := fn(record); err != nil { - return err - } - } - if err == io.EOF { - return nil - } - if err != nil { - return fmt.Errorf("read %s: %w", name, err) - } - } -} diff --git a/internal/sample/sample_test.go b/internal/sample/sample_test.go index 98fb0c0..97190a6 100644 --- a/internal/sample/sample_test.go +++ b/internal/sample/sample_test.go @@ -241,3 +241,56 @@ func writeInput(t *testing.T, content string) string { } return path } + +func TestInvertedRateSamplePartitionsRecords(t *testing.T) { + t.Parallel() + for _, stable := range []bool{false, true} { + for _, nul := range []bool{false, true} { + for _, rate := range []float64{0, 0.4, 1} { + cfg := Config{Rate: rate, RateSet: true, Stable: stable, Seed: 7, SeedSet: true, NUL: nul} + if stable { + cfg.Fields = prob.FieldOptions{Delimiter: "::", Field: 2} + } + delim := "\r\n" + if nul { + delim = "\x00" + } + var records []string + for i := range 100 { + records = append(records, fmt.Sprintf("%d::cohort-%d:: payload %s", i, i%10, delim)) + } + records[len(records)-1] = strings.TrimSuffix(records[len(records)-1], delim) + input := strings.Join(records, "") + var selected, rejected bytes.Buffer + if err := RunFrom(nil, cfg, &selected, strings.NewReader(input)); err != nil { + t.Fatal(err) + } + cfg.Invert = true + if err := RunFrom(nil, cfg, &rejected, strings.NewReader(input)); err != nil { + t.Fatal(err) + } + a, b := selected.String(), rejected.String() + for _, record := range records { + inA, inB := strings.HasPrefix(a, record), strings.HasPrefix(b, record) + if inA == inB { + t.Fatalf("stable=%v nul=%v rate=%v: record missing or duplicated: %q", stable, nul, rate, record) + } + if inA { + a = strings.TrimPrefix(a, record) + } else { + b = strings.TrimPrefix(b, record) + } + } + if a != "" || b != "" || (rate == 0 && selected.Len() != 0) || (rate == 1 && rejected.Len() != 0) { + t.Fatal("unexpected output in complementary sample") + } + if rate == 0.4 && (selected.Len() == 0 || rejected.Len() == 0) { + t.Fatal("split did not exercise both decisions") + } + } + } + } + if err := Validate(Config{Count: 2, Invert: true}); err == nil { + t.Fatal("accepted inverted reservoir sampling") + } +} From b9c32a82a99d5896d5665e24843d417eef30d746 Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:41:11 +0530 Subject: [PATCH 04/10] fix(card): preserve JSON value types --- internal/card/card.go | 59 +++++++++----------------------------- internal/card/card_test.go | 33 +++++++++++++++++++++ 2 files changed, 46 insertions(+), 46 deletions(-) diff --git a/internal/card/card.go b/internal/card/card.go index fbedae0..f20adf5 100644 --- a/internal/card/card.go +++ b/internal/card/card.go @@ -11,6 +11,7 @@ import ( "strings" "github.com/jo-cube/toolbox/internal/hll" + "github.com/jo-cube/toolbox/internal/prob" ) type Config struct { @@ -68,7 +69,7 @@ func runCSV(paths []string, cfg Config, stdin io.Reader) ([]Profile, error) { if err != nil { return nil, err } - return finish(counters, eachPath(paths, stdin, func(name string, r io.Reader) error { + return finish(counters, prob.EachFile(paths, stdin, func(name string, r io.Reader) error { cr := csv.NewReader(r) header, err := cr.Read() if err != nil { @@ -160,7 +161,9 @@ func runJSON(paths []string, cfg Config, stdin io.Reader) ([]Profile, error) { counters[i].total++ continue } - addAny(counters[i], value) + if err := addAny(counters[i], value); err != nil { + return err + } } return nil })) @@ -209,26 +212,22 @@ func addValue(c *counter, value string, ok bool) { c.sketch.Add([]byte(value)) } -func addAny(c *counter, value any) { +func addAny(c *counter, value any) error { c.total++ if value == nil { c.nulls++ - return + return nil } - if s, ok := value.(string); ok { - if s == "" { - c.empty++ - return - } - c.sketch.Add([]byte(s)) - return + if s, ok := value.(string); ok && s == "" { + c.empty++ + return nil } encoded, err := json.Marshal(value) if err != nil { - c.sketch.Add([]byte(fmt.Sprint(value))) - return + return err } c.sketch.Add(encoded) + return nil } func recordValue(record []string, idx int) (string, bool) { @@ -279,40 +278,8 @@ func lookup(value any, path []string) (any, bool) { return current, true } -func eachPath(paths []string, stdin io.Reader, fn func(string, io.Reader) error) error { - if len(paths) == 0 { - return fn("", stdin) - } - stdinUsed := false - for _, path := range paths { - if path == "-" { - if stdinUsed { - return fmt.Errorf("stdin may be read only once") - } - stdinUsed = true - if err := fn("", stdin); err != nil { - return err - } - continue - } - f, err := os.Open(path) - if err != nil { - return fmt.Errorf("open %s: %w", path, err) - } - err = fn(path, f) - closeErr := f.Close() - if err != nil { - return err - } - if closeErr != nil { - return fmt.Errorf("close %s: %w", path, closeErr) - } - } - return nil -} - func eachLine(paths []string, stdin io.Reader, fn func(string) error) error { - return eachPath(paths, stdin, func(name string, r io.Reader) error { + return prob.EachFile(paths, stdin, func(name string, r io.Reader) error { br := bufio.NewReader(r) line := 0 for { diff --git a/internal/card/card_test.go b/internal/card/card_test.go index e3377ca..eec3d40 100644 --- a/internal/card/card_test.go +++ b/internal/card/card_test.go @@ -112,3 +112,36 @@ func writeInput(t *testing.T, content string) string { } return path } + +func TestJSONProfileDistinguishesTypesAndNormalizesObjects(t *testing.T) { + t.Parallel() + input := `{"v":1} +{"v":"1"} +{"v":true} +{"v":"true"} +{"v":[]} +{"v":"[]"} +{"v":{"a":1,"b":2}} +{"v":{"b":2,"a":1}} +{"v":null} +{"v":""} +{} +` + profiles, err := RunFrom(nil, Config{Mode: "json", JSONPaths: []string{".v"}}, strings.NewReader(input)) + if err != nil { + t.Fatal(err) + } + got := profiles[0] + if got.ApproxUnique != 7 || got.Nulls != 1 || got.Empty != 1 || got.Missing != 1 || got.Total != 11 { + t.Fatalf("profile = %#v", got) + } +} + +func TestCSVFilesEachUseTheirOwnHeader(t *testing.T) { + t.Parallel() + path := writeInput(t, "id,country\nu1,US\n") + profiles, err := RunFrom([]string{path, "-"}, Config{Mode: "csv", Columns: []string{"id"}}, strings.NewReader("country,id\nCA,u2\nUS,u1\n")) + if err != nil || len(profiles) != 1 || profiles[0].ApproxUnique != 2 || profiles[0].Total != 3 { + t.Fatalf("profiles = %#v, error = %v", profiles, err) + } +} From 6fd8c749c139d636d83532d07cf0f59421810077 Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:41:35 +0530 Subject: [PATCH 05/10] docs(prob): document probabilistic tool extensions --- README.md | 12 ++++++++++-- docs/card.md | 2 +- docs/development.md | 5 +++-- docs/heavy.md | 11 ++++++++++- docs/hll.md | 11 ++++++++++- docs/probabilistic-tools.md | 6 ++++-- docs/sample.md | 20 +++++++++++++++++--- 7 files changed, 55 insertions(+), 12 deletions(-) diff --git a/README.md b/README.md index 9aae835..3986153 100644 --- a/README.md +++ b/README.md @@ -94,6 +94,13 @@ cat known-users.txt | bf build --expected-items 1000000 --false-positive-rate 0. cat candidates.txt | bf test users.bf ``` +Count or rank a tab-separated key without an extra extraction step: + +```sh +hll count -d $'\t' -f 2 events.tsv +heavy --top 20 -d $'\t' -f 2 events.tsv +``` + Profile JSON field cardinality: ```sh @@ -106,10 +113,11 @@ Find frequent API paths: awk '{print $7}' access.log | heavy --top 20 ``` -Take a stable sample: +Take a stable sample and its complementary cohort: ```sh -sample --rate 0.01 --stable events.jsonl +sample --rate 0.01 --stable events.jsonl > selected.jsonl +sample --rate 0.01 --stable --invert events.jsonl > remaining.jsonl ``` ## Behavior At A Glance diff --git a/docs/card.md b/docs/card.md index 21f2778..c0a23d7 100644 --- a/docs/card.md +++ b/docs/card.md @@ -131,7 +131,7 @@ Separate counters report: - `empty`: empty strings - `total`: total records observed for that field -For non-string JSON values, `card` hashes the JSON encoding of the value. JSON numbers retain their input spelling, so `1`, `1.0`, and `1e0` are distinct and large integers keep their full precision. +For JSON values, `card` hashes the JSON encoding, preserving type distinctions: `1` and `"1"`, or `true` and `"true"`, are different values. Object key order does not affect cardinality. JSON numbers retain their input spelling, so `1`, `1.0`, and `1e0` are distinct and large integers keep their full precision. Malformed CSV or malformed JSON fails the command with a line-aware error where possible. diff --git a/docs/development.md b/docs/development.md index a08c9fe..cfb4f06 100644 --- a/docs/development.md +++ b/docs/development.md @@ -48,7 +48,7 @@ The usual split is: - `cmd//main.go` owns flags, help text, stdout/stderr formatting, and exit codes. - `internal//` owns behavior that can be tested without shelling out. -- `internal/prob/` owns shared line/NUL input handling, literal field selection, and the stable hashing used by HLL and sampling. +- `internal/prob/` owns ordered file/stdin traversal, raw line/NUL records, value normalization and literal field selection, and the stable hashing used by HLL and sampling. Use `hello` as the minimal reference for that shape. @@ -101,7 +101,7 @@ After building, run the small local CLI smoke suite: TOOLBOX_BIN="$PWD/bin" sh scripts/smoke-local.sh ``` -It checks version aliases, help output, representative exit statuses, and the `hello` output for the locally built binaries. +It checks field-selection pipelines, complementary sampling, version aliases, help output, representative exit statuses, and the `hello` output for the locally built binaries. ## Implementation Notes @@ -131,6 +131,7 @@ It checks version aliases, help output, representative exit statuses, and the `h Probabilistic tools: - use the Go standard library, except for Bloom filtering's direct XXHash64 dependency +- share file/stdin traversal; sampling preserves raw records while HLL, Bloom filtering, and heavy-hitter counting normalize selected values - read streams without loading full inputs unless the selected algorithm requires it - estimate HLL cardinality from its register histogram using Ertl's improved raw estimator - keep state-file compatibility constants in package code diff --git a/docs/heavy.md b/docs/heavy.md index 7cffcab..db32aca 100644 --- a/docs/heavy.md +++ b/docs/heavy.md @@ -92,6 +92,8 @@ The true observed count is between `count_lower_bound` and `count_estimate`. In - `--exact`: use exact counts with unbounded memory - `--json`: write JSON output - `--tsv`: write tab-separated output +- `-d`, `--delimiter VALUE`: literal field delimiter +- `-f`, `--field N`: 1-based field to rank - `--trim`: trim surrounding whitespace - `--ignore-empty`: skip empty items - `-0`, `--nul`: read NUL-delimited items @@ -110,7 +112,14 @@ Defaults: - empty lines are counted as a value - no structured parsing is performed -Use tools such as `awk`, `cut`, or `jq` before `heavy` to select the field you want ranked. +Rank a literal-delimited field directly in either counting mode: + +```sh +heavy --top 20 -d $'\t' -f 2 events.tsv +heavy --top 20 --exact -d $'\t' -f 2 events.tsv +``` + +The delimiter and field must be supplied together. Output items are the selected values. Trimming and empty-value filtering apply after selection. Missing fields fail; empty fields count unless `--ignore-empty` is set. Use `jq` or a CSV parser upstream for structured data. ## Approximate Mode diff --git a/docs/hll.md b/docs/hll.md index 0514b55..73555c2 100644 --- a/docs/hll.md +++ b/docs/hll.md @@ -121,6 +121,8 @@ Input options for `count` and `build`: - `--precision N`: HLL precision from `4` to `20`; default is `14` - `--json`: write JSON output for `count` +- `-d`, `--delimiter VALUE`: literal field delimiter +- `-f`, `--field N`: 1-based field to count or insert - `--trim`: trim surrounding whitespace - `--ignore-empty`: skip empty items - `-0`, `--nul`: read NUL-delimited items @@ -141,7 +143,14 @@ Defaults: - empty lines are counted as a value - no structured parsing is performed -Use tools such as `awk`, `cut`, or `jq` before `hll` to select the value you want counted. +Select a literal-delimited field directly: + +```sh +hll count -d $'\t' -f 2 events.tsv +hll build -d $'\t' -f 2 events.tsv > users.hll +``` + +The delimiter and field must be supplied together. Trimming and empty-value filtering apply to the selected field. Missing fields fail; empty fields count unless `--ignore-empty` is set. Use `jq` or a CSV parser upstream for structured data. ## Accuracy And Memory diff --git a/docs/probabilistic-tools.md b/docs/probabilistic-tools.md index 567fdd0..b15870d 100644 --- a/docs/probabilistic-tools.md +++ b/docs/probabilistic-tools.md @@ -48,14 +48,16 @@ Defaults are conservative: ## Field Selection -`bf build`, `bf test`, `bf dedupe`, and `sample --stable` can use one field from a literal-delimited record: +`hll count`, `hll build`, `heavy`, `bf build`, `bf test`, `bf dedupe`, and `sample --stable` can use one field from a literal-delimited record: ```sh +hll count -d $'\t' -f 2 events.tsv +heavy -d $'\t' -f 2 events.tsv bf test -d $'\t' -f 2 users.bf events.tsv sample --rate 0.01 --stable -d $'\t' -f 2 events.tsv ``` -`-d`/`--delimiter` and `-f`/`--field` must be supplied together. Fields are 1-based. `bf test`, `bf dedupe`, and `sample` emit complete records even though the selected field controls the decision. This is intentionally not CSV or JSON parsing; use an upstream parser when quoting or structured data matters. +`-d`/`--delimiter` and `-f`/`--field` must be supplied together. Fields are 1-based. Missing fields are input errors; empty fields are valid values. `--trim` and `--ignore-empty`, where supported, apply after field selection. `bf test`, `bf dedupe`, and `sample` emit complete records even though the selected field controls the decision. This is intentionally not CSV or JSON parsing; use an upstream parser when quoting or structured data matters. ## Output diff --git a/docs/sample.md b/docs/sample.md index 4c4b6cf..e56d9c2 100644 --- a/docs/sample.md +++ b/docs/sample.md @@ -31,7 +31,7 @@ sample --version ## Synopsis ```sh -sample --rate

[--stable] [--seed n] [--delimiter value --field n] [-0] [file...] +sample --rate

[--stable] [--invert] [--seed n] [--delimiter value --field n] [-0] [file...] sample --count [--seed n] [-0] [file...] ``` @@ -87,6 +87,19 @@ The same input record, rate, and seed produce the same decision across runs. Stable mode hashes the full record without its trailing newline or NUL delimiter. With `--delimiter` and `--field`, it hashes only the selected field while emitting the complete record. Records with the same selected value therefore receive the same sampling decision. +### Complementary Rate Sampling + +`--invert` emits records excluded by the same rate sample. Create two disjoint cohorts with the same input, rate, seed, and field selection: + +```sh +sample --rate 0.2 --stable -d $'\t' -f 2 events.tsv > holdout.tsv +sample --rate 0.2 --stable --invert -d $'\t' -f 2 events.tsv > training.tsv +``` + +A selected key stays on the same side even if records are reordered or split across files. Random rate sampling also supports inversion; repeat the same input order and explicit `--seed` for complementary runs. Without a seed, separate random runs are independent. + +Inverted rate `0` emits everything; inverted rate `1` emits nothing. Inversion is unavailable with `--count`, which retains only the reservoir. + ### Reservoir Sampling `--count N` keeps up to `N` records from the stream without knowing the stream length in advance. @@ -112,6 +125,7 @@ With `-0` or `--nul`, records and emitted delimiters are NUL-separated instead. - `--rate P`: sample each record with probability `P`, from `0` to `1` - `--count N`: keep up to a positive `N` records using reservoir sampling - `--stable`: use deterministic hash sampling with `--rate` +- `--invert`: emit records excluded by the rate sample - `--seed N`: seed random modes or stable hashing - `-d`, `--delimiter VALUE`: literal field delimiter for stable sampling - `-f`, `--field N`: 1-based field used for stable sampling @@ -121,7 +135,7 @@ With `-0` or `--nul`, records and emitted delimiters are NUL-separated instead. Invalid combinations fail: - `--rate` with `--count` -- `--stable` with `--count` +- `--stable` or `--invert` with `--count` - a non-positive `--count` - neither `--rate` nor `--count` - field selection without `--stable` @@ -139,4 +153,4 @@ Invalid combinations fail: - CLI flags live in `cmd/sample/main.go`. - Sampling behavior lives in `internal/sample`. - Stable hashing uses `internal/prob.Hash64`. -- Preserve records exactly. Do not switch to the shared trimmed stream reader for this tool. +- Preserve records exactly through `internal/prob.EachRecordFrom`; normalization is for value-processing tools. From faf1ca93e63efe2f35dcb53006a127fed5894429 Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:41:51 +0530 Subject: [PATCH 06/10] test(cli): cover probabilistic extensions --- scripts/smoke-local.sh | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/scripts/smoke-local.sh b/scripts/smoke-local.sh index 15b6b05..ac7fb44 100644 --- a/scripts/smoke-local.sh +++ b/scripts/smoke-local.sh @@ -40,4 +40,13 @@ expect_status 2 "$bin/card" expect_status 2 "$bin/heavy" --top 0 expect_status 2 "$bin/sample" +[ "$(printf '1::a\n2::b\n3::a\n' | "$bin/hll" count -d :: -f 2 | head -n 1)" = "approx_unique=2" ] +[ "$(printf '1::a\n2::b\n3::a\n' | "$bin/hll" build -d :: -f 2 | "$bin/hll" estimate - | head -n 1)" = "approx_unique=2" ] +printf '1::a\n2::b\n3::a\n' | "$bin/heavy" --exact --top 1 -d :: -f 2 --tsv | grep -q '1[[:space:]]2[[:space:]]2[[:space:]]a' +[ "$(printf 'a\nb\n' | "$bin/sample" --rate 0 --invert)" = "$(printf 'a\nb')" ] +expect_status 2 "$bin/hll" count --field 2 +expect_status 2 "$bin/hll" build --delimiter :: +expect_status 2 "$bin/heavy" --field 2 +expect_status 2 "$bin/sample" --count 1 --invert + printf 'LOCAL CLI SMOKE TEST PASSED\n' From 300aa35e0c9788184cebb58e68ba96ce7578c09c Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:07:06 +0530 Subject: [PATCH 07/10] perf(card): reduce profiling allocations --- internal/card/card.go | 19 +++++++--- internal/card/card_benchmark_test.go | 54 ++++++++++++++++++++++++++++ internal/card/card_test.go | 40 +++++++++++++++++++++ 3 files changed, 109 insertions(+), 4 deletions(-) create mode 100644 internal/card/card_benchmark_test.go diff --git a/internal/card/card.go b/internal/card/card.go index f20adf5..4d1c4a9 100644 --- a/internal/card/card.go +++ b/internal/card/card.go @@ -7,6 +7,7 @@ import ( "fmt" "io" "os" + "sort" "strconv" "strings" @@ -71,6 +72,7 @@ func runCSV(paths []string, cfg Config, stdin io.Reader) ([]Profile, error) { } return finish(counters, prob.EachFile(paths, stdin, func(name string, r io.Reader) error { cr := csv.NewReader(r) + cr.ReuseRecord = true header, err := cr.Read() if err != nil { return fmt.Errorf("%s: read header: %w", name, err) @@ -116,11 +118,20 @@ func runDelimited(paths []string, cfg Config, stdin io.Reader) ([]Profile, error if err != nil { return nil, err } + order := make([]int, len(indexes)) + for i := range order { + order[i] = i + } + sort.Slice(order, func(i, j int) bool { return indexes[order[i]] < indexes[order[j]] }) return finish(counters, eachLine(paths, stdin, func(line string) error { - parts := strings.Split(line, cfg.Delimiter) - for i, idx := range indexes { - value, ok := recordValue(parts, idx) - addValue(counters[i], value, ok) + column, more := -1, true + var value string + for _, i := range order { + for more && column < indexes[i] { + value, line, more = strings.Cut(line, cfg.Delimiter) + column++ + } + addValue(counters[i], value, column == indexes[i]) } return nil })) diff --git a/internal/card/card_benchmark_test.go b/internal/card/card_benchmark_test.go new file mode 100644 index 0000000..aee8ac8 --- /dev/null +++ b/internal/card/card_benchmark_test.go @@ -0,0 +1,54 @@ +package card + +import ( + "fmt" + "strconv" + "strings" + "testing" +) + +func BenchmarkProfile(b *testing.B) { + for _, width := range []int{3, 100} { + fields := make([]string, width) + columns := make([]string, width) + for i := range fields { + fields[i] = fmt.Sprintf("value-%03d", i) + columns[i] = strconv.Itoa(i + 1) + } + var rows strings.Builder + for i := range 10000 { + fields[0] = strconv.Itoa(i) + rows.WriteString(strings.Join(fields, ",")) + rows.WriteByte('\n') + } + for _, mode := range []string{"csv", "delimiter"} { + for _, selection := range []string{"first", "last", "all"} { + b.Run(fmt.Sprintf("%s/%d/%s", mode, width, selection), func(b *testing.B) { + selected := columns + if selection == "first" { + selected = columns[:1] + } + if selection == "last" { + selected = columns[width-1:] + } + input := rows.String() + if mode == "csv" { + input = strings.Join(columns, ",") + "\n" + input + } + cfg := Config{Mode: mode, Columns: selected, Delimiter: ","} + b.SetBytes(int64(len(input))) + b.ReportAllocs() + for b.Loop() { + profiles, err := RunFrom(nil, cfg, strings.NewReader(input)) + if err != nil { + b.Fatal(err) + } + if profiles[0].Total != 10000 { + b.Fatal(profiles[0]) + } + } + }) + } + } + } +} diff --git a/internal/card/card_test.go b/internal/card/card_test.go index eec3d40..779df8c 100644 --- a/internal/card/card_test.go +++ b/internal/card/card_test.go @@ -4,6 +4,8 @@ import ( "bytes" "os" "path/filepath" + "reflect" + "strconv" "strings" "testing" ) @@ -145,3 +147,41 @@ func TestCSVFilesEachUseTheirOwnHeader(t *testing.T) { t.Fatalf("profiles = %#v, error = %v", profiles, err) } } + +func TestDelimitedSelectionMatchesSplit(t *testing.T) { + t.Parallel() + columns := []string{"4", "2", "1", "2", strconv.Itoa(int(^uint(0) >> 1)), "3"} + for _, delimiter := range []string{"::", "aa", "\t", "💠"} { + records := []string{ + strings.Join([]string{"a", "", "c", ""}, delimiter), + strings.Join([]string{"", "b"}, delimiter), + "solo", "", + "last" + delimiter, + strings.Join([]string{strings.Repeat("x", 64<<10), "z", ""}, delimiter), + } + cfg := Config{Mode: "delimiter", Delimiter: delimiter, Columns: columns, Precision: 8} + got, err := RunFrom(nil, cfg, strings.NewReader(strings.Join(records, "\r\n")+"\r")) + if err != nil { + t.Fatal(err) + } + counters, err := newCounters(columns, cfg.Precision) + if err != nil { + t.Fatal(err) + } + for _, record := range records { + parts := strings.Split(record, delimiter) + for i, col := range columns { + index, _ := strconv.Atoi(col) + value, ok := recordValue(parts, index-1) + addValue(counters[i], value, ok) + } + } + want, err := finish(counters, nil) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("delimiter %q: got %#v, want %#v", delimiter, got, want) + } + } +} From 5125cfb6ba496ed700c8d999a6777b0f029db73e Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:07:24 +0530 Subject: [PATCH 08/10] perf(heavy): repair heap root in place --- internal/heavy/heavy.go | 4 +-- internal/heavy/heavy_benchmark_test.go | 37 +++++++++++++++++++++ internal/heavy/heavy_test.go | 45 ++++++++++++++++++++++++++ 3 files changed, 84 insertions(+), 2 deletions(-) create mode 100644 internal/heavy/heavy_benchmark_test.go diff --git a/internal/heavy/heavy.go b/internal/heavy/heavy.go index 8f1b2f3..2078ec2 100644 --- a/internal/heavy/heavy.go +++ b/internal/heavy/heavy.go @@ -105,13 +105,13 @@ func approximate(paths []string, cfg Config) ([]Result, error) { return nil } - replaced := heap.Pop(&items).(*trackedItem) + replaced := items[0] delete(tracked, replaced.item) replaced.error = replaced.count replaced.item = key replaced.count++ tracked[key] = replaced - heap.Push(&items, replaced) + heap.Fix(&items, 0) return nil }); err != nil { return nil, err diff --git a/internal/heavy/heavy_benchmark_test.go b/internal/heavy/heavy_benchmark_test.go new file mode 100644 index 0000000..48ad780 --- /dev/null +++ b/internal/heavy/heavy_benchmark_test.go @@ -0,0 +1,37 @@ +package heavy + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "testing" +) + +func BenchmarkHighCardinality(b *testing.B) { + for _, skewed := range []bool{false, true} { + var input strings.Builder + for i := range 100000 { + if skewed && i%10 != 0 { + fmt.Fprintln(&input, "dominant") + } else { + fmt.Fprintf(&input, "key-%08d\n", i) + } + } + path := filepath.Join(b.TempDir(), "input") + if err := os.WriteFile(path, []byte(input.String()), 0600); err != nil { + b.Fatal(err) + } + for _, capacity := range []int{1000, 10000} { + b.Run(fmt.Sprintf("skewed=%t/capacity=%d", skewed, capacity), func(b *testing.B) { + b.SetBytes(int64(input.Len())) + b.ReportAllocs() + for b.Loop() { + if _, err := Run([]string{path}, Config{Top: 20, Capacity: capacity}); err != nil { + b.Fatal(err) + } + } + }) + } + } +} diff --git a/internal/heavy/heavy_test.go b/internal/heavy/heavy_test.go index 2057bf1..ed36f18 100644 --- a/internal/heavy/heavy_test.go +++ b/internal/heavy/heavy_test.go @@ -2,8 +2,10 @@ package heavy import ( "fmt" + "math/rand" "os" "path/filepath" + "reflect" "strings" "testing" @@ -119,3 +121,46 @@ func BenchmarkRepeatedValues(b *testing.B) { } } } + +func TestApproximateMatchesSpaceSavingReference(t *testing.T) { + t.Parallel() + rng := rand.New(rand.NewSource(1)) + input := make([]string, 2000) + for i := range input { + input[i] = fmt.Sprintf("key-%02d", rng.Intn(30)) + } + path := writeInput(t, strings.Join(input, "\n")) + for _, capacity := range []int{1, 2, 7, 40} { + tracked := map[string]Result{} + for _, item := range input { + entry, exists := tracked[item] + if !exists && len(tracked) == capacity { + var least Result + first := true + for _, candidate := range tracked { + if first || candidate.CountEstimate < least.CountEstimate || candidate.CountEstimate == least.CountEstimate && candidate.Item < least.Item { + least, first = candidate, false + } + } + delete(tracked, least.Item) + entry.CountEstimate = least.CountEstimate + } + entry.Item = item + entry.CountEstimate++ + entry.CountLowerBound++ + tracked[item] = entry + } + var want []Result + for _, entry := range tracked { + want = append(want, entry) + } + want = rank(want, capacity) + got, err := Run([]string{path}, Config{Top: capacity, Capacity: capacity}) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("capacity %d: got %#v, want %#v", capacity, got, want) + } + } +} From 5958e5b6479e9d6a2f33a221a3cc3740fecbb4da Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:08:07 +0530 Subject: [PATCH 09/10] perf(rdbsh): avoid unnecessary iterator copies --- internal/rdbsh/commands.go | 6 +- internal/rdbsh/export.go | 36 +++++------ internal/rdbsh/iter.go | 8 ++- internal/rdbsh/iter_test.go | 118 ++++++++++++++++++++++++++++++++++++ 4 files changed, 143 insertions(+), 25 deletions(-) create mode 100644 internal/rdbsh/iter_test.go diff --git a/internal/rdbsh/commands.go b/internal/rdbsh/commands.go index f856f60..e53b09f 100644 --- a/internal/rdbsh/commands.go +++ b/internal/rdbsh/commands.go @@ -131,7 +131,7 @@ func (s *Shell) cmdScan(args []string) error { } w := tabwriter.NewWriter(s.out, 0, 0, 2, ' ', 0) - result, err := s.iterate(prefix, limit, func(key, value []byte) error { + result, err := s.iterate(prefix, limit, true, func(key, value []byte) error { _, err := fmt.Fprintf(w, "%s\t%s\n", formatBytes(key), formatBytes(value)) return err }) @@ -157,7 +157,7 @@ func (s *Shell) cmdKeys(args []string) error { return err } - result, err := s.iterate(prefix, limit, func(key, _ []byte) error { + result, err := s.iterate(prefix, limit, false, func(key, _ []byte) error { _, err := fmt.Fprintln(s.out, formatBytes(key)) return err }) @@ -189,7 +189,7 @@ func (s *Shell) cmdCount(args []string) error { } } - result, err := s.iterate(prefix, 0, func(_, _ []byte) error { return nil }) + result, err := s.iterate(prefix, 0, false, func(_, _ []byte) error { return nil }) if err != nil { return err } diff --git a/internal/rdbsh/export.go b/internal/rdbsh/export.go index e46ea99..45d1186 100644 --- a/internal/rdbsh/export.go +++ b/internal/rdbsh/export.go @@ -1,6 +1,7 @@ package rdbsh import ( + "bufio" "encoding/csv" "encoding/json" "fmt" @@ -105,7 +106,7 @@ func (s *Shell) exportCSV(writer io.Writer, prefix []byte) (int, error) { return 0, err } - result, err := s.iterate(prefix, 0, func(key, value []byte) error { + result, err := s.iterate(prefix, 0, true, func(key, value []byte) error { return csvWriter.Write([]string{formatExportBytes(key), formatExportBytes(value)}) }) if err != nil { @@ -118,36 +119,28 @@ func (s *Shell) exportCSV(writer io.Writer, prefix []byte) (int, error) { return result.Count, nil } -func (s *Shell) exportJSON(writer io.Writer, prefix []byte) (int, error) { +func (s *Shell) exportJSON(out io.Writer, prefix []byte) (int, error) { + writer := bufio.NewWriter(out) + encoder := json.NewEncoder(writer) if _, err := fmt.Fprintln(writer, "["); err != nil { return 0, err } first := true - result, err := s.iterate(prefix, 0, func(key, value []byte) error { - entry, err := json.Marshal(struct { - Key string `json:"key"` - Value string `json:"value"` - }{ - Key: formatExportBytes(key), - Value: formatExportBytes(value), - }) - if err != nil { - return err - } + result, err := s.iterate(prefix, 0, true, func(key, value []byte) error { if !first { if _, err := fmt.Fprintln(writer, ","); err != nil { return err } } first = false - if _, err := writer.Write(entry); err != nil { - return err - } - if _, err := fmt.Fprintln(writer); err != nil { - return err - } - return nil + return encoder.Encode(struct { + Key string `json:"key"` + Value string `json:"value"` + }{ + Key: formatExportBytes(key), + Value: formatExportBytes(value), + }) }) if err != nil { return 0, err @@ -155,5 +148,8 @@ func (s *Shell) exportJSON(writer io.Writer, prefix []byte) (int, error) { if _, err := fmt.Fprintln(writer, "]"); err != nil { return 0, err } + if err := writer.Flush(); err != nil { + return 0, err + } return result.Count, nil } diff --git a/internal/rdbsh/iter.go b/internal/rdbsh/iter.go index 554ada8..20d12cc 100644 --- a/internal/rdbsh/iter.go +++ b/internal/rdbsh/iter.go @@ -11,7 +11,7 @@ type iterationResult struct { Limited bool } -func (s *Shell) iterate(prefix []byte, limit int, fn func(key, value []byte) error) (iterationResult, error) { +func (s *Shell) iterate(prefix []byte, limit int, readValues bool, fn func(key, value []byte) error) (iterationResult, error) { var result iterationResult it := s.newIterator() @@ -29,7 +29,11 @@ func (s *Shell) iterate(prefix []byte, limit int, fn func(key, value []byte) err break } - if err := fn(key, it.Value()); err != nil { + var value []byte + if readValues { + value = it.Value() + } + if err := fn(key, value); err != nil { return result, err } result.Count++ diff --git a/internal/rdbsh/iter_test.go b/internal/rdbsh/iter_test.go new file mode 100644 index 0000000..a7dd654 --- /dev/null +++ b/internal/rdbsh/iter_test.go @@ -0,0 +1,118 @@ +package rdbsh + +import ( + "bytes" + "errors" + "fmt" + "io" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" +) + +func newTestShell(t testing.TB) *Shell { + t.Helper() + ldb, err := exec.LookPath("ldb") + if err != nil { + ldb, err = exec.LookPath("rocksdb_ldb") + } + if err != nil { + t.Skip("RocksDB integration checks require ldb or rocksdb_ldb") + } + path := filepath.Join(t.TempDir(), "db") + if output, err := exec.Command(ldb, "--db="+path, "--create_if_missing", "put", "seed", "value").CombinedOutput(); err != nil { + t.Fatalf("create DB: %v: %s", err, output) + } + s, err := NewShell(Config{DBPath: path, Writable: true, Out: io.Discard, ErrOut: io.Discard}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(s.Close) + if err := s.delete([]byte("seed")); err != nil { + t.Fatal(err) + } + return s +} + +func BenchmarkReadCommands(b *testing.B) { + for _, size := range []int{32, 4096} { + s := newTestShell(b) + value := []byte(strings.Repeat("v", size)) + for i := range 10000 { + if err := s.put([]byte(fmt.Sprintf("key-%05d", i)), value); err != nil { + b.Fatal(err) + } + } + out, err := os.OpenFile(os.DevNull, os.O_WRONLY, 0) + if err != nil { + b.Fatal(err) + } + b.Cleanup(func() { out.Close() }) + s.out = out + for _, command := range []string{"count", "count key-0", "keys key-0 10000", "export - json"} { + b.Run(fmt.Sprintf("value=%d/%s", size, command), func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + if err := s.Exec(command); err != nil { + b.Fatal(err) + } + } + }) + } + } +} + +func TestReadCommandsPreserveOutput(t *testing.T) { + s := newTestShell(t) + for key, value := range map[string][]byte{"alpha": []byte("one"), "alpine": {}, "beta": {0, 10}} { + if err := s.put([]byte(key), value); err != nil { + t.Fatal(err) + } + } + for _, tt := range []struct{ command, want, diagnostic string }{ + {"count", "3 keys total\n", ""}, + {"count al", "2 keys (prefix: al)\n", ""}, + {"count z", "0 keys (prefix: z)\n", ""}, + {"keys al 1", "alpha\n... (limit 1 reached)\n(1 keys shown)\n", ""}, + {"keys al 2", "alpha\nalpine\n(2 keys shown)\n", ""}, + {"keys z", "(no keys found)\n", ""}, + {"scan b 1", "beta 0x000a\n", ""}, + {"export - csv al", "key,value\nalpha,one\nalpine,0x\n", "exported 2 entries to stdout (csv)\n"}, + {"export - json al", "[\n{\"key\":\"alpha\",\"value\":\"one\"}\n,\n{\"key\":\"alpine\",\"value\":\"0x\"}\n]\n", "exported 2 entries to stdout (json)\n"}, + {"export - json z", "[\n]\n", "exported 0 entries to stdout (json)\n"}, + } { + var out, diagnostic bytes.Buffer + s.out, s.errOut = &out, &diagnostic + if err := s.Exec(tt.command); err != nil { + t.Fatalf("%s: %v", tt.command, err) + } + if out.String() != tt.want || diagnostic.String() != tt.diagnostic { + t.Fatalf("%s: stdout=%q, stderr=%q", tt.command, out.String(), diagnostic.String()) + } + } + failure := errors.New("callback failed") + result, err := s.iterate([]byte("al"), 0, true, func(_, _ []byte) error { return failure }) + if !errors.Is(err, failure) || result.Count != 0 { + t.Fatalf("iterate: result=%#v, err=%v", result, err) + } +} + +type failedExportWriter struct{ err error } + +func (w failedExportWriter) Write([]byte) (int, error) { return 0, w.err } + +func TestJSONExportPropagatesWriteAndFlushErrors(t *testing.T) { + s := newTestShell(t) + failure := errors.New("output failed") + for _, size := range []int{32, 8192} { + if err := s.put([]byte("key"), []byte(strings.Repeat("v", size))); err != nil { + t.Fatal(err) + } + count, err := s.exportJSON(failedExportWriter{failure}, nil) + if !errors.Is(err, failure) || count != 0 { + t.Fatalf("size %d: count=%d, err=%v", size, count, err) + } + } +} From b7ec165df316fe553321ef3272b9f0d9a45eea97 Mon Sep 17 00:00:00 2001 From: jo-cube <56916509+jo-cube@users.noreply.github.com> Date: Mon, 14 Sep 2026 18:08:16 +0530 Subject: [PATCH 10/10] docs: document performance review --- docs/development.md | 81 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 81 insertions(+) diff --git a/docs/development.md b/docs/development.md index cfb4f06..787f466 100644 --- a/docs/development.md +++ b/docs/development.md @@ -145,6 +145,87 @@ Compatibility constants are the source of truth: Changing any of these can make old state files unreadable. Treat such changes as explicit file-format migrations. +## Performance Checks + +Run focused benchmarks serially so packages do not compete for CPU: + +```sh +go test -p 1 ./internal/card ./internal/heavy -run '^$' -bench 'Benchmark(Profile|HighCardinality|RepeatedValues)$' -benchmem -benchtime=300ms -count=3 +go test ./internal/rdbsh -run '^$' -bench '^BenchmarkReadCommands$' -benchmem -benchtime=300ms -count=3 +``` + +The RocksDB command needs the same compiler/linker flags printed by `make test` +when headers or libraries are outside the default search paths. Its integration +tests and benchmarks use `ldb` or `rocksdb_ldb` from PATH to create temporary +fixtures, and skip when neither is available. Go removes these fixtures after +checking them. They ran with Homebrew's `rocksdb_ldb` for the measurements below. + +### September 2026 Review + +Baseline: `faf1ca9`, before the performance changes in `card`, `heavy`, and +`rdbsh`. Measurements used an Apple M4, Go 1.26.4, darwin/arm64, and RocksDB +11.8.1. The table gives medians of three 300 ms benchmark runs. Allocation +volume is Go `B/op`, shown in decimal MB; it is not peak RSS or native RocksDB +memory. Inputs and databases are prepared outside the timed loops. RocksDB +output goes to `/dev/null`, so these measurements include write syscalls but +exclude database startup and persistent output storage. + +| Benchmark workload | Before → after (ms/op) | Before → after (MB allocated/op) | +| --- | --- | --- | +| CSV, 10k rows, all 3 columns | 0.872 → 0.741 | 0.774 → 0.294 | +| CSV, 10k rows, first of 100 columns | 10.929 → 10.274 | 26.911 → 8.991 | +| Delimited, 10k rows, first of 100 columns | 8.983 → 1.027 | 28.181 → 10.261 | +| Delimited, 10k rows, last of 100 columns | 9.126 → 8.739 | 28.181 → 10.261 | +| Delimited, 10k rows, all 100 columns | 19.136 → 16.806 | 29.822 → 11.903 | +| Heavy, 100k distinct values, capacity 1k | 19.186 → 17.358 | 1.820 → 1.820 | +| RocksDB count, 10k entries, 4 KiB values | 4.588 → 1.509 | 41.282 → 0.240 | +| RocksDB keys, 10k entries, 4 KiB values | 8.778 → 5.827 | 41.618 → 0.560 | +| RocksDB JSON export, 10k entries, 32-byte values | 14.944 → 3.274 | 2.082 → 1.446 | +| RocksDB JSON export, 10k entries, 4 KiB values | 53.053 → 41.801 | 131.659 → 82.882 | + +The changes remove work without changing CLI output or state formats: + +- CSV profiling enables the standard reader's record-slice reuse. Each row is + consumed before the next read; header indexes are resolved beforehand. +- Delimited profiling orders selectors once, scans forward with `strings.Cut`, + and stops after the last selected field. It preserves requested output order, + duplicate selectors, empty fields, and missing-field counts without building + a slice of every field for every row. Line reading still consumes each full row. +- RocksDB key-only commands skip value retrieval and C-to-Go copies. Scans and + exports still read values, and iterator errors and prefix/limit checks remain. +- JSON export uses a buffered writer and one standard JSON encoder, eliminating + per-entry encoded-byte copies and most small writes. Formatting and atomic + file publication remain unchanged, including propagation of flush failures. +- Heavy tracking replaces the heap root in place with `heap.Fix`, avoiding a + pop followed by a push. The existing logarithmic update algorithm, lexical + eviction tie-break, and count bounds remain intact. + +Whole CLI comparisons used identical pre-generated files, equal build flags, +startup-inclusive elapsed time, one warm-up pair, and five measured pairs with +alternating before/after execution order. All stdout and stderr matched: + +| CLI workload | Before → after median | +| --- | --- | +| CSV, 200k rows, user_id/country/plan from the release performance workload | 18.06 → 15.53 ms | +| Delimited, 200k rows, first of 100 columns | 205.36 → 46.27 ms | +| Delimited, 200k rows, last of 100 columns | 208.92 → 198.68 ms | +| Heavy, 500k sequential distinct values, top 20 | 103.56 → 93.93 ms | + +A `strings.SplitSeq` trial regressed late-column selection on wide rows by +about 13%; it was replaced with the smaller cursor loop. Heavy's skewed and +single-value controls showed little change, so its gains apply primarily to +replacement-heavy streams. No pooling, unsafe borrowing, custom parser, or +additional concurrency was needed. + +The review retained the shared borrowed-record reader, bounded sampling, +HLL histogram estimator, split-block Bloom layout, streaming Bloom union, +and buffered/streaming `kshape` artifact and report code. Stable hash and file +compatibility contracts rule out casually swapping hashes. Further work should +start with actual profiles of large `kshape` renders (which revisit sketches +at multiple resolutions), structured JSON profiling, or large-value exports. +Exact heavy counting still intentionally retains all distinct values. This +review did not benchmark live Kafka requests or published Linux release assets. + ## Release GitHub Actions: