From 010edd62864cf7bf23290370eda199f1b9118ce8 Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Tue, 14 Jul 2026 18:15:10 +0300 Subject: [PATCH 1/7] feat(orchestrator): add AWS S3 seed objects to LocalStack init --- internal/config/cloud.go | 45 +++++++++++++++++++ internal/config/cloud_test.go | 37 +++++++++++++++ internal/orchestrator/awsinit.go | 19 ++++++++ internal/orchestrator/compose_render_test.go | 47 ++++++++++++++++++++ 4 files changed, 148 insertions(+) diff --git a/internal/config/cloud.go b/internal/config/cloud.go index 40fd629..1090526 100644 --- a/internal/config/cloud.go +++ b/internal/config/cloud.go @@ -77,6 +77,30 @@ type AWSConfig struct { // Subscriptions wires SNS → SQS (RawMessageDelivery=true) so an // SNS-target case is observable by the `sqs` receiver mode. Subscriptions []AWSSubscription `yaml:"subscriptions"` + + // SeedObjects pre-uploads synthetic objects into a declared bucket at + // init, before the subject starts — for list-mode / first-run backlog + // source cases where the subject must find objects already present rather + // than a generator streaming them in during the run. + SeedObjects []AWSSeedObjects `yaml:"seed_objects"` +} + +// AWSSeedObjects pre-uploads Objects synthetic objects (Lines lines each, +// each line prefixed with Marker) under Prefix in Bucket during LocalStack +// init. +type AWSSeedObjects struct { + // Bucket must be one of the declared Buckets. + Bucket string `yaml:"bucket"` + // Prefix is the key prefix for seeded objects (default "seed/"). + Prefix string `yaml:"prefix"` + // Objects is the number of objects to create (> 0). + Objects int `yaml:"objects"` + // Lines is the number of lines per object (> 0). + Lines int `yaml:"lines"` + // Marker is the per-line content prefix (default "SEED"); validated + // against the cloud-name charset so it cannot inject into the init shell + // script. + Marker string `yaml:"marker"` } // AWSStream declares a Kinesis stream created at init. @@ -420,6 +444,27 @@ func (tc *TestCase) validateAWS() error { return err } } + for _, so := range tc.AWS.SeedObjects { + if _, ok := buckets[so.Bucket]; !ok { + return fmt.Errorf("case %q: seed_objects references undeclared bucket %q", tc.Name, so.Bucket) + } + if so.Prefix != "" { + if err := validateCloudName(tc.Name, "aws seed prefix", so.Prefix); err != nil { + return err + } + } + if so.Marker != "" { + if err := validateCloudName(tc.Name, "aws seed marker", so.Marker); err != nil { + return err + } + } + if so.Objects <= 0 { + return fmt.Errorf("case %q: seed_objects for bucket %q requires objects > 0, got %d", tc.Name, so.Bucket, so.Objects) + } + if so.Lines <= 0 { + return fmt.Errorf("case %q: seed_objects for bucket %q requires lines > 0, got %d", tc.Name, so.Bucket, so.Lines) + } + } return nil } diff --git a/internal/config/cloud_test.go b/internal/config/cloud_test.go index eb30a39..8089061 100644 --- a/internal/config/cloud_test.go +++ b/internal/config/cloud_test.go @@ -219,6 +219,43 @@ func TestAWSConfigDefaults(t *testing.T) { } } +func TestValidateAWSSeedObjects(t *testing.T) { + base := func(so AWSSeedObjects) *TestCase { + return &TestCase{ + Name: "seed-case", + Type: "correctness", + Duration: "10s", + AWS: &AWSConfig{ + Buckets: []string{"bench-in"}, + SeedObjects: []AWSSeedObjects{so}, + }, + Receiver: ReceiverConfig{Mode: "tcp", Listen: ":9001"}, + Correctness: CorrectnessConfig{}, + } + } + + tests := []struct { + name string + so AWSSeedObjects + wantErr bool + }{ + {name: "valid", so: AWSSeedObjects{Bucket: "bench-in", Objects: 10, Lines: 100}}, + {name: "valid with prefix and marker", so: AWSSeedObjects{Bucket: "bench-in", Prefix: "seed/", Objects: 1, Lines: 1, Marker: "SEED"}}, + {name: "undeclared bucket", so: AWSSeedObjects{Bucket: "nope", Objects: 1, Lines: 1}, wantErr: true}, + {name: "zero objects", so: AWSSeedObjects{Bucket: "bench-in", Objects: 0, Lines: 1}, wantErr: true}, + {name: "zero lines", so: AWSSeedObjects{Bucket: "bench-in", Objects: 1, Lines: 0}, wantErr: true}, + {name: "injecting marker", so: AWSSeedObjects{Bucket: "bench-in", Objects: 1, Lines: 1, Marker: "a'; rm -rf /"}, wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + err := base(tt.so).validateAWS() + if (err != nil) != tt.wantErr { + t.Fatalf("validateAWS() err = %v, wantErr %v", err, tt.wantErr) + } + }) + } +} + func TestMinioConfigDefaults(t *testing.T) { m := &MinioConfig{Buckets: []string{"bench-out"}} if got, want := m.ImageOrDefault(), "minio/minio:RELEASE.2025-04-22T22-12-26Z"; got != want { diff --git a/internal/orchestrator/awsinit.go b/internal/orchestrator/awsinit.go index d6d3754..19b5e54 100644 --- a/internal/orchestrator/awsinit.go +++ b/internal/orchestrator/awsinit.go @@ -56,6 +56,25 @@ func writeAWSInit(path string, aws *config.AWSConfig) error { "awslocal sns subscribe --topic-arn '%s' --protocol sqs --notification-endpoint '%s' --attributes RawMessageDelivery=true\n", aws.TopicARN(s.Topic), aws.QueueARN(s.Queue)) } + for _, so := range aws.SeedObjects { + prefix := so.Prefix + if prefix == "" { + prefix = "seed/" + } + marker := so.Marker + if marker == "" { + marker = "SEED" + } + // Build each object body with awk (busybox ships it), then upload — + // so the objects exist before LocalStack reports init "completed" and + // the subject's depends_on gate releases. All values are int-formatted + // or charset-validated, so single-quoting is defense in depth. + fmt.Fprintf(&b, "i=0; while [ \"$i\" -lt %d ]; do\n", so.Objects) + fmt.Fprintf(&b, " awk -v o=\"$i\" 'BEGIN{for(l=0;l<%d;l++) printf \"%s-OBJ%%d-LINE%%d\\n\", o, l}' > /tmp/pb-seed-obj\n", so.Lines, marker) + fmt.Fprintf(&b, " awslocal s3 cp /tmp/pb-seed-obj 's3://%s/%sobj-'\"$i\"'.log'\n", so.Bucket, prefix) + b.WriteString(" i=$((i+1))\n") + b.WriteString("done\n") + } b.WriteString("echo 'pipebench aws init complete'\n") return os.WriteFile(path, []byte(b.String()), 0o755) diff --git a/internal/orchestrator/compose_render_test.go b/internal/orchestrator/compose_render_test.go index 2920281..bcbd8ea 100644 --- a/internal/orchestrator/compose_render_test.go +++ b/internal/orchestrator/compose_render_test.go @@ -632,6 +632,53 @@ func TestComposeRendersAWS(t *testing.T) { mustContain(t, string(script), `"QueueArn":"arn:aws:sqs:us-east-1:000000000000:bench-events"`) } +// TestComposeRendersAWSSeedObjects verifies a `seed_objects:` entry renders a +// pre-upload loop into the LocalStack init script (objects exist before the +// subject starts, for list-mode backlog cases). +func TestComposeRendersAWSSeedObjects(t *testing.T) { + tc := &config.TestCase{ + Name: "aws-seed", + Type: "correctness", + Duration: "10s", + AWS: &config.AWSConfig{ + Buckets: []string{"bench-in"}, + SeedObjects: []config.AWSSeedObjects{ + {Bucket: "bench-in", Prefix: "seed/", Objects: 10, Lines: 100, Marker: "SEED"}, + }, + }, + Receiver: config.ReceiverConfig{Mode: "tcp", Listen: ":9001"}, + Correctness: config.CorrectnessConfig{MinReceived: 1000}, + } + if err := tc.Validate(); err != nil { + t.Fatalf("validate: %v", err) + } + subj := config.Subject{Name: "vmetric", Image: "vmetric/director", Version: "2.0.3", ConfigPath: "/config.yml"} + tmp, err := os.MkdirTemp("", "compose-seed-") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(tmp) + composePath := filepath.Join(tmp, "compose.yaml") + cfg := RunConfig{ + TestCase: tc, Subject: subj, ConfigName: "default", + ConfigSrcPath: composePath, TmpDir: tmp, + GeneratorImage: "img-gen", ReceiverImage: "img-recv", CollectorImage: "img-coll", + ReceiverHostPort: 19001, + } + if err := writeCompose(composePath, cfg); err != nil { + t.Fatalf("writeCompose: %v", err) + } + script, err := os.ReadFile(filepath.Join(tmp, "aws-init.sh")) + if err != nil { + t.Fatalf("aws-init.sh not written: %v", err) + } + mustContain(t, string(script), "awslocal s3 mb 's3://bench-in'") + mustContain(t, string(script), `while [ "$i" -lt 10 ]`) + mustContain(t, string(script), "for(l=0;l<100;l++)") + mustContain(t, string(script), "SEED-OBJ%d-LINE%d") + mustContain(t, string(script), "'s3://bench-in/seed/obj-'\"$i\"'.log'") +} + // TestComposeRendersAzure verifies an `azure:` case renders the Azurite // service plus the one-shot azure-init (receiver image), gates the subject on // init completion, and injects the connection string where needed. From 51126d6143b12b3292667319f78891f9ce747228 Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Tue, 14 Jul 2026 18:16:16 +0300 Subject: [PATCH 2/7] refactor(runner): make generic mid delivery action method --- internal/runner/runner.go | 88 +++++++++++++++++++++++++++++---------- 1 file changed, 67 insertions(+), 21 deletions(-) diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 7d51d32..252bc19 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -189,6 +189,14 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe if tc.Type == "kafka_inflight_crash_correctness" { return r.runKafkaInflightCrash(tc, subject) } + // Object-storage in-flight crash: same receiver-up mid-delivery flow as the + // kafka in-flight crash, for a poll-mode source (S3/Azure bucket). The + // bucket is the durable store, decoupled from the subject, so the generator + // keeps uploading across the SIGKILL+restart. Verifies no loss; duplicates + // (crash-resistance replay + cursor re-list) are reported, not failed. + if tc.Type == "persistence_inflight_crash_correctness" { + return r.runInflightCrashCorrectness(tc, subject) + } // Kafka offset-commit restart: receiver stays UP, ALL records are // delivered cleanly, then the subject is restarted gracefully. A // consumer whose offset commits actually persist resumes from the @@ -832,15 +840,20 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe // Re-deriving pass/fail from line-count loss + a strict over-delivery // cap here would wrongly flip a valid allow_overdelivery verifier pass. lossOK := lossPct <= tc.Correctness.ExpectedLossPct - // Kafka consumption is at-least-once: the consumer may re-deliver a - // fetch batch on its initial group join / rebalance, so allow bounded - // over-delivery (Correctness.MaxOverDeliveryPct, default 0 = exact). - // Non-kafka correctness stays strict. + // Over-delivery policy, in precedence order: + // - allow_overdelivery: the source is at-least-once by design (e.g. the + // list-poll since-cursor re-lists the tip object every idle cycle), so + // any over-delivery is accepted; only loss fails. This mirrors the + // carve-out the in-flight-crash / restart handlers already apply. + // - kafka: consumption is at-least-once — the consumer may re-deliver a + // fetch batch on its initial group join / rebalance, so allow bounded + // over-delivery (Correctness.MaxOverDeliveryPct, default 0 = exact). + // - otherwise: strict exact count. overCap := expectedOut if tc.IsKafkaType() { overCap += int64(float64(expectedOut) * tc.Correctness.MaxOverDeliveryPct / 100.0) } - overOK := recvMetrics.LinesReceived <= overCap + overOK := tc.Correctness.AllowOverDelivery || recvMetrics.LinesReceived <= overCap recvOK := result.Passed == nil || *result.Passed var failReasons []string @@ -1580,12 +1593,13 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject return result, nil } -// midDeliveryFlow parameterizes the shared kafka correctness driver -// (runKafkaMidDeliveryAction): produce to the broker with the receiver live, -// fire one disruptive action once the receiver has seen ~half of total_lines, -// then drain and assert no loss (over-delivery from at-least-once recovery is -// reported, not failed). The action is the only thing that varies between the -// flows — an in-flight subject crash, a broker cert rotation, etc. +// midDeliveryFlow parameterizes the shared mid-delivery correctness driver +// (runMidDeliveryAction): produce to a durable source (Kafka broker, S3/Azure +// bucket) with the receiver live, fire one disruptive action once the receiver +// has seen ~half of total_lines, then drain and assert no loss (over-delivery +// from at-least-once recovery is reported, not failed). The action is the only +// thing that varies between the flows — an in-flight subject crash, a broker +// cert rotation, etc. type midDeliveryFlow struct { // verdictLabel names the flow in the PASS/FAIL line, e.g. // "kafka cert rotation correctness". @@ -1609,12 +1623,15 @@ type midDeliveryFlow struct { action func(orch orchestrator.Orchestrator) error } -// runKafkaMidDeliveryAction is the shared driver behind the kafka in-flight -// crash and cert-rotation flows: both bring everything up with the receiver -// live, wait until the receiver has seen half the records, fire one disruptive -// action, then drain and apply the same no-loss / at-least-once verdict. Only -// the action (and a little setup/labelling) differs — see midDeliveryFlow. -func (r *Runner) runKafkaMidDeliveryAction(tc *config.TestCase, subject config.Subject, f midDeliveryFlow) (results.RunResult, error) { +// runMidDeliveryAction is the shared driver behind the receiver-up mid-delivery +// flows (kafka in-flight crash, kafka cert rotation, and object-storage +// in-flight crash): all bring everything up with the receiver live, wait until +// the receiver has seen half the records, fire one disruptive action, then +// drain and apply the same no-loss / at-least-once verdict. Only the action +// (and a little setup/labelling) differs — see midDeliveryFlow. The source is +// decoupled from the subject (broker or bucket), so the generator keeps +// producing across a subject restart. +func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subject, f midDeliveryFlow) (results.RunResult, error) { configName := r.opts.ConfigName subject = r.applySubjectOverrides(subject) @@ -1878,7 +1895,7 @@ func rotateAndReload(orch orchestrator.Orchestrator, service string, rotate func // offset-committed are re-consumed on restart. Verdict: no loss; duplicates are // reported, not failed. func (r *Runner) runKafkaInflightCrash(tc *config.TestCase, subject config.Subject) (results.RunResult, error) { - return r.runKafkaMidDeliveryAction(tc, subject, midDeliveryFlow{ + return r.runMidDeliveryAction(tc, subject, midDeliveryFlow{ verdictLabel: "kafka in-flight crash correctness", actionLog: "SIGKILL subject (no graceful shutdown), then restart", overDelivNote: "expected for a mid-delivery crash", @@ -1899,6 +1916,35 @@ func (r *Runner) runKafkaInflightCrash(tc *config.TestCase, subject config.Subje }) } +// runInflightCrashCorrectness SIGKILLs the subject WHILE it is actively +// delivering to a live receiver, then restarts it — the receiver-up +// mid-delivery worst case for any source whose store is decoupled from the +// subject (e.g. the S3/Azure bucket a poll-mode device lists). The since-cursor +// advances at cycle end with no per-object delivery commit, so recovery must +// still lose nothing; duplicates from the crash-resistance replay + re-list are +// reported, not failed. +func (r *Runner) runInflightCrashCorrectness(tc *config.TestCase, subject config.Subject) (results.RunResult, error) { + return r.runMidDeliveryAction(tc, subject, midDeliveryFlow{ + verdictLabel: "in-flight crash correctness", + actionLog: "SIGKILL subject (no graceful shutdown), then restart", + overDelivNote: "expected for a mid-delivery crash", + totalLinesErr: "persistence_inflight_crash_correctness requires generator.total_lines > 0", + action: func(orch orchestrator.Orchestrator) error { + if err := orch.KillServices("subject"); err != nil { + return fmt.Errorf("killing subject: %w", err) + } + // Settle before the subject restarts and re-lists from its cursor. + if err := sleepCtx(r.ctx, 3*time.Second); err != nil { + return fmt.Errorf("interrupted: %w", err) + } + if err := orch.UpServices("subject"); err != nil { + return fmt.Errorf("restarting subject: %w", err) + } + return nil + }, + }) +} + // runKafkaCertRotation verifies the subject's broker-cert handling over mTLS in // TWO halves, so the run fails if EITHER property breaks: // @@ -1919,7 +1965,7 @@ func (r *Runner) runKafkaCertRotation(tc *config.TestCase, subject config.Subjec // and prepare runs before action, so the capture is well-ordered. var certsDir string hosts := []string{"subject", "localhost", "redpanda"} - return r.runKafkaMidDeliveryAction(tc, subject, midDeliveryFlow{ + return r.runMidDeliveryAction(tc, subject, midDeliveryFlow{ verdictLabel: "kafka cert rotation correctness", actionLog: "rotating broker cert to an UNTRUSTED CA (must be rejected), then back to a trusted cert", overDelivNote: "expected across the broker reconnects", @@ -7178,7 +7224,7 @@ func (r *Runner) runSyslogVaultCertRotation(tc *config.TestCase, subject config. hosts := []string{"subject", "localhost"} var certsDir string - return r.runKafkaMidDeliveryAction(tc, subject, midDeliveryFlow{ + return r.runMidDeliveryAction(tc, subject, midDeliveryFlow{ verdictLabel: "syslog TLS vault cert rotation correctness", actionLog: "rotating syslog server cert to UNTRUSTED CA (generator TLS must fail), then restoring trusted cert", overDelivNote: "expected after the trusted cert is restored and the generator reconnects", @@ -7222,7 +7268,7 @@ func (r *Runner) runSyslogVaultCertRotation(tc *config.TestCase, subject config. mount := tc.Vault.MountOrDefault() token := tc.Vault.TokenOrDefault() - // Receiver metrics port is already forwarded by runKafkaMidDeliveryAction. + // Receiver metrics port is already forwarded by runMidDeliveryAction. metricsPort := orch.ReceiverMetricsPorts()["default"] // ---- Phase 1: UNTRUSTED cert — generator TLS must fail ---- From 00d90b5abbcb8213fb2c709ffa93668f60c3568f Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Tue, 14 Jul 2026 19:25:25 +0300 Subject: [PATCH 3/7] fix(runner): add missing container cleanup, correct function comment and correct formatting --- internal/runner/runner.go | 77 ++++++++++++++++++++++++--------------- 1 file changed, 48 insertions(+), 29 deletions(-) diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 252bc19..dc2e70d 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -804,7 +804,8 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe result.FailReason = fmt.Sprintf( "expect_failure: data path was NOT blocked — receiver observed %s line(s) (> %s); "+ "the control under test (e.g. auth) appears bypassed", - formatCount(recvMetrics.LinesReceived), formatCount(cap)) + formatCount(recvMetrics.LinesReceived), formatCount(cap), + ) } } else if tc.IsCorrectnessType() && !tc.HasGenerator() { // No generator: there's no expected line count to derive loss or @@ -824,7 +825,8 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe if !gotEnough { failReasons = append(failReasons, fmt.Sprintf( "expected >= %s received records, got %s", - formatCount(minRecv), formatCount(recvMetrics.LinesReceived))) + formatCount(minRecv), formatCount(recvMetrics.LinesReceived), + )) } passed := gotEnough && recvOK result.Passed = &passed @@ -863,13 +865,15 @@ func (r *Runner) Run(tc *config.TestCase, subject config.Subject) (results.RunRe if !lossOK { failReasons = append(failReasons, fmt.Sprintf( "expected loss <= %.2f%%, got %.2f%%", - tc.Correctness.ExpectedLossPct, lossPct)) + tc.Correctness.ExpectedLossPct, lossPct, + )) } if !overOK { extra := recvMetrics.LinesReceived - expectedOut failReasons = append(failReasons, fmt.Sprintf( "over-delivery: received %s lines but only %s were expected (%s extra/duplicate lines)", - formatCount(recvMetrics.LinesReceived), formatCount(expectedOut), formatCount(extra))) + formatCount(recvMetrics.LinesReceived), formatCount(expectedOut), formatCount(extra), + )) } passed := lossOK && overOK && recvOK @@ -1759,8 +1763,8 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec return results.RunResult{}, fmt.Errorf("never reached mid-delivery (%s) before timeout", formatCount(mid)) } - // The generator produces to Kafka independently of the subject; collect - // its final count. + // The generator produces to the durable source (broker or bucket) + // independently of the subject; collect its final count. duration := tc.DurationOrDefault(60 * time.Second) warmup := tc.WarmupOrDefault(30 * time.Second) genTimeout := min(duration+warmup+2*time.Minute, r.opts.Timeout) @@ -1929,6 +1933,7 @@ func (r *Runner) runInflightCrashCorrectness(tc *config.TestCase, subject config actionLog: "SIGKILL subject (no graceful shutdown), then restart", overDelivNote: "expected for a mid-delivery crash", totalLinesErr: "persistence_inflight_crash_correctness requires generator.total_lines > 0", + extraCleanup: []string{"bench-localstack", "bench-azurite", "bench-azure-init"}, action: func(orch orchestrator.Orchestrator) error { if err := orch.KillServices("subject"); err != nil { return fmt.Errorf("killing subject: %w", err) @@ -2715,7 +2720,8 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config if !blockedOK { errs = append(errs, fmt.Sprintf( "agent was NOT blocked before rotation (%s records during the %s baseline) — the initial allowlist already admits it; the recover transition is untested", - formatCount(rmBlocked.LinesReceived), baseline)) + formatCount(rmBlocked.LinesReceived), baseline, + )) } // Phase 1: rotate to the allow config; the director's refreshACL must pick @@ -2755,7 +2761,8 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config if !started { errs = append(errs, fmt.Sprintf( "delivery did not start after admitting the agent — %s records (expected >= %s); the ACL hot-reload did not take effect", - formatCount(finalCount), formatCount(minRecv))) + formatCount(finalCount), formatCount(minRecv), + )) } passed = blockedOK && started @@ -2788,7 +2795,8 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config return results.RunResult{}, fmt.Errorf( "agent never delivered the initial %s records — cannot test revocation (subject image %s:%s); "+ "if this is not the agent-capable image, re-run with VMETRIC_IMAGE=vmetric/director-enterprise", - formatCount(minRecv), subject.Image, subject.Version) + formatCount(minRecv), subject.Image, subject.Version, + ) } fmt.Printf(" delivery established — %s records before revocation ✓\n", formatCount(beforeCount)) @@ -2828,11 +2836,13 @@ func (r *Runner) runDirectorAgentACLRotation(tc *config.TestCase, subject config if advanced < 0 { errs = append(errs, fmt.Sprintf( "inconclusive: receiver count decreased after blocking the agent (after1=%s after2=%s) — the counter regressed, cannot confirm delivery stopped", - formatCount(rmAfter1.LinesReceived), formatCount(rmAfter2.LinesReceived))) + formatCount(rmAfter1.LinesReceived), formatCount(rmAfter2.LinesReceived), + )) } else if !stopped { errs = append(errs, fmt.Sprintf( "delivery did NOT stop after blocking the agent (%s new records across the block window) — the ACL was not enforced on the live data path", - formatCount(advanced))) + formatCount(advanced), + )) } passed = stopped @@ -3718,6 +3728,7 @@ func (st *fleetStatus) count(id, key string) int { } return d.Inbound[key].Count } + func (st *fleetStatus) lastData(id, key string) string { d, ok := st.Directors[id] if !ok { @@ -4474,7 +4485,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if leakedDuringBaseline > cfgUpdateLeak { errs = append(errs, fmt.Sprintf( "delivery was NOT suppressed by the BEFORE config (%s records during the %s baseline) — the update verdict would be vacuous", - formatCount(leakedDuringBaseline), baseline)) + formatCount(leakedDuringBaseline), baseline, + )) } // Phase 2: AFTER pipeline (configs/update.vmf) — enables delivery. @@ -4514,7 +4526,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if delivered < minRecv { errs = append(errs, fmt.Sprintf( "delivery did not start after the config update — %s new records (expected >= %s); the update did not take effect on the data plane", - formatCount(delivered), formatCount(minRecv))) + formatCount(delivered), formatCount(minRecv), + )) } else { fmt.Printf(" delivery started after the update — %s new records ✓\n", formatCount(delivered)) } @@ -4571,7 +4584,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if rmB.LinesReceived > leak { errs = append(errs, fmt.Sprintf( "delivery was not blocked before the change (%s received) — config A is not pointing at a dead target; the change is untested", - formatCount(rmB.LinesReceived))) + formatCount(rmB.LinesReceived), + )) } // Phase 1: push the changed VMF (config B → working target) over the fleet link. @@ -4626,7 +4640,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if finalCount < target { errs = append(errs, fmt.Sprintf( "delivery did not start after the pushed config change — %s received (expected >= %s = baseline %s + %s); the fleet-delivered VMF change did not reconcile onto the live data path", - formatCount(finalCount), formatCount(target), formatCount(rmB.LinesReceived), formatCount(minRecv))) + formatCount(finalCount), formatCount(target), formatCount(rmB.LinesReceived), formatCount(minRecv), + )) } else { fmt.Printf(" delivery started after the VMF config change (%s received, baseline %s) ✓\n", formatCount(finalCount), formatCount(rmB.LinesReceived)) } @@ -4844,7 +4859,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if finalCount < minRecv { errs = append(errs, fmt.Sprintf( "delivered pipeline produced too few records at the receiver — %s received (expected >= %s); the VMF-delivered pipeline/library did not apply on the data path", - formatCount(finalCount), formatCount(minRecv))) + formatCount(finalCount), formatCount(minRecv), + )) } else { fmt.Printf(" delivered pipeline produced %s record(s) at the receiver ✓\n", formatCount(finalCount)) } @@ -4964,8 +4980,8 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf // saveFleetResult records the verdict and, on failure, dumps short log tails of // the simulator and director to aid diagnosis. func (r *Runner) saveFleetResult(tc *config.TestCase, subject config.Subject, configName string, - startTime time.Time, finalCount int64, passed bool, errs []string, simContainer, subjectContainer string) (results.RunResult, error) { - + startTime time.Time, finalCount int64, passed bool, errs []string, simContainer, subjectContainer string, +) (results.RunResult, error) { elapsed := time.Since(startTime).Seconds() label := tc.Fleet.Scenario if passed { @@ -5340,8 +5356,8 @@ func (r *Runner) runCCFCorrectness(tc *config.TestCase, subject config.Subject) // saveCCFResult records the CCF verdict and, on failure, dumps short log tails of // the mock API and director to aid diagnosis. func (r *Runner) saveCCFResult(tc *config.TestCase, subject config.Subject, configName string, - startTime time.Time, finalCount int64, passed bool, errs []string, apiContainer, subjectContainer string) (results.RunResult, error) { - + startTime time.Time, finalCount int64, passed bool, errs []string, apiContainer, subjectContainer string, +) (results.RunResult, error) { elapsed := time.Since(startTime).Seconds() label := tc.CCF.Scenario if passed { @@ -5603,8 +5619,8 @@ func (r *Runner) runHTTPSourceCorrectness(tc *config.TestCase, subject config.Su } func (r *Runner) saveHTTPSourceResult(tc *config.TestCase, subject config.Subject, configName string, - startTime time.Time, finalCount int64, passed bool, errs []string, senderContainer, subjectContainer string) (results.RunResult, error) { - + startTime time.Time, finalCount int64, passed bool, errs []string, senderContainer, subjectContainer string, +) (results.RunResult, error) { elapsed := time.Since(startTime).Seconds() if passed { fmt.Printf(" http source (%s): PASSED ✓\n", tc.HTTPSource.Scenario) @@ -5759,7 +5775,8 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con // (tabseparatedraw/parquet/native write only `message`; the rest DEFAULT). createTable := fmt.Sprintf( "CREATE TABLE IF NOT EXISTS %s (message String, `@timestamp` String DEFAULT '', ts DateTime DEFAULT now()) ENGINE = MergeTree ORDER BY ts", - fqTable) + fqTable, + ) if ch.CreateTableSQL != "" { createTable = strings.ReplaceAll(ch.CreateTableSQL, "{{TABLE}}", fqTable) } @@ -5848,8 +5865,8 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con } func (r *Runner) saveClickHouseTargetResult(tc *config.TestCase, subject config.Subject, configName string, - startTime time.Time, finalCount int64, passed bool, errs []string, chContainer, subjectContainer string) (results.RunResult, error) { - + startTime time.Time, finalCount int64, passed bool, errs []string, chContainer, subjectContainer string, +) (results.RunResult, error) { elapsed := time.Since(startTime).Seconds() if passed { fmt.Println(" clickhouse target: PASSED ✓") @@ -6231,8 +6248,8 @@ func (r *Runner) setupAuxRun(tc *config.TestCase, subject config.Subject, config // saveAuxResult records the verdict for the small aux-container drivers. func (r *Runner) saveAuxResult(tc *config.TestCase, subject config.Subject, configName, label string, - startTime time.Time, finalCount int64, passed bool, errs []string, auxContainer, subjectContainer string) (results.RunResult, error) { - + startTime time.Time, finalCount int64, passed bool, errs []string, auxContainer, subjectContainer string, +) (results.RunResult, error) { elapsed := time.Since(startTime).Seconds() if passed { fmt.Printf(" %s: PASSED ✓\n", label) @@ -6707,7 +6724,8 @@ func (r *Runner) runKafkaOffsetCommitRestart(tc *config.TestCase, subject config if overPct > tc.Correctness.MaxOverDeliveryPct { errors = append(errors, fmt.Sprintf( "expected over-delivery <= %.2f%%, got %.2f%% (%s duplicate lines) — restart re-consumed records whose offsets should have been committed", - tc.Correctness.MaxOverDeliveryPct, overPct, formatCount(extra))) + tc.Correctness.MaxOverDeliveryPct, overPct, formatCount(extra), + )) } passed := len(errors) == 0 @@ -7322,7 +7340,8 @@ func (r *Runner) runSyslogVaultCertRotation(tc *config.TestCase, subject config. "SECURITY: receiver count never stalled after wrong-CA rotation "+ "(last count: %d) — director did not serve the untrusted cert, "+ "or generator ignored cert validation; check debug.console.status: true logs", - lastCount) + lastCount, + ) } // ---- Phase 2: TRUSTED cert restored — recovery ---- From fc4fbb415149df4aa294a645b939d8a729ab083c Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Thu, 16 Jul 2026 18:33:01 +0300 Subject: [PATCH 4/7] feat(runner): aggregate resource metrics for crash/restart --- internal/results/compare.go | 14 +++ internal/results/compare_test.go | 103 ++++++++++++++++++++++ internal/runner/runner.go | 143 ++++++++++++++++++++++++++++++- 3 files changed, 256 insertions(+), 4 deletions(-) create mode 100644 internal/results/compare_test.go diff --git a/internal/results/compare.go b/internal/results/compare.go index 8d17b3c..e4cba83 100644 --- a/internal/results/compare.go +++ b/internal/results/compare.go @@ -279,6 +279,20 @@ func AggregateAllMetricsFromCSVWindow(csvPath string, startNs, endNs int64) (Agg } } + // A sample with zero CPU AND zero memory is a Docker stats snapshot + // of a stopped container (the crash/restart drivers stop and restart + // the subject mid-run, and Docker returns zeroed stats for an exited + // container). A running process never has 0 RSS, so real idle samples + // (cpu 0, mem > 0) are kept; counting stopped-window rows would + // dilute the averages with time the subject wasn't running. + if cpuIdx >= 0 && cpuIdx < len(record) && memIdx >= 0 && memIdx < len(record) { + cpuV, _ := strconv.ParseFloat(record[cpuIdx], 64) + memV, _ := strconv.ParseFloat(record[memIdx], 64) + if cpuV == 0 && memV == 0 { + continue + } + } + if cpuIdx >= 0 && cpuIdx < len(record) { v, _ := strconv.ParseFloat(record[cpuIdx], 64) cpuSum += v diff --git a/internal/results/compare_test.go b/internal/results/compare_test.go new file mode 100644 index 0000000..3a162d5 --- /dev/null +++ b/internal/results/compare_test.go @@ -0,0 +1,103 @@ +package results + +import ( + "os" + "path/filepath" + "testing" + "time" +) + +func writeMetricsCSV(t *testing.T, rows string) string { + t.Helper() + path := filepath.Join(t.TempDir(), "metrics.csv") + header := "epoch,cpu_usr,mem_used,net_recv,net_send,dsk_read,dsk_writ,load_avg1,load_avg5,load_avg15\n" + if err := os.WriteFile(path, []byte(header+rows), 0o644); err != nil { + t.Fatalf("writing csv: %v", err) + } + return path +} + +func TestAggregateAllMetricsFromCSV(t *testing.T) { + t.Parallel() + + // Stopped-container samples (cpu 0 AND mem 0) are dropped so the + // crash/restart down window doesn't dilute the averages. + t.Run("all-zero rows skipped", func(t *testing.T) { + t.Parallel() + csv := writeMetricsCSV(t, + "100,0,0,0,0,0,0,0,0,0\n"+ + "101,10,104857600,0,0,0,0,0,0,0\n"+ + "102,0,0,0,0,0,0,0,0,0\n"+ + "103,20,209715200,0,0,0,0,0,0,0\n") + m, err := AggregateAllMetricsFromCSV(csv) + if err != nil { + t.Fatalf("aggregate: %v", err) + } + if m.Samples != 2 { + t.Fatalf("Samples = %d, want 2", m.Samples) + } + if m.CPUAvg != 15 || m.CPUMax != 20 { + t.Fatalf("cpu avg/max = %v/%v, want 15/20", m.CPUAvg, m.CPUMax) + } + if m.MemAvgMB != 150 || m.MemMaxMB != 200 { + t.Fatalf("mem avg/max = %v/%v, want 150/200", m.MemAvgMB, m.MemMaxMB) + } + }) + + // A running-but-idle sample (cpu 0, mem > 0) is a real sample and stays. + t.Run("idle row kept", func(t *testing.T) { + t.Parallel() + csv := writeMetricsCSV(t, + "100,0,104857600,0,0,0,0,0,0,0\n"+ + "101,10,104857600,0,0,0,0,0,0,0\n") + m, err := AggregateAllMetricsFromCSV(csv) + if err != nil { + t.Fatalf("aggregate: %v", err) + } + if m.Samples != 2 { + t.Fatalf("Samples = %d, want 2", m.Samples) + } + if m.CPUAvg != 5 { + t.Fatalf("CPUAvg = %v, want 5", m.CPUAvg) + } + }) + + // Net/disk totals accumulate only over surviving rows. + t.Run("io totals over kept rows", func(t *testing.T) { + t.Parallel() + csv := writeMetricsCSV(t, + "100,0,0,999,999,999,999,0,0,0\n"+ + "101,10,104857600,100,200,300,400,0,0,0\n") + m, err := AggregateAllMetricsFromCSV(csv) + if err != nil { + t.Fatalf("aggregate: %v", err) + } + if m.NetRecv != 100 || m.NetSend != 200 || m.DiskRead != 300 || m.DiskWrite != 400 { + t.Fatalf("io = net %d/%d disk %d/%d, want 100/200 300/400", + m.NetRecv, m.NetSend, m.DiskRead, m.DiskWrite) + } + }) +} + +func TestAggregateAllMetricsFromCSVWindow(t *testing.T) { + t.Parallel() + + // The epoch window filter still applies alongside the zero-row skip. + csv := writeMetricsCSV(t, + "100,10,104857600,0,0,0,0,0,0,0\n"+ + "200,0,0,0,0,0,0,0,0,0\n"+ + "201,30,314572800,0,0,0,0,0,0,0\n"+ + "300,50,524288000,0,0,0,0,0,0,0\n") + startNs := int64(200) * int64(time.Second) + endNs := int64(250) * int64(time.Second) + m, err := AggregateAllMetricsFromCSVWindow(csv, startNs, endNs) + if err != nil { + t.Fatalf("aggregate: %v", err) + } + if m.Samples != 1 { + t.Fatalf("Samples = %d, want 1 (window keeps 200-201, zero row dropped)", m.Samples) + } + if m.CPUAvg != 30 || m.MemMaxMB != 300 { + t.Fatalf("cpu/mem = %v/%v, want 30/300", m.CPUAvg, m.MemMaxMB) + } +} diff --git a/internal/runner/runner.go b/internal/runner/runner.go index dc2e70d..487aaed 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -1194,6 +1194,9 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S return results.RunResult{}, fmt.Errorf("querying receiver metrics: %w", err) } + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() + elapsed := time.Since(startTime).Seconds() // Compute results @@ -1258,14 +1261,31 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S LinesOut: recvMetrics.LinesReceived, BytesIn: genStats.BytesSent, BytesOut: recvMetrics.BytesReceived, + LinesPerSec: receiveWindowRate(recvMetrics), LossPercent: lossPct, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, Passed: &passed, } if !passed { result.FailReason = strings.Join(errors, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -1273,6 +1293,11 @@ func (r *Runner) runPersistenceCorrectness(tc *config.TestCase, subject config.S fmt.Printf(" done. results → %s\n", dir) fmt.Printf(" lines sent: %s lines received: %s loss: %.2f%%\n", formatCount(genStats.LinesSent), formatCount(recvMetrics.LinesReceived), lossPct) + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) + if metrics.IOThroughputAvg > 0 { + fmt.Printf(" io throughput: avg %.1f MB/s\n", metrics.IOThroughputAvg/(1024*1024)) + } if tc.Correctness.ValidateDedup { fmt.Printf(" unique lines: %s duplicates: %s\n", formatCount(recvMetrics.UniqueLines), formatCount(recvMetrics.Duplicates)) @@ -1491,6 +1516,9 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject return results.RunResult{}, fmt.Errorf("querying receiver metrics: %w", err) } + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() + elapsed := time.Since(startTime).Seconds() lossPct := 0.0 @@ -1554,14 +1582,31 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject LinesOut: recvMetrics.LinesReceived, BytesIn: genStats.BytesSent, BytesOut: recvMetrics.BytesReceived, + LinesPerSec: receiveWindowRate(recvMetrics), LossPercent: lossPct, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, Passed: &passed, } if !passed { result.FailReason = strings.Join(errors, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -1569,6 +1614,11 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject fmt.Printf(" done. results → %s\n", dir) fmt.Printf(" lines sent: %s lines received: %s loss: %.2f%%\n", formatCount(genStats.LinesSent), formatCount(recvMetrics.LinesReceived), lossPct) + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) + if metrics.IOThroughputAvg > 0 { + fmt.Printf(" io throughput: avg %.1f MB/s\n", metrics.IOThroughputAvg/(1024*1024)) + } if tc.Correctness.ValidateDedup { fmt.Printf(" unique lines: %s duplicates: %s\n", formatCount(recvMetrics.UniqueLines), formatCount(recvMetrics.Duplicates)) @@ -1597,6 +1647,41 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject return result, nil } +// harvestResourceMetrics copies the collector's CSV, stops the collector, and +// aggregates resource usage over the whole run. Used by the crash/restart +// drivers, whose subject stops and restarts mid-run: the collector samples by +// container name each tick, so it rides across the restart; samples taken +// while the subject is down are all-zero rows the aggregator drops. Returns +// the aggregated metrics and the local CSV path ("" if unavailable) for +// saveResult. +func (r *Runner) harvestResourceMetrics(orch orchestrator.Orchestrator, tmpDir string) (results.AggregatedMetrics, string) { + csvPath := filepath.Join(tmpDir, "metrics.csv") + fmt.Println(" collecting metrics…") + if err := orch.CopyMetricsCSV(csvPath); err != nil { + fmt.Fprintf(os.Stderr, " warning: metrics CSV not available: %v\n", err) + csvPath = "" + } + fmt.Println(" stopping collector…") + if err := orch.StopCollector(); err != nil { + fmt.Fprintf(os.Stderr, " warning: stopping collector: %v\n", err) + } + var metrics results.AggregatedMetrics + if csvPath != "" { + metrics, _ = results.AggregateAllMetricsFromCSV(csvPath) + } + return metrics, csvPath +} + +// receiveWindowRate returns lines/sec over the receiver's active window +// (first→last received). Crash/restart delivery is a burst after the +// disruption, so lines/total-run-time would understate to near zero. +func receiveWindowRate(rm ReceiverMetrics) float64 { + if rm.LastReceivedNs > rm.FirstReceivedNs { + return float64(rm.LinesReceived) / (float64(rm.LastReceivedNs-rm.FirstReceivedNs) / 1e9) + } + return 0 +} + // midDeliveryFlow parameterizes the shared mid-delivery correctness driver // (runMidDeliveryAction): produce to a durable source (Kafka broker, S3/Azure // bucket) with the receiver live, fire one disruptive action once the receiver @@ -1810,6 +1895,9 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec return results.RunResult{}, fmt.Errorf("querying receiver metrics: %w", err) } + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() + elapsed := time.Since(startTime).Seconds() lossPct := 0.0 if genStats.LinesSent > 0 { @@ -1835,6 +1923,11 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec fmt.Printf(" lines sent: %s lines received: %s loss: %.2f%%\n", formatCount(genStats.LinesSent), formatCount(recvMetrics.LinesReceived), lossPct) + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) + if metrics.IOThroughputAvg > 0 { + fmt.Printf(" io throughput: avg %.1f MB/s\n", metrics.IOThroughputAvg/(1024*1024)) + } if passed { fmt.Printf(" %s: PASSED ✓\n", f.verdictLabel) } else { @@ -1857,7 +1950,24 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec LinesOut: recvMetrics.LinesReceived, BytesIn: genStats.BytesSent, BytesOut: recvMetrics.BytesReceived, + LinesPerSec: receiveWindowRate(recvMetrics), LossPercent: lossPct, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, Passed: &passed, } if !passed { @@ -1866,7 +1976,7 @@ func (r *Runner) runMidDeliveryAction(tc *config.TestCase, subject config.Subjec // Persist the result like every other run path — Run's contract is to // return the *persisted* result. - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -6954,6 +7064,9 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje lastCount = rm.LinesReceived } + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() + // Evaluate. elapsed := time.Since(startTime).Seconds() lossPct := 0.0 @@ -6997,14 +7110,31 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje LinesOut: recvMetrics.LinesReceived, BytesIn: genStats.BytesSent, BytesOut: recvMetrics.BytesReceived, + LinesPerSec: receiveWindowRate(recvMetrics), LossPercent: lossPct, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, Passed: &passed, } if !passed { result.FailReason = strings.Join(perrs, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -7012,6 +7142,11 @@ func (r *Runner) runPersistenceFileRestartCorrectness(tc *config.TestCase, subje fmt.Printf(" done. results → %s\n", dir) fmt.Printf(" lines sent: %s lines received: %s loss: %.2f%%\n", formatCount(genStats.LinesSent), formatCount(recvMetrics.LinesReceived), lossPct) + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) + if metrics.IOThroughputAvg > 0 { + fmt.Printf(" io throughput: avg %.1f MB/s\n", metrics.IOThroughputAvg/(1024*1024)) + } if tc.Correctness.ValidateDedup { fmt.Printf(" unique lines: %s duplicates: %s\n", formatCount(recvMetrics.UniqueLines), formatCount(recvMetrics.Duplicates)) From 5210813441c954fe83a730cbe8ed5712eb6d456a Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Thu, 16 Jul 2026 21:20:07 +0300 Subject: [PATCH 5/7] fix(results): skip only truly idle containers in aggregation, not malformed rows --- internal/results/compare.go | 11 +++++++---- internal/results/compare_test.go | 20 ++++++++++++++++++++ internal/runner/runner.go | 9 ++++++--- 3 files changed, 33 insertions(+), 7 deletions(-) diff --git a/internal/results/compare.go b/internal/results/compare.go index e4cba83..67092e9 100644 --- a/internal/results/compare.go +++ b/internal/results/compare.go @@ -284,11 +284,14 @@ func AggregateAllMetricsFromCSVWindow(csvPath string, startNs, endNs int64) (Agg // the subject mid-run, and Docker returns zeroed stats for an exited // container). A running process never has 0 RSS, so real idle samples // (cpu 0, mem > 0) are kept; counting stopped-window rows would - // dilute the averages with time the subject wasn't running. + // dilute the averages with time the subject wasn't running. An + // unparseable field is not evidence of a stopped container, so rows + // with malformed cpu/mem values are kept (they contribute 0 to the + // cpu/mem sums below, as they always have). if cpuIdx >= 0 && cpuIdx < len(record) && memIdx >= 0 && memIdx < len(record) { - cpuV, _ := strconv.ParseFloat(record[cpuIdx], 64) - memV, _ := strconv.ParseFloat(record[memIdx], 64) - if cpuV == 0 && memV == 0 { + cpuV, cpuErr := strconv.ParseFloat(record[cpuIdx], 64) + memV, memErr := strconv.ParseFloat(record[memIdx], 64) + if cpuErr == nil && memErr == nil && cpuV == 0 && memV == 0 { continue } } diff --git a/internal/results/compare_test.go b/internal/results/compare_test.go index 3a162d5..0e14a2c 100644 --- a/internal/results/compare_test.go +++ b/internal/results/compare_test.go @@ -62,6 +62,26 @@ func TestAggregateAllMetricsFromCSV(t *testing.T) { } }) + // An unparseable cpu/mem field is not evidence of a stopped container: + // the row is kept and its valid net/disk columns still count. + t.Run("malformed cpu field kept", func(t *testing.T) { + t.Parallel() + csv := writeMetricsCSV(t, + "100,x,0,100,200,300,400,0,0,0\n"+ + "101,0,0,0,0,0,0,0,0,0\n") + m, err := AggregateAllMetricsFromCSV(csv) + if err != nil { + t.Fatalf("aggregate: %v", err) + } + if m.Samples != 1 { + t.Fatalf("Samples = %d, want 1 (malformed row kept, all-zero row dropped)", m.Samples) + } + if m.NetRecv != 100 || m.NetSend != 200 || m.DiskRead != 300 || m.DiskWrite != 400 { + t.Fatalf("io = net %d/%d disk %d/%d, want 100/200 300/400", + m.NetRecv, m.NetSend, m.DiskRead, m.DiskWrite) + } + }) + // Net/disk totals accumulate only over surviving rows. t.Run("io totals over kept rows", func(t *testing.T) { t.Parallel() diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 487aaed..58f4984 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -1651,9 +1651,12 @@ func (r *Runner) runPersistenceShutdownCorrectness(tc *config.TestCase, subject // aggregates resource usage over the whole run. Used by the crash/restart // drivers, whose subject stops and restarts mid-run: the collector samples by // container name each tick, so it rides across the restart; samples taken -// while the subject is down are all-zero rows the aggregator drops. Returns -// the aggregated metrics and the local CSV path ("" if unavailable) for -// saveResult. +// while the subject is down are all-zero rows the aggregator drops. The +// copy-then-stop order is deliberate and matches the standard driver: the +// collector fsyncs every row as it writes it (see CopyMetricsCSV) and has no +// shutdown flush, so copying first loses nothing and excludes post-cutoff +// idle samples from the stop grace. Returns the aggregated metrics and the +// local CSV path ("" if unavailable) for saveResult. func (r *Runner) harvestResourceMetrics(orch orchestrator.Orchestrator, tmpDir string) (results.AggregatedMetrics, string) { csvPath := filepath.Join(tmpDir, "metrics.csv") fmt.Println(" collecting metrics…") From ff5c9205b1165fc1825f951cd653b4bafe672e6a Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Fri, 17 Jul 2026 00:35:19 +0300 Subject: [PATCH 6/7] fix: support multi-container metric collection --- containers/collector/main.go | 61 ++++++++++++--- containers/collector/main_test.go | 20 +++++ containers/receiver/main.go | 34 ++++++++- containers/receiver/main_test.go | 85 +++++++++++++++++++++ internal/orchestrator/docker.go | 22 +++++- internal/runner/runner.go | 123 ++++++++++++++++++++++-------- 6 files changed, 296 insertions(+), 49 deletions(-) create mode 100644 containers/receiver/main_test.go diff --git a/containers/collector/main.go b/containers/collector/main.go index bbf8aea..220c1e5 100644 --- a/containers/collector/main.go +++ b/containers/collector/main.go @@ -100,11 +100,25 @@ func main() { } func runDockerMode(ctx context.Context, ticker *time.Ticker, emit func(MetricsRow)) { - container := mustEnv("COLLECTOR_TARGET_CONTAINER") + // COLLECTOR_TARGET_CONTAINER is a comma-separated container list. Single- + // subject cases pass one name; director-cluster cases pass every node so + // each row captures the whole cluster's footprint (per-tick sum across + // nodes). Deltas (net/blkio) are computed per container against its own + // previous sample, then summed. + var targets []string + for _, t := range strings.Split(mustEnv("COLLECTOR_TARGET_CONTAINER"), ",") { + if t = strings.TrimSpace(t); t != "" { + targets = append(targets, t) + } + } + if len(targets) == 0 { + fmt.Fprintln(os.Stderr, "collector: COLLECTOR_TARGET_CONTAINER has no container names") + os.Exit(1) + } dockerHost := getEnv("DOCKER_HOST", "unix:///var/run/docker.sock") httpClient := buildClient(dockerHost) - var prev *dockerStats + prev := make([]*dockerStats, len(targets)) var errCount int for { @@ -112,24 +126,47 @@ func runDockerMode(ctx context.Context, ticker *time.Ticker, emit func(MetricsRo case <-ctx.Done(): return case <-ticker.C: - stats, err := fetchDockerStats(httpClient, dockerHost, container) - if err != nil { - errCount++ - // Log the first error immediately, then every 30 ticks, - // so a silent permission/DNS/container-missing failure is - // diagnosable but we don't spam the run. - if errCount == 1 || errCount%30 == 0 { - fmt.Fprintf(os.Stderr, "collector: fetch stats (err #%d): %v\n", errCount, err) + var row MetricsRow + sampled := false + for i, container := range targets { + stats, err := fetchDockerStats(httpClient, dockerHost, container) + if err != nil { + errCount++ + // Log the first error immediately, then every 30 ticks, + // so a silent permission/DNS/container-missing failure is + // diagnosable but we don't spam the run. + if errCount == 1 || errCount%30 == 0 { + fmt.Fprintf(os.Stderr, "collector: fetch stats %s (err #%d): %v\n", container, errCount, err) + } + continue } + addRow(&row, dockerStatsToRow(stats, prev[i])) + prev[i] = stats + sampled = true + } + if !sampled { continue } - row := dockerStatsToRow(stats, prev) + row.Epoch = time.Now().Unix() + row.CpuIdl = 100.0 - row.CpuUsr emit(row) - prev = stats } } } +// addRow accumulates one container's sample into the tick's combined row. +// Epoch and CpuIdl are derived once per tick by the caller, not summed. +func addRow(dst *MetricsRow, s MetricsRow) { + dst.CpuUsr += s.CpuUsr + dst.MemUsed += s.MemUsed + dst.MemCach += s.MemCach + dst.MemFree += s.MemFree + dst.NetRecv += s.NetRecv + dst.NetSend += s.NetSend + dst.DskRead += s.DskRead + dst.DskWrit += s.DskWrit +} + // --- Docker Stats API --- type dockerStats struct { diff --git a/containers/collector/main_test.go b/containers/collector/main_test.go index 5dc7c20..8b6b2f4 100644 --- a/containers/collector/main_test.go +++ b/containers/collector/main_test.go @@ -100,3 +100,23 @@ func TestMemUsage(t *testing.T) { }) } } + +func TestAddRow(t *testing.T) { + var dst MetricsRow + addRow(&dst, MetricsRow{CpuUsr: 40, MemUsed: 100, MemCach: 10, MemFree: 50, NetRecv: 1, NetSend: 2, DskRead: 3, DskWrit: 4}) + addRow(&dst, MetricsRow{CpuUsr: 15, MemUsed: 200, MemCach: 20, MemFree: 70, NetRecv: 10, NetSend: 20, DskRead: 30, DskWrit: 40}) + + if dst.CpuUsr != 55 { + t.Errorf("CpuUsr = %v, want 55", dst.CpuUsr) + } + if dst.MemUsed != 300 || dst.MemCach != 30 || dst.MemFree != 120 { + t.Errorf("mem = %d/%d/%d, want 300/30/120", dst.MemUsed, dst.MemCach, dst.MemFree) + } + if dst.NetRecv != 11 || dst.NetSend != 22 || dst.DskRead != 33 || dst.DskWrit != 44 { + t.Errorf("io = %d/%d/%d/%d, want 11/22/33/44", dst.NetRecv, dst.NetSend, dst.DskRead, dst.DskWrit) + } + // Epoch and CpuIdl are derived per tick by the caller, never summed. + if dst.Epoch != 0 || dst.CpuIdl != 0 { + t.Errorf("Epoch/CpuIdl = %d/%v, want 0/0", dst.Epoch, dst.CpuIdl) + } +} diff --git a/containers/receiver/main.go b/containers/receiver/main.go index d968451..c218f2d 100644 --- a/containers/receiver/main.go +++ b/containers/receiver/main.go @@ -642,10 +642,34 @@ func receiveTCP(cfg config, cnt *counters, val *validator) error { } } +// stampingReader refreshes the shard's receive-window timestamps on every +// successful read from the connection. recordLine samples time.Now() only +// every 1024 lines (and finish() only runs at connection close), so a burst +// smaller than 1024 lines on a connection the sender keeps open would leave +// lastNs == firstNs and the runner's receive-window rate at 0. Stamping per +// read costs one time.Now() per socket read — already syscall-scale — and +// bounds the window by actual socket activity. +type stampingReader struct { + r io.Reader + shard *connStats +} + +func (sr *stampingReader) Read(p []byte) (int, error) { + n, err := sr.r.Read(p) + if n > 0 { + now := time.Now().UnixNano() + sr.shard.firstNs.CompareAndSwap(0, now) + sr.shard.lastNs.Store(now) + } + return n, err +} + func handleConn(conn net.Conn, shard *connStats, val *validator, cfg config) { defer conn.Close() - defer shard.finish() - scanner := bufio.NewScanner(conn) + // No shard.finish() here: the stampingReader already recorded the exact + // time of the last data read. finish() would overwrite it with the close + // time, inflating the receive window when the sender idles before closing. + scanner := bufio.NewScanner(&stampingReader{r: conn, shard: shard}) scanner.Buffer(make([]byte, 1024*1024), 1024*1024) needsValidation := cfg.ValidateDedup || cfg.ValidateContent || cfg.ValidateJSON || cfg.RequiredSubstring != "" // Pass the scanner's internal slice directly. shard.recordLine + validator @@ -660,6 +684,12 @@ func handleConn(conn net.Conn, shard *connStats, val *validator, cfg config) { val.recordLine(b, cfg) } } + // A scan error (oversized line, connection reset) silently truncates the + // count for this connection — surface it so a lossy-looking run is + // diagnosable from the receiver log. EOF is not reported by Err(). + if err := scanner.Err(); err != nil { + fmt.Fprintf(os.Stderr, "receiver: conn %s: %v\n", conn.RemoteAddr(), err) + } } func receiveFile(cfg config, cnt *counters, val *validator) error { diff --git a/containers/receiver/main_test.go b/containers/receiver/main_test.go new file mode 100644 index 0000000..83bc6b6 --- /dev/null +++ b/containers/receiver/main_test.go @@ -0,0 +1,85 @@ +package main + +import ( + "bufio" + "io" + "strings" + "testing" + "time" +) + +// chunkedReader yields one fixed chunk per Read call with a small delay +// between chunks, so consecutive reads land on distinct timestamps. +type chunkedReader struct { + chunks []string + delay time.Duration +} + +func (c *chunkedReader) Read(p []byte) (int, error) { + if len(c.chunks) == 0 { + return 0, io.EOF + } + if c.delay > 0 { + time.Sleep(c.delay) + } + n := copy(p, c.chunks[0]) + if n == len(c.chunks[0]) { + c.chunks = c.chunks[1:] + } else { + c.chunks[0] = c.chunks[0][n:] + } + return n, nil +} + +func TestStampingReader(t *testing.T) { + // A burst far below recordLine's 1024-line sampling threshold must still + // produce a non-empty receive window (lastNs > firstNs) when the data + // arrives across multiple reads — the regression behind EPS 0 on + // 1,000-line crash/restart cases. + shard := &connStats{} + src := &chunkedReader{ + chunks: []string{"a\nb\n", "c\nd\n"}, + delay: 5 * time.Millisecond, + } + scanner := bufio.NewScanner(&stampingReader{r: src, shard: shard}) + lines := 0 + for scanner.Scan() { + shard.recordLine(int64(len(scanner.Bytes())) + 1) + lines++ + } + if err := scanner.Err(); err != nil { + t.Fatalf("scan: %v", err) + } + if lines != 4 { + t.Fatalf("lines = %d, want 4", lines) + } + first, last := shard.firstNs.Load(), shard.lastNs.Load() + if first == 0 { + t.Fatal("firstNs not stamped") + } + if last <= first { + t.Fatalf("lastNs (%d) <= firstNs (%d): receive window empty", last, first) + } +} + +func TestStampingReaderFirstNsStable(t *testing.T) { + // firstNs is stamped once on the first read and never moves. + shard := &connStats{} + sr := &stampingReader{r: strings.NewReader("x\n"), shard: shard} + buf := make([]byte, 16) + if _, err := sr.Read(buf); err != nil { + t.Fatalf("read: %v", err) + } + first := shard.firstNs.Load() + time.Sleep(2 * time.Millisecond) + sr.r = strings.NewReader("y\n") + if _, err := sr.Read(buf); err != nil { + t.Fatalf("read: %v", err) + } + if got := shard.firstNs.Load(); got != first { + t.Fatalf("firstNs moved: %d → %d", first, got) + } + if shard.lastNs.Load() <= first { + t.Fatal("lastNs not refreshed on second read") + } +} diff --git a/internal/orchestrator/docker.go b/internal/orchestrator/docker.go index e3fa046..1ff3934 100644 --- a/internal/orchestrator/docker.go +++ b/internal/orchestrator/docker.go @@ -607,7 +607,7 @@ services: user: "0:0" {{- end }} environment: - COLLECTOR_TARGET_CONTAINER: "{{ .SubjectContainer }}" + COLLECTOR_TARGET_CONTAINER: "{{ .CollectorTargets }}" COLLECTOR_INTERVAL_SECS: "1" COLLECTOR_OUTPUT: "/results/metrics.csv" restart: "no" @@ -1726,7 +1726,11 @@ type composeVars struct { // template ranges over them instead of emitting the singular subject block. // Shared fields (.SubjectImage, .SubjectEntrypoint, .SubjectCmd, .ConfigDst, // .CaseCertsHost, .TmpDir) are referenced via $ inside the range. - Subjects []clusterNode + Subjects []clusterNode + // CollectorTargets is the comma-separated container list the collector + // samples: the singular subject container normally, every cluster node + // container for cluster cases. + CollectorTargets string SubjectEnv map[string]string HasResourceLimits bool CPULimit string @@ -2033,6 +2037,19 @@ func writeCompose(path string, cfg RunConfig) error { } } + // The collector samples targets by container name. Cluster cases have no + // singular subject container — point the collector at every node so each + // CSV row carries the whole cluster's footprint (the collector sums the + // comma-separated targets per tick). + collectorTargets := subjectContainer + if len(clusterNodes) > 0 { + names := make([]string, len(clusterNodes)) + for i, n := range clusterNodes { + names[i] = n.Container + } + collectorTargets = strings.Join(names, ",") + } + // cluster_ip_failover: pin the bench subnet (so the configured virtual IP is in // range and won't collide with an auto-assigned container address) and grant the // node containers NET_ADMIN (so the elected leader can add the IP to its iface). @@ -2174,6 +2191,7 @@ func writeCompose(path string, cfg RunConfig) error { vars := composeVars{ SubjectImage: s.ImageRef(), SubjectContainer: subjectContainer, + CollectorTargets: collectorTargets, ConfigSrc: configSrc, ConfigDst: configDst, ConfigMountOpts: configMountOpts(s), diff --git a/internal/runner/runner.go b/internal/runner/runner.go index 58f4984..dd5feb8 100644 --- a/internal/runner/runner.go +++ b/internal/runner/runner.go @@ -3709,6 +3709,11 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf finalCount = rm.LinesReceived } + // Cluster resource usage: the collector samples every node container and + // sums them per tick, so these figures are the whole cluster's footprint. + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() + elapsed := time.Since(startTime).Seconds() if passed { fmt.Printf(" director cluster (%s): PASSED ✓\n", clusterActionLabel(tc.Cluster.Action)) @@ -3720,23 +3725,41 @@ func (r *Runner) runDirectorClusterCorrectness(tc *config.TestCase, subject conf fmt.Printf(" --- %s (tail) ---\n%s\n", c, string(out)) } } + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB (cluster total)\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) result := results.RunResult{ - TestName: tc.Name, - Config: configName, - Subject: subject.Name, - Version: subject.Version, - Hardware: hardwareID(), - Timestamp: startTime, - DurationSec: elapsed, - LinesOut: finalCount, - Passed: &passed, + TestName: tc.Name, + Config: configName, + Subject: subject.Name, + Version: subject.Version, + Hardware: hardwareID(), + Timestamp: startTime, + DurationSec: elapsed, + LinesOut: finalCount, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, + Passed: &passed, } if !passed { result.FailReason = strings.Join(errs, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -4177,7 +4200,7 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf fmt.Println(" director did not connect with a bad token, as expected ✓") } passed = len(errs) == 0 - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } // All other scenarios: wait for the director to connect first. @@ -4194,7 +4217,7 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if !fleetWaitConnected(r.ctx, simContainer, dirID, connDeadline) { errs = append(errs, "director never connected to the fleet simulator") passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } fmt.Println(" director connected to the fleet simulator ✓") } @@ -4212,14 +4235,14 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if e != nil { errs = append(errs, "cannot read deliver_config "+fc.DeliverConfig+": "+e.Error()) passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } fmt.Printf(" delivering operational config (%s, %d bytes vmf)…\n", fc.DeliverConfig, len(vmfBytes)) b64 := base64.StdEncoding.EncodeToString(vmfBytes) if _, se := fleetSimSend(simContainer, dirID, "config", map[string]any{"data_b64": b64}); se != nil { errs = append(errs, "failed to deliver operational config: "+se.Error()) passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } // Confirm the director APPLIED the config via its fleet-link // acknowledgment (a rep.config reply with executed=true), not by grepping @@ -4242,7 +4265,7 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if _, ok := fleetWaitCount(r.ctx, simContainer, dirID, "rep.config", 1, applyDeadline); !ok { errs = append(errs, "director did not acknowledge the delivered config over the fleet link") passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } last := "" if st, e := fleetSimStatus(simContainer); e == nil { @@ -4251,7 +4274,7 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if !strings.Contains(last, "\"executed\":true") { errs = append(errs, fmt.Sprintf("director replied to the delivered config but did not report it executed: %.200s", last)) passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } fmt.Println(" director applied the delivered config (executed) ✓") } @@ -4274,7 +4297,7 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf if fc.Scenario != "stats" { errs = append(errs, fmt.Sprintf("restart_mid_run is only supported for the stats scenario, not %q", fc.Scenario)) passed = false - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } fmt.Println(" letting the agent stream before the restart…") if err := sleepCtx(r.ctx, 25*time.Second); err != nil { @@ -5087,14 +5110,17 @@ func (r *Runner) runFleetAutomationCorrectness(tc *config.TestCase, subject conf } passed = len(errs) == 0 - return r.saveFleetResult(tc, subject, configName, startTime, finalCount, passed, errs, simContainer, subjectContainer) + return r.saveFleetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, simContainer, subjectContainer) } // saveFleetResult records the verdict and, on failure, dumps short log tails of // the simulator and director to aid diagnosis. func (r *Runner) saveFleetResult(tc *config.TestCase, subject config.Subject, configName string, + orch orchestrator.Orchestrator, tmpDir string, startTime time.Time, finalCount int64, passed bool, errs []string, simContainer, subjectContainer string, ) (results.RunResult, error) { + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() elapsed := time.Since(startTime).Seconds() label := tc.Fleet.Scenario if passed { @@ -5106,22 +5132,40 @@ func (r *Runner) saveFleetResult(tc *config.TestCase, subject config.Subject, co fmt.Printf(" --- %s (tail) ---\n%s\n", c, string(out)) } } + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) result := results.RunResult{ - TestName: tc.Name, - Config: configName, - Subject: subject.Name, - Version: subject.Version, - Hardware: hardwareID(), - Timestamp: startTime, - DurationSec: elapsed, - LinesOut: finalCount, - Passed: &passed, + TestName: tc.Name, + Config: configName, + Subject: subject.Name, + Version: subject.Version, + Hardware: hardwareID(), + Timestamp: startTime, + DurationSec: elapsed, + LinesOut: finalCount, + AvgCPUPercent: metrics.CPUAvg, + MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, + MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, + DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, + NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, + LoadAvg5: metrics.LoadAvg5, + LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, + SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, + SubjectMemLimit: r.opts.MemLimit, + Passed: &passed, } if !passed { result.FailReason = strings.Join(errs, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } @@ -5874,11 +5918,11 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con } if !ready { errs = append(errs, "clickhouse-server never became reachable") - return r.saveClickHouseTargetResult(tc, subject, configName, startTime, 0, false, errs, chContainer, subjectContainer) + return r.saveClickHouseTargetResult(tc, subject, configName, orch, tmpDir, startTime, 0, false, errs, chContainer, subjectContainer) } if _, e := chExec(chContainer, "CREATE DATABASE IF NOT EXISTS "+db); e != nil { errs = append(errs, "create database: "+e.Error()) - return r.saveClickHouseTargetResult(tc, subject, configName, startTime, 0, false, errs, chContainer, subjectContainer) + return r.saveClickHouseTargetResult(tc, subject, configName, orch, tmpDir, startTime, 0, false, errs, chContainer, subjectContainer) } // Flexible log table: the director sends {"message": "...", "@timestamp": "..."} // (raw) or normalized JSON; skip_unknown_fields tolerates extra keys, so a @@ -5895,7 +5939,7 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con } if _, e := chExec(chContainer, createTable); e != nil { errs = append(errs, "create table: "+e.Error()) - return r.saveClickHouseTargetResult(tc, subject, configName, startTime, 0, false, errs, chContainer, subjectContainer) + return r.saveClickHouseTargetResult(tc, subject, configName, orch, tmpDir, startTime, 0, false, errs, chContainer, subjectContainer) } fmt.Printf(" clickhouse ready; table %s created ✓\n", fqTable) @@ -5974,12 +6018,15 @@ func (r *Runner) runClickHouseTargetCorrectness(tc *config.TestCase, subject con } passed = len(errs) == 0 - return r.saveClickHouseTargetResult(tc, subject, configName, startTime, finalCount, passed, errs, chContainer, subjectContainer) + return r.saveClickHouseTargetResult(tc, subject, configName, orch, tmpDir, startTime, finalCount, passed, errs, chContainer, subjectContainer) } func (r *Runner) saveClickHouseTargetResult(tc *config.TestCase, subject config.Subject, configName string, + orch orchestrator.Orchestrator, tmpDir string, startTime time.Time, finalCount int64, passed bool, errs []string, chContainer, subjectContainer string, ) (results.RunResult, error) { + metrics, metricsCSVSrc := r.harvestResourceMetrics(orch, tmpDir) + sysCPUs, sysMemMB := getSystemInfo() elapsed := time.Since(startTime).Seconds() if passed { fmt.Println(" clickhouse target: PASSED ✓") @@ -5990,14 +6037,24 @@ func (r *Runner) saveClickHouseTargetResult(tc *config.TestCase, subject config. fmt.Printf(" --- %s (tail) ---\n%s\n", c, string(out)) } } + fmt.Printf(" cpu: avg %.1f%% max %.1f%% mem: avg %.0f MB max %.0f MB\n", + metrics.CPUAvg, metrics.CPUMax, metrics.MemAvgMB, metrics.MemMaxMB) result := results.RunResult{ TestName: tc.Name, Config: configName, Subject: subject.Name, Version: subject.Version, Hardware: hardwareID(), Timestamp: startTime, DurationSec: elapsed, LinesOut: finalCount, Passed: &passed, + AvgCPUPercent: metrics.CPUAvg, MaxCPUPercent: metrics.CPUMax, + AvgMemMB: metrics.MemAvgMB, MaxMemMB: metrics.MemMaxMB, + DiskReadBytes: metrics.DiskRead, DiskWriteBytes: metrics.DiskWrite, + NetRecvBytes: metrics.NetRecv, NetSendBytes: metrics.NetSend, + IOThroughputAvg: metrics.IOThroughputAvg, + LoadAvg1: metrics.LoadAvg1, LoadAvg5: metrics.LoadAvg5, LoadAvg15: metrics.LoadAvg15, + SystemCPUs: sysCPUs, SystemMemMB: sysMemMB, + SubjectCPULimit: r.opts.CPULimit, SubjectMemLimit: r.opts.MemLimit, } if !passed { result.FailReason = strings.Join(errs, "; ") } - dir, err := r.saveResult(result, "") + dir, err := r.saveResult(result, metricsCSVSrc) if err != nil { return result, fmt.Errorf("saving results: %w", err) } From fb36f090408d79ceb4b628fb82df28d73acd283f Mon Sep 17 00:00:00 2001 From: Eren Aslan <16862833+erenaslandev@users.noreply.github.com> Date: Fri, 17 Jul 2026 00:51:06 +0300 Subject: [PATCH 7/7] fix(collector): floor CPU idle at 0 when summed usage exceeds 100% --- containers/collector/main.go | 11 ++++++----- containers/collector/main_test.go | 21 +++++++++++++++++++++ 2 files changed, 27 insertions(+), 5 deletions(-) diff --git a/containers/collector/main.go b/containers/collector/main.go index 220c1e5..8c57d19 100644 --- a/containers/collector/main.go +++ b/containers/collector/main.go @@ -148,7 +148,9 @@ func runDockerMode(ctx context.Context, ticker *time.Ticker, emit func(MetricsRo continue } row.Epoch = time.Now().Unix() - row.CpuIdl = 100.0 - row.CpuUsr + // Summed usage regularly exceeds 100 (multi-core containers, + // multi-node clusters); idle is a percentage, so floor it at 0. + row.CpuIdl = max(100.0-row.CpuUsr, 0) emit(row) } } @@ -235,9 +237,10 @@ func fetchDockerStats(client *http.Client, dockerHost, container string) (*docke return &s, nil } +// dockerStatsToRow converts one container's stats snapshot into a partial +// row: per-tick fields (Epoch, CpuIdl) are derived by the caller after the +// per-target rows are summed, so they stay zero here. func dockerStatsToRow(cur *dockerStats, prev *dockerStats) MetricsRow { - now := time.Now().Unix() - var cpuPct float64 cpuDelta := float64(cur.CPUStats.CPUUsage.TotalUsage) - float64(cur.PreCPUStats.CPUUsage.TotalUsage) sysDelta := float64(cur.CPUStats.SystemUsage) - float64(cur.PreCPUStats.SystemUsage) @@ -306,9 +309,7 @@ func dockerStatsToRow(cur *dockerStats, prev *dockerStats) MetricsRow { } return MetricsRow{ - Epoch: now, CpuUsr: cpuPct, - CpuIdl: 100.0 - cpuPct, MemUsed: memUsed, MemCach: cache, MemFree: memFree, diff --git a/containers/collector/main_test.go b/containers/collector/main_test.go index 8b6b2f4..0392e92 100644 --- a/containers/collector/main_test.go +++ b/containers/collector/main_test.go @@ -120,3 +120,24 @@ func TestAddRow(t *testing.T) { t.Errorf("Epoch/CpuIdl = %d/%v, want 0/0", dst.Epoch, dst.CpuIdl) } } + +func TestDockerStatsToRowLeavesPerTickFieldsZero(t *testing.T) { + // Epoch and CpuIdl are derived once per tick by runDockerMode after the + // per-target rows are summed (CpuIdl floored at 0 there, since summed + // usage can exceed 100). A per-container row must leave them zero, or a + // stale per-container idle would leak into the combined row's contract. + var s dockerStats + s.CPUStats.CPUUsage.TotalUsage = 2_000_000 + s.CPUStats.SystemUsage = 10_000_000 + s.CPUStats.OnlineCPUs = 4 + s.PreCPUStats.CPUUsage.TotalUsage = 1_000_000 + s.PreCPUStats.SystemUsage = 9_000_000 + + row := dockerStatsToRow(&s, nil) + if row.CpuUsr <= 0 { + t.Fatalf("CpuUsr = %v, want > 0 (sanity: stats deltas present)", row.CpuUsr) + } + if row.Epoch != 0 || row.CpuIdl != 0 { + t.Fatalf("Epoch/CpuIdl = %d/%v, want 0/0 (caller-derived per tick)", row.Epoch, row.CpuIdl) + } +}