diff --git a/HOWTO.md b/HOWTO.md index d25ac8b..105ee7e 100644 --- a/HOWTO.md +++ b/HOWTO.md @@ -453,6 +453,96 @@ The practical checks, for any comparison of large data: - **Do not validate candidates one at a time** before comparing them. Use `ValidatePair`, or `Compare`, which does. +## One process is one observation + +Everything `Compare` reports — the interval, the noise floor, `Resolved` — +describes the noise **within one run of your program**. Run the same program +again and you get a new process, and a new process lays its data out in memory +differently: different addresses, different pages, different cache sets. For +small data that barely matters. For data larger than the caches, or full of +pointers (trees, linked structures, maps of heap objects), it can move the +difference between A and B by several points, and that shift is fixed for the +whole life of the process. No amount of `Repeats` inside the process sees it, +and the A/A validation cannot either, because both halves of an A/A run share +the same layout. + +How big this gets was measured on a 1M-key ordered index +([issue #109](https://github.com/TomTonic/rtcompare/issues/109)): the same +comparison, repeated in separate processes, scattered 4 to 10 times more +widely than each process's own interval said it should. One process reported +`+14.3% [+13.2, +15.2]`, the next `+24.8%`. Merely adding an unrelated +structure to the fixture build moved another comparison from `+2.9%` to +`-27.3%`, both resolved with narrow intervals. Even for data that fits in the +cache, the intervals were about twice too narrow. + +So the rule is: **when your data is large or pointer-heavy, one process is one +observation.** Run several and pool them. `Compare` reminds you: when the +program holds more than 16 MB of live data, its warnings say so. + +### How + +The `multiproc` package does all of it. You say how to build each candidate; +it does the rest: + +```go +func main() { + multiproc.Main(multiproc.Options{}, multiproc.Pair{ + Name: "lookup", + A: func() rtcompare.Candidate { return lookupIn(buildTreeA()) }, + B: func() rtcompare.Candidate { return lookupIn(buildTreeB()) }, + }) +} +``` + +In a test, `multiproc.RunTest(t, multiproc.Options{}, pair)` does the same and +returns the pooled results for your assertions. + +What happens behind that call, so that you don't have to remember any of it: + +- **Your program is started again as child processes**, one after another, so + they don't disturb each other. In a test, the children run only that test. +- **Each child gets a different heap layout.** Before anything is built, it + fills the heap with a seeded amount of filler, so that your data lands at + different addresses in each process (`rtcompare.PerturbHeap`). Otherwise + every process would repeat the same layout, and its bias with it. +- **The build order alternates.** Even-numbered processes build A's data + first, odd-numbered ones B's. This matters more than the heap: in the + reproduction in `cmd/rtcompare-aa`, whichever of two identical 1M-node lists + was built second was about 3% faster in every process, however the heap was + perturbed. That is also why the builders are functions: they have to run + inside each process, after the perturbation, in the right order. Build + everything the measurement depends on inside them. +- **It stops when the answer is precise enough.** At least 5 processes, then + more until every comparison's pooled interval is within ±2 percentage + points, or within ±10% of the difference itself, and at most 20. It only + stops after an even number, so both build orders count equally. In the + measurements behind issue #109, five were enough for data in the cache and + 13 to 15 were needed for data far out of it. + +Keep the machine awake while this runs: a laptop that goes to sleep pauses the +measurement for as long as it sleeps. rtcompare notices that +(`Report.Suspended`) but cannot give you the time back. + +For anything the pairs don't cover, `multiproc.Run` takes a suite function of +your own, and still perturbs the heap and pools for you; alternate the build +order by `p.Index` yourself there. If you run the processes by some other +means, `rtcompare.Combine(reports, 0)` does the pooling part on its own. + +### Reading a pooled result + +`Combine` treats each process as one number, its delta, and puts a Student t +interval around their mean. With few processes that interval is wide, and it +should be: five observations are five observations. Next to it you get: + +- **`Inflation`** — how many times more the processes scatter than one + process's interval implies. Near 1, a single process would have told you the + truth. Well above 1, it would not have, and the pooled result is the one to + quote. +- **`I2`** — the share of the scatter that the per-process intervals do not + explain. Above about 0.5, the differences between processes dominate. +- **Warnings** — among them, when processes resolved the difference with + opposite signs, each of them confident. + ## Troubleshooting: what to do, when This is the part of the "long protocol" that's normally invisible — the @@ -526,6 +616,11 @@ whether the result keeps its sign. **Symptom:** `DriftReport.Drifted(...)` returns true, or `Report.Warnings` mentions a candidate that "drifted during the run." +This warning only appears when the trend is both significant and larger than +what the result can resolve (the noise floor, or half the interval, whichever +is larger). A long run finds shifts of a tenth of a percent significant, and a +warning that fires on nearly every run tells you nothing. + **What it means:** the machine changed behavior over the course of the run — usually getting slower, most often from thermal throttling as sustained load heats up the CPU. Because measurement order is interleaved by default (ABBA), @@ -542,6 +637,20 @@ beyond what the reported interval shows. result from it with a bit more skepticism than the headline confidence suggests. +### A warning says the machine was suspended + +**Symptom:** `Report.Suspended` is non-zero, or a warning says the machine +"appears to have been suspended." + +**What it means:** the wall clock moved further than the monotonic clock +during the comparison, which on Linux and macOS happens when the machine +sleeps. The samples straddle a pause, possibly a long one, after which caches, +clock speeds and everything else started cold. + +**What to do:** repeat the comparison with the machine kept awake — plugged +in, and with `caffeinate -i` on macOS, `systemd-inhibit --what=idle:sleep` on +Linux, or the power settings on Windows. + ### The autocorrelation is high **Symptom:** `HarnessValidation.Autocorrelation` (or `Report.Autocorrelation`) @@ -656,6 +765,10 @@ package-level sink variable. - **Autocorrelation** — how much each measurement resembles the one right before it. High values mean measurements aren't fully independent of one another. +- **Layout effect** — the part of a measured difference that comes from where + the data happens to lie in memory in this particular process, not from the + code. It is fixed for the life of a process and changes with the next one; + see [One process is one observation](#one-process-is-one-observation). - **Attenuation** — the true difference between two pieces of code getting diluted in the measured number because of fixed overhead (a loop, an accumulator) that both candidates carry equally. See diff --git a/README.md b/README.md index c4164d9..a6fb7b6 100644 --- a/README.md +++ b/README.md @@ -23,6 +23,7 @@ Keywords: benchmarking, performance, bootstrap, runtime comparison, statistics, - Measure the harness against itself, so that a result can be compared with the difference the same setup reports between two runs of identical code. - Detect a trend across a measurement run, which resampling structurally cannot see because it discards the order the samples arrived in. - Resample in blocks when the measurements are correlated enough that treating them as independent would overstate confidence. +- Run a comparison in several processes, each with its own heap layout, and pool the results into an interval that covers the scatter between processes (`multiproc`, `Combine`, `PerturbHeap`). - Deterministic PRNG for reproducible inputs, and a crypto/rand-backed one where unpredictability is wanted. ## What this cannot tell you @@ -31,6 +32,8 @@ Two limits are worth knowing before the first run, because neither is visible in **Attenuation.** What is measured is the loop, not the function. Whatever fixed per-operation cost the batch body carries — the loop itself, an accumulator, regenerating an input the candidate mutates — is present in both candidates and shrinks the difference between them. In a controlled experiment where the true difference was exactly 50%, the measured difference was 35%, because 1.81 ns/op of loop overhead sat on top of 2.13 ns/op of real work. Subtracting an empty-loop baseline does not repair it: the compiler optimizes an empty loop differently, and that correction recovered 2 of the 15 missing percentage points. Read a result as the speedup of the measured region, not of the isolated function. +**The process.** Every interval a single run reports covers the noise within that one process. A process fixes its memory layout for its whole life, and for data larger than the caches or full of pointers that layout alone has moved a difference by several points — 4 to 10 times the reported interval, and sometimes past zero. Where that applies, one process is one observation: run several with the `multiproc` package and read the pooled result. See [One process is one observation](HOWTO.md#one-process-is-one-observation). + **The noise floor.** Resampling quantifies how much an estimate would move if the same measurements were drawn again. It cannot see a bias that affected all of them equally, and will report a tight confidence around one. Measured on identical code, this package has seen apparent differences from a few tenths of a percent to well over one, carried with high confidence. `ValidateHarness` exists to measure that floor for your machine and your options; a result below it has resolved nothing, however confident the number looks. ## Install @@ -156,10 +159,18 @@ Checking the measurement itself: - `ValidateHarness(candidate, ValidationOptions)` — runs a candidate against itself and reports the noise floor, the tie rate, the drift rate and the autocorrelation. `Resolves(difference)` answers whether a result clears that floor. The floor is the 90th percentile of the differences observed on identical code, not their maximum, so that it converges as you validate longer instead of growing; roughly one A/A run in ten exceeds it. Validate both candidates and use the worse floor. - `ValidatePair(a, b, ValidationOptions)` — validates two candidates together, interleaved batch by batch, and returns one `HarnessValidation` each. Use it instead of two `ValidateHarness` calls before comparing the two: a candidate validated on its own last would start the comparison with a warm cache. - `DetectDrift(samples)` — tests a sample series for a trend across the run. +- `Report.Suspended` — how long the machine slept during a `Compare`, from the gap between the wall clock and the monotonic clock. +- `Report.LiveHeap` — the program's live data; above 16 MB, `Compare` warns that one process is not enough and points to `multiproc`. + +Across processes: + +- `multiproc.Main(Options, pairs...)` / `multiproc.RunTest(t, Options, pairs...)` — the whole job in one call: each `Pair` says how to build candidate A and B, and the program is re-executed as child processes, one at a time, each with its own heap layout and with the build order alternating, until every pooled interval is within ±2 points or ±10% of the difference (at least 5, at most 20 processes). `multiproc.Run(Options, suite)` does the same for a suite function of your own. +- `Combine(reports, level)` — pools per-process reports of one comparison into a `Pooled` result: a t interval over the per-process deltas, plus how far the processes scatter beyond their own intervals (`Inflation`, Cochran's Q, I²). +- `PerturbHeap(seed)` — allocates seeded filler in every small size class and one large block, so that data built afterwards lands at different addresses in each process. Primitives: -- `DPRNG` / `CPRNG` — deterministic and cryptographic generators with `Uint64`, `Float64` and `Uint32N`. +- `DPRNG` / `CPRNG` — deterministic and cryptographic generators with `Uint64`, `Float64` and `Uint32N`. `DPRNG.Shuffle` permutes, e.g. the order in which fixtures are built. - `SampleTime()` / `DiffTimeStamps()` — high-resolution timestamps, and `GetSampleTimePrecision()` for the smallest interval they can resolve here. - `Median` / `QuickMedian` / `Statistics` — small statistics helpers. diff --git a/cmd/rtcompare-aa/main.go b/cmd/rtcompare-aa/main.go index 7be54cf..298d4b3 100644 --- a/cmd/rtcompare-aa/main.go +++ b/cmd/rtcompare-aa/main.go @@ -10,8 +10,14 @@ // measurement (issue #111), which -mode prefix isolates and -mode compare // shows end to end; // - scatter between processes far beyond the interval a single process -// reports (issue #109), for which the command is meant to be run several -// times. +// reports (issue #109). Run -mode compare several times to see it, and add +// -multi to have the multiproc package run the processes, perturb each +// one's heap and pool their results. +// +// Two workloads are available. -workload chase loads a random node per +// operation and follows two random pointers from it. -workload scan walks 100 +// consecutive elements of a linked list whose nodes were allocated in random +// order, which is what a range scan over a pointer-based index does. // // Pick -n so that one instance is roughly 0.5 to 1 times the size of the // machine's last-level cache; a node is 64 bytes, so the default of 262,144 @@ -26,6 +32,7 @@ import ( "time" "github.com/TomTonic/rtcompare" + "github.com/TomTonic/rtcompare/multiproc" ) // node is one cache line: a value, a pointer to follow and padding. @@ -35,23 +42,143 @@ type node struct { _ [6]uint64 } -// fixture is a random pointer chase: nodes allocated in random order with -// random successors, and a fixed sequence of probes over them. +// fixture is a set of nodes allocated in random order, linked either at random +// or in key order, and a fixed sequence of probes into them. type fixture struct { nodes []*node probe []int32 } +// scanLength is how many consecutive elements one scan operation visits. +const scanLength = 100 + // sink keeps the chases observable to the compiler. var sink uint64 -// build allocates both fixtures from the same seed, so that they hold the same -// content, in the order the flag asks for. Whatever is allocated last is also -// what the caches hold when the measurement starts, which is one of the head -// starts this command is about. -func build(n int, order string) (a, b *fixture, err error) { - fa, fb := newFixture(n), newFixture(n) - perm, next, probe := layout(n) +type config struct { + n int + mode string + workload string + order string + seed uint64 + perturb bool + multi bool + minProcs int + maxProcs int + repeats int + validation int + warmup int + warmupDur time.Duration + skipVal bool +} + +func main() { + var c config + flag.IntVar(&c.n, "n", 1<<18, "nodes per instance (64 bytes each)") + flag.StringVar(&c.mode, "mode", "compare", "compare: Compare in both role assignments; prefix: Collect after validating none, A, B, or A then B") + flag.StringVar(&c.workload, "workload", "chase", "chase: random node and two random pointers per operation; scan: 100 consecutive list elements per operation") + flag.StringVar(&c.order, "build", "ab", "allocation order of the two fixtures: ab, ba, mixed, or random (ab or ba by seed); -multi alternates ab and ba by process") + flag.Uint64Var(&c.seed, "seed", 1, "seed for -perturb and -build random when not running under -multi") + flag.BoolVar(&c.perturb, "perturb", false, "perturb the heap layout from the seed before building the fixtures") + flag.BoolVar(&c.multi, "multi", false, "run -mode compare in several processes via the multiproc package, perturbing each one's heap and alternating the build order between processes") + flag.IntVar(&c.minProcs, "minprocs", 0, "multiproc.Options.MinProcesses (0: default)") + flag.IntVar(&c.maxProcs, "maxprocs", 0, "multiproc.Options.MaxProcesses (0: default)") + flag.IntVar(&c.repeats, "repeats", 0, "CollectOptions.Repeats (0: default)") + flag.IntVar(&c.validation, "validation", 0, "CompareOptions.ValidationRuns (0: default)") + flag.IntVar(&c.warmup, "warmup", 0, "CollectOptions.Warmup (0: default)") + flag.DurationVar(&c.warmupDur, "warmupdur", 0, "CollectOptions.WarmupDuration (0: default)") + flag.BoolVar(&c.skipVal, "skipvalidation", false, "CompareOptions.SkipValidation") + flag.Parse() + + var err error + if c.multi { + err = runMulti(c) + } else { + err = run(c) + } + if err != nil { + fmt.Fprintln(os.Stderr, "rtcompare-aa:", err) + os.Exit(1) + } + + // sink is only ever written, so that the compiler cannot drop the chases; + // in a package main the linter can see that nothing reads it. This read + // tells it what the writes already told the compiler, as in + // cmd/rtcompare-example. + _ = sink +} + +func run(c config) error { + if c.perturb { + defer rtcompare.PerturbHeap(c.seed).KeepAlive() + } + rng := rtcompare.NewDPRNG(c.seed | 1) + fa, fb, err := build(c, &rng) + if err != nil { + return err + } + switch c.mode { + case "compare": + return compare(c, fa, fb, func(name string, rep rtcompare.Report) { printReport(c, name, rep) }) + case "prefix": + return prefix(c, fa, fb) + default: + return fmt.Errorf("unknown mode %q, want compare or prefix", c.mode) + } +} + +// runMulti runs -mode compare in several processes. multiproc perturbs each +// child's heap from its own seed, so that the processes sample different +// layouts instead of repeating one, and the build order alternates between +// processes. It uses Run with its own suite rather than Pairs, because it +// builds the two fixtures once and compares them in both role assignments. +// Alternating rather than drawing it at random matters: whichever fixture is +// built second is consistently a few percent faster here, and a random draw +// over a handful of processes is rarely balanced. +func runMulti(c config) error { + res, err := multiproc.Run(multiproc.Options{ + MinProcesses: c.minProcs, + MaxProcesses: c.maxProcs, + Stdout: os.Stdout, + Progress: func(r multiproc.Results) { + fmt.Fprintf(os.Stderr, "process %d done\n", r.Processes) + }, + }, func(p *multiproc.Process) error { + // multiproc has already perturbed this process's heap. + c.order = [2]string{"ab", "ba"}[p.Index%2] + fa, fb, err := build(c, p.Rand()) + if err != nil { + return err + } + fmt.Printf("process %d, seed %d:\n", p.Index, p.Seed) + return compare(c, fa, fb, func(name string, rep rtcompare.Report) { + printReport(c, name, rep) + p.Record(name, rep) + }) + }) + if err != nil || res.Child { + return err + } + fmt.Printf("\nn=%d workload=%s, pooled:\n%s\n", c.n, c.workload, res) + return nil +} + +// build allocates both fixtures with the same content, in the order the +// configuration asks for. Whatever is allocated last is also what the caches +// hold when the measurement starts, and it lands in different memory. +func build(c config, rng *rtcompare.DPRNG) (a, b *fixture, err error) { + if c.workload != "chase" && c.workload != "scan" { + return nil, nil, fmt.Errorf("unknown workload %q, want chase or scan", c.workload) + } + if c.workload == "scan" && c.n <= scanLength { + return nil, nil, fmt.Errorf("-workload scan needs more than %d nodes, got %d", scanLength, c.n) + } + fa, fb := newFixture(c.n), newFixture(c.n) + perm, next, probe := layout(c.n, c.workload) + order := c.order + if order == "random" { + order = [2]string{"ab", "ba"}[rng.Uint32N(2)] + } switch order { case "ab": fa.fill(perm, next, probe) @@ -62,7 +189,7 @@ func build(n int, order string) (a, b *fixture, err error) { case "mixed": fillMixed(fa, fb, perm, next, probe) default: - return nil, nil, fmt.Errorf("unknown build order %q, want ab, ba or mixed", order) + return nil, nil, fmt.Errorf("unknown build order %q, want ab, ba, mixed or random", c.order) } return fa, fb, nil } @@ -72,22 +199,31 @@ func newFixture(n int) *fixture { } // layout draws the shared content of both fixtures: the order nodes are -// allocated in, each node's successor and the probe sequence. -func layout(n int) (perm, next []int, probe []int32) { +// allocated in, each node's successor and the probe sequence. It is seeded +// with a constant, so the content is the same in every process and only the +// layout differs. +func layout(n int, workload string) (perm, next []int, probe []int32) { rng := rtcompare.NewDPRNG(1) perm = make([]int, n) for i := range perm { perm[i] = i } - for i := n - 1; i > 0; i-- { - j := int(rng.Uint32N(uint32(i + 1))) - perm[i], perm[j] = perm[j], perm[i] - } + rng.Shuffle(n, func(i, j int) { perm[i], perm[j] = perm[j], perm[i] }) next = make([]int, n) + probe = make([]int32, 1<<20) + if workload == "scan" { + // Linked in key order, with scans that never run off the end. + for i := range next { + next[i] = (i + 1) % n + } + for i := range probe { + probe[i] = int32(rng.Uint32N(uint32(n - scanLength))) + } + return perm, next, probe + } for i := range next { next[i] = int(rng.Uint32N(uint32(n))) } - probe = make([]int32, 1<<20) even := uint32(n &^ 1) // probe^1 has to stay in range for i := range probe { probe[i] = int32(rng.Uint32N(even)) @@ -120,11 +256,18 @@ func fillMixed(a, b *fixture, perm, next []int, probe []int32) { b.link(next, probe) } -// candidate loads a node per operation and follows two pointers from it. Each -// operation depends on the one before, like a lookup whose branches depend on -// the data it loads, so every cache miss is paid in full. Every candidate keeps -// its own cursor into the probe sequence. -func candidate(name string, f *fixture) rtcompare.Candidate { +func candidate(workload, name string, f *fixture) rtcompare.Candidate { + if workload == "scan" { + return scanCandidate(name, f) + } + return chaseCandidate(name, f) +} + +// chaseCandidate loads a node per operation and follows two pointers from it. +// Each operation depends on the one before, like a lookup whose branches +// depend on the data it loads, so every cache miss is paid in full. Every +// candidate keeps its own cursor into the probe sequence. +func chaseCandidate(name string, f *fixture) rtcompare.Candidate { j := 0 return rtcompare.Candidate{Name: name, Batch: func(n uint64) { var acc uint64 @@ -139,63 +282,34 @@ func candidate(name string, f *fixture) rtcompare.Candidate { }} } -type config struct { - n int - mode string - order string - repeats int - validation int - warmup int - warmupDur time.Duration - skipVal bool -} - -func main() { - var c config - flag.IntVar(&c.n, "n", 1<<18, "nodes per instance (64 bytes each)") - flag.StringVar(&c.mode, "mode", "compare", "compare: Compare in both role assignments; prefix: Collect after validating none, A, B, or A then B") - flag.StringVar(&c.order, "build", "ab", "allocation order of the two fixtures: ab, ba or mixed") - flag.IntVar(&c.repeats, "repeats", 0, "CollectOptions.Repeats (0: default)") - flag.IntVar(&c.validation, "validation", 0, "CompareOptions.ValidationRuns (0: default)") - flag.IntVar(&c.warmup, "warmup", 0, "CollectOptions.Warmup (0: default)") - flag.DurationVar(&c.warmupDur, "warmupdur", 0, "CollectOptions.WarmupDuration (0: default)") - flag.BoolVar(&c.skipVal, "skipvalidation", false, "CompareOptions.SkipValidation") - flag.Parse() - - if err := run(c); err != nil { - fmt.Fprintln(os.Stderr, "rtcompare-aa:", err) - os.Exit(1) - } - - // sink is only ever written, so that the compiler cannot drop the chases; - // in a package main the linter can see that nothing reads it. This read - // tells it what the writes already told the compiler, as in - // cmd/rtcompare-example. - _ = sink -} - -func run(c config) error { - fa, fb, err := build(c.n, c.order) - if err != nil { - return err - } - switch c.mode { - case "compare": - return compare(c, fa, fb) - case "prefix": - return prefix(c, fa, fb) - default: - return fmt.Errorf("unknown mode %q, want compare or prefix", c.mode) - } +// scanCandidate walks scanLength consecutive elements per operation, starting +// at a random one: a range scan over a list whose nodes are scattered in +// memory. +func scanCandidate(name string, f *fixture) rtcompare.Candidate { + j := 0 + return rtcompare.Candidate{Name: name, Batch: func(n uint64) { + var acc uint64 + for range n { + x := f.nodes[f.probe[j]] + for range scanLength { + acc += x.val + x = x.next + } + if j++; j == len(f.probe) { + j = 0 + } + } + sink += acc + }} } func collectOptions(c config) rtcompare.CollectOptions { return rtcompare.CollectOptions{MaxQuantizationError: 1e-4, Repeats: c.repeats, Warmup: c.warmup, WarmupDuration: c.warmupDur} } -// compare runs Compare with each fixture in each role. Delta is 1 - A/B, so a -// negative delta means A measured slower. -func compare(c config, fa, fb *fixture) error { +// compare runs Compare with each fixture in each role and hands every report +// to report. Delta is 1 - A/B, so a negative delta means A measured slower. +func compare(c config, fa, fb *fixture, report func(name string, rep rtcompare.Report)) error { opt := rtcompare.CompareOptions{ Collect: collectOptions(c), ValidationRuns: c.validation, @@ -205,23 +319,27 @@ func compare(c config, fa, fb *fixture) error { name string a, b *fixture }{{"A=first B=second", fa, fb}, {"A=second B=first", fb, fa}} { - rep, err := rtcompare.Compare(candidate("A", roles.a), candidate("B", roles.b), opt) + rep, err := rtcompare.Compare(candidate(c.workload, "A", roles.a), candidate(c.workload, "B", roles.b), opt) if err != nil { return err } - fmt.Printf("n=%d build=%s %s: A %.2f ns, B %.2f ns, delta %+.2f%% [%+.2f, %+.2f] floor %.2f%% resolved=%v\n", - c.n, c.order, roles.name, rep.NsPerOpA, rep.NsPerOpB, - 100*rep.Estimate.Delta, 100*rep.Estimate.Low, 100*rep.Estimate.High, 100*rep.NoiseFloor, rep.Resolved) - fmt.Printf(" drift A %+.1f%% (p %.3f), B %+.1f%% (p %.3f), B/A %+.1f%% (p %.3f)\n", - 100*rep.DriftA.RelativeShift, rep.DriftA.PValue, 100*rep.DriftB.RelativeShift, rep.DriftB.PValue, - 100*rep.DriftRatio.RelativeShift, rep.DriftRatio.PValue) - for _, w := range rep.Warnings { - fmt.Println(" warning:", w) - } + report(roles.name, rep) } return nil } +func printReport(c config, name string, rep rtcompare.Report) { + fmt.Printf("n=%d workload=%s build=%s %s: A %.2f ns, B %.2f ns, delta %+.2f%% [%+.2f, %+.2f] floor %.2f%% resolved=%v\n", + c.n, c.workload, c.order, name, rep.NsPerOpA, rep.NsPerOpB, + 100*rep.Estimate.Delta, 100*rep.Estimate.Low, 100*rep.Estimate.High, 100*rep.NoiseFloor, rep.Resolved) + fmt.Printf(" drift A %+.1f%% (p %.3f), B %+.1f%% (p %.3f), B/A %+.1f%% (p %.3f)\n", + 100*rep.DriftA.RelativeShift, rep.DriftA.PValue, 100*rep.DriftB.RelativeShift, rep.DriftB.PValue, + 100*rep.DriftRatio.RelativeShift, rep.DriftRatio.PValue) + for _, w := range rep.Warnings { + fmt.Println(" warning:", w) + } +} + // prefix runs Collect at a fixed, calibrated batch size after validating no // candidate, only A, only B, or A and then B, which is what Compare did before // v0.7.0, and prints the ratio of the medians B/A. 1.00 is the truth. The @@ -230,14 +348,14 @@ func compare(c config, fa, fb *fixture) error { // Collect's warm-up washes out on its own. func prefix(c config, fa, fb *fixture) error { co := collectOptions(c) - cal, err := rtcompare.CalibrateInnerLoops(candidate("A", fa), rtcompare.CalibrationOptions{MaxQuantizationError: co.MaxQuantizationError}) + cal, err := rtcompare.CalibrateInnerLoops(candidate(c.workload, "A", fa), rtcompare.CalibrationOptions{MaxQuantizationError: co.MaxQuantizationError}) if err != nil { return err } co.InnerLoops = cal.InnerLoops vo := rtcompare.ValidationOptions{Collect: co, Runs: c.validation} for _, pre := range []string{"", "A", "B", "AB"} { - a, b := candidate("A", fa), candidate("B", fb) + a, b := candidate(c.workload, "A", fa), candidate(c.workload, "B", fb) for _, who := range pre { switch who { case 'A': diff --git a/combine.go b/combine.go new file mode 100644 index 0000000..994b633 --- /dev/null +++ b/combine.go @@ -0,0 +1,367 @@ +package rtcompare + +// This file combines the same comparison run in several processes, which is +// the only way to put an interval on a difference that includes the scatter +// between processes, and not just the noise within one. + +import ( + "fmt" + "math" + "slices" + "strings" + "time" +) + +// Pooled is one comparison combined across several processes by [Combine]. +// +// It treats each process as a single observation of the difference. Its +// interval therefore covers what a single [Report]'s interval cannot: that each +// process lays its data out in memory differently, and that the layout can move +// the difference by more than the noise within any one of them. +type Pooled struct { + // Processes is the number of per-process results that were combined. + Processes int + + // Delta is the mean of the per-process deltas, each 1 - median(A)/median(B). + // Positive means A is faster. The mean is unweighted on purpose: the + // per-process intervals understate each process's real uncertainty, which is + // the problem this type exists to solve, so weighting by them would give the + // processes that happened to be least noisy inside the most say. + Delta float64 + + // Low and High bound Delta at Level: a Student t interval over the + // per-process deltas with Processes-1 degrees of freedom. With few processes + // it is wide, and honestly so. + Low, High float64 + + // Level is the coverage level of the interval. + Level float64 + + // SpreadBetween is the standard deviation of the per-process deltas. + SpreadBetween float64 + + // SpreadWithin is the median standard error that a single process's own + // interval implies, (High-Low)/(2z) at that interval's level. + SpreadWithin float64 + + // Inflation is SpreadBetween/SpreadWithin: how many times more the + // processes scatter than one process's interval says they should. Near one, + // a single process could be trusted on its own; measured on large, + // pointer-heavy data it has been 5 and more. Zero when SpreadWithin is zero. + Inflation float64 + + // Q is Cochran's heterogeneity statistic, the sum of the squared deviations + // of the per-process deltas from their inverse-variance mean, in units of + // each process's own standard error. Without variance between processes it + // is about Processes-1. NaN when a process reported an interval of zero + // width, since its standard error is then unknown. + Q float64 + + // I2 is the share of the scatter between processes that their own intervals + // do not explain, max(0, (Q-(k-1))/Q), in [0,1]. Above about 0.5 the + // variance between processes dominates. NaN when Q is. + I2 float64 + + // Validated reports whether every combined process ran its A/A validation. + Validated bool + + // NoiseFloor is the median of the validated processes' noise floors, zero + // when none was validated. Resolved requires Delta to clear it, as for a + // single Report. + NoiseFloor float64 + + // Resolved is the short answer: the pooled interval excludes zero and Delta + // clears the noise floor. + Resolved bool + + // Warnings lists everything that undermines the pooled result, in plain + // sentences. + Warnings []string +} + +// Precise reports whether the interval is narrow enough for the question: its +// half-width is at most abs, or at most rel times |Delta|. The multiproc +// package uses it as its rule for when to stop starting processes, with +// defaults of 0.02 and 0.10. +func (p Pooled) Precise(abs, rel float64) bool { + half := (p.High - p.Low) / 2 + return half <= abs || half <= rel*math.Abs(p.Delta) +} + +// String renders the pooled result as a short multi-line summary. +func (p Pooled) String() string { + var b strings.Builder + fmt.Fprintf(&b, "difference %+.2f%% [%+.2f%%, %+.2f%%] at %.0f%% confidence, pooled over %d processes\n", + p.Delta*100, p.Low*100, p.High*100, p.Level*100, p.Processes) + fmt.Fprintf(&b, "scatter between processes %.3f%%, within one %.3f%% (%.1fx), I² %.2f\n", + p.SpreadBetween*100, p.SpreadWithin*100, p.Inflation, p.I2) + switch { + case p.Resolved && p.Delta > 0: + b.WriteString("resolved: A is faster than B\n") + case p.Resolved: + b.WriteString("resolved: A is slower than B\n") + default: + b.WriteString("not resolved: these processes did not establish a difference\n") + } + for _, w := range p.Warnings { + fmt.Fprintf(&b, " warning: %s\n", w) + } + return strings.TrimRight(b.String(), "\n") +} + +// Combine pools the reports of the same comparison, run once in each of +// several processes, into one estimate whose interval includes the scatter +// between processes. +// +// Parameters: reports holds one [Report] per process, all of the same +// comparison, with the same candidates in the same roles; at least three are +// required, and five or more are advisable. level is the coverage level of the +// pooled interval, zero selecting [DefaultConfidenceLevel]. +// +// It returns the [Pooled] result, or an error if there are fewer than three +// reports, if level is not strictly between zero and one, or if a report holds +// a non-finite estimate. +// +// Use it whenever the data of a comparison does not fit in the caches or is +// full of pointers. A single process reports an interval that covers only the +// noise within that process. The process's memory layout is fixed for its +// whole lifetime, and it can move the difference by several points, so the +// next process reports a different, equally narrow interval somewhere else. +// Observed on 1M-key trees, the processes scattered 4 to 10 times more widely +// than their own intervals implied, and two processes resolved the same +// comparison with opposite signs. Combine treats each process as one +// observation: Delta is the mean of the per-process deltas and the interval is +// a t interval over them, which is honest for few processes where random-effects +// weighting is known to be too narrow. Q, I2 and Inflation say how much the +// processes disagreed beyond their own intervals. +// +// Pooling only helps if the processes sample different layouts. A process that +// allocates the same things in the same order gets much the same layout every +// time, and its bias repeats rather than averages out; see [PerturbHeap]. The +// multiproc package starts the processes, perturbs each one's heap and calls +// Combine for you. +// +// var reports []rtcompare.Report // one per process, e.g. read back from files +// pooled, err := rtcompare.Combine(reports, 0) +// if err != nil { ... } +// fmt.Println(pooled) +func Combine(reports []Report, level float64) (Pooled, error) { + k := len(reports) + if k < 3 { + return Pooled{}, fmt.Errorf("rtcompare: need at least 3 per-process reports to combine, got %d", k) + } + if level == 0 { + level = DefaultConfidenceLevel + } + if math.IsNaN(level) || level <= 0 || level >= 1 { + return Pooled{}, fmt.Errorf("rtcompare: level must be strictly between 0 and 1, got %v", level) + } + deltas := make([]float64, k) + ses := make([]float64, k) + for i, r := range reports { + e := r.Estimate + if !finite(e.Delta) || !finite(e.Low) || !finite(e.High) { + return Pooled{}, fmt.Errorf("rtcompare: report %d holds a non-finite estimate %s", i, e) + } + deltas[i] = e.Delta + ses[i] = standardErrorOf(e) + } + + mean, _, sd := Statistics(deltas) + half := studentTQuantile((1+level)/2, float64(k-1)) * sd / math.Sqrt(float64(k)) + p := Pooled{ + Processes: k, + Delta: mean, + Low: mean - half, + High: mean + half, + Level: level, + SpreadBetween: sd, + SpreadWithin: Median(slices.Clone(ses)), + } + if p.SpreadWithin > 0 { + p.Inflation = p.SpreadBetween / p.SpreadWithin + } + p.Q, p.I2 = heterogeneity(deltas, ses) + p.Validated, p.NoiseFloor = pooledFloor(reports) + p.Resolved = (p.Low > 0 || p.High < 0) && math.Abs(p.Delta) > p.NoiseFloor + p.Warnings = p.warnings(reports) + return p, nil +} + +func finite(x float64) bool { return !math.IsNaN(x) && !math.IsInf(x, 0) } + +// standardErrorOf recovers the standard error a single process's interval +// implies, treating it as a normal interval at its own level. A report that +// carries no level is read at the default. +func standardErrorOf(e Estimate) float64 { + level := e.Level + if !(level > 0 && level < 1) { + level = DefaultConfidenceLevel + } + z := math.Sqrt2 * math.Erfinv(level) + return (e.High - e.Low) / (2 * z) +} + +// heterogeneity computes Cochran's Q and I² from the per-process deltas and +// their standard errors. Both are NaN when a standard error is zero, because +// the weight of that process is then undefined. +func heterogeneity(deltas, ses []float64) (q, i2 float64) { + var sw, swd float64 + for i, se := range ses { + if !(se > 0) { + return math.NaN(), math.NaN() + } + w := 1 / (se * se) + sw += w + swd += w * deltas[i] + } + fixed := swd / sw + for i, se := range ses { + d := deltas[i] - fixed + q += d * d / (se * se) + } + if q > 0 { + i2 = math.Max(0, (q-float64(len(deltas)-1))/q) + } + return q, i2 +} + +// pooledFloor returns whether every report was validated and the median noise +// floor of those that were. +func pooledFloor(reports []Report) (all bool, floor float64) { + var floors []float64 + for _, r := range reports { + if r.Validated { + floors = append(floors, r.NoiseFloor) + } + } + return len(floors) == len(reports), medianOrZero(floors) +} + +// warnings lists the things that undermine a pooled result. +func (p Pooled) warnings(reports []Report) []string { + var w []string + if p.Processes < 5 { + w = append(w, fmt.Sprintf( + "only %d processes were combined; below five the interval is wide and the heterogeneity figures are rough", p.Processes)) + } + if !p.Validated { + w = append(w, "not every process ran its A/A validation, so the noise floor is taken from those that did, or is unknown") + } else if math.Abs(p.Delta) <= p.NoiseFloor { + w = append(w, fmt.Sprintf( + "the pooled difference of %.2f%% does not clear the %.2f%% median noise floor of the processes", + p.Delta*100, p.NoiseFloor*100)) + } + if !(p.Low > 0 || p.High < 0) { + w = append(w, fmt.Sprintf( + "the pooled interval [%.2f%%, %.2f%%] includes zero", p.Low*100, p.High*100)) + } + if p.Inflation > 2 { + w = append(w, fmt.Sprintf( + "the processes scatter %.1f times as widely as one process's interval implies, so a single process's result for this comparison is not to be trusted on its own", + p.Inflation)) + } + faster, slower := 0, 0 + for _, r := range reports { + switch { + case r.Resolved && r.Estimate.Delta > 0: + faster++ + case r.Resolved && r.Estimate.Delta < 0: + slower++ + } + } + if faster > 0 && slower > 0 { + w = append(w, fmt.Sprintf( + "%d processes resolved A as faster and %d as slower; each was confident, and the disagreement between them is the layout of memory, not the code", + faster, slower)) + } + for i, r := range reports { + if r.Suspended > 0 { + w = append(w, fmt.Sprintf("process %d was suspended for %s during its comparison", i, r.Suspended.Round(time.Second))) + } + } + return w +} + +// studentTQuantile returns the p-quantile of Student's t distribution with df +// degrees of freedom, by bisection on its distribution function. It exists so +// that Combine needs no dependency for the one quantile it uses. +func studentTQuantile(p, df float64) float64 { + if p == 0.5 { + return 0 + } + if p < 0.5 { + return -studentTQuantile(1-p, df) + } + lo, hi := 0.0, 1.0 + for studentTCDF(hi, df) < p { + lo, hi = hi, 2*hi + } + for range 200 { + mid := (lo + hi) / 2 + if studentTCDF(mid, df) < p { + lo = mid + } else { + hi = mid + } + } + return (lo + hi) / 2 +} + +// studentTCDF is the distribution function of Student's t for t >= 0, written +// through the regularized incomplete beta function. +func studentTCDF(t, df float64) float64 { + return 1 - 0.5*regularizedIncompleteBeta(df/(df+t*t), df/2, 0.5) +} + +// regularizedIncompleteBeta is I_x(a, b), evaluated by its continued fraction +// on whichever side of the mean converges fastest. +func regularizedIncompleteBeta(x, a, b float64) float64 { + if x <= 0 { + return 0 + } + if x >= 1 { + return 1 + } + la, _ := math.Lgamma(a) + lb, _ := math.Lgamma(b) + lab, _ := math.Lgamma(a + b) + front := math.Exp(a*math.Log(x) + b*math.Log1p(-x) + lab - la - lb) + if x < (a+1)/(a+b+2) { + return front * betaContinuedFraction(x, a, b) / a + } + return 1 - front*betaContinuedFraction(1-x, b, a)/b +} + +// betaContinuedFraction evaluates the continued fraction of the incomplete +// beta function by the modified Lentz method. +func betaContinuedFraction(x, a, b float64) float64 { + const ( + maxIter = 300 + eps = 1e-15 + tiny = 1e-300 + ) + guard := func(v float64) float64 { + if math.Abs(v) < tiny { + return tiny + } + return v + } + c, d := 1.0, 1/guard(1-(a+b)*x/(a+1)) + h := d + for m := 1.0; m <= maxIter; m++ { + even := m * (b - m) * x / ((a + 2*m - 1) * (a + 2*m)) + d = 1 / guard(1+even*d) + c = guard(1 + even/c) + h *= d * c + odd := -(a + m) * (a + b + m) * x / ((a + 2*m) * (a + 2*m + 1)) + d = 1 / guard(1+odd*d) + c = guard(1 + odd/c) + step := d * c + h *= step + if math.Abs(step-1) < eps { + break + } + } + return h +} diff --git a/combine_test.go b/combine_test.go new file mode 100644 index 0000000..071f3f9 --- /dev/null +++ b/combine_test.go @@ -0,0 +1,203 @@ +package rtcompare + +import ( + "math" + "strings" + "testing" + "time" +) + +// perProcess builds a validated report whose estimate is delta with a +// symmetric 95% interval of the given half-width, as one process would report. +func perProcess(delta, half float64) Report { + return Report{ + Estimate: Estimate{Delta: delta, Low: delta - half, High: delta + half, Level: 0.95}, + Validated: true, + NoiseFloor: 0.005, + Resolved: math.Abs(delta) > half && math.Abs(delta) > 0.005, + } +} + +// TestStudentTQuantile checks the one distribution quantile the pooled +// interval depends on. Combine computes it without a statistics dependency, so +// it is compared against published two-sided 95% and 90% critical values of +// Student's t, which it has to reproduce to five decimals. +func TestStudentTQuantile(t *testing.T) { + cases := []struct { + p, df, want float64 + }{ + {0.975, 1, 12.706205}, + {0.975, 2, 4.302653}, + {0.975, 4, 2.776445}, + {0.975, 9, 2.262157}, + {0.975, 19, 2.093024}, + {0.975, 100, 1.983972}, + {0.95, 4, 2.131847}, + {0.025, 4, -2.776445}, + {0.5, 7, 0}, + } + for _, c := range cases { + if got := studentTQuantile(c.p, c.df); math.Abs(got-c.want) > 1e-5 { + t.Errorf("t quantile p=%v df=%v: got %.6f, want %.6f", c.p, c.df, got, c.want) + } + } +} + +// TestCombinePoolsProcessesAsObservations checks the central promise of +// pooling: the interval reflects how much the processes actually disagree, not +// how confident each of them was. It belongs to Combine, which turns several +// single-process reports of the same comparison into one estimate. Five +// processes with tight intervals that scatter widely have to produce the mean +// of their deltas, a t interval over them, an inflation factor well above one +// and a warning that says so. +func TestCombinePoolsProcessesAsObservations(t *testing.T) { + deltas := []float64{0.10, 0.14, 0.12, 0.25, 0.20} + var reports []Report + for _, d := range deltas { + reports = append(reports, perProcess(d, 0.01)) + } + p, err := Combine(reports, 0) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + t.Log("\n" + p.String()) + + mean, _, sd := Statistics(deltas) + half := 2.776445 * sd / math.Sqrt(5) + if math.Abs(p.Delta-mean) > 1e-12 || math.Abs(p.Low-(mean-half)) > 1e-5 || math.Abs(p.High-(mean+half)) > 1e-5 { + t.Errorf("pooled %+.4f [%+.4f, %+.4f], want %+.4f [%+.4f, %+.4f]", p.Delta, p.Low, p.High, mean, mean-half, mean+half) + } + wantWithin := 0.01 / 1.959964 + if math.Abs(p.SpreadWithin-wantWithin) > 1e-6 { + t.Errorf("within-process standard error: got %v, want %v", p.SpreadWithin, wantWithin) + } + if p.Inflation < 5 { + t.Errorf("inflation %v, want well above one for processes this far apart", p.Inflation) + } + if !(p.I2 > 0.9 && p.I2 <= 1) { + t.Errorf("I² %v, want close to one", p.I2) + } + if !p.Resolved { + t.Errorf("a pooled interval well above zero should resolve") + } + if got := strings.Join(p.Warnings, "\n"); !strings.Contains(got, "times as widely") { + t.Errorf("warnings do not mention the inflation:\n%s", got) + } +} + +// TestCombineAgreeingProcessesShowNoHeterogeneity checks the reassuring case: +// processes that agree within their own intervals. Combine is expected to +// report an inflation near one, an I² of zero and no warning about scatter, so +// that a clean result reads clean. +func TestCombineAgreeingProcessesShowNoHeterogeneity(t *testing.T) { + reports := []Report{ + perProcess(0.050, 0.02), perProcess(0.052, 0.02), perProcess(0.049, 0.02), + perProcess(0.051, 0.02), perProcess(0.048, 0.02), + } + p, err := Combine(reports, 0.95) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if p.I2 != 0 { + t.Errorf("I² %v, want 0 for processes that agree within their intervals", p.I2) + } + if p.Inflation > 1 { + t.Errorf("inflation %v, want below one", p.Inflation) + } + if got := strings.Join(p.Warnings, "\n"); strings.Contains(got, "times as widely") { + t.Errorf("unexpected scatter warning:\n%s", got) + } +} + +// TestCombineWarnings checks the plain-language fine print of a pooled result. +// Each case sets up one condition that undermines it and expects a sentence +// that names it. +func TestCombineWarnings(t *testing.T) { + opposite := []Report{perProcess(0.10, 0.01), perProcess(-0.10, 0.01), perProcess(0.02, 0.01)} + unvalidated := []Report{perProcess(0.1, 0.01), perProcess(0.1, 0.01), {Estimate: Estimate{Delta: 0.1, Low: 0.09, High: 0.11, Level: 0.95}}} + suspended := []Report{perProcess(0.1, 0.01), perProcess(0.1, 0.01), perProcess(0.1, 0.01)} + suspended[1].Suspended = 20 * time.Minute + cases := []struct { + name string + reports []Report + want string + }{ + {"warns about too few processes", opposite, "only 3 processes"}, + {"warns about processes resolving in opposite directions", opposite, "2 processes resolved A as faster and 1 as slower"}, + {"warns about an interval that includes zero", opposite, "includes zero"}, + {"warns about processes without validation", unvalidated, "not every process"}, + {"warns about a suspended process", suspended, "process 1 was suspended for 20m0s"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + p, err := Combine(c.reports, 0) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if got := strings.Join(p.Warnings, "\n"); !strings.Contains(got, c.want) { + t.Errorf("warnings do not mention %q:\n%s", c.want, got) + } + }) + } +} + +// TestCombineRejectsBadInput checks that pooling fails loudly on input it +// cannot honestly pool: too few processes for a t interval to mean anything, a +// level outside (0,1), and an estimate that is not a number. +func TestCombineRejectsBadInput(t *testing.T) { + ok := perProcess(0.1, 0.01) + nan := perProcess(0.1, 0.01) + nan.Estimate.Delta = math.NaN() + cases := []struct { + name string + reports []Report + level float64 + want string + }{ + {"returns error for two processes", []Report{ok, ok}, 0, "at least 3"}, + {"returns error for a level of one", []Report{ok, ok, ok}, 1, "level"}, + {"returns error for a NaN level", []Report{ok, ok, ok}, math.NaN(), "level"}, + {"returns error for a non-finite estimate", []Report{ok, nan, ok}, 0, "report 1"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + _, err := Combine(c.reports, c.level) + if err == nil || !strings.Contains(err.Error(), c.want) { + t.Errorf("got %v, want an error mentioning %q", err, c.want) + } + }) + } +} + +// TestHeterogeneityWithoutStandardErrors checks that Q and I² are reported as +// unknown, not as zero, when a process gave an interval of zero width. Its +// weight in the heterogeneity statistic would be infinite, so any number would +// be made up. +func TestHeterogeneityWithoutStandardErrors(t *testing.T) { + q, i2 := heterogeneity([]float64{0.1, 0.2, 0.3}, []float64{0.01, 0, 0.01}) + if !math.IsNaN(q) || !math.IsNaN(i2) { + t.Errorf("got Q %v, I² %v, want NaN for both", q, i2) + } +} + +// TestPooledPrecise checks the stop rule the multi-process driver uses: a +// pooled interval is precise enough when its half-width is within an absolute +// bound or within a share of the difference itself, whichever is looser. +func TestPooledPrecise(t *testing.T) { + cases := []struct { + name string + p Pooled + want bool + }{ + {"is precise within two points", Pooled{Delta: 0.01, Low: -0.005, High: 0.025}, true}, + {"is precise within ten percent of a large difference", Pooled{Delta: 0.5, Low: 0.46, High: 0.54}, true}, + {"is not precise when wide and small", Pooled{Delta: 0.05, Low: 0.0, High: 0.1}, false}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if got := c.p.Precise(0.02, 0.10); got != c.want { + t.Errorf("Precise: got %v, want %v", got, c.want) + } + }) + } +} diff --git a/compare.go b/compare.go index c50aa0d..271ac1f 100644 --- a/compare.go +++ b/compare.go @@ -3,7 +3,9 @@ package rtcompare import ( "fmt" "math" + "runtime/metrics" "strings" + "time" ) // AutocorrelationThreshold is the lag-1 autocorrelation above which [Compare] @@ -17,6 +19,29 @@ import ( // tables. const AutocorrelationThreshold = 0.2 +// SuspendThreshold is how far the wall clock may run ahead of the monotonic +// clock during a [Compare] before the machine is taken to have been suspended. +// +// On Linux and macOS the monotonic clock stops while the machine sleeps and the +// wall clock does not, so the gap between them is the time spent asleep. A +// second is far above what clock adjustments produce in the span of one +// comparison and far below the idle sleeps that have actually paused +// comparisons, which lasted 15 to 25 minutes. +const SuspendThreshold = time.Second + +// LargeHeapThreshold is the live heap above which [Compare] warns that a +// single process's result is not to be trusted on its own. +// +// Above roughly the size of a last-level cache, where the data lies in memory +// starts to decide how fast it is, and that placement is fixed for a process +// and different in the next. Measured on 1M-node lists, processes scattered 4.5 +// times as widely as their own intervals said. The threshold sits at the lower +// end of common last-level caches, so that the warning comes early rather than +// late; it is a hint drawn from the heap as a whole, not a measurement of what +// the candidates touch, and can fire for a program that holds much memory the +// candidates never read. +const LargeHeapThreshold = 16 << 20 + // CompareOptions configures [Compare]. The zero value is usable and selects the // documented defaults throughout. type CompareOptions struct { @@ -70,6 +95,13 @@ type Report struct { // Estimate is how much smaller A is than B, with an interval around it. // Positive means A is faster. + // + // The interval covers the noise within this one process and nothing else. + // Each process gets its own memory layout, and for data larger than the + // caches or full of pointers that layout alone can shift the difference by + // several points, far beyond this interval, and differently in the next + // process. Where that matters, run the comparison in several processes and + // read [Combine]'s pooled interval instead; see the multiproc package. Estimate Estimate // Confidence maps each requested threshold to the confidence that it is @@ -100,6 +132,20 @@ type Report struct { // means the test could not be run. DriftA, DriftB DriftReport + // LiveHeap is the live heap in bytes as of the last garbage collection + // during the comparison, zero if the runtime did not report it. Above + // [LargeHeapThreshold] the warnings recommend running the comparison in + // several processes. + LiveHeap uint64 + + // Suspended is how much longer the wall clock ran than the monotonic clock + // during the comparison, when that exceeded [SuspendThreshold], and zero + // otherwise. Non-zero means the machine was asleep for about that long in + // the middle of the measurement, which neither the samples nor the drift + // tests can show reliably. See SuspendThreshold for the platforms this + // works on. + Suspended time.Duration + // DriftRatio tests the ratio B/A of each pair of neighbouring batches for a // trend across the run. A zero N means the test could not be run. // @@ -242,6 +288,7 @@ func sortedKeys(m map[float64]float64) []float64 { // fails, if any threshold is NaN, or if the underlying measurement or // validation fails. See [Collect] for the option-validation errors. func Compare(a, b Candidate, opt CompareOptions) (Report, error) { + start := time.Now() if a.Batch == nil { return Report{}, fmt.Errorf("rtcompare: candidate %s has a nil Batch function", a.label("A")) } @@ -336,10 +383,41 @@ func Compare(a, b Candidate, opt CompareOptions) (Report, error) { // zero, and it must be larger than what this setup invents from identical // code. Neither implies the other. r.Resolved = est.Excludes(0) && math.Abs(est.Delta) > r.NoiseFloor + r.Suspended = suspendedSince(start) + r.LiveHeap = liveHeap() r.Warnings = r.warnings() return r, nil } +// suspendedSince returns how much further the wall clock has moved than the +// monotonic clock since start, if that exceeds SuspendThreshold. Round(0) +// strips the monotonic reading, which leaves the two subtractions measuring +// the same interval on the two clocks. +func suspendedSince(start time.Time) time.Duration { + return suspendGap(time.Now().Round(0).Sub(start.Round(0)), time.Since(start)) +} + +// liveHeap reads the live heap as of the last garbage collection. It is used +// rather than the heap in use, which includes garbage not yet collected, and +// rather than forcing a collection, which would be a side effect of Compare. +func liveHeap() uint64 { + sample := []metrics.Sample{{Name: "/gc/heap/live:bytes"}} + metrics.Read(sample) + if sample[0].Value.Kind() != metrics.KindUint64 { + return 0 + } + return sample[0].Value.Uint64() +} + +// suspendGap is the decision behind suspendedSince, apart from the clocks so +// that it can be tested. +func suspendGap(wall, monotonic time.Duration) time.Duration { + if gap := wall - monotonic; gap > SuspendThreshold { + return gap + } + return 0 +} + // pairRatios returns b[i]/a[i] for each pair of samples taken next to each // other, which is the series a head start of one candidate shows up in. A zero // in a produces a non-finite ratio, which DetectDrift then declines. @@ -377,14 +455,29 @@ func (r Report) warnings() []string { r.Estimate.Low*100, r.Estimate.High*100)) } + if r.LiveHeap > LargeHeapThreshold { + w = append(w, fmt.Sprintf( + "the program holds %d MB of live data, more than the caches of many machines; for data that size, where it lies in memory can shift the result by several points in ways the interval does not cover, and differently in the next process; run the comparison in several processes with the multiproc package", + r.LiveHeap>>20)) + } + + if r.Suspended > 0 { + w = append(w, fmt.Sprintf( + "the machine appears to have been suspended for %s during the comparison; its samples straddle the pause, so repeat it with the machine kept awake", + r.Suspended.Round(time.Second))) + } + // Drift is a warning rather than a veto: interleaving the measurement order // means a trend hits both candidates about equally, so it inflates the - // spread more than it biases the comparison. + // spread more than it biases the comparison. It is only worth a warning + // when it is large enough to matter to this result; a long run finds + // shifts of a tenth of a percent significant, and a warning that fires on + // nearly every run tells nobody anything. for _, d := range []struct { name string rep DriftReport }{{"A", r.DriftA}, {"B", r.DriftB}} { - if d.rep.N > 0 && d.rep.Drifted(DriftLevel) { + if d.rep.N > 0 && d.rep.Drifted(DriftLevel) && math.Abs(d.rep.RelativeShift) > r.resolution() { w = append(w, fmt.Sprintf( "candidate %s drifted during the run, shifting %+.2f%% from its first half to its second; the machine did not hold still", d.name, d.rep.RelativeShift*100)) diff --git a/compare_test.go b/compare_test.go index c0ba6ce..e7d3465 100644 --- a/compare_test.go +++ b/compare_test.go @@ -2,6 +2,7 @@ package rtcompare import ( "math" + "runtime" "strings" "testing" "time" @@ -436,14 +437,41 @@ func TestReportWarnings(t *testing.T) { want: "includes zero", }, { - name: "drift in one series", + name: "drift in one series larger than the resolution", report: Report{ Validated: true, NoiseFloor: 0.01, Estimate: Estimate{Delta: 0.5, Low: 0.4, High: 0.6}, - DriftB: DriftReport{N: 101, PValue: 0.0001, RelativeShift: -0.07}, + DriftB: DriftReport{N: 101, PValue: 0.0001, RelativeShift: -0.17}, }, want: "candidate B drifted", }, + { + name: "significant drift smaller than the resolution", + report: Report{ + Validated: true, NoiseFloor: 0.01, + Estimate: Estimate{Delta: 0.05, Low: 0.04, High: 0.06}, + DriftA: DriftReport{N: 101, PValue: 0.0001, RelativeShift: 0.006}, + }, + absent: "drifted", + }, + { + name: "large heap", + report: Report{ + Validated: true, NoiseFloor: 0.01, + Estimate: Estimate{Delta: 0.5, Low: 0.4, High: 0.6}, + LiveHeap: 64 << 20, + }, + want: "holds 64 MB of live data", + }, + { + name: "machine suspended", + report: Report{ + Validated: true, NoiseFloor: 0.01, + Estimate: Estimate{Delta: 0.5, Low: 0.4, High: 0.6}, + Suspended: 17 * time.Minute, + }, + want: "suspended for 17m0s", + }, { name: "ratio trend larger than the resolution", report: Report{ @@ -553,3 +581,45 @@ func TestPairRatios(t *testing.T) { t.Error("a zero sample in A should make the ratio series unusable for DetectDrift") } } + +// TestSuspendedSince checks that a comparison notices a machine that went to +// sleep in the middle of it. Compare reads both the wall clock and the +// monotonic clock, and a gap above SuspendThreshold between them is the time +// spent asleep. A gap of an hour has to be reported as an hour, small gaps from +// clock adjustments and a wall clock stepped backwards as zero, and a real +// start time a moment ago as zero. +func TestSuspendedSince(t *testing.T) { + cases := []struct { + name string + wall, monotonic time.Duration + want time.Duration + }{ + {"reports an hour asleep", time.Hour + time.Minute, time.Minute, time.Hour}, + {"ignores a gap below the threshold", time.Minute + 200*time.Millisecond, time.Minute, 0}, + {"ignores a wall clock stepped backwards", time.Minute - 5*time.Second, time.Minute, 0}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if got := suspendGap(c.wall, c.monotonic); got != c.want { + t.Errorf("got %v, want %v", got, c.want) + } + }) + } + if got := suspendedSince(time.Now()); got != 0 { + t.Errorf("no pause: got %v, want 0", got) + } +} + +// TestCompareRecordsTheLiveHeap checks that a comparison records how much +// live data the program holds, which is what the warning about single-process +// results is based on. The runtime reports it after its first collection, so a +// collection is forced first; the value then has to be positive and must not +// exceed what the runtime has obtained from the operating system. +func TestCompareRecordsTheLiveHeap(t *testing.T) { + runtime.GC() + var ms runtime.MemStats + runtime.ReadMemStats(&ms) + if got := liveHeap(); got == 0 || got > ms.Sys { + t.Errorf("live heap %d bytes, want between 1 and %d", got, ms.Sys) + } +} diff --git a/diagnostics.go b/diagnostics.go index 1243cb4..89913a6 100644 --- a/diagnostics.go +++ b/diagnostics.go @@ -276,7 +276,10 @@ func (e Estimate) String() string { // uncertainty only. It says nothing about a bias that affected every // measurement, and a machine that drifted during the run will produce a tight // interval around the wrong number. Read it alongside [ValidateHarness] and -// [DetectDrift]. +// [DetectDrift]. Nor does it cover the next process: measurements from one run +// of a program share that run's memory layout, which for large or +// pointer-heavy data can shift the difference far beyond this interval. Pool +// several processes with [Combine] where that matters. // // Level zero selects [DefaultConfidenceLevel]. Resamples zero selects // [DefaultResamples]. An error is returned if either input holds fewer than diff --git a/dprng.go b/dprng.go index 677140f..17b2f97 100644 --- a/dprng.go +++ b/dprng.go @@ -121,3 +121,30 @@ func (thisState *DPRNG) Uint32N(n uint32) uint32 { func (thisState *DPRNG) UInt32N(n uint32) uint32 { return thisState.Uint32N(n) } + +// Shuffle puts n elements in a pseudo-random order by calling swap(i, j) for +// the pairs a Fisher-Yates shuffle exchanges, like math/rand's Shuffle. +// +// Parameters: n is the number of elements, at most 2^32; swap exchanges the +// elements at two indices. Zero or one element leaves nothing to do, and a +// negative n does nothing either. +// +// Its intended use is the build order of benchmark fixtures. Whatever is +// allocated last lands in different memory, and is the last data the caches +// saw, so a fixed order hands one candidate the same advantage in every +// process. Shuffling the order per process, with a seed per process, turns that +// into scatter that pooling across processes can average out; see [Combine]. +// +// rng := rtcompare.NewDPRNG(seed) +// builders := []func(){buildA, buildB} +// rng.Shuffle(len(builders), func(i, j int) { builders[i], builders[j] = builders[j], builders[i] }) +// for _, build := range builders { build() } +// +// The bias of Uint32N for n that are not powers of two carries over, and it is +// far below anything a benchmark could notice. +func (thisState *DPRNG) Shuffle(n int, swap func(i, j int)) { + for i := n - 1; i > 0; i-- { + j := int(thisState.Uint32N(uint32(i + 1))) + swap(i, j) + } +} diff --git a/layout.go b/layout.go new file mode 100644 index 0000000..6ff2847 --- /dev/null +++ b/layout.go @@ -0,0 +1,125 @@ +package rtcompare + +// This file holds PerturbHeap, which gives each process a different heap layout +// so that the layout becomes random scatter between processes rather than a +// bias that every process repeats. + +import "runtime" + +// Spacers holds the allocations [PerturbHeap] made. Its only purpose is to stay +// reachable: while it is, the memory it occupies is not handed to anything +// else. +type Spacers struct { + noscan [][]byte + scan [][]*byte + large []byte +} + +// Bytes returns how much memory the spacers occupy, for reporting. +func (s *Spacers) Bytes() int { + total := len(s.large) + for _, b := range s.noscan { + total += len(b) + } + for _, p := range s.scan { + total += 8 * len(p) + } + return total +} + +// KeepAlive keeps the spacers reachable up to the point of the call. Call it +// after the last measurement, e.g. with defer right after PerturbHeap. +func (s *Spacers) KeepAlive() { runtime.KeepAlive(s) } + +// perturbClassBudget bounds the spacers per allocation size, and so the total, +// which comes to a few megabytes across all sizes. +const perturbClassBudget = 64 << 10 + +// perturbLargeMax bounds the one large spacer block. Large objects get their +// own pages, so this block moves where the next large allocations begin. +const perturbLargeMax = 16 << 20 + +// PerturbHeap allocates a pseudo-random amount of filler in every small +// allocation size and one large block, so that data allocated afterwards lands +// at addresses and in cache sets that differ from one seed to the next. +// +// Parameters: seed selects the perturbation; zero is a valid seed like any +// other. Give each process a different one, and record it, so that a process +// can be repeated exactly. +// +// It returns the spacers, which must stay reachable until the measurement is +// over: `defer rtcompare.PerturbHeap(seed).KeepAlive()` does that. If they +// were collected, their memory could be handed to the code under test and the +// perturbation would partly undo itself. +// +// The multiproc package calls it in every child process before the suite +// runs, so that code using multiproc never needs to. Call it yourself only +// when running processes by other means, at the start of a process, before +// building the data the candidates work on. The problem it addresses is that a Go program's heap layout is +// deterministic: the same allocations in the same order give much the same +// addresses in every run. For data larger than the caches or full of pointers, +// that layout moves a measured difference by several points, and a program +// that repeats its layout repeats its bias in every process, which pooling +// across processes cannot remove (issue #109). Perturbing the heap per process +// turns the layout into a random effect, which [Combine] can average over. +// It does not replace varying the order in which the fixtures are built: in +// the reproduction in cmd/rtcompare-aa, whichever of two identical lists was +// built second stayed about 3% faster whatever the perturbation, so the build +// order has to alternate between processes as well; see the multiproc package. The idea follows Curtsinger and Berger, "STABILIZER: +// Statistically Sound Performance Evaluation", ASPLOS 2013, which randomizes +// code, stack and heap layout for the same reason. +// +// It only moves the heap. Code layout, stack addresses and the size of the +// environment, which also shift performance (Mytkowicz et al., "Producing Wrong +// Data Without Doing Anything Obviously Wrong!", ASPLOS 2009), are untouched. +// +// The filler costs a few megabytes of spacers across the size classes, both +// with and without pointers, plus a large block of up to 16 MB whose pages are +// mostly never touched, and about a millisecond. +func PerturbHeap(seed uint64) *Spacers { + rng := NewDPRNG(mixSeed(seed)) + sizes := smallSizes() + rng.Shuffle(len(sizes), func(i, j int) { sizes[i], sizes[j] = sizes[j], sizes[i] }) + + s := &Spacers{} + for _, size := range sizes { + limit := uint32(perturbClassBudget / size) + // Objects with and without pointers live in different spans, so both + // kinds are shifted. + for range rng.Uint32N(limit + 1) { + s.noscan = append(s.noscan, make([]byte, size)) + } + for range rng.Uint32N(limit + 1) { + s.scan = append(s.scan, make([]*byte, size/8)) + } + } + s.large = make([]byte, rng.Uint32N(perturbLargeMax)) + return s +} + +// smallSizes returns allocation sizes from 8 bytes to 32 KB spaced about an +// eighth apart, which is roughly how Go spaces its size classes, so that every +// class receives some filler without depending on the runtime's exact table. +func smallSizes() []int { + const largest = 32 << 10 + var sizes []int + for size := 8; size < largest; size += max(8, size/8&^7) { + sizes = append(sizes, size) + } + return append(sizes, largest) +} + +// mixSeed spreads a seed over all 64 bits with one round of splitmix64, so that +// small seeds such as a process index do not start the xorshift generator in a +// low-entropy state, and so that zero stays a fixed seed instead of the request +// for a random one that NewDPRNG takes it to be. +func mixSeed(seed uint64) uint64 { + z := seed + 0x9E3779B97F4A7C15 + z = (z ^ (z >> 30)) * 0xBF58476D1CE4E5B9 + z = (z ^ (z >> 27)) * 0x94D049BB133111EB + z ^= z >> 31 + if z == 0 { + return 1 + } + return z +} diff --git a/layout_test.go b/layout_test.go new file mode 100644 index 0000000..06a3788 --- /dev/null +++ b/layout_test.go @@ -0,0 +1,104 @@ +package rtcompare + +import ( + "slices" + "testing" +) + +// TestDPRNGShuffleIsASeededPermutation checks the helper for randomizing the +// build order of benchmark fixtures per process. A shuffle has to keep every +// element, has to repeat exactly for the same seed so that a process can be +// reproduced, and has to leave nothing to do for zero, one or a negative +// number of elements. +func TestDPRNGShuffleIsASeededPermutation(t *testing.T) { + shuffled := func(seed uint64, n int) []int { + xs := make([]int, n) + for i := range xs { + xs[i] = i + } + rng := NewDPRNG(seed) + rng.Shuffle(len(xs), func(i, j int) { xs[i], xs[j] = xs[j], xs[i] }) + return xs + } + a, b := shuffled(7, 50), shuffled(7, 50) + if !slices.Equal(a, b) { + t.Error("the same seed produced two different orders") + } + if slices.Equal(a, shuffled(8, 50)) { + t.Error("two seeds produced the same order of 50 elements") + } + sorted := slices.Clone(a) + slices.Sort(sorted) + for i, v := range sorted { + if v != i { + t.Fatalf("the shuffle lost or duplicated elements: %v", a) + } + } + for _, n := range []int{-1, 0, 1} { + rng := NewDPRNG(1) + rng.Shuffle(n, func(i, j int) { t.Errorf("swap called for n=%d", n) }) + } +} + +// TestDPRNGShuffleIsUniform checks that no build order is favoured. Over many +// shuffles of three elements, each of the six orders has to come up about a +// sixth of the time; the bound is five standard errors wide. +func TestDPRNGShuffleIsUniform(t *testing.T) { + const trials = 60000 + counts := map[[3]int]int{} + rng := NewDPRNG(42) + for range trials { + xs := [3]int{0, 1, 2} + rng.Shuffle(3, func(i, j int) { xs[i], xs[j] = xs[j], xs[i] }) + counts[xs]++ + } + if len(counts) != 6 { + t.Fatalf("expected all 6 orders, got %d", len(counts)) + } + want := trials / 6.0 + for order, c := range counts { + if d := float64(c) - want; d > 5*91.3 || d < -5*91.3 { // sd = sqrt(n p (1-p)) + t.Errorf("order %v came up %d times, want about %.0f", order, c, want) + } + } +} + +// TestPerturbHeapIsSeededAndBounded checks the heap perturbation that gives +// each process of a multi-process comparison its own memory layout. It has to +// be reproducible from its seed, differ between seeds, and stay within the few +// megabytes of filler plus one large block that its documentation promises. +func TestPerturbHeapIsSeededAndBounded(t *testing.T) { + a, b, c := PerturbHeap(0), PerturbHeap(0), PerturbHeap(1) + defer a.KeepAlive() + defer b.KeepAlive() + defer c.KeepAlive() + if a.Bytes() != b.Bytes() { + t.Errorf("the same seed allocated %d and %d bytes", a.Bytes(), b.Bytes()) + } + if a.Bytes() == c.Bytes() { + t.Errorf("two seeds allocated the same %d bytes", a.Bytes()) + } + limit := 2*perturbClassBudget*len(smallSizes()) + perturbLargeMax + for _, s := range []*Spacers{a, c} { + if s.Bytes() <= 0 || s.Bytes() > limit { + t.Errorf("spacers occupy %d bytes, want between 1 and %d", s.Bytes(), limit) + } + } +} + +// TestSmallSizesCoverTheSizeClasses checks that the filler reaches every +// small allocation size class: sizes have to start at 8 bytes, end at 32 KB, +// rise strictly, stay multiples of 8 and never step by more than an eighth +// plus 8 bytes, which is at least as fine as Go's own class spacing. +func TestSmallSizesCoverTheSizeClasses(t *testing.T) { + sizes := smallSizes() + if sizes[0] != 8 || sizes[len(sizes)-1] != 32<<10 { + t.Errorf("sizes run from %d to %d, want 8 to 32768", sizes[0], sizes[len(sizes)-1]) + } + for i := 1; i < len(sizes); i++ { + step := sizes[i] - sizes[i-1] + if step <= 0 || sizes[i]%8 != 0 || step > sizes[i-1]/8+8 { + t.Errorf("step from %d to %d", sizes[i-1], sizes[i]) + } + } +} diff --git a/multiproc/multiproc.go b/multiproc/multiproc.go new file mode 100644 index 0000000..59d976b --- /dev/null +++ b/multiproc/multiproc.go @@ -0,0 +1,495 @@ +// Package multiproc runs rtcompare comparisons in several processes and pools +// their results, so that the reported interval covers the scatter between +// processes and not only the noise within one. +// +// A single process fixes its memory layout for its whole lifetime, and for data +// larger than the caches or full of pointers that layout can move a measured +// difference by several points. The interval of a single [rtcompare.Report] +// cannot see this: the next process reports a different, equally narrow +// interval somewhere else (issue #109). The remedy is to treat one process as +// one observation. +// +// Everything that takes is done here, so that none of it has to be remembered: +// the program is started again as child processes, one at a time; each child +// perturbs its heap from its own seed before anything is built; the two +// candidates' data is built in an order that alternates between processes, +// since whichever is built second can be consistently a few percent faster; +// the reports are pooled per comparison with [rtcompare.Combine]; and +// processes keep being started until every pooled interval is precise enough. +// +// A program needs nothing but the comparisons: +// +// func main() { +// multiproc.Main(multiproc.Options{}, multiproc.Pair{ +// Name: "lookup", +// A: func() rtcompare.Candidate { return lookupIn(buildTreeA()) }, +// B: func() rtcompare.Candidate { return lookupIn(buildTreeB()) }, +// }) +// } +// +// In a test, [RunTest] does the same and returns the results for assertions. +// [Run] takes an arbitrary suite for anything the pairs do not cover. +package multiproc + +import ( + "crypto/rand" + "encoding/binary" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "time" + + "github.com/TomTonic/rtcompare" +) + +// The environment variables through which a parent tells a child that it is +// one, and what it is to do. +const ( + envOut = "RTCOMPARE_MULTIPROC_OUT" + envSeed = "RTCOMPARE_MULTIPROC_SEED" + envIndex = "RTCOMPARE_MULTIPROC_INDEX" +) + +// Defaults for [Options], taken from a downstream suite that pooled 46 +// comparisons this way. Five processes were enough while the data fit in the +// caches; out of cache the stop rule needed 13 to 15 before every interval was +// within two points. +const ( + DefaultMinProcesses = 5 + DefaultMaxProcesses = 20 + DefaultRotation = 2 + DefaultAbsPrecision = 0.02 + DefaultRelPrecision = 0.10 +) + +// Options configures [Run]. The zero value is usable and selects the +// documented defaults. +type Options struct { + // MinProcesses is the least number of processes to run before the stop + // rule is consulted. Zero selects [DefaultMinProcesses]. Must be at least 3, + // the minimum [rtcompare.Combine] pools. + MinProcesses int + + // MaxProcesses bounds the number of processes. Zero selects + // [DefaultMaxProcesses]. Must not be below MinProcesses. + MaxProcesses int + + // Rotation is the number of build orders the suite cycles through by + // Process.Index, such as 2 for building A's data last in even processes and + // B's in odd ones. The stop rule is only consulted after a whole rotation, + // so that every order is represented equally in the pooled result. Zero + // selects [DefaultRotation], which matches that two-way alternation and + // costs a suite that does not alternate at most one extra process; 1 + // consults the rule after every process. + Rotation int + + // AbsPrecision and RelPrecision are the stop rule: no more processes are + // started once every comparison's pooled interval has a half-width of at + // most AbsPrecision, or of at most RelPrecision times its |Delta|. Zero + // selects [DefaultAbsPrecision] and [DefaultRelPrecision]. + AbsPrecision, RelPrecision float64 + + // Level is the coverage level of the pooled intervals. Zero selects + // [rtcompare.DefaultConfidenceLevel]. + Level float64 + + // Seed determines the per-process seeds, so that a whole run can be + // repeated. Zero selects a random one; Results.Seeds records what each + // process got either way. + Seed uint64 + + // Executable is the program to start as a child. Empty selects the running + // binary, from os.Executable. + Executable string + + // Args are the arguments for the children. Nil selects os.Args[1:], which + // repeats the parent's own command line. In a test, run only the calling + // test in the children, or every test of the package runs in every child: + // + // Args: []string{"-test.run=^" + regexp.QuoteMeta(t.Name()) + "$"} + Args []string + + // Stdout and Stderr receive the children's output. Nil discards stdout and + // passes stderr through to the parent's. + Stdout, Stderr io.Writer + + // Progress, when set, is called after each process with the results so + // far, for example to print a line per process. + Progress func(Results) +} + +// Process is what a child knows about itself. The suite receives it and +// records its results through it. +type Process struct { + // Index counts the processes of a run from zero. + Index int + + // Seed is this process's seed. The heap perturbation Run applies before + // the suite and Rand are derived from it, so a process can be repeated by + // running the suite again with the same seed. + Seed uint64 + + records []record +} + +// Record hands a comparison's report to the parent under a name. The same name +// in every process identifies the same comparison; recording a name twice in +// one process records two observations of it. +func (p *Process) Record(name string, r rtcompare.Report) { + p.records = append(p.records, newRecord(name, r)) +} + +// Rand returns a generator seeded from this process's seed, independent of the +// heap perturbation, for shuffling the order in which many fixtures are built; +// see [rtcompare.DPRNG.Shuffle]. For two fixtures, alternate by Index instead, +// which balances exactly where a random draw over a few processes rarely does. +func (p *Process) Rand() *rtcompare.DPRNG { + rng := rtcompare.NewDPRNG(splitmix(p.Seed^0x5DEECE66D) | 1) // zero would ask NewDPRNG for a random seed + return &rng +} + +// Comparison collects one named comparison across processes. +type Comparison struct { + // Name is the name the suite recorded it under. + Name string + + // Reports holds one report per process that recorded it, in process order. + // They carry the estimate, noise floor, verdict, suspension and warnings of + // each process, but not the raw samples or validations, which stay in the + // child. + Reports []rtcompare.Report + + // Pooled is the result of [rtcompare.Combine] over Reports. It is the zero + // value until three processes have recorded the comparison. + Pooled rtcompare.Pooled +} + +// Results is what [Run] found. +type Results struct { + // Child is true in a child process, where the other fields are empty. The + // caller should return without reporting anything. + Child bool + + // Processes is how many processes ran, and Seeds the seed each one got. + Processes int + Seeds []uint64 + + // Comparisons holds every named comparison in the order it was first + // recorded. + Comparisons []Comparison + + // Precise reports whether the stop rule was met; false means the run + // stopped at MaxProcesses with at least one interval still wider than + // requested. + Precise bool +} + +// String renders the pooled results, one block per comparison. +func (r Results) String() string { + var b strings.Builder + fmt.Fprintf(&b, "%d processes", r.Processes) + if !r.Precise { + b.WriteString(", stopped before every interval was as precise as requested") + } + b.WriteString("\n") + for _, c := range r.Comparisons { + fmt.Fprintf(&b, "\n%s:\n", c.Name) + if c.Pooled.Processes == 0 { + fmt.Fprintf(&b, " recorded by %d processes, too few to pool\n", len(c.Reports)) + continue + } + for _, line := range strings.Split(c.Pooled.String(), "\n") { + fmt.Fprintf(&b, " %s\n", line) + } + } + return strings.TrimRight(b.String(), "\n") +} + +// Run runs suite in several child processes and pools what they record. +// +// Parameters: opt sets how many processes to run and when to stop, see +// [Options]; suite is the measurement one process performs. It receives a +// [Process] with that process's seed, and records each comparison's report +// under a name with Process.Record. +// +// Run behaves differently depending on where it is called. In the parent, the +// program as started by the user, it never calls suite: it starts the current +// binary again as a child, with the same arguments unless Options.Args says +// otherwise, waits for it, and repeats, one process at a time so that they do +// not disturb each other. Once MinProcesses have run it stops as soon as every +// comparison's pooled interval meets the precision in Options at the end of a +// whole Rotation, and in any case after MaxProcesses. It returns the pooled +// results. In a child, Run perturbs the heap from the process's seed (see +// [rtcompare.PerturbHeap]), calls suite once, hands its records to the parent +// through a file, and returns Results with Child set; the caller should then +// return without doing anything else, since the parent does the reporting. +// [Main] and [RunTest] take care of that. The children find out which they are +// from environment variables that Run sets for them. +// +// Use it for any comparison whose data is larger than the caches or full of +// pointers, where a single process's interval is known to be several times too +// narrow; see [rtcompare.Combine] for the numbers. It costs a process start and +// the whole suite per process, and out of cache that has taken 13 to 15 +// processes. Keep the machine awake while it runs: each Report notices a +// suspension, but the time is lost. +// +// An error is returned for invalid options, if a child cannot be started, +// exits with a failure, returns an error from suite or records nothing, and if +// pooling fails. In a child, the error is suite's. +func Run(opt Options, suite func(*Process) error) (Results, error) { + if path := os.Getenv(envOut); path != "" { + return Results{Child: true}, runChild(path, suite) + } + opt, err := opt.resolve() + if err != nil { + return Results{}, err + } + dir, err := os.MkdirTemp("", "rtcompare-multiproc-*") + if err != nil { + return Results{}, fmt.Errorf("multiproc: creating a directory for the children's results: %w", err) + } + // A leftover directory of small result files in the system's temporary + // directory is not worth failing a finished run for. + defer func() { _ = os.RemoveAll(dir) }() + + var res Results + index := map[string]int{} + for i := range opt.MaxProcesses { + seed := splitmix(opt.Seed + uint64(i)) + recs, err := runProcess(opt, filepath.Join(dir, strconv.Itoa(i)+".json"), i, seed) + if err != nil { + return res, err + } + res.Processes++ + res.Seeds = append(res.Seeds, seed) + for _, rec := range recs { + j, ok := index[rec.Name] + if !ok { + j = len(res.Comparisons) + index[rec.Name] = j + res.Comparisons = append(res.Comparisons, Comparison{Name: rec.Name}) + } + res.Comparisons[j].Reports = append(res.Comparisons[j].Reports, rec.report()) + } + if err := res.pool(opt.Level); err != nil { + return res, err + } + if opt.Progress != nil { + opt.Progress(res) + } + if res.Processes >= opt.MinProcesses && res.Processes%opt.Rotation == 0 && res.precise(opt.AbsPrecision, opt.RelPrecision) { + res.Precise = true + break + } + } + return res, nil +} + +// resolve checks the options and fills in their defaults. +func (opt Options) resolve() (Options, error) { + if opt.MinProcesses == 0 { + opt.MinProcesses = DefaultMinProcesses + } + if opt.MaxProcesses == 0 { + opt.MaxProcesses = max(DefaultMaxProcesses, opt.MinProcesses) + } + if opt.MinProcesses < 3 { + return opt, fmt.Errorf("multiproc: MinProcesses must be at least 3, got %d", opt.MinProcesses) + } + if opt.MaxProcesses < opt.MinProcesses { + return opt, fmt.Errorf("multiproc: MaxProcesses (%d) must not be below MinProcesses (%d)", opt.MaxProcesses, opt.MinProcesses) + } + if opt.Rotation == 0 { + opt.Rotation = DefaultRotation + } + if opt.Rotation < 0 { + return opt, fmt.Errorf("multiproc: Rotation must not be negative, got %d", opt.Rotation) + } + if opt.AbsPrecision == 0 { + opt.AbsPrecision = DefaultAbsPrecision + } + if opt.RelPrecision == 0 { + opt.RelPrecision = DefaultRelPrecision + } + if !(opt.AbsPrecision > 0) || !(opt.RelPrecision > 0) { + return opt, fmt.Errorf("multiproc: AbsPrecision and RelPrecision must be positive, got %v and %v", opt.AbsPrecision, opt.RelPrecision) + } + if opt.Seed == 0 { + var b [8]byte + if _, err := rand.Read(b[:]); err != nil { + return opt, fmt.Errorf("multiproc: drawing a seed: %w", err) + } + opt.Seed = binary.LittleEndian.Uint64(b[:]) + } + if opt.Executable == "" { + exe, err := os.Executable() + if err != nil { + return opt, fmt.Errorf("multiproc: finding the running binary to start as a child: %w", err) + } + opt.Executable = exe + } + if opt.Args == nil { + opt.Args = os.Args[1:] + } + if opt.Stdout == nil { + opt.Stdout = io.Discard + } + if opt.Stderr == nil { + opt.Stderr = os.Stderr + } + return opt, nil +} + +// pool combines every comparison that at least three processes recorded. +func (r *Results) pool(level float64) error { + for i := range r.Comparisons { + c := &r.Comparisons[i] + if len(c.Reports) < 3 { + continue + } + p, err := rtcompare.Combine(c.Reports, level) + if err != nil { + return fmt.Errorf("multiproc: pooling %q: %w", c.Name, err) + } + c.Pooled = p + } + return nil +} + +// precise reports whether every comparison is pooled and meets the stop rule. +func (r *Results) precise(abs, rel float64) bool { + for _, c := range r.Comparisons { + if c.Pooled.Processes == 0 || !c.Pooled.Precise(abs, rel) { + return false + } + } + return len(r.Comparisons) > 0 +} + +// runProcess starts one child and reads back what it recorded. +func runProcess(opt Options, out string, i int, seed uint64) ([]record, error) { + cmd := exec.Command(opt.Executable, opt.Args...) + cmd.Env = append(os.Environ(), + envOut+"="+out, + envSeed+"="+strconv.FormatUint(seed, 10), + envIndex+"="+strconv.Itoa(i), + ) + cmd.Stdout, cmd.Stderr = opt.Stdout, opt.Stderr + start := time.Now() + if err := cmd.Run(); err != nil { + return nil, fmt.Errorf("multiproc: process %d (seed %d) failed after %s: %w", i, seed, time.Since(start).Round(time.Millisecond), childError(out, err)) + } + data, err := os.ReadFile(out) + if err != nil { + return nil, fmt.Errorf("multiproc: process %d (seed %d) left no results; does the program call multiproc.Run? %w", i, seed, err) + } + var f childFile + if err := json.Unmarshal(data, &f); err != nil { + return nil, fmt.Errorf("multiproc: reading the results of process %d: %w", i, err) + } + if f.Error != "" { + return nil, fmt.Errorf("multiproc: process %d (seed %d): %s", i, seed, f.Error) + } + if len(f.Records) == 0 { + return nil, fmt.Errorf("multiproc: process %d (seed %d) recorded no comparisons", i, seed) + } + return f.Records, nil +} + +// childError prefers the suite's own error, if the child got as far as writing +// it, over the bare exit status. +func childError(out string, exitErr error) error { + data, err := os.ReadFile(out) + if err != nil { + return exitErr + } + var f childFile + if json.Unmarshal(data, &f) != nil || f.Error == "" { + return exitErr + } + return errors.New(f.Error) +} + +// runChild runs the suite once as a child and writes what it recorded, or the +// error it returned, to the file the parent named. +func runChild(out string, suite func(*Process) error) error { + seed, err := strconv.ParseUint(os.Getenv(envSeed), 10, 64) + if err != nil { + return fmt.Errorf("multiproc: child started without a valid %s: %w", envSeed, err) + } + index, err := strconv.Atoi(os.Getenv(envIndex)) + if err != nil { + return fmt.Errorf("multiproc: child started without a valid %s: %w", envIndex, err) + } + p := &Process{Index: index, Seed: seed} + // Before anything the suite builds, so that its data lands at addresses + // that differ from one process to the next; kept alive until the suite is + // done, so that the filler's memory is not handed to the code under test. + spacers := rtcompare.PerturbHeap(seed) + suiteErr := suite(p) + spacers.KeepAlive() + f := childFile{Records: p.records} + if suiteErr != nil { + f.Error = suiteErr.Error() + } + data, err := json.Marshal(f) + if err != nil { + return fmt.Errorf("multiproc: encoding the results for the parent: %w", err) + } + if err := os.WriteFile(out, data, 0o600); err != nil { + return fmt.Errorf("multiproc: writing the results for the parent: %w", err) + } + return suiteErr +} + +// childFile is what a child hands to its parent. +type childFile struct { + Records []record `json:"records"` + Error string `json:"error,omitempty"` +} + +// record is the part of a Report that crosses the process boundary. The full +// Report cannot: its Confidence map has float keys, which JSON does not allow, +// and its samples and validations are not needed for pooling. +type record struct { + Name string `json:"name"` + NsPerOpA float64 `json:"ns_a"` + NsPerOpB float64 `json:"ns_b"` + Estimate rtcompare.Estimate `json:"estimate"` + NoiseFloor float64 `json:"noise_floor"` + Validated bool `json:"validated"` + Resolved bool `json:"resolved"` + Suspended time.Duration `json:"suspended"` + Warnings []string `json:"warnings,omitempty"` +} + +func newRecord(name string, r rtcompare.Report) record { + return record{ + Name: name, NsPerOpA: r.NsPerOpA, NsPerOpB: r.NsPerOpB, Estimate: r.Estimate, + NoiseFloor: r.NoiseFloor, Validated: r.Validated, Resolved: r.Resolved, + Suspended: r.Suspended, Warnings: r.Warnings, + } +} + +func (rec record) report() rtcompare.Report { + return rtcompare.Report{ + NsPerOpA: rec.NsPerOpA, NsPerOpB: rec.NsPerOpB, Estimate: rec.Estimate, + NoiseFloor: rec.NoiseFloor, Validated: rec.Validated, Resolved: rec.Resolved, + Suspended: rec.Suspended, Warnings: rec.Warnings, + } +} + +// splitmix derives well-spread seeds from consecutive integers, one round of +// splitmix64. +func splitmix(x uint64) uint64 { + z := x + 0x9E3779B97F4A7C15 + z = (z ^ (z >> 30)) * 0xBF58476D1CE4E5B9 + z = (z ^ (z >> 27)) * 0x94D049BB133111EB + return z ^ (z >> 31) +} diff --git a/multiproc/multiproc_test.go b/multiproc/multiproc_test.go new file mode 100644 index 0000000..d53178d --- /dev/null +++ b/multiproc/multiproc_test.go @@ -0,0 +1,213 @@ +package multiproc + +import ( + "errors" + "regexp" + "slices" + "strings" + "testing" + + "github.com/TomTonic/rtcompare" +) + +// onlyThisTest makes the children run nothing but the calling test. +func onlyThisTest(t *testing.T) []string { + return []string{"-test.run=^" + regexp.QuoteMeta(t.Name()) + "$"} +} + +// fakeReport stands in for a real comparison, so that these tests exercise +// the process machinery without spending seconds on measurement. The delta +// varies with the seed like a layout effect would. +func fakeReport(p *Process) rtcompare.Report { + d := 0.05 + float64(p.Seed%7)*0.002 + return rtcompare.Report{ + NsPerOpA: float64(p.Index), + Estimate: rtcompare.Estimate{Delta: d, Low: d - 0.004, High: d + 0.004, Level: 0.95}, + Validated: true, + NoiseFloor: 0.004, + Resolved: true, + } +} + +// TestRunPoolsChildProcesses checks the whole round trip a user relies on: the +// suite runs in separate child processes, each with its own seed, and the +// parent gets back one pooled result per named comparison. It belongs to the +// multiproc driver around rtcompare.Combine. With a loose stop rule the run is +// expected to stop at MinProcesses, with every child's report in process +// order, distinct seeds, and a pooled interval over them. +func TestRunPoolsChildProcesses(t *testing.T) { + res, err := Run(Options{MinProcesses: 3, MaxProcesses: 6, Rotation: 1, AbsPrecision: 0.5, Seed: 99, Args: onlyThisTest(t)}, + func(p *Process) error { + p.Record("lookup", fakeReport(p)) + p.Record("insert", fakeReport(p)) + return nil + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Child { + return + } + t.Log("\n" + res.String()) + if res.Processes != 3 || !res.Precise { + t.Errorf("expected to stop precisely after 3 processes, got %d (precise %v)", res.Processes, res.Precise) + } + if len(res.Seeds) != 3 || res.Seeds[0] == res.Seeds[1] || res.Seeds[1] == res.Seeds[2] { + t.Errorf("expected three distinct seeds, got %v", res.Seeds) + } + if len(res.Comparisons) != 2 || res.Comparisons[0].Name != "lookup" || res.Comparisons[1].Name != "insert" { + t.Fatalf("expected the comparisons lookup and insert in that order, got %+v", res.Comparisons) + } + for _, c := range res.Comparisons { + if len(c.Reports) != 3 || c.Pooled.Processes != 3 { + t.Errorf("%s: %d reports, pooled over %d", c.Name, len(c.Reports), c.Pooled.Processes) + } + for i, r := range c.Reports { + if r.NsPerOpA != float64(i) { + t.Errorf("%s: report %d came from process %v", c.Name, i, r.NsPerOpA) + } + } + } +} + +// TestRunIsReproducibleFromItsSeed checks that a whole multi-process run can +// be repeated: the same Options.Seed has to hand the same seeds to the same +// processes, and a different one different seeds. +func TestRunIsReproducibleFromItsSeed(t *testing.T) { + seeds := func(seed uint64) []uint64 { + res, err := Run(Options{MinProcesses: 3, MaxProcesses: 3, Rotation: 1, AbsPrecision: 0.5, Seed: seed, Args: onlyThisTest(t)}, + func(p *Process) error { + p.Record("x", fakeReport(p)) + return nil + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + return res.Seeds + } + a := seeds(1) + if a == nil { + return // in a child + } + if b := seeds(1); !slices.Equal(a, b) { + t.Errorf("the same seed gave %v and %v", a, b) + } + if c := seeds(2); a[0] == c[0] { + t.Errorf("seeds 1 and 2 gave the same first process seed %d", a[0]) + } +} + +// TestRunKeepsGoingUntilPrecise checks the adaptive stop rule: when the +// processes disagree more than the requested precision allows, the driver has +// to keep starting processes up to MaxProcesses and then say that it stopped +// short. +func TestRunKeepsGoingUntilPrecise(t *testing.T) { + res, err := Run(Options{MinProcesses: 3, MaxProcesses: 4, Rotation: 1, AbsPrecision: 1e-9, RelPrecision: 1e-9, Args: onlyThisTest(t)}, + func(p *Process) error { + r := fakeReport(p) + r.Estimate.Delta += float64(p.Index%2) * 0.05 + p.Record("scattered", r) + return nil + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Child { + return + } + if res.Processes != 4 || res.Precise { + t.Errorf("expected to run all 4 processes and stop imprecise, got %d (precise %v)", res.Processes, res.Precise) + } + if !strings.Contains(res.String(), "stopped before") { + t.Errorf("the summary should say the run stopped short:\n%s", res) + } +} + +// TestRunReportsChildFailures checks that a failing suite does not vanish into +// a child's exit status: its error has to reach the caller in the parent, with +// the process named, and so does a suite that recorded nothing. +func TestRunReportsChildFailures(t *testing.T) { + cases := []struct { + name string + suite func(*Process) error + want string + }{ + {"returns the suite's error", func(*Process) error { return errors.New("fixture exploded") }, "fixture exploded"}, + {"returns error when nothing was recorded", func(*Process) error { return nil }, "recorded no comparisons"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + res, err := Run(Options{MinProcesses: 3, Args: onlyThisTest(t), Stderr: &strings.Builder{}}, c.suite) + if res.Child { + return + } + if err == nil || !strings.Contains(err.Error(), c.want) || !strings.Contains(err.Error(), "process 0") { + t.Errorf("got %v, want an error about process 0 mentioning %q", err, c.want) + } + }) + } +} + +// TestOptionsResolve checks that the driver rejects settings under which +// pooling cannot mean anything, and fills in its documented defaults. +func TestOptionsResolve(t *testing.T) { + bad := []struct { + name string + opt Options + want string + }{ + {"returns error for fewer than three processes", Options{MinProcesses: 2}, "at least 3"}, + {"returns error for a maximum below the minimum", Options{MinProcesses: 6, MaxProcesses: 5}, "must not be below"}, + {"returns error for a negative precision", Options{AbsPrecision: -1}, "positive"}, + {"returns error for a negative rotation", Options{Rotation: -2}, "Rotation"}, + } + for _, c := range bad { + t.Run(c.name, func(t *testing.T) { + if _, err := c.opt.resolve(); err == nil || !strings.Contains(err.Error(), c.want) { + t.Errorf("got %v, want an error mentioning %q", err, c.want) + } + }) + } + opt, err := Options{MinProcesses: 30}.resolve() + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if opt.MaxProcesses != 30 || opt.Rotation != DefaultRotation || opt.AbsPrecision != DefaultAbsPrecision || opt.RelPrecision != DefaultRelPrecision || + opt.Seed == 0 || opt.Executable == "" || opt.Args == nil || opt.Stdout == nil || opt.Stderr == nil { + t.Errorf("defaults not filled in: %+v", opt) + } +} + +// TestProcessRandIsSeeded checks that the generator a child uses to shuffle +// its build order is reproducible from the process seed, including seed zero, +// which rtcompare.NewDPRNG would otherwise take as a request for a random one. +func TestProcessRandIsSeeded(t *testing.T) { + for _, seed := range []uint64{0, 1, 12345} { + a, b := (&Process{Seed: seed}).Rand(), (&Process{Seed: seed}).Rand() + if a.Uint64() != b.Uint64() { + t.Errorf("seed %d gave two different streams", seed) + } + } +} + +// TestRunStopsOnlyAfterAWholeRotation checks that a suite alternating its +// build order between processes gets every order equally often in the pooled +// result. With a stop rule that is met at once, MinProcesses 3 and the default +// rotation of 2, the driver has to run a fourth process rather than stop with +// one order represented twice and the other once. +func TestRunStopsOnlyAfterAWholeRotation(t *testing.T) { + res, err := Run(Options{MinProcesses: 3, AbsPrecision: 0.5, Args: onlyThisTest(t)}, + func(p *Process) error { + p.Record("x", fakeReport(p)) + return nil + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Child { + return + } + if res.Processes != 4 || !res.Precise { + t.Errorf("expected to stop precisely after 4 processes, got %d (precise %v)", res.Processes, res.Precise) + } +} diff --git a/multiproc/pairs.go b/multiproc/pairs.go new file mode 100644 index 0000000..1a50780 --- /dev/null +++ b/multiproc/pairs.go @@ -0,0 +1,171 @@ +package multiproc + +import ( + "fmt" + "io" + "os" + "regexp" + "testing" + + "github.com/TomTonic/rtcompare" +) + +// Pair is one comparison for [Main], [RunTest] or [Pairs]: how to build each +// candidate, and the options to compare them with. +type Pair struct { + // Name identifies the comparison in the results. It must be unique among + // the pairs of one run. + Name string + + // A and B build their candidate: they create the data the candidate works + // on and return the candidate. Build everything the measurement depends + // on here, not before, so that it is built after the heap perturbation and + // in the alternating order. Neither may be nil. + A, B func() rtcompare.Candidate + + // Options are passed to [rtcompare.Compare] unchanged. + Options rtcompare.CompareOptions +} + +// build calls the two builders in an order that alternates with the process +// index: A first in even processes, B first in odd ones. Whichever is built +// second lands in different memory, and has been measured consistently a few +// percent faster for it; alternating lets pooling average that out, and the +// default Options.Rotation of 2 keeps the two orders equally represented. +func (pair Pair) build(index int) (a, b rtcompare.Candidate) { + if index%2 == 0 { + a = pair.A() + b = pair.B() + } else { + b = pair.B() + a = pair.A() + } + return a, b +} + +// checkPairs rejects pairs that cannot run, before any process is started. +func checkPairs(pairs []Pair) error { + if len(pairs) == 0 { + return fmt.Errorf("multiproc: no pairs to compare") + } + seen := map[string]bool{} + for i, pair := range pairs { + if pair.Name == "" { + return fmt.Errorf("multiproc: pair %d has no Name", i) + } + if seen[pair.Name] { + return fmt.Errorf("multiproc: two pairs are named %q", pair.Name) + } + seen[pair.Name] = true + if pair.A == nil || pair.B == nil { + return fmt.Errorf("multiproc: pair %q needs both an A and a B builder", pair.Name) + } + } + return nil +} + +// Pairs returns a suite for [Run] that builds and compares each pair in turn, +// and records each report under the pair's name. +// +// Parameters: pairs are the comparisons, see [Pair]. +// +// Use it when Run's other options are needed with pairs, or to combine pairs +// with measurements of your own in one suite. [Main] and [RunTest] use it. The +// suite returns an error, naming the pair, if a comparison fails; it does not +// check the pairs, which Main and RunTest do before starting any process. +func Pairs(pairs ...Pair) func(*Process) error { + return func(p *Process) error { + for _, pair := range pairs { + a, b := pair.build(p.Index) + rep, err := rtcompare.Compare(a, b, pair.Options) + if err != nil { + return fmt.Errorf("%s: %w", pair.Name, err) + } + p.Record(pair.Name, rep) + } + return nil + } +} + +// Main runs the pairs in several processes and prints the pooled results. It +// is meant to be the whole of a benchmark program's main function. +// +// Parameters: opt configures the processes, see [Options]; pairs are the +// comparisons, see [Pair]. +// +// In the parent, the process the user started, Main prints a line per finished +// process to standard error and the pooled results to standard output, then +// returns. In a child it runs the pairs once and exits, so that nothing after +// Main runs there. On an error it prints it to standard error and exits with +// status 1, in either role. +// +// func main() { +// multiproc.Main(multiproc.Options{}, multiproc.Pair{Name: "lookup", A: buildA, B: buildB}) +// } +func Main(opt Options, pairs ...Pair) { + if code, exit := runMain(opt, pairs, os.Stdout, os.Stderr); exit { + os.Exit(code) + } +} + +// runMain is Main without the exit, so that it can be tested. It reports the +// status to exit with, and whether to exit at all. Its output goes to the +// terminal, where a failed write leaves nothing better to do than carry on, so +// write errors are deliberately ignored. +func runMain(opt Options, pairs []Pair, stdout, stderr io.Writer) (code int, exit bool) { + if err := checkPairs(pairs); err != nil { + _, _ = fmt.Fprintln(stderr, err) + return 1, true + } + if opt.Progress == nil { + opt.Progress = func(r Results) { + _, _ = fmt.Fprintf(stderr, "multiproc: process %d done\n", r.Processes) + } + } + res, err := Run(opt, Pairs(pairs...)) + if err != nil { + _, _ = fmt.Fprintln(stderr, err) + return 1, true + } + if res.Child { + return 0, true + } + _, _ = fmt.Fprintln(stdout, res) + return 0, false +} + +// RunTest runs the pairs in several processes from within a test and returns +// the pooled results. +// +// Parameters: t is the calling test; opt configures the processes, see +// [Options]; pairs are the comparisons, see [Pair]. +// +// It does for a test what [Main] does for a program. Unless opt.Args is set, +// the children run only the calling test, since otherwise every test of the +// package would run in every child. In a child, RunTest runs the pairs once and +// skips the rest of the test, which the parent reports; in the parent, an error +// fails the test. The children re-run the test function up to the call, so +// build the candidates' data inside the pairs' builders rather than before. +// +// func TestLookup(t *testing.T) { +// res := multiproc.RunTest(t, multiproc.Options{}, multiproc.Pair{Name: "lookup", A: buildA, B: buildB}) +// if p := res.Comparisons[0].Pooled; p.Resolved { ... } +// } +func RunTest(t testing.TB, opt Options, pairs ...Pair) Results { + t.Helper() + if err := checkPairs(pairs); err != nil { + t.Fatal(err) + } + if opt.Args == nil { + opt.Args = []string{"-test.run=^" + regexp.QuoteMeta(t.Name()) + "$"} + } + res, err := Run(opt, Pairs(pairs...)) + if res.Child { + // The parent reads any error from the child's results file. + t.SkipNow() + } + if err != nil { + t.Fatal(err) + } + return res +} diff --git a/multiproc/pairs_test.go b/multiproc/pairs_test.go new file mode 100644 index 0000000..1dae8d7 --- /dev/null +++ b/multiproc/pairs_test.go @@ -0,0 +1,138 @@ +package multiproc + +import ( + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/TomTonic/rtcompare" +) + +var pairsSink uint64 + +// spin is a cheap candidate with unremovable work, for pairs that must run a +// real comparison quickly. +func spin(name string, cost uint64) rtcompare.Candidate { + return rtcompare.Candidate{Name: name, Batch: func(n uint64) { + var acc uint64 + for i := range n * cost { + acc = acc*31 + i + } + pairsSink ^= acc + }} +} + +// quick keeps a real comparison to a few milliseconds. +var quick = rtcompare.CompareOptions{ + Collect: rtcompare.CollectOptions{Repeats: 11, InnerLoops: 2000, WarmupDuration: time.Millisecond}, + SkipValidation: true, + Resamples: 300, +} + +// TestPairBuildOrderAlternates checks the automatic protection against the +// build-order effect: whichever candidate's data is built second can be +// consistently faster, so the driver has to build A first in even processes +// and B first in odd ones without the user having to arrange it. +func TestPairBuildOrderAlternates(t *testing.T) { + var order []string + pair := Pair{ + Name: "x", + A: func() rtcompare.Candidate { order = append(order, "A"); return spin("a", 1) }, + B: func() rtcompare.Candidate { order = append(order, "B"); return spin("b", 1) }, + } + for index := range 4 { + a, b := pair.build(index) + if a.Name != "a" || b.Name != "b" { + t.Errorf("process %d: the candidates swapped roles", index) + } + } + if got := strings.Join(order, ""); got != "ABBAABBA" { + t.Errorf("build order over four processes: got %s, want ABBAABBA", got) + } +} + +// TestPairsRecordsEachComparison checks the suite behind Main and RunTest in a +// single process: every pair is built, compared and recorded under its name, +// in order, and a failing comparison is reported with the pair's name. +func TestPairsRecordsEachComparison(t *testing.T) { + p := &Process{Index: 1} + err := Pairs( + Pair{Name: "first", A: func() rtcompare.Candidate { return spin("a", 1) }, B: func() rtcompare.Candidate { return spin("b", 2) }, Options: quick}, + Pair{Name: "second", A: func() rtcompare.Candidate { return spin("a", 2) }, B: func() rtcompare.Candidate { return spin("b", 1) }, Options: quick}, + )(p) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(p.records) != 2 || p.records[0].Name != "first" || p.records[1].Name != "second" { + t.Fatalf("expected records first and second, got %+v", p.records) + } + if p.records[0].Estimate.Delta <= 0 || p.records[1].Estimate.Delta >= 0 { + t.Errorf("the cheaper candidate should come out faster: %+v", p.records) + } + + broken := Pair{Name: "broken", A: func() rtcompare.Candidate { return rtcompare.Candidate{} }, B: func() rtcompare.Candidate { return spin("b", 1) }} + if err := Pairs(broken)(p); err == nil || !strings.Contains(err.Error(), "broken") { + t.Errorf("got %v, want an error naming the pair", err) + } +} + +// TestCheckPairs checks that pairs which cannot run are rejected before any +// process is started, with the reason named. +func TestCheckPairs(t *testing.T) { + build := func() rtcompare.Candidate { return spin("x", 1) } + cases := []struct { + name string + pairs []Pair + want string + }{ + {"returns error for no pairs", nil, "no pairs"}, + {"returns error for a missing name", []Pair{{A: build, B: build}}, "no Name"}, + {"returns error for a duplicate name", []Pair{{Name: "x", A: build, B: build}, {Name: "x", A: build, B: build}}, "two pairs"}, + {"returns error for a missing builder", []Pair{{Name: "x", A: build}}, "both an A and a B"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if err := checkPairs(c.pairs); err == nil || !strings.Contains(err.Error(), c.want) { + t.Errorf("got %v, want an error mentioning %q", err, c.want) + } + }) + } +} + +// TestRunTestPoolsPairsAcrossProcesses checks the one-call path for tests: +// the pairs run in real child processes, restricted to this test, and come +// back pooled, with nothing to remember about child processes or build order. +func TestRunTestPoolsPairsAcrossProcesses(t *testing.T) { + res := RunTest(t, Options{MinProcesses: 3, MaxProcesses: 4, AbsPrecision: 0.5}, + Pair{Name: "cheap vs costly", A: func() rtcompare.Candidate { return spin("a", 1) }, B: func() rtcompare.Candidate { return spin("b", 2) }, Options: quick}) + if res.Processes != 4 || len(res.Comparisons) != 1 { + t.Fatalf("expected one comparison over 4 processes, got %d over %d", len(res.Comparisons), res.Processes) + } + if p := res.Comparisons[0].Pooled; p.Processes != 4 || p.Delta <= 0 { + t.Errorf("pooled result %+v, want A faster over 4 processes", p) + } +} + +// TestRunMainRoles checks what Main does in each role without exiting the +// test binary: invalid pairs end the program with status 1, and a child runs +// the pairs, writes its results for the parent and ends with status 0. +func TestRunMainRoles(t *testing.T) { + var out, errOut strings.Builder + if code, exit := runMain(Options{}, nil, &out, &errOut); code != 1 || !exit || !strings.Contains(errOut.String(), "no pairs") { + t.Errorf("invalid pairs: code %d, exit %v, stderr %q", code, exit, errOut.String()) + } + + file := filepath.Join(t.TempDir(), "child.json") + t.Setenv(envOut, file) + t.Setenv(envSeed, "7") + t.Setenv(envIndex, "0") + pair := Pair{Name: "x", A: func() rtcompare.Candidate { return spin("a", 1) }, B: func() rtcompare.Candidate { return spin("b", 1) }, Options: quick} + if code, exit := runMain(Options{}, []Pair{pair}, &out, &errOut); code != 0 || !exit { + t.Errorf("child: code %d, exit %v, stderr %q", code, exit, errOut.String()) + } + if data, err := os.ReadFile(file); err != nil || !strings.Contains(string(data), `"name":"x"`) { + t.Errorf("child results file: %s, %v", data, err) + } +}