From 0f48b182c680220b32e83c55571fb8ce4b8a6389 Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Fri, 21 Aug 2026 22:47:17 +0300 Subject: [PATCH 1/3] feat(runner): add resource limits and assertions Allow `TestCase` to override global CPU and memory limits using `SubjectCPULimit` and `SubjectMemLimit`. Add `ExpectResources` to `FleetConfig` to assert on `DeviceResource` gauges (e.g., cpu.container) with exact values or ranges via the new `ResourceBound` type. Implement `fleetWaitResources` to poll for conditions during fleet tests. --- internal/config/case.go | 52 ++++++- internal/config/resource_bound_test.go | 68 +++++++++ internal/runner/runner.go | 184 +++++++++++++++++++++---- 3 files changed, 273 insertions(+), 31 deletions(-) create mode 100644 internal/config/resource_bound_test.go diff --git a/internal/config/case.go b/internal/config/case.go index dcfab7e..a3c87cb 100644 --- a/internal/config/case.go +++ b/internal/config/case.go @@ -132,8 +132,20 @@ type TestCase struct { // Leave empty to use the registry default (or whatever --image/--version // specifies). Non-strict YAML decode means older harness binaries silently // ignore these fields — they fall back to the registry default. - SubjectImage string `yaml:"subject_image"` - SubjectVersion string `yaml:"subject_version"` + SubjectImage string `yaml:"subject_image"` + + // SubjectCPULimit and SubjectMemLimit pin the subject container's cgroup + // ceilings for this case, overriding the --cpu-limit / --mem-limit flags. + // Same syntax as those flags ("2", "0.5"; "512m", "4g"). + // + // A case needs these when the limits are part of what it asserts rather + // than part of how it is being benchmarked — container-awareness cases + // have to run constrained or they assert nothing, and relying on the + // operator to remember the flags would make a plain run report a product + // failure that is really a missing argument. + SubjectCPULimit string `yaml:"subject_cpu_limit"` + SubjectMemLimit string `yaml:"subject_mem_limit"` + SubjectVersion string `yaml:"subject_version"` Subjects []string `yaml:"subjects"` Configurations map[string]Configuration `yaml:"configurations"` @@ -1435,6 +1447,19 @@ type FleetConfig struct { // Example: {route.in: {events_in: 500, dropped_count: 500, events_out: 0}}. ExpectStats map[string]map[string]int64 `yaml:"expect_stats"` + // ExpectResources (stats scenario only) asserts on the DeviceResource + // (inputtype=4) rows the simulator decodes: CPU, memory, volumes and the + // runtime rows that say which environment the subject measured itself in. + // Keyed "." (e.g. "cpu.container", + // "runtime.container") → field ("count" | "total" | "used" | "sockets" | + // "cores" | "threads" | "samples") → bound. + // + // Bounds rather than ExpectStats' exact match, because these are GAUGES: + // a CPU row says how busy the box was during one 250 ms sample, so the + // only honest assertions are "equals the configured ceiling" for a limit + // and "within [0, ceiling]" for a reading. + ExpectResources map[string]map[string]ResourceBound `yaml:"expect_resources"` + // BaselineSeconds (config_update data-plane mode only) is how long the driver // confirms delivery is suppressed after the BEFORE config is delivered, before // pushing the AFTER config. It proves the case really starts suppressed @@ -2812,3 +2837,26 @@ func ListCases(casesDir string) ([]string, error) { } return names, nil } + +// ResourceBound constrains one DeviceResource gauge field. Set Eq for an exact +// value, or Min/Max (either or both) for a range. An empty bound matches +// anything, which is how a case asserts only that the row exists. +type ResourceBound struct { + Eq *int64 `yaml:"eq"` + Min *int64 `yaml:"min"` + Max *int64 `yaml:"max"` +} + +// Check reports why v fails the bound, or "" when it passes. +func (b ResourceBound) Check(v int64) string { + if b.Eq != nil && v != *b.Eq { + return fmt.Sprintf("= %d, want exactly %d", v, *b.Eq) + } + if b.Min != nil && v < *b.Min { + return fmt.Sprintf("= %d, want >= %d", v, *b.Min) + } + if b.Max != nil && v > *b.Max { + return fmt.Sprintf("= %d, want <= %d", v, *b.Max) + } + return "" +} diff --git a/internal/config/resource_bound_test.go b/internal/config/resource_bound_test.go new file mode 100644 index 0000000..c755515 --- /dev/null +++ b/internal/config/resource_bound_test.go @@ -0,0 +1,68 @@ +package config + +import ( + "testing" + + "gopkg.in/yaml.v3" +) + +// TestResourceBoundYAMLDecode guards against a silently vacuous assertion. +// +// ResourceBound's fields are pointers so "unset" is distinguishable from zero. +// If a yaml tag ever stops matching, decoding yields all-nil pointers, Check +// returns "pass" for every value, and a case that looks like it asserts exact +// cgroup ceilings would assert only that the row exists — while still printing +// a tick. This test fails instead. +func TestResourceBoundYAMLDecode(t *testing.T) { + var fc struct { + ExpectResources map[string]map[string]ResourceBound `yaml:"expect_resources"` + } + src := ` +expect_resources: + runtime.container: + count: {eq: 1} + total: {eq: 524288} + cpu.container: + used: {min: 0, max: 2000} +` + if err := yaml.Unmarshal([]byte(src), &fc); err != nil { + t.Fatal(err) + } + + total := fc.ExpectResources["runtime.container"]["total"] + if total.Eq == nil || *total.Eq != 524288 { + t.Fatalf("total.eq did not decode: %+v", total) + } + if why := total.Check(524288); why != "" { + t.Errorf("Check(524288) = %q, want pass", why) + } + // The exact value a pre-cgroup director reported: host RAM in KB. + if why := total.Check(32528416); why == "" { + t.Error("Check(host-sized value) passed — the bound is vacuous") + } + + used := fc.ExpectResources["cpu.container"]["used"] + if used.Min == nil || used.Max == nil { + t.Fatalf("min/max did not decode: %+v", used) + } + if why := used.Check(2000); why != "" { + t.Errorf("Check(2000) = %q, want pass at the ceiling", why) + } + if why := used.Check(2001); why == "" { + t.Error("Check(over max) passed — max is vacuous") + } + if why := used.Check(-1); why == "" { + t.Error("Check(under min) passed — min is vacuous") + } +} + +// TestResourceBoundEmptyMatchesAnything pins the documented "row must exist" +// form: a bound with nothing set accepts any value. +func TestResourceBoundEmptyMatchesAnything(t *testing.T) { + var b ResourceBound + for _, v := range []int64{-1, 0, 1 << 40} { + if why := b.Check(v); why != "" { + t.Errorf("empty bound rejected %d: %s", v, why) + } + } +} diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 4be2247..a95ecec 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -432,8 +432,8 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe VerifierImage: r.opts.VerifierImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), TLSCertsHost: tlsCertsHost, } @@ -1051,9 +1051,9 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe if tc.Correctness.ValidateContent { fmt.Printf(" malformed lines: %s\n", formatCount(recvMetrics.MalformedLines)) } - if r.opts.CPULimit != "" || r.opts.MemLimit != "" { + if cpuLim, memLim := subjectCPULimit(tc, r.opts.CPULimit), subjectMemLimit(tc, r.opts.MemLimit); cpuLim != "" || memLim != "" { fmt.Printf(" subject limits: cpu=%s mem=%s\n", - defaultVal(r.opts.CPULimit, "unlimited"), defaultVal(r.opts.MemLimit, "unlimited")) + defaultVal(cpuLim, "unlimited"), defaultVal(memLim, "unlimited")) } recvWindow := 0.0 if recvMetrics.FirstReceivedNs > 0 && recvMetrics.LastReceivedNs > recvMetrics.FirstReceivedNs { @@ -1161,8 +1161,8 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -1465,8 +1465,8 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -1844,8 +1844,8 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } // Per-flow setup (e.g. generate TLS certs and set rc.TLSCertsHost) before @@ -2371,8 +2371,8 @@ func (r *Runner) runDirectorAgentCertRotation(tc *config.TestCase, subject confi CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), TLSCertsHost: certsDir, } @@ -2788,8 +2788,8 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -3457,8 +3457,8 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -4025,6 +4025,36 @@ type fleetStatBucket struct { ExecTimeNs int64 `json:"exec_time_ns"` // pipeline.* only: per-processor execution time } +// subjectCPULimit and subjectMemLimit resolve the subject container's cgroup +// ceilings: the case's own pin when it sets one, otherwise the --cpu-limit / +// --mem-limit flags. The case wins because a case that pins a limit is +// asserting on it — see Case.SubjectCPULimit. +func subjectCPULimit(tc *config.TestCase, flag string) string { + if tc != nil && tc.SubjectCPULimit != "" { + return tc.SubjectCPULimit + } + return flag +} + +func subjectMemLimit(tc *config.TestCase, flag string) string { + if tc != nil && tc.SubjectMemLimit != "" { + return tc.SubjectMemLimit + } + return flag +} + +// fleetResourceSample mirrors the simulator's per-item DeviceResource gauge. +type fleetResourceSample struct { + Samples int64 `json:"samples"` + Count int64 `json:"count"` + Total int64 `json:"total"` + Used int64 `json:"used"` + Cores int64 `json:"cores"` + Threads int64 `json:"threads"` + Sockets int64 `json:"sockets"` + DeviceID int64 `json:"device_id"` +} + type fleetStatus struct { Directors map[string]struct { Connected bool `json:"connected"` @@ -4037,9 +4067,14 @@ type fleetStatus struct { // Stats are decoded from the director's forwarded VMF metric frames by // the simulator (see fleetsim/vmfstats.go), keyed by "." // e.g. "route.in", "target.out". - Stats map[string]fleetStatBucket `json:"stats"` - StatsFrames int `json:"stats_frames"` - StatsRecords int `json:"stats_records"` + Stats map[string]fleetStatBucket `json:"stats"` + // Resources are the DeviceResource (inputtype=4) gauge rows, keyed + // "." e.g. "cpu.container", + // "runtime.container". LAST value seen, not summed — see + // fleetsim/vmfstats.go resourceSample. + Resources map[string]fleetResourceSample `json:"resources"` + StatsFrames int `json:"stats_frames"` + StatsRecords int `json:"stats_records"` } `json:"directors"` } @@ -4076,6 +4111,85 @@ func (st *fleetStatus) statField(id, bucket, field string) (int64, bool) { return 0, false } +// resourceField returns the named field of a decoded DeviceResource item. +func (st *fleetStatus) resourceField(id, key, field string) (int64, bool) { + d, ok := st.Directors[id] + if !ok { + return 0, false + } + r, ok := d.Resources[key] + if !ok { + return 0, false + } + switch field { + case "samples": + return r.Samples, true + case "count": + return r.Count, true + case "total": + return r.Total, true + case "used": + return r.Used, true + case "cores": + return r.Cores, true + case "threads": + return r.Threads, true + case "sockets": + return r.Sockets, true + case "device_id": + return r.DeviceID, true + } + return 0, false +} + +// fleetWaitResources polls until every expected DeviceResource item exists and +// satisfies its bounds. +// +// Unlike fleetWaitStats it never fails early on a high reading: these are +// gauges, and the value the subject reports depends on what it was doing during +// that 250 ms sample. A bound violation is only fatal once the deadline passes, +// so a case can wait for a container to warm up without a flaky first tick +// deciding the verdict. +func fleetWaitResources(ctx context.Context, simContainer, dirID string, expect map[string]map[string]config.ResourceBound, deadline time.Time) error { + var lastErr error + for { + st, err := fleetSimStatus(simContainer) + if err != nil { + lastErr = err + } else { + lastErr = nil + for key, fields := range expect { + for field, bound := range fields { + got, ok := st.resourceField(dirID, key, field) + if !ok { + lastErr = fmt.Errorf("resource %s.%s not reported", key, field) + break + } + if why := bound.Check(got); why != "" { + lastErr = fmt.Errorf("resource %s.%s %s", key, field, why) + break + } + } + if lastErr != nil { + break + } + } + if lastErr == nil { + return nil + } + } + if !time.Now().Before(deadline) { + if lastErr == nil { + lastErr = fmt.Errorf("resource expectations not reached by deadline") + } + return lastErr + } + if err := sleepCtx(ctx, 3*time.Second); err != nil { + return fmt.Errorf("interrupted: %w", err) + } + } +} + func (st *fleetStatus) connected(id string) bool { d, ok := st.Directors[id] return ok && d.Connected @@ -4356,8 +4470,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf VerifierImage: r.opts.VerifierImage, // pipeline_verify reuses the DuckDB verifier ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -5082,6 +5196,18 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf } } + // Same idea for the DeviceResource rows: assert the subject reported its + // OWN ceilings, not the host's. Without this a containerized director + // reporting the host's 32 cores and 32 GB looks identical to one + // reporting its 2-core / 512 MB share. + if len(fc.ExpectResources) > 0 { + if rerr := fleetWaitResources(r.ctx, simContainer, dirID, fc.ExpectResources, scenarioDeadline()); rerr != nil { + errs = append(errs, "resource rows mismatch: "+rerr.Error()) + } else { + fmt.Println(" decoded resource rows match expectations ✓") + } + } + case "reconnect": st0, _ := fleetSimStatus(simContainer) before := 0 @@ -5546,8 +5672,8 @@ func (r *Runner) runCCFCorrectness(tc *config.TestCase, subject config.Subject) CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -5858,8 +5984,8 @@ func (r *Runner) runHTTPSourceCorrectness(tc *config.TestCase, subject config.Su CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -6963,8 +7089,8 @@ func (r *Runner) runKafkaOffsetCommitRestart(tc *config.TestCase, subject config CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -7243,8 +7369,8 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: r.opts.CPULimit, - MemLimit: r.opts.MemLimit, + CPULimit: subjectCPULimit(tc, r.opts.CPULimit), + MemLimit: subjectMemLimit(tc, r.opts.MemLimit), } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) From 4c04f9eda67c33c35b0e4831fe66a3c0966874bc Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Sat, 22 Aug 2026 00:58:43 +0300 Subject: [PATCH 2/3] fix(runner): use pre-resolved limits and return errors --- internal/config/case.go | 12 +-- internal/config/resource_bound_test.go | 18 ++-- internal/runner/runner.go | 111 +++++++++++++++---------- 3 files changed, 81 insertions(+), 60 deletions(-) diff --git a/internal/config/case.go b/internal/config/case.go index a3c87cb..b2ae38d 100644 --- a/internal/config/case.go +++ b/internal/config/case.go @@ -2847,16 +2847,16 @@ type ResourceBound struct { Max *int64 `yaml:"max"` } -// Check reports why v fails the bound, or "" when it passes. -func (b ResourceBound) Check(v int64) string { +// Check reports why v fails the bound, or nil when it passes. +func (b ResourceBound) Check(v int64) error { if b.Eq != nil && v != *b.Eq { - return fmt.Sprintf("= %d, want exactly %d", v, *b.Eq) + return fmt.Errorf("= %d, want exactly %d", v, *b.Eq) } if b.Min != nil && v < *b.Min { - return fmt.Sprintf("= %d, want >= %d", v, *b.Min) + return fmt.Errorf("= %d, want >= %d", v, *b.Min) } if b.Max != nil && v > *b.Max { - return fmt.Sprintf("= %d, want <= %d", v, *b.Max) + return fmt.Errorf("= %d, want <= %d", v, *b.Max) } - return "" + return nil } diff --git a/internal/config/resource_bound_test.go b/internal/config/resource_bound_test.go index c755515..eae0054 100644 --- a/internal/config/resource_bound_test.go +++ b/internal/config/resource_bound_test.go @@ -33,11 +33,11 @@ expect_resources: if total.Eq == nil || *total.Eq != 524288 { t.Fatalf("total.eq did not decode: %+v", total) } - if why := total.Check(524288); why != "" { - t.Errorf("Check(524288) = %q, want pass", why) + if err := total.Check(524288); err != nil { + t.Errorf("Check(524288) = %v, want pass", err) } // The exact value a pre-cgroup director reported: host RAM in KB. - if why := total.Check(32528416); why == "" { + if err := total.Check(32528416); err == nil { t.Error("Check(host-sized value) passed — the bound is vacuous") } @@ -45,13 +45,13 @@ expect_resources: if used.Min == nil || used.Max == nil { t.Fatalf("min/max did not decode: %+v", used) } - if why := used.Check(2000); why != "" { - t.Errorf("Check(2000) = %q, want pass at the ceiling", why) + if err := used.Check(2000); err != nil { + t.Errorf("Check(2000) = %v, want pass at the ceiling", err) } - if why := used.Check(2001); why == "" { + if err := used.Check(2001); err == nil { t.Error("Check(over max) passed — max is vacuous") } - if why := used.Check(-1); why == "" { + if err := used.Check(-1); err == nil { t.Error("Check(under min) passed — min is vacuous") } } @@ -61,8 +61,8 @@ expect_resources: func TestResourceBoundEmptyMatchesAnything(t *testing.T) { var b ResourceBound for _, v := range []int64{-1, 0, 1 << 40} { - if why := b.Check(v); why != "" { - t.Errorf("empty bound rejected %d: %s", v, why) + if err := b.Check(v); err != nil { + t.Errorf("empty bound rejected %d: %v", v, err) } } } diff --git a/internal/runner/runner.go b/internal/runner/runner.go index a95ecec..6acaa76 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -116,6 +116,22 @@ type Runner struct { ctx context.Context opts Options store *results.Store + + // runCPULimit and runMemLimit are the subject's cgroup ceilings for the + // CURRENT run: the case's own pin when it sets one, otherwise the CLI flag. + // Resolved once at the top of Run and read everywhere — compose config, + // console output, and the persisted result — so those three cannot + // disagree. They did: the compose file honoured a case pin while the saved + // result recorded the empty CLI flag, so a run under a 2-core/512MB ceiling + // was filed as unrestricted, which is worse than unrecorded because it + // looks recorded. + // + // A Runner is reused across (case, subject) pairs, so Run assigns BOTH + // unconditionally on entry — never conditionally, or one case's pin leaks + // into the next. Safe without a lock: Run is the only exported method and + // cmd/harness calls it sequentially. + runCPULimit string + runMemLimit string } // hardwareID returns the BENCH_HARDWARE env var or "custom" when unset. @@ -226,6 +242,11 @@ func (r *Runner) resolveValues(tc *config.TestCase, subject config.Subject) erro // Run executes the test and returns the persisted result. func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunResult, error) { + // Resolve the subject's ceilings once for this run. Unconditional: a + // Runner is reused across cases and a stale pin must not survive. + r.runCPULimit = subjectCPULimit(tc, r.opts.CPULimit) + r.runMemLimit = subjectMemLimit(tc, r.opts.MemLimit) + if err := r.resolveValues(tc, subject); err != nil { return results.RunResult{}, err } @@ -432,8 +453,8 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe VerifierImage: r.opts.VerifierImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, TLSCertsHost: tlsCertsHost, } @@ -814,8 +835,8 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, LatencyP50Ms: recvMetrics.LatencyP50Ms, LatencyP95Ms: recvMetrics.LatencyP95Ms, LatencyP99Ms: recvMetrics.LatencyP99Ms, @@ -1051,7 +1072,7 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe if tc.Correctness.ValidateContent { fmt.Printf(" malformed lines: %s\n", formatCount(recvMetrics.MalformedLines)) } - if cpuLim, memLim := subjectCPULimit(tc, r.opts.CPULimit), subjectMemLimit(tc, r.opts.MemLimit); cpuLim != "" || memLim != "" { + if cpuLim, memLim := r.runCPULimit, r.runMemLimit; cpuLim != "" || memLim != "" { fmt.Printf(" subject limits: cpu=%s mem=%s\n", defaultVal(cpuLim, "unlimited"), defaultVal(memLim, "unlimited")) } @@ -1161,8 +1182,8 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -1347,8 +1368,8 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -1465,8 +1486,8 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -1668,8 +1689,8 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -1844,8 +1865,8 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } // Per-flow setup (e.g. generate TLS certs and set rc.TLSCertsHost) before @@ -2039,8 +2060,8 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -2371,8 +2392,8 @@ func (r *Runner) runDirectorAgentCertRotation(tc *config.TestCase, subject confi CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, TLSCertsHost: certsDir, } @@ -2788,8 +2809,8 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -3457,8 +3478,8 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -3975,8 +3996,8 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -4165,8 +4186,8 @@ func fleetWaitResources(ctx context.Context, simContainer, dirID string, expect lastErr = fmt.Errorf("resource %s.%s not reported", key, field) break } - if why := bound.Check(got); why != "" { - lastErr = fmt.Errorf("resource %s.%s %s", key, field, why) + if err := bound.Check(got); err != nil { + lastErr = fmt.Errorf("resource %s.%s %w", key, field, err) break } } @@ -4470,8 +4491,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf VerifierImage: r.opts.VerifierImage, // pipeline_verify reuses the DuckDB verifier ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -5508,8 +5529,8 @@ func (r *Runner) saveFleetResult(tc *config.TestCase, subject config.Subject, co LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -5672,8 +5693,8 @@ func (r *Runner) runCCFCorrectness(tc *config.TestCase, subject config.Subject) CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -5984,8 +6005,8 @@ func (r *Runner) runHTTPSourceCorrectness(tc *config.TestCase, subject config.Su CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -6215,7 +6236,7 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con runCfg := orchestrator.RunConfig{ TestCase: tc, Subject: subject, ConfigName: configName, ConfigSrcPath: srcCfg, CaseDir: caseDir, TmpDir: tmpDir, GeneratorImage: r.opts.GeneratorImage, ReceiverImage: r.opts.ReceiverImage, CollectorImage: r.opts.CollectorImage, - ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.opts.CPULimit, MemLimit: r.opts.MemLimit, + ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.runCPULimit, MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) if err != nil { @@ -6399,7 +6420,7 @@ func (r *Runner) saveClickHouseTargetResult(tc *config.TestCase, subject config. IOThroughputAvg: metrics.IOThroughputAvg, LoadAvg1: metrics.LoadAvg1, LoadAvg5: metrics.LoadAvg5, LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, SubjectMemLimit: r.runMemLimit, } if !passed { result.FailReason = strings.Join(errs, "; ") @@ -6757,7 +6778,7 @@ func (r *Runner) setupAuxRun(tc *config.TestCase, subject config.Subject, config runCfg := orchestrator.RunConfig{ TestCase: tc, Subject: subject, ConfigName: configName, ConfigSrcPath: srcCfg, CaseDir: caseDir, TmpDir: tmpDir, GeneratorImage: r.opts.GeneratorImage, ReceiverImage: r.opts.ReceiverImage, CollectorImage: r.opts.CollectorImage, - ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.opts.CPULimit, MemLimit: r.opts.MemLimit, + ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.runCPULimit, MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) if err != nil { @@ -6890,7 +6911,7 @@ func (r *Runner) runHTTPVaultCertRotation(tc *config.TestCase, subject config.Su runCfg := orchestrator.RunConfig{ TestCase: tc, Subject: subject, ConfigName: configName, ConfigSrcPath: srcCfg, CaseDir: caseDir, TmpDir: tmpDir, GeneratorImage: r.opts.GeneratorImage, ReceiverImage: r.opts.ReceiverImage, CollectorImage: r.opts.CollectorImage, - ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.opts.CPULimit, MemLimit: r.opts.MemLimit, + ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, CPULimit: r.runCPULimit, MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) if err != nil { @@ -7089,8 +7110,8 @@ func (r *Runner) runKafkaOffsetCommitRestart(tc *config.TestCase, subject config CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } orch, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -7369,8 +7390,8 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje CollectorImage: r.opts.CollectorImage, ReceiverHostPort: r.opts.ReceiverHostPort, ExtraSubjectEnv: extraEnv, - CPULimit: subjectCPULimit(tc, r.opts.CPULimit), - MemLimit: subjectMemLimit(tc, r.opts.MemLimit), + CPULimit: r.runCPULimit, + MemLimit: r.runMemLimit, } cr, err := orchestrator.NewComposeRunner(r.ctx, runCfg) @@ -7536,8 +7557,8 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.opts.CPULimit, - SubjectMemLimit: r.opts.MemLimit, + SubjectCPULimit: r.runCPULimit, + SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { From 9267ef8421656813458825c24a3c38878674d4f5 Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Sat, 22 Aug 2026 14:59:52 +0300 Subject: [PATCH 3/3] feat(runner): stamp resource limits in saveResult Centralize the assignment of SubjectCPULimit and SubjectMemLimit in saveResult to ensure all run paths consistently record limits. This ensures that case-pinned limits take precedence over CLI flags and prevents runs from being recorded as unconstrained when they were actually restricted. --- internal/config/case.go | 24 ++++-- internal/runner/runner.go | 39 +++++---- internal/runner/save_result_test.go | 123 ++++++++++++++++++++++++++++ 3 files changed, 164 insertions(+), 22 deletions(-) create mode 100644 internal/runner/save_result_test.go diff --git a/internal/config/case.go b/internal/config/case.go index b2ae38d..e85f839 100644 --- a/internal/config/case.go +++ b/internal/config/case.go @@ -135,14 +135,24 @@ type TestCase struct { SubjectImage string `yaml:"subject_image"` // SubjectCPULimit and SubjectMemLimit pin the subject container's cgroup - // ceilings for this case, overriding the --cpu-limit / --mem-limit flags. - // Same syntax as those flags ("2", "0.5"; "512m", "4g"). + // ceilings for this case. Same syntax as the flags ("2", "0.5"; "512m", + // "4g"). // - // A case needs these when the limits are part of what it asserts rather - // than part of how it is being benchmarked — container-awareness cases - // have to run constrained or they assert nothing, and relying on the - // operator to remember the flags would make a plain run report a product - // failure that is really a missing argument. + // PRECEDENCE IS THE INVERSE of SubjectImage above — case pin > CLI flag, + // not CLI flag > case pin — and the inversion is deliberate. An image pin + // says "this case was written against this build", which an operator + // testing a different build should be able to override. A limits pin says + // "these ceilings are part of what this case ASSERTS", which an operator + // flag must not silently change. + // + // Concretely: director_container_resource_stats_correctness asserts the + // director reports exactly 2000 millicores. If --cpu-limit 8 won, a routine + // suite run would fail that case and report a product bug that is really a + // harness argument — the exact failure the pin exists to prevent. + // + // The rule for a new pinnable field: if a case sets it to make an assertion + // true, the case wins; if it only describes the environment the case was + // developed in, the flag wins. SubjectCPULimit string `yaml:"subject_cpu_limit"` SubjectMemLimit string `yaml:"subject_mem_limit"` SubjectVersion string `yaml:"subject_version"` diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 6acaa76..eb404dc 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -835,8 +835,6 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, LatencyP50Ms: recvMetrics.LatencyP50Ms, LatencyP95Ms: recvMetrics.LatencyP95Ms, LatencyP99Ms: recvMetrics.LatencyP99Ms, @@ -1368,8 +1366,6 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -1689,8 +1685,6 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -2060,8 +2054,6 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -3996,8 +3988,6 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -5529,8 +5519,6 @@ func (r *Runner) saveFleetResult(tc *config.TestCase, subject config.Subject, co LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -6420,7 +6408,6 @@ func (r *Runner) saveClickHouseTargetResult(tc *config.TestCase, subject config. IOThroughputAvg: metrics.IOThroughputAvg, LoadAvg1: metrics.LoadAvg1, LoadAvg5: metrics.LoadAvg5, LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, SubjectMemLimit: r.runMemLimit, } if !passed { result.FailReason = strings.Join(errs, "; ") @@ -7557,8 +7544,6 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje LoadAvg15: metrics.LoadAvg15, SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, - SubjectCPULimit: r.runCPULimit, - SubjectMemLimit: r.runMemLimit, Passed: &passed, } if !passed { @@ -7606,10 +7591,34 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje // saveResult persists a run result unless the run was interrupted — a verdict // computed after cancellation reflects a half-finished run (waits bail out // early on cancel) and must never land in the store. +// +// It is also where the subject's effective cgroup ceilings are stamped onto the +// record. RunResult is hand-built in fourteen places, and six of them never set +// these two fields, so those runs were filed as unrestricted while the container +// really ran under a ceiling — worse than unrecorded, because it looks recorded. +// Setting them at each constructor is how they drifted apart in the first place: +// every field common to all runs eventually gets retrofitted into fourteen +// literals and misses some. This function is the single choke point — it wraps +// the only store.Save call in the package — so stamping here cannot be forgotten +// by a new run path, and a future common field belongs here too. +// +// Safe to do centrally because no flow runs unconstrained: every +// orchestrator.RunConfig in this file carries CPULimit. When nothing is pinned +// the values are empty and `omitempty` keeps them out of the JSON, so an +// unconstrained run is still never recorded as constrained. +// +// The stamp lands on the local copy, so it reaches the STORE but not the +// RunResult returned up to cmd/harness. Nothing reads these fields off the +// returned value — the console summary ignores them, and the only consumer is +// results.Store.Save copying them into the persisted entry. func (r *Runner) saveResult(result results.RunResult, metricsCSVSrc string) (string, error) { if err := r.ctx.Err(); err != nil { return "", fmt.Errorf("interrupted: %w", err) } + + result.SubjectCPULimit = r.runCPULimit + result.SubjectMemLimit = r.runMemLimit + return r.store.Save(result, metricsCSVSrc) } diff --git a/internal/runner/save_result_test.go b/internal/runner/save_result_test.go new file mode 100644 index 0000000..d039217 --- /dev/null +++ b/internal/runner/save_result_test.go @@ -0,0 +1,123 @@ +package runner + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "testing" + + "github.com/VirtualMetric/PipeBench/internal/results" +) + +// newTestRunner builds a Runner writing to a throwaway store. The store field is +// unexported, which is why this test lives in package runner rather than +// runner_test — saveResult is the thing under test and it is unexported too. +func newTestRunner(t *testing.T, ctx context.Context, cpuLimit, memLimit string) (*Runner, string) { + t.Helper() + dir := t.TempDir() + return &Runner{ + ctx: ctx, + store: results.NewStore(dir), + runCPULimit: cpuLimit, + runMemLimit: memLimit, + }, dir +} + +// readEntry returns the single persisted result entry as a raw map, so the test +// asserts on the JSON that actually lands on disk rather than on a Go struct +// that would happily zero-fill a missing key. +func readEntry(t *testing.T, baseDir, hardware, subject string) map[string]any { + t.Helper() + path := filepath.Join(baseDir, hardware, subject+".json") + data, err := os.ReadFile(path) + if err != nil { + t.Fatalf("reading %s: %v", path, err) + } + + var file struct { + Results []map[string]any `json:"results"` + } + if err := json.Unmarshal(data, &file); err != nil { + t.Fatalf("unmarshalling %s: %v", path, err) + } + if len(file.Results) != 1 { + t.Fatalf("expected 1 persisted entry, got %d — %s", len(file.Results), data) + } + return file.Results[0] +} + +// TestSaveResultStampsEffectiveLimits is the regression for the defect this +// function's stamp exists to close. +// +// The RunResult handed in has EMPTY limit fields — the shape six of the fourteen +// result constructors produce (runDirectorAgentCertRotation, +// runDirectorAgentACLRotation, runKafkaOffsetCommitRestart, saveCCFResult, +// saveHTTPSourceResult, saveAuxResult). Those runs were filed as unrestricted +// while the container ran under a ceiling. Because every one of the fourteen +// funnels through saveResult, asserting here covers all of them at once. +func TestSaveResultStampsEffectiveLimits(t *testing.T) { + r, dir := newTestRunner(t, context.Background(), "2", "512m") + + if _, err := r.saveResult(results.RunResult{ + TestName: "some_correctness", + Subject: "vmetric", + Hardware: "custom", + // SubjectCPULimit / SubjectMemLimit deliberately unset. + }, ""); err != nil { + t.Fatalf("saveResult: %v", err) + } + + entry := readEntry(t, dir, "custom", "vmetric") + if got := entry["subject_cpu_limit"]; got != "2" { + t.Errorf("subject_cpu_limit = %v, want %q", got, "2") + } + if got := entry["subject_mem_limit"]; got != "512m" { + t.Errorf("subject_mem_limit = %v, want %q", got, "512m") + } +} + +// TestSaveResultOmitsAbsentLimits is the other half: the stamp must never invent +// a ceiling. An unconstrained run has to stay recorded as unconstrained, or the +// fix reintroduces the original defect with the sign flipped. +func TestSaveResultOmitsAbsentLimits(t *testing.T) { + r, dir := newTestRunner(t, context.Background(), "", "") + + if _, err := r.saveResult(results.RunResult{ + TestName: "some_correctness", + Subject: "vmetric", + Hardware: "custom", + }, ""); err != nil { + t.Fatalf("saveResult: %v", err) + } + + entry := readEntry(t, dir, "custom", "vmetric") + if _, ok := entry["subject_cpu_limit"]; ok { + t.Error("subject_cpu_limit present for an unconstrained run — omitempty was defeated") + } + if _, ok := entry["subject_mem_limit"]; ok { + t.Error("subject_mem_limit present for an unconstrained run — omitempty was defeated") + } +} + +// TestSaveResultInterruptedWritesNothing keeps the pre-existing guard intact: a +// verdict computed after cancellation reflects a half-finished run, and the +// stamp must not have moved the early return. +func TestSaveResultInterruptedWritesNothing(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + r, dir := newTestRunner(t, ctx, "2", "512m") + + if _, err := r.saveResult(results.RunResult{ + TestName: "some_correctness", + Subject: "vmetric", + Hardware: "custom", + }, ""); err == nil { + t.Fatal("saveResult on a cancelled context returned nil error") + } + + if _, err := os.Stat(filepath.Join(dir, "custom", "vmetric.json")); !os.IsNotExist(err) { + t.Errorf("an interrupted run wrote to the store (stat err = %v)", err) + } +}