From e7a3fa43e26983d3c227686245e4412a014df25e Mon Sep 17 00:00:00 2001 From: fcostaoliveira Date: Mon, 13 Jul 2026 22:15:33 +0100 Subject: [PATCH] fix: make reply capture opt-in so it doesn't inflate query latency (#117) #114 changed the per-command receiver from a discarded nil interface to new(interface{}) so RxBytes could be populated. But that fully unmarshals every RESP reply into nested Go values via reflection INSIDE client.Do -- i.e. inside the measured latency window. For FT.SEARCH / FT.AGGREGATE (O(docs*fields) allocations per query) this materially inflates the very read/vector latencies these benchmarks exist to measure, and RxBytes has no consumer in the specs. Make capture opt-in via --capture-replies (default false): the receiver is a nil interface again, so radix reads-and-discards the reply with no allocation or reflection -- restoring latency fidelity by default. With the flag on, replies are decoded and getRxLen populates RxBytes. Command errors are surfaced by radix regardless of the receiver, so error accounting is unaffected either way. Validated E2E: default run -> RxBytes=0, TxBytes unchanged; --capture-replies -> RxBytes>0, same TxBytes. Adds TestFTSBCaptureRepliesControlsRxBytes and updates the pipeline-tail test to the new default. --- benchmark_runner/redisearch_test.go | 41 +++++++++++++++++++++++++--- cmd/ftsb_redisearch/cmd_processor.go | 15 ++++++++-- cmd/ftsb_redisearch/main.go | 3 ++ 3 files changed, 53 insertions(+), 6 deletions(-) diff --git a/benchmark_runner/redisearch_test.go b/benchmark_runner/redisearch_test.go index 711a12f..c5db359 100644 --- a/benchmark_runner/redisearch_test.go +++ b/benchmark_runner/redisearch_test.go @@ -145,13 +145,46 @@ func TestFTSBPipelineTailIsFlushedAndCounted(t *testing.T) { if parsed.Totals.TotalOps != 100 { t.Errorf("TotalOps = %d, want 100 (pipeline=3 must flush and count the trailing window of a 100-row input)", parsed.Totals.TotalOps) } - // Byte accounting (#111/#112/#114): sent bytes land in TxBytes, and reply - // bytes are actually captured (RxBytes>0 — the 100 HSET integer replies). + // Byte accounting (#111/#112/#114): sent bytes land in TxBytes. if parsed.Totals.TxBytes == 0 { t.Errorf("TxBytes = 0, want > 0 (sent bytes must be counted)") } - if parsed.Totals.RxBytes == 0 { - t.Errorf("RxBytes = 0, want > 0 (reply bytes must be captured)") + // Reply capture is off by default (#117), so RxBytes is 0 here; see + // TestFTSBCaptureRepliesControlsRxBytes for the --capture-replies path. + if parsed.Totals.RxBytes != 0 { + t.Errorf("RxBytes = %d, want 0 without --capture-replies", parsed.Totals.RxBytes) + } +} + +// #117: --capture-replies controls whether reply bytes are decoded and counted. +// Default off (RxBytes==0, no client-side unmarshal on the latency hot path); +// on -> RxBytes>0 (the HSET integer replies are decoded). +func TestFTSBCaptureRepliesControlsRxBytes(t *testing.T) { + startRedisContainer(t) + + rx := func(args ...string) uint64 { + data := runFTSBReadJSON(t, "../testdata/results.capture.json", + append([]string{"--input", "../testdata/minimal.csv"}, args...)...) + var parsed struct { + Totals struct { + TotalOps int `json:"TotalOps"` + RxBytes uint64 `json:"RxBytes"` + } `json:"Totals"` + } + if err := json.Unmarshal(data, &parsed); err != nil { + t.Fatalf("parse: %v", err) + } + if parsed.Totals.TotalOps != 100 { + t.Fatalf("TotalOps = %d, want 100", parsed.Totals.TotalOps) + } + return parsed.Totals.RxBytes + } + + if got := rx(); got != 0 { + t.Errorf("default RxBytes = %d, want 0 (capture off)", got) + } + if got := rx("--capture-replies"); got == 0 { + t.Errorf("RxBytes with --capture-replies = 0, want > 0") } } diff --git a/cmd/ftsb_redisearch/cmd_processor.go b/cmd/ftsb_redisearch/cmd_processor.go index bf87ee9..d1dfd72 100644 --- a/cmd/ftsb_redisearch/cmd_processor.go +++ b/cmd/ftsb_redisearch/cmd_processor.go @@ -286,13 +286,24 @@ type pendingCmd struct { } func sendFlatCmd(p *processor, client radix.Client, cmdType, cmdQueryId, cmd string, docfields []string, txBytesCount uint64, pending []pendingCmd) ([]pendingCmd, bool) { - reply := new(interface{}) + // By default use a nil receiver: radix reads and DISCARDS the reply (no + // allocation, no reflection) so the measured latency isn't inflated by + // client-side unmarshalling -- which is significant for large FT.SEARCH / + // FT.AGGREGATE replies (issue #117). --capture-replies opts into decoding the + // reply so getRxLen can populate RxBytes. Command errors are surfaced by radix + // regardless of the receiver. + var reply *interface{} + var rcv interface{} // nil interface -> discard + if captureReplies { + reply = new(interface{}) + rcv = reply + } key := "" if len(docfields) > 0 { key = docfields[0] } pending = append(pending, pendingCmd{ - action: radix.Cmd(reply, cmd, docfields...), + action: radix.Cmd(rcv, cmd, docfields...), reply: reply, cmdType: cmdType, cmdQueryId: cmdQueryId, diff --git a/cmd/ftsb_redisearch/main.go b/cmd/ftsb_redisearch/main.go index d5407df..f90e53a 100644 --- a/cmd/ftsb_redisearch/main.go +++ b/cmd/ftsb_redisearch/main.go @@ -21,6 +21,7 @@ var ( pipeline int clusterMode bool continueOnErr bool + captureReplies bool timeout time.Duration versionFlag bool logFile string @@ -34,6 +35,7 @@ func init() { flag.StringVar(&password, "a", "", "Password for Redis Auth.") flag.IntVar(&debug, "debug", 0, "Debug printing (choices: 0, 1, 2). (default 0)") flag.BoolVar(&continueOnErr, "continue-on-error", true, "If set to true, it will continue the benchmark and print the error message to stderr.") + flag.BoolVar(&captureReplies, "capture-replies", false, "If true, decode each command's reply so RxBytes is populated. Off by default: capturing fully unmarshals every reply on the client hot path (allocation + reflection inside the measured latency window), which inflates FT.SEARCH/FT.AGGREGATE latency. Command errors are detected regardless of this setting.") flag.BoolVar(&clusterMode, "cluster-mode", false, "If set to true, it will run the client in cluster mode.") flag.IntVar(&pipeline, "pipeline", 1, "Pipeline requests. Default 1 (no pipeline).") flag.IntVar(&timeoutSeconds, "timeout", 60, "Redis connection timeout in seconds.") @@ -64,6 +66,7 @@ func (b *benchmark) GetConfigurationParametersMap() map[string]interface{} { configs["host"] = host configs["clusterMode"] = clusterMode configs["continueOnError"] = continueOnErr + configs["captureReplies"] = captureReplies configs["debug"] = debug configs["pipeline"] = pipeline configs["logFile"] = logFile