From 4cceb468f41c2a9f7f98edc296887a2d75684b4e Mon Sep 17 00:00:00 2001 From: Max Smythe Date: Sat, 3 Oct 2026 11:22:23 -0700 Subject: [PATCH] benchmarking(boomer): let each user class own its slice of the runtime config dynconfig had become one struct of every knob every user class ever needed, with a payload mirror, a merge, a validator, and a log line that all had to grow by hand for each new field, and every class saw and was checked against every other class's knobs. The payload is now an arbitrary JSON object, and a worker runs exactly one user class, known at launch. So each class declares the keys it reads as a struct with json tags, wraps it in a dynconfig.Typed codec with its defaults and its validation func, and registers the codec with its userclass.Entry. The worker builds the holder from that codec; the holder stores the typed value, folds each payload into it with one decode over the current value (so an absent or null key keeps what it had), refuses a payload the class cannot run on in favor of the last good config, and hands the class its knobs back through dynconfig.Get. Keys the class does not name are ignored. The wait window and the resume and lifecycle modes are shared pieces a class embeds; durations travel as seconds, the unit the locust flags use. The GluttonUser defaults the worker used to seed from main.go (a 500ms max wait, one ping per wake) live in that class's codec now, and the durdir read modes move into the glutton package with the class that reads them. --- benchmarking/locust/common/boomer_config.py | 6 +- cmd/benchmarking/boomer-worker/main.go | 31 +- .../boomer/agentsession/agentsession.go | 44 +- .../boomer/agentsession/agentsession_test.go | 50 +- .../boomer/dynconfig/dynconfig.go | 607 ++++++++---------- .../boomer/dynconfig/dynconfig_test.go | 390 ++++++----- .../boomer/glutton/cpuload_test.go | 7 +- .../benchmarking/boomer/glutton/durdir.go | 33 +- .../boomer/glutton/durdir_test.go | 29 +- .../boomer/glutton/fixture_test.go | 2 +- internal/benchmarking/boomer/glutton/knobs.go | 150 +++++ .../benchmarking/boomer/glutton/lifecycle.go | 34 +- .../boomer/glutton/lifecycle_test.go | 24 +- .../boomer/glutton/memfill_test.go | 20 +- internal/benchmarking/boomer/glutton/spawn.go | 22 +- .../benchmarking/boomer/glutton/spawn_test.go | 6 +- .../benchmarking/boomer/glutton/wait_test.go | 19 +- .../benchmarking/boomer/sweperf/sweperf.go | 49 +- .../boomer/sweperf/sweperf_test.go | 29 +- .../benchmarking/boomer/userclass/config.go | 7 +- .../benchmarking/boomer/userclass/registry.go | 8 + 21 files changed, 813 insertions(+), 754 deletions(-) create mode 100644 internal/benchmarking/boomer/glutton/knobs.go diff --git a/benchmarking/locust/common/boomer_config.py b/benchmarking/locust/common/boomer_config.py index 17902d78c9..77c36619c1 100644 --- a/benchmarking/locust/common/boomer_config.py +++ b/benchmarking/locust/common/boomer_config.py @@ -35,7 +35,11 @@ * serve_config_headless(): the same /boomer-config payload from a plain HTTP server, for a headless run whose values change while it runs. -Keep _FLAGS aligned with internal/benchmarking/boomer/dynconfig.payload. +_FLAGS is the set of keys the payload carries. The Go side takes the payload +as an arbitrary JSON object, and the user class a worker runs reads and +validates the keys it documents (the dynconfig.Typed codec it registers with +its userclass.Entry), so a flag added here needs a matching field in the +class that consumes it. """ import argparse diff --git a/cmd/benchmarking/boomer-worker/main.go b/cmd/benchmarking/boomer-worker/main.go index b7e48a7ec0..84d53174d7 100644 --- a/cmd/benchmarking/boomer-worker/main.go +++ b/cmd/benchmarking/boomer-worker/main.go @@ -47,7 +47,7 @@ func main() { routerURL = flag.String("router-url", "http://atenet-router.ate-system.svc.cluster.local", "atenet HTTP router base URL (no trailing slash).") atespace = flag.String("atespace", "benchmark", "Atespace every actor this worker creates lives in. Ensured (CreateAtespace, AlreadyExists is ok) at startup.") promAddr = flag.String("prometheus-addr", ":8001", "Address for the Prometheus /metrics endpoint.") - configJSON = flag.String("config-json", "", "Initial dynconfig as a JSON object (keys: trace_probability, min_wait_time, max_wait_time, min_live_time, max_live_time in seconds, durdir_file_size_bytes, resume_mode, lifecycle_mode, durdir_read_mode, durdir_template, mem_target, mem_churn, mem_read, cpu_cores, cpu_duty_cycle, max_pings_per_wake). Unset fields keep their built-in defaults.") + configJSON = flag.String("config-json", "", "Initial dynconfig as a JSON object keyed by locust flag name in snake_case (trace_probability, min_wait_time, ...); the selected user class reads the keys it documents and ignores the rest. Unset keys keep the class's built-in defaults.") masterWebPort = flag.Int("master-web-port", 0, "If non-zero, fetch dynconfig from http://{master-host}:{master-web-port}/boomer-config on each spawn message. Exits if the first fetch fails; later failures keep the last fetched values. {master-host} comes from boomer's existing --master-host flag.") configPollInterval = flag.Duration("config-poll-interval", 10*time.Second, "With --master-web-port, also fetch dynconfig on this interval. A spawn message comes only when the number of users or the spawn rate changes, thus a load shape that changes the sample rate alone needs this. Zero stops the polling.") userClass = flag.String("user-class", "glutton", fmt.Sprintf("Locust user class to run, lowercase; one of %s.", strings.Join(userclass.Names(), "|"))) @@ -102,17 +102,25 @@ func main() { return } - initialCfg, err := dynconfig.Parse([]byte(*configJSON), dynconfig.Config{ - MaxWait: 500 * time.Millisecond, - MaxPingsPerWake: 1, - }) - if err != nil { + entry, ok := userclass.Lookup(class) + if !ok { + slog.Error("fatal: unknown --user-class value", + slog.String("user_class", *userClass), + slog.String("known", strings.Join(userclass.Names(), ","))) + os.Exit(1) + } + + // The holder is built from the class's own codec, so a --config-json + // the class cannot run on fails here, and a bad value from the master + // later is refused in favor of the last good config. + dyn := dynconfig.NewHolder(entry.Config) + if _, err := dyn.Apply([]byte(*configJSON)); err != nil { slog.Error("failed to parse --config-json", slog.String("err", err.Error())) os.Exit(1) } ctx := context.Background() - sampler := btrace.NewUpdatableSampler(initialCfg.TraceProbability) + sampler := btrace.NewUpdatableSampler(dyn.Common().TraceProbability) tp, err := btrace.Init(ctx, "substrate-boomer", sampler) if err != nil { slog.Error("failed to initialize tracing", slog.String("err", err.Error())) @@ -144,8 +152,6 @@ func main() { transport.IdleConnTimeout = 5 * time.Minute httpClient := &http.Client{Timeout: 30 * time.Second, Transport: transport} - dyn := dynconfig.NewHolder(initialCfg) - if *masterWebPort > 0 { masterHost := flag.Lookup("master-host").Value.String() configURL := fmt.Sprintf("http://%s:%d/boomer-config", masterHost, *masterWebPort) @@ -195,13 +201,6 @@ func main() { ActorDeadline: time.Duration(*actorDeadline * float64(time.Second)), } - entry, ok := userclass.Lookup(class) - if !ok { - slog.Error("fatal: unknown --user-class value", - slog.String("user_class", *userClass), - slog.String("known", strings.Join(userclass.Names(), ","))) - os.Exit(1) - } taskFn, shutdownFn := entry.Init(cfg) slog.Info("registered boomer task", diff --git a/internal/benchmarking/boomer/agentsession/agentsession.go b/internal/benchmarking/boomer/agentsession/agentsession.go index b386128def..c555044770 100644 --- a/internal/benchmarking/boomer/agentsession/agentsession.go +++ b/internal/benchmarking/boomer/agentsession/agentsession.go @@ -96,10 +96,36 @@ func init() { Name: "agentsession", LocustFile: "agentsession.py", UserClass: agentSessionUserClass, + Config: knobsCodec, Init: initAgentSession, }) } +// knobs is the AgentSessionUser slice of the runtime config. The master +// populates the keys from the --agentsession-* locust flags +// (common/agentsession_config.py). +type knobs struct { + dynconfig.Lifecycle + // Script is the built-in script variant; "" falls back to the default. + Script string `json:"agentsession_script"` + // ScriptFile is a script YAML on the worker; wins over Script when set. + ScriptFile string `json:"agentsession_script_file"` + // ThinkScale multiplies the script's per-step think times; 0 reads as 1. + ThinkScale float64 `json:"agentsession_think_scale"` +} + +var knobsCodec = dynconfig.Typed[knobs]{ + Validate: func(k knobs) error { + if err := k.Lifecycle.Validate(); err != nil { + return err + } + if k.ThinkScale < 0 { + return fmt.Errorf("agentsession_think_scale cannot be negative: %f", k.ThinkScale) + } + return nil + }, +} + func initAgentSession(cfg *userclass.Config) (taskFn func(), shutdown func(context.Context)) { if cfg.Tracer == nil { cfg.Tracer = otel.Tracer("substrate-boomer/agentsession") @@ -177,12 +203,12 @@ type loadedScript struct { // scriptSource is what the knobs select: a file named by // --agentsession-script-file wins, else the built-in variant named by // --agentsession-script, else the default. -func scriptSource(dyn dynconfig.Config) (source string, fromFile bool) { - if dyn.AgentSessionScriptFile != "" { - return dyn.AgentSessionScriptFile, true +func scriptSource(dyn knobs) (source string, fromFile bool) { + if dyn.ScriptFile != "" { + return dyn.ScriptFile, true } - if dyn.AgentSessionScript != "" { - return dyn.AgentSessionScript, false + if dyn.Script != "" { + return dyn.Script, false } return DefaultScript, false } @@ -198,7 +224,7 @@ func scriptSource(dyn dynconfig.Config) (source string, fromFile bool) { func (r *runtime) loadScript() (*loadedScript, error) { r.scriptMu.Lock() defer r.scriptMu.Unlock() - source, fromFile := scriptSource(r.cfg.Dyn.Load()) + source, fromFile := scriptSource(dynconfig.Get[knobs](r.cfg.Dyn)) data, err := readScript(source, fromFile) if err != nil { return nil, err @@ -292,7 +318,7 @@ func memoryLimit(tmpl *ateapipb.ActorTemplate) (limit int64, found bool, err err // script's think time scaled by --agentsession-think-scale (0 reads as 1.0), // with ±20% jitter so a fleet of sessions doesn't move in lockstep. func (r *runtime) think(s Step) time.Duration { - scale := r.cfg.Dyn.Load().AgentSessionThinkScale + scale := dynconfig.Get[knobs](r.cfg.Dyn).ThinkScale if scale <= 0 { scale = 1.0 } @@ -517,7 +543,7 @@ func (u *sessionUser) runStep(ctx context.Context, step Step) bool { return false } - if u.cfg.Dyn.Load().ResumeMode == dynconfig.ResumeModeExplicit { + if dynconfig.Get[knobs](u.cfg.Dyn).ResumeMode == dynconfig.ResumeModeExplicit { if err := u.resume(ctx); err != nil { u.noteFailure(err) return false @@ -736,7 +762,7 @@ func (u *sessionUser) resume(ctx context.Context) error { // it instead of waking a stranded actor. func (u *sessionUser) hibernate(ctx context.Context) { var err error - if u.cfg.Dyn.Load().LifecycleMode == dynconfig.LifecycleModePause { + if dynconfig.Get[knobs](u.cfg.Dyn).LifecycleMode == dynconfig.LifecycleModePause { err = u.tracedCall(ctx, "PauseActor", func(callCtx context.Context, tr *metadata.MD) error { _, err := u.cfg.APIStub.PauseActor(callCtx, &ateapipb.PauseActorRequest{Actor: u.ref()}, grpc.Trailer(tr)) return err diff --git a/internal/benchmarking/boomer/agentsession/agentsession_test.go b/internal/benchmarking/boomer/agentsession/agentsession_test.go index 7da28b1258..aafc12ced2 100644 --- a/internal/benchmarking/boomer/agentsession/agentsession_test.go +++ b/internal/benchmarking/boomer/agentsession/agentsession_test.go @@ -210,7 +210,7 @@ func TestExecOpAgainstFake(t *testing.T) { HTTPClient: http.DefaultClient, RouterURL: ts.URL, Atespace: "benchmark", - Dyn: dynconfig.NewHolder(dynconfig.Config{}), + Dyn: dynconfig.Static(knobs{}), }, actorName: "agent-test", } @@ -256,7 +256,7 @@ func TestExecOpAgainstFake(t *testing.T) { // error that leaves the previous script in place, and a changed knob loads // the new script for sessions that start afterwards. func TestLoadScriptFollowsTheKnob(t *testing.T) { - rt := &runtime{cfg: &userclass.Config{Dyn: dynconfig.NewHolder(dynconfig.Config{AgentSessionScript: "no-such-script"})}} + rt := &runtime{cfg: &userclass.Config{Dyn: dynconfig.Static(knobs{Script: "no-such-script"})}} if _, err := rt.loadScript(); err == nil { t.Fatal("loadScript accepted an unknown script name") } @@ -264,7 +264,7 @@ func TestLoadScriptFollowsTheKnob(t *testing.T) { t.Fatal("a failed load must not cache a script") } - rt.cfg.Dyn.Store(dynconfig.Config{}) + rt.cfg.Dyn = dynconfig.Static(knobs{}) s, err := rt.loadScript() if err != nil { t.Fatal(err) @@ -281,7 +281,7 @@ func TestLoadScriptFollowsTheKnob(t *testing.T) { // A knob that names something that does not load is an error, and the // previous script stays available to sessions that already run on it. - rt.cfg.Dyn.Store(dynconfig.Config{AgentSessionScript: "no-such-script"}) + rt.cfg.Dyn = dynconfig.Static(knobs{Script: "no-such-script"}) if _, err := rt.loadScript(); err == nil { t.Error("loadScript accepted an unknown name after a successful load") } @@ -294,7 +294,7 @@ func TestLoadScriptFollowsTheKnob(t *testing.T) { if err := os.WriteFile(path, []byte(validScript), 0o600); err != nil { t.Fatal(err) } - rt.cfg.Dyn.Store(dynconfig.Config{AgentSessionScriptFile: path}) + rt.cfg.Dyn = dynconfig.Static(knobs{ScriptFile: path}) next, err := rt.loadScript() if err != nil { t.Fatal(err) @@ -329,7 +329,7 @@ func TestLoadScriptFollowsTheKnob(t *testing.T) { func TestShutdownFansOut(t *testing.T) { const sessions = 40 ctl := &fakeControlClient{deleteDelay: 20 * time.Millisecond} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{}) rt := &runtime{cfg: u.cfg} for i := range sessions { rt.users.Store(int64(i), &sessionUser{cfg: u.cfg, actorName: "agent-" + strconv.Itoa(i)}) @@ -389,9 +389,9 @@ func TestLoadScriptPrefersFile(t *testing.T) { if err := os.WriteFile(path, []byte(validScript), 0o600); err != nil { t.Fatal(err) } - rt := &runtime{cfg: &userclass.Config{Dyn: dynconfig.NewHolder(dynconfig.Config{ - AgentSessionScript: DefaultScript, - AgentSessionScriptFile: path, + rt := &runtime{cfg: &userclass.Config{Dyn: dynconfig.Static(knobs{ + Script: DefaultScript, + ScriptFile: path, })}} s, err := rt.loadScript() if err != nil { @@ -422,7 +422,7 @@ func TestStartUserChecksTemplateMemory(t *testing.T) { } { t.Run(tc.name, func(t *testing.T) { ctl := &fakeControlClient{templateMemory: tc.memory} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{}) rt := &runtime{cfg: u.cfg} started, err := rt.startUser(context.Background(), &loadedScript{Script: script, ingestBuf: makeIngestBuf(script.Steps)}) @@ -455,7 +455,7 @@ func TestTemplateMemoryRefusalExpires(t *testing.T) { t.Fatal(err) } ctl := &fakeControlClient{templateMemory: "512Mi"} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{}) rt := &runtime{cfg: u.cfg} if err := rt.checkTemplateMemory(context.Background(), script); err == nil { @@ -481,7 +481,7 @@ func TestTemplateMemoryRefusalExpires(t *testing.T) { // start afterwards, not one already walking its steps. func TestRunningSessionKeepsItsScript(t *testing.T) { ctl := &fakeControlClient{templateMemory: "1Gi"} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{}) rt := &runtime{cfg: u.cfg} first, err := rt.loadScript() if err != nil { @@ -496,7 +496,7 @@ func TestRunningSessionKeepsItsScript(t *testing.T) { if err := os.WriteFile(path, []byte(validScript), 0o600); err != nil { t.Fatal(err) } - rt.cfg.Dyn.Store(dynconfig.Config{AgentSessionScriptFile: path}) + rt.cfg.Dyn = dynconfig.Static(knobs{ScriptFile: path}) second, err := rt.loadScript() if err != nil { t.Fatal(err) @@ -510,7 +510,7 @@ func TestRunningSessionKeepsItsScript(t *testing.T) { func TestThinkScaling(t *testing.T) { r := &runtime{ cfg: &userclass.Config{ - Dyn: dynconfig.NewHolder(dynconfig.Config{AgentSessionThinkScale: 0.5}), + Dyn: dynconfig.Static(knobs{ThinkScale: 0.5}), }, } s := Step{Think: 10 * time.Second} @@ -522,7 +522,7 @@ func TestThinkScaling(t *testing.T) { } // Zero scale reads as 1.0. - r.cfg.Dyn.Store(dynconfig.Config{}) + r.cfg.Dyn = dynconfig.Static(knobs{}) for i := 0; i < 100; i++ { got := r.think(s) if got < 8*time.Second || got > 12*time.Second { @@ -645,7 +645,7 @@ func countCalls(calls []string, name string) int { return n } -func newTestUser(t *testing.T, srv *fake.Server, ctl *fakeControlClient, dyn dynconfig.Config) *sessionUser { +func newTestUser(t *testing.T, srv *fake.Server, ctl *fakeControlClient, dyn knobs) *sessionUser { t.Helper() ts := srv.Start(t) return &sessionUser{ @@ -654,7 +654,7 @@ func newTestUser(t *testing.T, srv *fake.Server, ctl *fakeControlClient, dyn dyn HTTPClient: ts.Client(), RouterURL: ts.URL, Atespace: "benchmark", - Dyn: dynconfig.NewHolder(dyn), + Dyn: dynconfig.Static(dyn), Tracer: otel.Tracer("test"), }, actorName: "agent-test", @@ -669,7 +669,7 @@ var pingStep = Step{Name: "01_test", Agent: "testing", Think: time.Second, Ops: func TestRunStep_RetriesStrandedHibernate(t *testing.T) { srv := &fake.Server{} ctl := &fakeControlClient{suspendErrs: []error{status.Error(codes.Unavailable, "ate-api-server restarting")}} - u := newTestUser(t, srv, ctl, dynconfig.Config{}) + u := newTestUser(t, srv, ctl, knobs{}) if !u.runStep(context.Background(), pingStep) { t.Fatal("runStep = false; the step's ops succeeded and must count") @@ -702,7 +702,7 @@ func TestRunStep_RetriesStrandedHibernate(t *testing.T) { // session must replace it on the first failure, not the third. func TestRunStep_ReplacesCrashedActorImmediately(t *testing.T) { ctl := &fakeControlClient{resumeErrs: []error{status.Error(codes.Aborted, "actor benchmark/agent-test crashed")}} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{ResumeMode: dynconfig.ResumeModeExplicit}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{Lifecycle: dynconfig.Lifecycle{ResumeMode: dynconfig.ResumeModeExplicit}}) if u.runStep(context.Background(), pingStep) { t.Fatal("runStep = true with a crashed actor") @@ -717,7 +717,7 @@ func TestRunStep_ReplacesCrashedActorImmediately(t *testing.T) { // a state nothing the driver can call moves the actor out of. func TestRunStep_ReplacesStuckActorAfterFailedWake(t *testing.T) { ctl := &fakeControlClient{suspendErrs: []error{status.Error(codes.FailedPrecondition, "MarkSuspending prerequisite not met (got: ACTOR_STATE_CRASHED)")}} - u := newTestUser(t, &fake.Server{Status: 503}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{Status: 503}, ctl, knobs{}) if u.runStep(context.Background(), pingStep) { t.Fatal("runStep = true with a failing router") @@ -736,7 +736,7 @@ func TestRunStep_KeepsActorThroughCapacityShortage(t *testing.T) { errs[i] = status.Error(codes.ResourceExhausted, "no free workers available") } ctl := &fakeControlClient{resumeErrs: errs} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{ResumeMode: dynconfig.ResumeModeExplicit}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{Lifecycle: dynconfig.Lifecycle{ResumeMode: dynconfig.ResumeModeExplicit}}) for range rounds { if u.runStep(context.Background(), pingStep) { @@ -755,7 +755,7 @@ func TestRunStep_KeepsActorThroughRouterCapacityErrors(t *testing.T) { for _, status := range []int{503, 504, 429} { t.Run(strconv.Itoa(status), func(t *testing.T) { ctl := &fakeControlClient{} - u := newTestUser(t, &fake.Server{Status: status}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{Status: status}, ctl, knobs{}) for range maxConsecutiveStepFailures + 2 { if u.runStep(context.Background(), pingStep) { t.Fatal("runStep = true with a failing router") @@ -772,7 +772,7 @@ func TestRunStep_KeepsActorThroughRouterCapacityErrors(t *testing.T) { // brings it back, so it is replaced on the first failure. func TestRunStep_ReplacesActorOnRouterNotFound(t *testing.T) { ctl := &fakeControlClient{} - u := newTestUser(t, &fake.Server{Status: 404}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{Status: 404}, ctl, knobs{}) u.runStep(context.Background(), pingStep) if !u.broken { t.Error("broken = false after a 404 wake; want immediate replacement") @@ -784,7 +784,7 @@ func TestRunStep_ReplacesActorOnRouterNotFound(t *testing.T) { // after maxConsecutiveStepFailures steps. func TestRunStep_ReplacesActorAfterRepeatedStepFailures(t *testing.T) { ctl := &fakeControlClient{} - u := newTestUser(t, &fake.Server{Status: 502}, ctl, dynconfig.Config{}) + u := newTestUser(t, &fake.Server{Status: 502}, ctl, knobs{}) for i := 1; i <= maxConsecutiveStepFailures; i++ { u.runStep(context.Background(), pingStep) @@ -796,7 +796,7 @@ func TestRunStep_ReplacesActorAfterRepeatedStepFailures(t *testing.T) { func TestControlRPCsCarryADeadline(t *testing.T) { ctl := &fakeControlClient{} - u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{ResumeMode: dynconfig.ResumeModeExplicit}) + u := newTestUser(t, &fake.Server{}, ctl, knobs{Lifecycle: dynconfig.Lifecycle{ResumeMode: dynconfig.ResumeModeExplicit}}) u.runStep(context.Background(), pingStep) u.suspendAndDelete(context.Background()) diff --git a/internal/benchmarking/boomer/dynconfig/dynconfig.go b/internal/benchmarking/boomer/dynconfig/dynconfig.go index 6280dddd9d..50ccfd7c6f 100644 --- a/internal/benchmarking/boomer/dynconfig/dynconfig.go +++ b/internal/benchmarking/boomer/dynconfig/dynconfig.go @@ -13,341 +13,303 @@ // limitations under the License. // Package dynconfig fetches and holds the boomer worker's runtime-mutable -// settings — the subset of locust flags the operator can change in the web -// UI form. The boomer wire protocol only carries num_users + spawn_rate, so -// these come over an HTTP side channel from the master's /boomer-config -// endpoint (common/boomer_config.py). +// settings: the locust flags the operator can change in the web UI form. +// The boomer wire protocol only carries num_users + spawn_rate, so these +// come over an HTTP side channel from the master's /boomer-config endpoint +// (common/boomer_config.py). +// +// The payload is a JSON object keyed by locust flag name in snake_case. +// This package does not know what the keys mean. A worker runs exactly one +// user class, known at launch, so the class hands the Holder its Codec at +// construction: the typed struct it reads its knobs from, with the +// defaults a key left unset falls back to and the rules a fetched value +// must pass. The Holder then stores that typed value, folds each payload +// into it, refuses one the class cannot run on in favor of the last good +// config, and hands the class its config back through Get. Keys the class +// does not name are ignored. package dynconfig import ( + "bytes" "context" "encoding/json" "fmt" + "io" "log/slog" - "math" + "math/rand/v2" "net/http" + "reflect" "sync/atomic" "time" "github.com/myzhan/boomer" ) -// Resume modes. Explicit issues a ResumeActor RPC before sending traffic. -// Implicit issues no wake request at all: the actor stays suspended until a -// request reaches the atenet router, which wakes it while the request is -// parked. -const ( - ResumeModeExplicit = "explicit" - ResumeModeImplicit = "implicit" +// Codec is how a user class's config comes to be: the value a worker +// starts from, and how a payload folds into the current value. Typed is +// the implementation a class declares over its knobs struct. +type Codec interface { + // Initial is the config before any payload: the class's defaults. + Initial() any + // Merge folds update, a JSON object, into current and validates the + // result. A key absent from update, or null in it, keeps current's + // value: the master serves every flag it knows, null for the ones the + // operator left blank. + Merge(current any, update []byte) (any, error) +} - ReadModeData = "data" - ReadModeDigest = "digest" +// Typed is the Codec over a class's knobs struct T, whose json tags name +// the keys it reads. Merge decodes the update over a copy of the current +// value, so untouched keys keep what they had, and runs Validate on the +// result. +type Typed[T any] struct { + Defaults T + Validate func(T) error +} - LifecycleModeSuspend = "suspend" - LifecycleModePause = "pause" -) +func (c Typed[T]) Initial() any { return c.Defaults } -// Config is the dynamic-mutable subset of boomer's behavior. Holder swaps -// it atomically so task goroutines read a consistent snapshot. -type Config struct { - MinWait time.Duration // gap between one actor's suspend and the VU's next resume, lower bound - MaxWait time.Duration // upper bound of the same gap +func (c Typed[T]) Merge(current any, update []byte) (any, error) { + next, ok := current.(T) + if !ok { + return nil, fmt.Errorf("dynconfig: current config is %T, codec expects %T", current, c.Defaults) + } + if len(bytes.TrimSpace(update)) > 0 { + if err := json.Unmarshal(update, &next); err != nil { + return nil, fmt.Errorf("decode config: %w", err) + } + } + if c.Validate != nil { + if err := c.Validate(next); err != nil { + return nil, err + } + } + return next, nil +} - MinLive time.Duration // time a GluttonUser actor stays resumed between its first ping and suspend, lower bound - MaxLive time.Duration // upper bound of the live window; zero (the default) suspends right after the ping +// Common is the slice of the payload the worker itself reads, whatever the +// user class. Every Holder folds it alongside the class's config. +type Common struct { + TraceProbability float64 `json:"trace_probability"` +} - TraceProbability float64 +var commonCodec = Typed[Common]{ + Validate: func(c Common) error { + if c.TraceProbability < 0 || c.TraceProbability > 1 { + return fmt.Errorf("trace_probability must be between 0.0 and 1.0, got: %f", c.TraceProbability) + } + return nil + }, +} - ResumeMode string // ResumeModeExplicit | ResumeModeImplicit - LifecycleMode string // LifecycleModeSuspend | LifecycleModePause +// Holder holds the running class's config and the worker's Common slice, +// swapped together atomically so task goroutines read a consistent pair. +type Holder struct { + codec Codec + v atomic.Pointer[state] +} - DurDirFileSize int64 // bytes - DurDirReadMode string // ReadModeData | ReadModeDigest - DurDirTemplate string // ActorTemplate name +type state struct { + class any + common Common +} + +// NewHolder returns a holder at codec's defaults. A nil codec is a class +// that reads no runtime config. +func NewHolder(codec Codec) *Holder { + if codec == nil { + codec = Typed[struct{}]{} + } + h := &Holder{codec: codec} + h.v.Store(&state{class: codec.Initial()}) + return h +} - MemTarget string // resident RAM the GluttonUser fills via WriteRAM, suffixed (e.g. "2Gi"); "" disables - MemChurn string // RAM re-randomized in place each cycle via WriteRAM rotate, suffixed (e.g. "64Mi"); "" disables - MemRead string // RAM walked (one byte per page) via ReadRAM after each resume, suffixed (e.g. "1Gi") or "all"; "" disables - MaxPingsPerWake int // cap on pings a GluttonUser sends during one resume/suspend cycle; values < 1 read as 1 +// Static is for tests: a holder fixed at cfg, with no validation. +func Static[T any](cfg T) *Holder { + return NewHolder(Typed[T]{Defaults: cfg}) +} - SweperfTemplate string // ActorTemplate name for the sweperf workload; "" falls back to default - SweperfTotalSteps int // total steps in trace; 0 falls back to default - SweperfNumCycles int // number of cycles to partition steps into; 0 falls back to default - SweperfPollIntervalMs int // /status poll interval in ms; 0 falls back to default +// Apply folds a payload into the current config. A payload the class's +// codec or Common refuses leaves the holder as it was and comes back as +// the error, naming the key. changed is false when the payload left every +// value as it was, so a caller can keep an unchanged poll out of the log. +func (h *Holder) Apply(update []byte) (changed bool, err error) { + prev := h.v.Load() + class, err := h.codec.Merge(prev.class, update) + if err != nil { + return false, err + } + common, err := commonCodec.Merge(prev.common, update) + if err != nil { + return false, err + } + next := &state{class: class, common: common.(Common)} + if next.common == prev.common && reflect.DeepEqual(next.class, prev.class) { + return false, nil + } + h.v.Store(next) + return true, nil +} - CPUCores int // goroutines each GluttonUser's actor spins via UseCPU; 0 disables - CPUDutyCycle float64 // fraction of one core each of those goroutines consumes, in [0, 1] +// Load returns the class's current config as the codec's type; Get is the +// typed form. +func (h *Holder) Load() any { return h.v.Load().class } - AgentSessionScript string // built-in agent-session script variant; "" falls back to the default - AgentSessionScriptFile string // path to a script YAML on the worker; wins over AgentSessionScript when set - AgentSessionThinkScale float64 // multiplier on the script's per-step think times; 0 reads as 1.0 +// Common returns the worker's current Common slice. +func (h *Holder) Common() Common { return h.v.Load().common } - TotalActors int // spawn batch size; 0 keeps --total-actors - SpawnConcurrency int // actors the spawn batch creates concurrently; 0 keeps --spawn-concurrency - ActorDeadline time.Duration // per-actor timeout in the spawn batch; 0 keeps --actor-deadline +// Get is the class's current config. A nil holder reads as T's zero +// value; a holder built for another type is a programming error and +// panics with both types named. +func Get[T any](h *Holder) T { + var zero T + if h == nil { + return zero + } + cfg, ok := h.Load().(T) + if !ok { + panic(fmt.Sprintf("dynconfig: holder carries %T, Get asked for %T", h.Load(), zero)) + } + return cfg } -// Holder lets readers Load() the current Config and writers Store() a new -// one. Backed by atomic.Pointer for lock-free reads on the hot path. -type Holder struct { - v atomic.Pointer[Config] -} +// Seconds is a duration carried in the payload as a JSON number of seconds, +// which is how the locust flags spell every time value. +type Seconds time.Duration -func NewHolder(initial Config) *Holder { - h := &Holder{} - h.v.Store(&initial) - return h +// Duration converts to the time package's unit. +func (s Seconds) Duration() time.Duration { return time.Duration(s) } + +func (s Seconds) MarshalJSON() ([]byte, error) { + return json.Marshal(time.Duration(s).Seconds()) } -func (h *Holder) Load() Config { return *h.v.Load() } +func (s *Seconds) UnmarshalJSON(data []byte) error { + if bytes.Equal(bytes.TrimSpace(data), []byte("null")) { + return nil + } + var secs float64 + if err := json.Unmarshal(data, &secs); err != nil { + return fmt.Errorf("seconds: %w", err) + } + *s = Seconds(secs * float64(time.Second)) + return nil +} -func (h *Holder) Store(c Config) { h.v.Store(&c) } +// WaitTime is the gap between one iteration of a user and the next, drawn +// uniformly from [MinWait, MaxWait]. Embed it in a class's knobs struct +// and call its Validate from the class's. +type WaitTime struct { + MinWait Seconds `json:"min_wait_time"` + MaxWait Seconds `json:"max_wait_time"` +} -// ProbabilityUpdater is the subset of trace.UpdatableSampler we touch here; -// kept as an interface so this package doesn't depend on the trace package. -type ProbabilityUpdater interface { - UpdateProbability(p float64) +func (w WaitTime) Validate() error { + if w.MinWait < 0 { + return fmt.Errorf("min_wait_time cannot be negative: %v", w.MinWait.Duration()) + } + if w.MaxWait < 0 { + return fmt.Errorf("max_wait_time cannot be negative: %v", w.MaxWait.Duration()) + } + if w.MaxWait < w.MinWait { + return fmt.Errorf("max_wait_time (%v) cannot be less than min_wait_time (%v)", w.MaxWait.Duration(), w.MinWait.Duration()) + } + return nil } -// payload mirrors the master's /boomer-config JSON. Fields are pointers so -// we can distinguish "absent" (leave current value) from "explicitly zero". -// The same shape is used for static --config-json input and the live -// /boomer-config endpoint, so master + Python runner + Go worker share one -// vocabulary for the boomer's runtime-tunable knobs. -type payload struct { - TraceProbability *float64 `json:"trace_probability"` - MinWaitTime *float64 `json:"min_wait_time"` - MaxWaitTime *float64 `json:"max_wait_time"` - MinLiveTime *float64 `json:"min_live_time"` - MaxLiveTime *float64 `json:"max_live_time"` - DurDirFileSize *float64 `json:"durdir_file_size_bytes"` - ResumeMode *string `json:"resume_mode"` - LifecycleMode *string `json:"lifecycle_mode"` - DurDirReadMode *string `json:"durdir_read_mode"` - DurDirTemplate *string `json:"durdir_template"` - MemTarget *string `json:"mem_target"` - MemChurn *string `json:"mem_churn"` - MemRead *string `json:"mem_read"` - CPUCores *float64 `json:"cpu_cores"` - CPUDutyCycle *float64 `json:"cpu_duty_cycle"` - MaxPingsPerWake *float64 `json:"max_pings_per_wake"` - SweperfTemplate *string `json:"sweperf_template"` - SweperfTotalSteps *float64 `json:"sweperf_total_steps"` - SweperfNumCycles *float64 `json:"sweperf_num_cycles"` - SweperfPollIntervalMs *float64 `json:"sweperf_poll_interval_ms"` - - AgentSessionScript *string `json:"agentsession_script"` - AgentSessionScriptFile *string `json:"agentsession_script_file"` - AgentSessionThinkScale *float64 `json:"agentsession_think_scale"` - - TotalActors *float64 `json:"total_actors"` - SpawnConcurrency *float64 `json:"spawn_concurrency"` - ActorDeadline *float64 `json:"actor_deadline"` +// Draw returns a wait from [MinWait, MaxWait]; an inverted or empty range +// yields MinWait. +func (w WaitTime) Draw() time.Duration { + return Uniform(w.MinWait.Duration(), w.MaxWait.Duration()) } -// Parse decodes a JSON blob (typically from a CLI flag) and merges its -// fields into `current`. Returns the merged Config — unset fields preserve -// `current`'s existing values, matching Fetch's behavior. -func Parse(jsonBytes []byte, current Config) (Config, error) { - if len(jsonBytes) == 0 { - return current, nil +// Uniform draws from [lo, hi]; an inverted or empty range yields lo. +func Uniform(lo, hi time.Duration) time.Duration { + if hi <= lo { + return lo } - var p payload - if err := json.Unmarshal(jsonBytes, &p); err != nil { - return current, fmt.Errorf("decode config json: %w", err) + return lo + time.Duration(rand.Float64()*float64(hi-lo)) +} + +// Resume and lifecycle modes. Explicit issues a ResumeActor RPC before +// sending traffic; implicit issues no wake request at all: the actor stays +// suspended until a request reaches the atenet router, which wakes it +// while the request is parked. Suspend writes the snapshot to durable +// storage; pause keeps it on the node. +const ( + ResumeModeExplicit = "explicit" + ResumeModeImplicit = "implicit" + + LifecycleModeSuspend = "suspend" + LifecycleModePause = "pause" +) + +// Lifecycle is how a class takes its actors off a worker and brings them +// back. The empty string is each class's own default; see the class. +type Lifecycle struct { + ResumeMode string `json:"resume_mode"` + LifecycleMode string `json:"lifecycle_mode"` +} + +func (l Lifecycle) Validate() error { + if l.ResumeMode != "" && l.ResumeMode != ResumeModeExplicit && l.ResumeMode != ResumeModeImplicit { + return fmt.Errorf("invalid resume_mode %q: must be %q or %q", l.ResumeMode, ResumeModeExplicit, ResumeModeImplicit) } - merged := p.merge(current) - if err := merged.Validate(); err != nil { - return current, fmt.Errorf("validate config: %w", err) + if l.LifecycleMode != "" && l.LifecycleMode != LifecycleModeSuspend && l.LifecycleMode != LifecycleModePause { + return fmt.Errorf("invalid lifecycle_mode %q: must be %q or %q", l.LifecycleMode, LifecycleModeSuspend, LifecycleModePause) } - return merged, nil + return nil } -// Fetch GETs `url` and merges any returned fields into `current`. Returns -// the merged Config (or current unchanged on a soft no-op response). -func Fetch(ctx context.Context, url string, current Config) (Config, error) { +// Fetch GETs url and returns the payload it serves. +func Fetch(ctx context.Context, url string) ([]byte, error) { req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { - return current, err + return nil, err } resp, err := http.DefaultClient.Do(req) if err != nil { - return current, fmt.Errorf("GET %s: %w", url, err) + return nil, fmt.Errorf("GET %s: %w", url, err) } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return current, fmt.Errorf("GET %s: status %d", url, resp.StatusCode) + return nil, fmt.Errorf("GET %s: status %d", url, resp.StatusCode) } - var p payload - if err := json.NewDecoder(resp.Body).Decode(&p); err != nil { - return current, fmt.Errorf("decode %s: %w", url, err) - } - merged := p.merge(current) - if err := merged.Validate(); err != nil { - return current, fmt.Errorf("validate %s: %w", url, err) + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("read %s: %w", url, err) } - return merged, nil + return body, nil } -// Validate checks that the config values are within legal ranges. -func (c Config) Validate() error { - if c.MinWait < 0 { - return fmt.Errorf("min_wait_time cannot be negative: %v", c.MinWait) - } - if c.MaxWait < 0 { - return fmt.Errorf("max_wait_time cannot be negative: %v", c.MaxWait) - } - if c.MaxWait < c.MinWait { - return fmt.Errorf("max_wait_time (%v) cannot be less than min_wait_time (%v)", c.MaxWait, c.MinWait) - } - if c.MinLive < 0 { - return fmt.Errorf("min_live_time cannot be negative: %v", c.MinLive) - } - if c.MaxLive < 0 { - return fmt.Errorf("max_live_time cannot be negative: %v", c.MaxLive) - } - if c.MaxLive < c.MinLive { - return fmt.Errorf("max_live_time (%v) cannot be less than min_live_time (%v)", c.MaxLive, c.MinLive) - } - if c.TraceProbability < 0 || c.TraceProbability > 1 { - return fmt.Errorf("trace_probability must be between 0.0 and 1.0, got: %f", c.TraceProbability) - } - if c.DurDirFileSize < 0 { - return fmt.Errorf("durdir_file_size_bytes cannot be negative: %d", c.DurDirFileSize) - } - if c.DurDirFileSize > math.MaxInt32 { - return fmt.Errorf("durdir_file_size_bytes cannot exceed %d (2 GiB), got: %d", math.MaxInt32, c.DurDirFileSize) - } - if c.ResumeMode != "" && c.ResumeMode != ResumeModeExplicit && c.ResumeMode != ResumeModeImplicit { - return fmt.Errorf("invalid resume_mode %q: must be %q or %q", c.ResumeMode, ResumeModeExplicit, ResumeModeImplicit) - } - if c.LifecycleMode != "" && c.LifecycleMode != LifecycleModeSuspend && c.LifecycleMode != LifecycleModePause { - return fmt.Errorf("invalid lifecycle_mode %q: must be %q or %q", c.LifecycleMode, LifecycleModeSuspend, LifecycleModePause) - } - if c.DurDirReadMode != "" && c.DurDirReadMode != ReadModeData && c.DurDirReadMode != ReadModeDigest { - return fmt.Errorf("invalid durdir_read_mode %q: must be %q or %q", c.DurDirReadMode, ReadModeData, ReadModeDigest) - } - if c.CPUCores < 0 { - return fmt.Errorf("cpu_cores cannot be negative: %d", c.CPUCores) - } - if c.CPUDutyCycle < 0 || c.CPUDutyCycle > 1 { - return fmt.Errorf("cpu_duty_cycle must be between 0.0 and 1.0, got: %f", c.CPUDutyCycle) - } - if c.SweperfTotalSteps < 0 { - return fmt.Errorf("sweperf_total_steps cannot be negative: %d", c.SweperfTotalSteps) - } - if c.SweperfNumCycles < 0 { - return fmt.Errorf("sweperf_num_cycles cannot be negative: %d", c.SweperfNumCycles) - } - if c.SweperfPollIntervalMs < 0 { - return fmt.Errorf("sweperf_poll_interval_ms cannot be negative: %d", c.SweperfPollIntervalMs) - } - if c.AgentSessionThinkScale < 0 { - return fmt.Errorf("agentsession_think_scale cannot be negative: %f", c.AgentSessionThinkScale) - } - if c.TotalActors < 0 { - return fmt.Errorf("total_actors cannot be negative: %d", c.TotalActors) - } - if c.SpawnConcurrency < 0 { - return fmt.Errorf("spawn_concurrency cannot be negative: %d", c.SpawnConcurrency) - } - if c.ActorDeadline < 0 { - return fmt.Errorf("actor_deadline cannot be negative: %v", c.ActorDeadline) - } - // MaxPingsPerWake < 1 is treated as 1 at read time (see iterate() in - // glutton/lifecycle.go), so Config's zero value stays usable — no - // validate rejection here. - // MemTarget, MemChurn, and MemRead are passed to glutton verbatim - // (MemRead's "all" excepted, which the driver maps to an empty - // whole-array walk), which owns the parse; invalid values fail loudly - // there as GluttonFillRAM / GluttonChurnRAM / GluttonReadRAM errors. - return nil +// ProbabilityUpdater is the subset of trace.UpdatableSampler we touch here; +// kept as an interface so this package doesn't depend on the trace package. +type ProbabilityUpdater interface { + UpdateProbability(p float64) } -// merge folds the payload's set fields into `current`, leaving unset fields -// at their existing values. Used by both Parse (CLI input) and Fetch (HTTP -// pull) so the merge semantics are identical. -func (p payload) merge(current Config) Config { - out := current - if p.TraceProbability != nil { - out.TraceProbability = *p.TraceProbability - } - if p.MinWaitTime != nil { - out.MinWait = time.Duration(*p.MinWaitTime * float64(time.Second)) - } - if p.MaxWaitTime != nil { - out.MaxWait = time.Duration(*p.MaxWaitTime * float64(time.Second)) - } - if p.MinLiveTime != nil { - out.MinLive = time.Duration(*p.MinLiveTime * float64(time.Second)) - } - if p.MaxLiveTime != nil { - out.MaxLive = time.Duration(*p.MaxLiveTime * float64(time.Second)) - } - if p.DurDirFileSize != nil { - out.DurDirFileSize = int64(*p.DurDirFileSize) - } - if p.ResumeMode != nil { - out.ResumeMode = *p.ResumeMode - } - if p.LifecycleMode != nil { - out.LifecycleMode = *p.LifecycleMode - } - if p.DurDirReadMode != nil { - out.DurDirReadMode = *p.DurDirReadMode - } - if p.DurDirTemplate != nil { - out.DurDirTemplate = *p.DurDirTemplate - } - if p.MemTarget != nil { - out.MemTarget = *p.MemTarget - } - if p.MemChurn != nil { - out.MemChurn = *p.MemChurn - } - if p.MemRead != nil { - out.MemRead = *p.MemRead - } - if p.CPUCores != nil { - out.CPUCores = int(*p.CPUCores) - } - if p.CPUDutyCycle != nil { - out.CPUDutyCycle = *p.CPUDutyCycle - } - if p.MaxPingsPerWake != nil { - out.MaxPingsPerWake = int(*p.MaxPingsPerWake) - } - if p.SweperfTemplate != nil { - out.SweperfTemplate = *p.SweperfTemplate - } - if p.SweperfTotalSteps != nil { - out.SweperfTotalSteps = int(*p.SweperfTotalSteps) - } - if p.SweperfNumCycles != nil { - out.SweperfNumCycles = int(*p.SweperfNumCycles) - } - if p.SweperfPollIntervalMs != nil { - out.SweperfPollIntervalMs = int(*p.SweperfPollIntervalMs) - } - if p.AgentSessionScript != nil { - out.AgentSessionScript = *p.AgentSessionScript - } - if p.AgentSessionScriptFile != nil { - out.AgentSessionScriptFile = *p.AgentSessionScriptFile - } - if p.AgentSessionThinkScale != nil { - out.AgentSessionThinkScale = *p.AgentSessionThinkScale - } - if p.TotalActors != nil { - out.TotalActors = int(*p.TotalActors) +// fetchAndApply fetches url into holder and, on a change, pushes the trace +// probability to the sampler and logs the config now in force. The error +// is a failed fetch or the holder's refusal, with the key it names. +func fetchAndApply(ctx context.Context, trigger, url string, holder *Holder, sampler ProbabilityUpdater) error { + body, err := Fetch(ctx, url) + if err != nil { + return err } - if p.SpawnConcurrency != nil { - out.SpawnConcurrency = int(*p.SpawnConcurrency) + changed, err := holder.Apply(body) + if err != nil { + return fmt.Errorf("%s: %w", url, err) } - if p.ActorDeadline != nil { - out.ActorDeadline = time.Duration(*p.ActorDeadline * float64(time.Second)) + if !changed { + return nil } - return out + sampler.UpdateProbability(holder.Common().TraceProbability) + slog.Info("dynconfig applied", + slog.String("trigger", trigger), + slog.Float64("trace_probability", holder.Common().TraceProbability), + slog.String("config", fmt.Sprintf("%+v", holder.Load()))) + return nil } // StartPoll fetches `url` every `interval` until `ctx` is done, and applies @@ -361,9 +323,11 @@ func (p payload) merge(current Config) Config { // it. The sample-rate sweep of benchmarking/observability.md is one such // shape: each of its steps holds 10 users. // -// `onError` gets each failed fetch. A caller must not exit the process there, -// as it does for a spawn: the worker holds the last good value, and one -// failed poll of a long run is not a reason to lose the run. +// `onError` gets each failed fetch and each refused payload. A caller must +// not exit the process there, as it does for a spawn: the worker holds the +// last good value, and one failed poll of a long run is not a reason to +// lose the run. Only a change goes to the log: a poll of each few seconds +// for the length of a soak would otherwise make a log that hides the run. func StartPoll( ctx context.Context, url string, @@ -384,7 +348,7 @@ func StartPoll( return case <-ticker.C: fetchCtx, cancel := context.WithTimeout(ctx, fetchTimeout) - next, err := Fetch(fetchCtx, url, holder.Load()) + err := fetchAndApply(fetchCtx, "poll", url, holder, sampler) cancel() if err != nil { // The end of the run stops a fetch that is in @@ -394,44 +358,7 @@ func StartPoll( return } onError(err) - continue - } - // Only a change goes to the log. A poll of each few seconds - // for the length of a soak makes a log that hides the run. - if next == holder.Load() { - continue } - holder.Store(next) - sampler.UpdateProbability(next.TraceProbability) - slog.Info("dynconfig applied", - slog.String("trigger", "poll"), - slog.Float64("trace_probability", next.TraceProbability), - slog.Duration("min_wait", next.MinWait), - slog.Duration("max_wait", next.MaxWait), - slog.Duration("min_live", next.MinLive), - slog.Duration("max_live", next.MaxLive), - slog.Int64("durdir_file_size_bytes", next.DurDirFileSize), - slog.String("resume_mode", next.ResumeMode), - slog.String("lifecycle_mode", next.LifecycleMode), - slog.String("durdir_read_mode", next.DurDirReadMode), - slog.String("durdir_template", next.DurDirTemplate), - slog.String("mem_target", next.MemTarget), - slog.String("mem_churn", next.MemChurn), - slog.String("mem_read", next.MemRead), - slog.Int("cpu_cores", next.CPUCores), - slog.Float64("cpu_duty_cycle", next.CPUDutyCycle), - slog.Int("max_pings_per_wake", next.MaxPingsPerWake), - slog.String("sweperf_template", next.SweperfTemplate), - slog.Int("sweperf_total_steps", next.SweperfTotalSteps), - slog.Int("sweperf_num_cycles", next.SweperfNumCycles), - slog.Int("sweperf_poll_interval_ms", next.SweperfPollIntervalMs), - slog.String("agentsession_script", next.AgentSessionScript), - slog.String("agentsession_script_file", next.AgentSessionScriptFile), - slog.Float64("agentsession_think_scale", next.AgentSessionThinkScale), - slog.Int("total_actors", next.TotalActors), - slog.Int("spawn_concurrency", next.SpawnConcurrency), - slog.Duration("actor_deadline", next.ActorDeadline), - ) } } }() @@ -440,52 +367,22 @@ func StartPoll( // SubscribeSpawn registers a boomer Events handler that fetches `url` on // each spawn message and applies the result to `holder` + `sampler`. Locust // sends a spawn message for every ramp step, so a long ramp fetches once per -// second. `onError` is invoked when a fetch fails, with `fetched` true if an -// earlier spawn fetch succeeded: the holder then still has a value from the -// master, and the caller can keep running on it. With `fetched` false the -// worker has only its command-line defaults, and callers typically exit. -// Returns an error if the event subscription itself fails (handler signature -// mismatch), which is a programmer error and should be treated as fatal too. +// second. `onError` is invoked when a fetch fails or the holder refuses the +// payload, with `fetched` true if an earlier spawn fetch was applied: the +// holder then still has a value from the master, and the caller can keep +// running on it. With `fetched` false the worker has only its command-line +// values, and callers typically exit. Returns an error if the event +// subscription itself fails (handler signature mismatch), which is a +// programmer error and should be treated as fatal too. func SubscribeSpawn(url string, holder *Holder, sampler ProbabilityUpdater, fetchTimeout time.Duration, onError func(err error, fetched bool)) error { var fetched atomic.Bool return boomer.Events.Subscribe("boomer:spawn", func(spawnCount int, spawnRate float64) { ctx, cancel := context.WithTimeout(context.Background(), fetchTimeout) defer cancel() - next, err := Fetch(ctx, url, holder.Load()) - if err != nil { + if err := fetchAndApply(ctx, "spawn", url, holder, sampler); err != nil { onError(err, fetched.Load()) return } fetched.Store(true) - holder.Store(next) - sampler.UpdateProbability(next.TraceProbability) - slog.Info("dynconfig applied", - slog.Float64("trace_probability", next.TraceProbability), - slog.Duration("min_wait", next.MinWait), - slog.Duration("max_wait", next.MaxWait), - slog.Duration("min_live", next.MinLive), - slog.Duration("max_live", next.MaxLive), - slog.Int64("durdir_file_size_bytes", next.DurDirFileSize), - slog.String("resume_mode", next.ResumeMode), - slog.String("lifecycle_mode", next.LifecycleMode), - slog.String("durdir_read_mode", next.DurDirReadMode), - slog.String("durdir_template", next.DurDirTemplate), - slog.String("mem_target", next.MemTarget), - slog.String("mem_churn", next.MemChurn), - slog.String("mem_read", next.MemRead), - slog.Int("cpu_cores", next.CPUCores), - slog.Float64("cpu_duty_cycle", next.CPUDutyCycle), - slog.Int("max_pings_per_wake", next.MaxPingsPerWake), - slog.String("sweperf_template", next.SweperfTemplate), - slog.Int("sweperf_total_steps", next.SweperfTotalSteps), - slog.Int("sweperf_num_cycles", next.SweperfNumCycles), - slog.Int("sweperf_poll_interval_ms", next.SweperfPollIntervalMs), - slog.String("agentsession_script", next.AgentSessionScript), - slog.String("agentsession_script_file", next.AgentSessionScriptFile), - slog.Float64("agentsession_think_scale", next.AgentSessionThinkScale), - slog.Int("total_actors", next.TotalActors), - slog.Int("spawn_concurrency", next.SpawnConcurrency), - slog.Duration("actor_deadline", next.ActorDeadline), - ) }) } diff --git a/internal/benchmarking/boomer/dynconfig/dynconfig_test.go b/internal/benchmarking/boomer/dynconfig/dynconfig_test.go index 2f1390d888..b9a8eb9d05 100644 --- a/internal/benchmarking/boomer/dynconfig/dynconfig_test.go +++ b/internal/benchmarking/boomer/dynconfig/dynconfig_test.go @@ -19,6 +19,7 @@ import ( "fmt" "net/http" "net/http/httptest" + "strings" "sync" "sync/atomic" "testing" @@ -27,232 +28,170 @@ import ( "github.com/myzhan/boomer" ) -func TestParseValid(t *testing.T) { - jsonBlob := []byte(`{ - "trace_probability": 0.5, - "min_wait_time": 0.1, - "max_wait_time": 0.5, - "min_live_time": 9, - "max_live_time": 14, - "durdir_file_size_bytes": 1048576, - "resume_mode": "explicit", - "lifecycle_mode": "pause", - "durdir_read_mode": "data", - "durdir_template": "glutton-durdir-data", - "cpu_cores": 2, - "cpu_duty_cycle": 0.1, - "sweperf_template": "swebench-astropy-7336", - "sweperf_total_steps": 21, - "sweperf_num_cycles": 4, - "sweperf_poll_interval_ms": 100, - "agentsession_script": "coding-session", - "agentsession_script_file": "/etc/agentsession/script.yaml", - "total_actors": 50, - "spawn_concurrency": 5, - "actor_deadline": 60.0 - }`) +// knobs is the kind of struct a user class declares: its slice of the +// payload with its own defaults and rules. +type knobs struct { + WaitTime + Lifecycle + Template string `json:"test_template"` + Workers int `json:"test_workers"` +} + +var knobsCodec = Typed[knobs]{ + Defaults: knobs{Template: "stock", Workers: 1}, + Validate: func(k knobs) error { + if err := k.WaitTime.Validate(); err != nil { + return err + } + if err := k.Lifecycle.Validate(); err != nil { + return err + } + if k.Workers < 1 { + return fmt.Errorf("test_workers must be positive: %d", k.Workers) + } + return nil + }, +} - cfg, err := Parse(jsonBlob, Config{}) +func mustApply(t *testing.T, h *Holder, payload string) bool { + t.Helper() + changed, err := h.Apply([]byte(payload)) if err != nil { - t.Fatalf("Parse failed: %v", err) + t.Fatalf("Apply(%s): %v", payload, err) } + return changed +} - if cfg.TraceProbability != 0.5 { - t.Errorf("TraceProbability: got %f, want 0.5", cfg.TraceProbability) - } - if cfg.MinWait != 100*time.Millisecond { - t.Errorf("MinWait: got %v, want 100ms", cfg.MinWait) - } - if cfg.MaxWait != 500*time.Millisecond { - t.Errorf("MaxWait: got %v, want 500ms", cfg.MaxWait) - } - if cfg.MinLive != 9*time.Second { - t.Errorf("MinLive: got %v, want 9s", cfg.MinLive) - } - if cfg.MaxLive != 14*time.Second { - t.Errorf("MaxLive: got %v, want 14s", cfg.MaxLive) - } - if cfg.DurDirFileSize != 1048576 { - t.Errorf("DurDirFileSize: got %d, want 1048576", cfg.DurDirFileSize) - } - if cfg.ResumeMode != ResumeModeExplicit { - t.Errorf("ResumeMode: got %q, want %q", cfg.ResumeMode, ResumeModeExplicit) - } - if cfg.LifecycleMode != LifecycleModePause { - t.Errorf("LifecycleMode: got %q, want %q", cfg.LifecycleMode, LifecycleModePause) - } - if cfg.DurDirReadMode != ReadModeData { - t.Errorf("DurDirReadMode: got %q, want %q", cfg.DurDirReadMode, ReadModeData) - } - if cfg.DurDirTemplate != "glutton-durdir-data" { - t.Errorf("DurDirTemplate: got %q, want glutton-durdir-data", cfg.DurDirTemplate) - } - if cfg.CPUCores != 2 { - t.Errorf("CPUCores: got %d, want 2", cfg.CPUCores) +// A payload sets the keys it carries, leaves the ones it nulls or omits at +// their current values, converts seconds to durations, and ignores keys +// the class does not name. +func TestApplyMerges(t *testing.T) { + h := NewHolder(knobsCodec) + if got := Get[knobs](h); got != knobsCodec.Defaults { + t.Fatalf("fresh holder = %+v, want the defaults", got) + } + mustApply(t, h, `{ + "min_wait_time": 0.5, "max_wait_time": 2, + "resume_mode": "implicit", + "test_workers": 3, + "someone_elses_key": "ignored" + }`) + want := knobs{ + WaitTime: WaitTime{MinWait: Seconds(500 * time.Millisecond), MaxWait: Seconds(2 * time.Second)}, + Lifecycle: Lifecycle{ResumeMode: ResumeModeImplicit}, + Template: "stock", + Workers: 3, } - if cfg.CPUDutyCycle != 0.1 { - t.Errorf("CPUDutyCycle: got %f, want 0.1", cfg.CPUDutyCycle) + if got := Get[knobs](h); got != want { + t.Errorf("after first payload =\n %+v, want\n %+v", got, want) } - if cfg.SweperfTemplate != "swebench-astropy-7336" { - t.Errorf("SweperfTemplate: got %q, want swebench-astropy-7336", cfg.SweperfTemplate) + mustApply(t, h, `{"test_template": "big", "test_workers": null, "max_wait_time": null}`) + want.Template = "big" + if got := Get[knobs](h); got != want { + t.Errorf("after nulls =\n %+v, want\n %+v", got, want) } - if cfg.SweperfTotalSteps != 21 { - t.Errorf("SweperfTotalSteps: got %d, want 21", cfg.SweperfTotalSteps) + if changed := mustApply(t, h, `{"test_template": "big"}`); changed { + t.Error("a payload that changes nothing reported a change") } - if cfg.SweperfNumCycles != 4 { - t.Errorf("SweperfNumCycles: got %d, want 4", cfg.SweperfNumCycles) + if changed := mustApply(t, h, ``); changed { + t.Error("an empty payload reported a change") } - if cfg.SweperfPollIntervalMs != 100 { - t.Errorf("SweperfPollIntervalMs: got %d, want 100", cfg.SweperfPollIntervalMs) +} + +// A payload the class's rules or Common refuse leaves the holder as it was, +// so a bad value from the master cannot take a run down. +func TestApplyRefuses(t *testing.T) { + h := NewHolder(knobsCodec) + mustApply(t, h, `{"test_workers": 2, "trace_probability": 0.5}`) + before := Get[knobs](h) + for name, payload := range map[string]string{ + "class rule": `{"test_workers": 0}`, + "embedded": `{"max_wait_time": 1, "min_wait_time": 2}`, + "mode": `{"lifecycle_mode": "hibernate"}`, + "wrong type": `{"test_workers": "three"}`, + "bad seconds": `{"min_wait_time": "soon"}`, + "common rule": `{"trace_probability": 1.5}`, + "not json": ``, + } { + t.Run(name, func(t *testing.T) { + changed, err := h.Apply([]byte(payload)) + if err == nil || changed { + t.Fatalf("Apply(%s) = changed %v, err %v; want a refusal", payload, changed, err) + } + if got := Get[knobs](h); got != before { + t.Errorf("a refused payload changed the config to %+v", got) + } + if got := h.Common().TraceProbability; got != 0.5 { + t.Errorf("a refused payload changed trace_probability to %v", got) + } + }) } - if cfg.AgentSessionScript != "coding-session" { - t.Errorf("AgentSessionScript: got %q, want coding-session", cfg.AgentSessionScript) +} + +// Get on the wrong type is a programming error that must not pass +// silently; a nil holder is the zero value, for classes that run without +// one in tests. +func TestGetTypeChecks(t *testing.T) { + if got := Get[knobs](nil); got != (knobs{}) { + t.Errorf("Get on a nil holder = %+v, want zero", got) + } + defer func() { + if recover() == nil { + t.Error("Get with the wrong type did not panic") + } + }() + Get[Common](NewHolder(knobsCodec)) +} + +func TestStaticAndNilCodec(t *testing.T) { + h := Static(knobs{Template: "t", WaitTime: WaitTime{MaxWait: Seconds(3 * time.Second)}}) + if got := Get[knobs](h); got.Template != "t" || got.MaxWait != Seconds(3*time.Second) { + t.Errorf("Static round trip = %+v", got) } - if cfg.AgentSessionScriptFile != "/etc/agentsession/script.yaml" { - t.Errorf("AgentSessionScriptFile: got %q", cfg.AgentSessionScriptFile) + // A class without knobs still gets Common. + h = NewHolder(nil) + mustApply(t, h, `{"trace_probability": 0.25, "anything": 1}`) + if got := h.Common().TraceProbability; got != 0.25 { + t.Errorf("trace_probability = %v, want 0.25", got) } - if cfg.TotalActors != 50 { - t.Errorf("TotalActors: got %d, want 50", cfg.TotalActors) +} + +func TestWaitTimeDraw(t *testing.T) { + window := WaitTime{MinWait: Seconds(10 * time.Millisecond), MaxWait: Seconds(50 * time.Millisecond)} + for i := 0; i < 100; i++ { + if got := window.Draw(); got < 10*time.Millisecond || got > 50*time.Millisecond { + t.Fatalf("Draw = %v, outside the window", got) + } } - if cfg.SpawnConcurrency != 5 { - t.Errorf("SpawnConcurrency: got %d, want 5", cfg.SpawnConcurrency) + inverted := WaitTime{MinWait: Seconds(100 * time.Millisecond), MaxWait: Seconds(50 * time.Millisecond)} + if got := inverted.Draw(); got != 100*time.Millisecond { + t.Errorf("inverted Draw = %v, want the lower bound", got) } - if cfg.ActorDeadline != 60*time.Second { - t.Errorf("ActorDeadline: got %v, want 60s", cfg.ActorDeadline) + if got := Uniform(0, 0); got != 0 { + t.Errorf("Uniform(0, 0) = %v, want 0", got) } -} - -func TestParseInvalidValues(t *testing.T) { - tests := []struct { - name string - json string - }{ - { - name: "negative trace probability", - json: `{"trace_probability": -0.1}`, - }, - { - name: "trace probability > 1.0", - json: `{"trace_probability": 1.5}`, - }, - { - name: "negative min wait", - json: `{"min_wait_time": -1.0}`, - }, - { - name: "negative max wait", - json: `{"max_wait_time": -1.0}`, - }, - { - name: "max wait less than min wait", - json: `{"min_wait_time": 2.0, "max_wait_time": 1.0}`, - }, - { - name: "negative min live", - json: `{"min_live_time": -1.0}`, - }, - { - name: "negative max live", - json: `{"max_live_time": -1.0}`, - }, - { - name: "max live less than min live", - json: `{"min_live_time": 14.0, "max_live_time": 9.0}`, - }, - { - name: "negative file size", - json: `{"durdir_file_size_bytes": -100}`, - }, - { - name: "file size exceeds 2 GiB", - json: `{"durdir_file_size_bytes": 2147483648}`, - }, - { - name: "invalid resume mode", - json: `{"resume_mode": "invalid_mode"}`, - }, - { - name: "invalid lifecycle mode", - json: `{"lifecycle_mode": "invalid_lifecycle"}`, - }, - { - name: "invalid read mode", - json: `{"durdir_read_mode": "invalid_read"}`, - }, - { - name: "negative cpu cores", - json: `{"cpu_cores": -1}`, - }, - { - name: "negative cpu duty cycle", - json: `{"cpu_duty_cycle": -0.1}`, - }, - { - name: "cpu duty cycle > 1.0", - json: `{"cpu_duty_cycle": 1.5}`, - }, - { - name: "negative sweperf total steps", - json: `{"sweperf_total_steps": -1}`, - }, - { - name: "negative sweperf num cycles", - json: `{"sweperf_num_cycles": -1}`, - }, - { - name: "negative sweperf poll interval", - json: `{"sweperf_poll_interval_ms": -1}`, - }, - { - name: "negative total actors", - json: `{"total_actors": -1}`, - }, - { - name: "negative spawn concurrency", - json: `{"spawn_concurrency": -1}`, - }, - { - name: "negative actor deadline", - json: `{"actor_deadline": -1.0}`, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - _, err := Parse([]byte(tt.json), Config{}) - if err == nil { - t.Errorf("expected Parse to fail for %s, got nil error", tt.name) - } - }) + if got := Uniform(3*time.Second, 3*time.Second); got != 3*time.Second { + t.Errorf("Uniform(3s, 3s) = %v, want 3s", got) } } -func TestFetchValidAndInvalid(t *testing.T) { +func TestFetchReportsErrors(t *testing.T) { mux := http.NewServeMux() mux.HandleFunc("/valid", func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - _, _ = w.Write([]byte(`{"resume_mode": "implicit", "durdir_read_mode": "digest"}`)) + _, _ = w.Write([]byte(`{"resume_mode": "implicit"}`)) }) - mux.HandleFunc("/invalid", func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Content-Type", "application/json") - _, _ = w.Write([]byte(`{"resume_mode": "bogus"}`)) + mux.HandleFunc("/busy", func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "busy", http.StatusServiceUnavailable) }) ts := httptest.NewServer(mux) defer ts.Close() - ctx := context.Background() - - cfg, err := Fetch(ctx, ts.URL+"/valid", Config{}) - if err != nil { - t.Fatalf("Fetch valid failed: %v", err) + body, err := Fetch(context.Background(), ts.URL+"/valid") + if err != nil || string(body) != `{"resume_mode": "implicit"}` { + t.Errorf("Fetch valid = %q, %v", body, err) } - if cfg.ResumeMode != ResumeModeImplicit || cfg.DurDirReadMode != ReadModeDigest { - t.Errorf("Fetch valid values mismatch: got %+v", cfg) - } - - _, err = Fetch(ctx, ts.URL+"/invalid", Config{}) - if err == nil { - t.Errorf("expected Fetch invalid to fail, got nil") + if _, err := Fetch(context.Background(), ts.URL+"/busy"); err == nil { + t.Error("Fetch of a 503 returned no error") } } @@ -293,7 +232,7 @@ func TestStartPollAppliesAChange(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() - holder := NewHolder(Config{}) + holder := NewHolder(knobsCodec) sampler := &fakeSampler{} StartPoll(ctx, ts.URL, holder, sampler, 10*time.Millisecond, time.Second, func(err error) { t.Errorf("unexpected poll error: %v", err) }) @@ -306,7 +245,7 @@ func TestStartPollAppliesAChange(t *testing.T) { } time.Sleep(10 * time.Millisecond) } - if got := holder.Load().TraceProbability; got != 0.5 { + if got := holder.Common().TraceProbability; got != 0.5 { t.Fatalf("holder trace_probability = %v, want 0.5", got) } @@ -318,17 +257,47 @@ func TestStartPollAppliesAChange(t *testing.T) { } } +// A poll that brings a value the class refuses goes to onError and leaves +// the holder as it was. +func TestStartPollReportsRefusedConfig(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(`{"test_workers": 0}`)) + })) + defer ts.Close() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + holder := NewHolder(knobsCodec) + errs := make(chan error, 1) + StartPoll(ctx, ts.URL, holder, &fakeSampler{}, 10*time.Millisecond, time.Second, func(err error) { + select { + case errs <- err: + default: + } + }) + select { + case err := <-errs: + if !strings.Contains(err.Error(), "test_workers") { + t.Errorf("onError got %v, want the refusing key named", err) + } + case <-time.After(2 * time.Second): + t.Fatal("a refused poll never reached onError") + } + if got := Get[knobs](holder); got != knobsCodec.Defaults { + t.Errorf("a refused poll changed the holder to %+v", got) + } +} + func TestStartPollStopsWithTheContext(t *testing.T) { var hits atomic.Int64 ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { hits.Add(1) - w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(`{}`)) })) defer ts.Close() ctx, cancel := context.WithCancel(context.Background()) - StartPoll(ctx, ts.URL, NewHolder(Config{}), &fakeSampler{}, + StartPoll(ctx, ts.URL, NewHolder(nil), &fakeSampler{}, 10*time.Millisecond, time.Second, func(error) {}) time.Sleep(60 * time.Millisecond) cancel() @@ -351,7 +320,6 @@ func TestSubscribeSpawnReportsPriorSuccess(t *testing.T) { http.Error(w, "busy", http.StatusServiceUnavailable) return } - w.Header().Set("Content-Type", "application/json") _, _ = w.Write([]byte(`{"trace_probability": 0.25}`)) })) defer ts.Close() @@ -361,7 +329,7 @@ func TestSubscribeSpawnReportsPriorSuccess(t *testing.T) { fetched bool } var calls []call - holder := NewHolder(Config{}) + holder := NewHolder(nil) // boomer.Events is a process-wide bus and SubscribeSpawn owns the handler // value, so the subscription outlives this test. No other test publishes // boomer:spawn, and each test uses its own holder, so that is harmless. @@ -373,8 +341,8 @@ func TestSubscribeSpawnReportsPriorSuccess(t *testing.T) { fail.Store(true) boomer.Events.Publish("boomer:spawn", 1, 1.0) - if len(calls) != 1 || calls[0].fetched { - t.Fatalf("after a failure with no prior success: calls = %+v, want one with fetched=false", calls) + if len(calls) != 1 || calls[0].fetched || calls[0].err == nil { + t.Fatalf("after a failure with no prior success: calls = %+v, want one error with fetched=false", calls) } fail.Store(false) @@ -382,7 +350,7 @@ func TestSubscribeSpawnReportsPriorSuccess(t *testing.T) { if len(calls) != 1 { t.Fatalf("a successful fetch invoked onError: %+v", calls) } - if got := holder.Load().TraceProbability; got != 0.25 { + if got := holder.Common().TraceProbability; got != 0.25 { t.Fatalf("holder trace_probability = %v, want 0.25", got) } @@ -391,7 +359,7 @@ func TestSubscribeSpawnReportsPriorSuccess(t *testing.T) { if len(calls) != 2 || !calls[1].fetched { t.Fatalf("after a failure with a prior success: calls = %+v, want a second with fetched=true", calls) } - if got := holder.Load().TraceProbability; got != 0.25 { + if got := holder.Common().TraceProbability; got != 0.25 { t.Fatalf("failed fetch changed holder trace_probability to %v", got) } } @@ -404,7 +372,7 @@ func TestStartPollZeroInterval(t *testing.T) { })) defer ts.Close() - StartPoll(context.Background(), ts.URL, NewHolder(Config{}), &fakeSampler{}, + StartPoll(context.Background(), ts.URL, NewHolder(nil), &fakeSampler{}, 0, time.Second, func(error) {}) time.Sleep(50 * time.Millisecond) if hits.Load() != 0 { diff --git a/internal/benchmarking/boomer/glutton/cpuload_test.go b/internal/benchmarking/boomer/glutton/cpuload_test.go index aacf4eaddd..91adc77a9f 100644 --- a/internal/benchmarking/boomer/glutton/cpuload_test.go +++ b/internal/benchmarking/boomer/glutton/cpuload_test.go @@ -19,13 +19,12 @@ import ( "net/http" "testing" - "github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig" "github.com/agent-substrate/substrate/internal/benchmarking/glutton/fake" ) func TestEnsureCPULoadRequestsConfiguredLoad(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{CPUCores: 2, CPUDutyCycle: 0.1}) + u := newTestGluttonActor(t, srv, gluttonKnobs{CPUCores: 2, CPUDutyCycle: 0.1}) u.ensureCPULoad(context.Background()) @@ -49,7 +48,7 @@ func TestEnsureCPULoadRequestsConfiguredLoad(t *testing.T) { func TestEnsureCPULoadDisabledByDefault(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{}) + u := newTestGluttonActor(t, srv, gluttonKnobs{}) u.ensureCPULoad(context.Background()) @@ -63,7 +62,7 @@ func TestEnsureCPULoadDisabledByDefault(t *testing.T) { func TestEnsureCPULoadRetriesAfterFailure(t *testing.T) { srv := &fake.Server{Status: http.StatusServiceUnavailable} - u := newTestGluttonActor(t, srv, dynconfig.Config{CPUCores: 1, CPUDutyCycle: 0.5}) + u := newTestGluttonActor(t, srv, gluttonKnobs{CPUCores: 1, CPUDutyCycle: 0.5}) ctx := context.Background() u.ensureCPULoad(ctx) diff --git a/internal/benchmarking/boomer/glutton/durdir.go b/internal/benchmarking/boomer/glutton/durdir.go index ff9e072f8a..6017dce8c5 100644 --- a/internal/benchmarking/boomer/glutton/durdir.go +++ b/internal/benchmarking/boomer/glutton/durdir.go @@ -22,7 +22,6 @@ import ( "fmt" "io" "log/slog" - "math/rand/v2" "net/http" "strings" "sync" @@ -65,6 +64,7 @@ func init() { Name: "durdir", LocustFile: "durdir.py", UserClass: durDirUserClass, + Config: durDirCodec, Init: initDurDir, }) } @@ -85,19 +85,14 @@ type durDirRuntime struct { } func (r *durDirRuntime) dynamicWait() time.Duration { - cfg := r.cfg.Dyn.Load() - if cfg.MaxWait <= cfg.MinWait { - return cfg.MinWait - } - jitter := cfg.MaxWait - cfg.MinWait - return cfg.MinWait + time.Duration(rand.Float64()*float64(jitter)) + return dynconfig.Get[durDirKnobs](r.cfg.Dyn).WaitTime.Draw() } func (r *durDirRuntime) iterate() { gid := boomerutil.GoroutineID() val, loaded := r.users.Load(gid) if !loaded { - dynCfg := r.cfg.Dyn.Load() + dynCfg := dynconfig.Get[durDirKnobs](r.cfg.Dyn) u, err := r.startUser(context.Background(), dynCfg) if err != nil { slog.Warn("durdir on_start failed; goroutine will retry next iter", @@ -109,15 +104,15 @@ func (r *durDirRuntime) iterate() { } user := val.(*durDirUser) - dynCfg := r.cfg.Dyn.Load() + dynCfg := dynconfig.Get[durDirKnobs](r.cfg.Dyn) ctx := context.Background() user.step(ctx, dynCfg) time.Sleep(r.dynamicWait()) } -func (r *durDirRuntime) startUser(ctx context.Context, dynCfg dynconfig.Config) (*durDirUser, error) { - tmpl := dynCfg.DurDirTemplate +func (r *durDirRuntime) startUser(ctx context.Context, dynCfg durDirKnobs) (*durDirUser, error) { + tmpl := dynCfg.Template if tmpl == "" { tmpl = defaultDurTemplate } @@ -146,7 +141,7 @@ func (r *durDirRuntime) startUser(ctx context.Context, dynCfg dynconfig.Config) } func (r *durDirRuntime) shutdown(ctx context.Context) { - dynCfg := r.cfg.Dyn.Load() + dynCfg := dynconfig.Get[durDirKnobs](r.cfg.Dyn) r.users.Range(func(_, val any) bool { u := val.(*durDirUser) u.hibernateAndDelete(ctx, dynCfg) @@ -231,7 +226,7 @@ func (u *durDirUser) suspend(ctx context.Context) { }) } -func (u *durDirUser) hibernate(ctx context.Context, dynCfg dynconfig.Config) { +func (u *durDirUser) hibernate(ctx context.Context, dynCfg durDirKnobs) { if dynCfg.LifecycleMode == dynconfig.LifecycleModePause { u.pause(ctx) } else { @@ -242,7 +237,7 @@ func (u *durDirUser) hibernate(ctx context.Context, dynCfg dynconfig.Config) { // hibernateAndDelete hibernates (suspends or pauses) the actor before deleting it. // The hibernate call is unmetered (teardown precondition, not benchmark latency), // while the delete is metered so true leaks still surface in failures.csv. -func (u *durDirUser) hibernateAndDelete(ctx context.Context, dynCfg dynconfig.Config) { +func (u *durDirUser) hibernateAndDelete(ctx context.Context, dynCfg durDirKnobs) { if dynCfg.LifecycleMode == dynconfig.LifecycleModePause { _, _ = u.cfg.APIStub.PauseActor(ctx, &ateapipb.PauseActorRequest{ Actor: u.ref(), @@ -287,19 +282,19 @@ func (u *durDirUser) tracedCall(ctx context.Context, name string, do func(contex return nil } -func (u *durDirUser) params(dynCfg dynconfig.Config) (int64, gluttonpb.ReadMode) { - fileSize := dynCfg.DurDirFileSize +func (u *durDirUser) params(dynCfg durDirKnobs) (int64, gluttonpb.ReadMode) { + fileSize := dynCfg.FileSize if fileSize <= 0 { fileSize = defaultFileSize } readMode := gluttonpb.ReadMode_READ_MODE_DATA - if dynCfg.DurDirReadMode == dynconfig.ReadModeDigest { + if dynCfg.ReadMode == readModeDigest { readMode = gluttonpb.ReadMode_READ_MODE_DIGEST_ONLY } return fileSize, readMode } -func (u *durDirUser) step(ctx context.Context, dynCfg dynconfig.Config) { +func (u *durDirUser) step(ctx context.Context, dynCfg durDirKnobs) { fileSize, readMode := u.params(dynCfg) // 1. Suspend or pause actor @@ -326,7 +321,7 @@ func (u *durDirUser) step(ctx context.Context, dynCfg dynconfig.Config) { } } -func (u *durDirUser) bootstrap(ctx context.Context, dynCfg dynconfig.Config) error { +func (u *durDirUser) bootstrap(ctx context.Context, dynCfg durDirKnobs) error { fileSize, readMode := u.params(dynCfg) if !u.resume(ctx, dynCfg.ResumeMode) { diff --git a/internal/benchmarking/boomer/glutton/durdir_test.go b/internal/benchmarking/boomer/glutton/durdir_test.go index dd1b5d9410..d0149ed4ee 100644 --- a/internal/benchmarking/boomer/glutton/durdir_test.go +++ b/internal/benchmarking/boomer/glutton/durdir_test.go @@ -66,16 +66,14 @@ func TestDurDirLoopSequence(t *testing.T) { fakeCtrl := &fakeControlClient{} cfg := &userclass.Config{ APIStub: fakeCtrl, - Dyn: dynconfig.NewHolder(dynconfig.Config{ - ResumeMode: tc.resumeMode, - LifecycleMode: tc.lifecycleMode, + Dyn: dynconfig.Static(durDirKnobs{ + Lifecycle: dynconfig.Lifecycle{ResumeMode: tc.resumeMode, LifecycleMode: tc.lifecycleMode}, }), } du := newTestDurDirUser(t, srv, cfg) du.expectedDigest = srv.HexDigest() - dynCfg := cfg.Dyn.Load() - du.step(context.Background(), dynCfg) + du.step(context.Background(), dynconfig.Get[durDirKnobs](cfg.Dyn)) if got := fakeCtrl.recordedCalls(); !reflect.DeepEqual(got, tc.wantGRPCCall) { t.Errorf("gRPC calls: got %v, want %v", got, tc.wantGRPCCall) @@ -195,14 +193,14 @@ func TestDurDirBootstrapUsesConfiguredResumeMode(t *testing.T) { fakeCtrl := &fakeControlClient{} cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, - Dyn: dynconfig.NewHolder(dynconfig.Config{ - DurDirFileSize: int64(len(srv.Data)), - ResumeMode: tc.resumeMode, + Dyn: dynconfig.Static(durDirKnobs{ + FileSize: int64(len(srv.Data)), + Lifecycle: dynconfig.Lifecycle{ResumeMode: tc.resumeMode}, }), }) rt := &durDirRuntime{cfg: cfg} - _, err := rt.startUser(context.Background(), cfg.Dyn.Load()) + _, err := rt.startUser(context.Background(), dynconfig.Get[durDirKnobs](cfg.Dyn)) if err != nil { t.Fatalf("startUser failed: %v", err) } @@ -221,14 +219,14 @@ func TestDurDirBootstrapFailureSuspendsBeforeDelete(t *testing.T) { fakeCtrl := &fakeControlClient{} cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, - Dyn: dynconfig.NewHolder(dynconfig.Config{ - DurDirFileSize: 1024, - ResumeMode: dynconfig.ResumeModeExplicit, + Dyn: dynconfig.Static(durDirKnobs{ + FileSize: 1024, + Lifecycle: dynconfig.Lifecycle{ResumeMode: dynconfig.ResumeModeExplicit}, }), }) rt := &durDirRuntime{cfg: cfg} - _, err := rt.startUser(context.Background(), cfg.Dyn.Load()) + _, err := rt.startUser(context.Background(), dynconfig.Get[durDirKnobs](cfg.Dyn)) if err == nil { t.Fatalf("startUser expected error on failing server, got nil") } @@ -243,6 +241,7 @@ func TestDurDirShutdownSuspendsBeforeDelete(t *testing.T) { fakeCtrl := &fakeControlClient{} cfg := &userclass.Config{ APIStub: fakeCtrl, + Dyn: dynconfig.Static(durDirKnobs{}), } du := newTestDurDirUser(t, &fake.Server{}, cfg) @@ -264,8 +263,8 @@ func TestDurDirShutdownPausesBeforeDelete(t *testing.T) { fakeCtrl := &fakeControlClient{} cfg := &userclass.Config{ APIStub: fakeCtrl, - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModePause, + Dyn: dynconfig.Static(durDirKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModePause}, }), } du := newTestDurDirUser(t, &fake.Server{}, cfg) diff --git a/internal/benchmarking/boomer/glutton/fixture_test.go b/internal/benchmarking/boomer/glutton/fixture_test.go index 170afe8cca..cbdad69701 100644 --- a/internal/benchmarking/boomer/glutton/fixture_test.go +++ b/internal/benchmarking/boomer/glutton/fixture_test.go @@ -140,7 +140,7 @@ func newTestConfig(t *testing.T, srv *fake.Server, cfg *userclass.Config) *userc cfg.Tracer = otel.Tracer("test") } if cfg.Dyn == nil { - cfg.Dyn = dynconfig.NewHolder(dynconfig.Config{}) + cfg.Dyn = dynconfig.Static(gluttonKnobs{}) } cfg.HTTPClient = ts.Client() cfg.RouterURL = ts.URL diff --git a/internal/benchmarking/boomer/glutton/knobs.go b/internal/benchmarking/boomer/glutton/knobs.go new file mode 100644 index 0000000000..a3396140bd --- /dev/null +++ b/internal/benchmarking/boomer/glutton/knobs.go @@ -0,0 +1,150 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package glutton + +import ( + "fmt" + "math" + "time" + + "github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig" +) + +// The three user classes in this package each read their own slice of the +// runtime config. The structs below name the keys; the sections carry the +// defaults a missing key falls back to and the rules a fetched value must +// pass before the worker takes it on. + +// gluttonKnobs is the GluttonUser slice: the wait and live windows, the +// lifecycle modes, the working-set sizes, the CPU load, and the ping cap. +type gluttonKnobs struct { + dynconfig.WaitTime + dynconfig.Lifecycle + // MinLive and MaxLive bound how long an actor stays resumed between its + // first ping and its suspend; a zero window suspends right after the ping. + MinLive dynconfig.Seconds `json:"min_live_time"` + MaxLive dynconfig.Seconds `json:"max_live_time"` + // MemTarget is the resident RAM filled via WriteRAM, suffixed ("2Gi"); + // MemChurn the RAM re-randomized in place each cycle; MemRead the RAM + // walked after each resume, or "all". Empty disables each. Glutton + // parses the sizes, so a bad one fails loudly as a GluttonFillRAM, + // GluttonChurnRAM, or GluttonReadRAM error rather than here. + MemTarget string `json:"mem_target"` + MemChurn string `json:"mem_churn"` + MemRead string `json:"mem_read"` + // CPUCores goroutines each burn CPUDutyCycle of one core via UseCPU; 0 + // disables. + CPUCores int `json:"cpu_cores"` + CPUDutyCycle float64 `json:"cpu_duty_cycle"` + // MaxPingsPerWake caps the pings sent during one resume/suspend cycle; + // values below 1 read as 1. + MaxPingsPerWake int `json:"max_pings_per_wake"` +} + +var gluttonCodec = dynconfig.Typed[gluttonKnobs]{ + Defaults: gluttonKnobs{ + WaitTime: dynconfig.WaitTime{MaxWait: dynconfig.Seconds(500 * time.Millisecond)}, + MaxPingsPerWake: 1, + }, + Validate: func(k gluttonKnobs) error { + if err := k.WaitTime.Validate(); err != nil { + return err + } + if err := k.Lifecycle.Validate(); err != nil { + return err + } + if k.MinLive < 0 { + return fmt.Errorf("min_live_time cannot be negative: %v", k.MinLive.Duration()) + } + if k.MaxLive < 0 { + return fmt.Errorf("max_live_time cannot be negative: %v", k.MaxLive.Duration()) + } + if k.MaxLive < k.MinLive { + return fmt.Errorf("max_live_time (%v) cannot be less than min_live_time (%v)", k.MaxLive.Duration(), k.MinLive.Duration()) + } + if k.CPUCores < 0 { + return fmt.Errorf("cpu_cores cannot be negative: %d", k.CPUCores) + } + if k.CPUDutyCycle < 0 || k.CPUDutyCycle > 1 { + return fmt.Errorf("cpu_duty_cycle must be between 0.0 and 1.0, got: %f", k.CPUDutyCycle) + } + return nil + }, +} + +// DurDir read modes: whether the actor ships the file's bytes back or only +// its digest. +const ( + readModeData = "data" + readModeDigest = "digest" +) + +// durDirKnobs is the DurdirUser slice. +type durDirKnobs struct { + dynconfig.WaitTime + dynconfig.Lifecycle + // FileSize is the file written and re-read each cycle, in bytes; 0 + // falls back to defaultFileSize. + FileSize int64 `json:"durdir_file_size_bytes"` + // ReadMode is readModeData or readModeDigest; empty reads as data. + ReadMode string `json:"durdir_read_mode"` + // Template is the ActorTemplate name; empty falls back to + // defaultDurTemplate. + Template string `json:"durdir_template"` +} + +var durDirCodec = dynconfig.Typed[durDirKnobs]{ + Validate: func(k durDirKnobs) error { + if err := k.WaitTime.Validate(); err != nil { + return err + } + if err := k.Lifecycle.Validate(); err != nil { + return err + } + if k.FileSize < 0 { + return fmt.Errorf("durdir_file_size_bytes cannot be negative: %d", k.FileSize) + } + if k.FileSize > math.MaxInt32 { + return fmt.Errorf("durdir_file_size_bytes cannot exceed %d (2 GiB), got: %d", math.MaxInt32, k.FileSize) + } + if k.ReadMode != "" && k.ReadMode != readModeData && k.ReadMode != readModeDigest { + return fmt.Errorf("invalid durdir_read_mode %q: must be %q or %q", k.ReadMode, readModeData, readModeDigest) + } + return nil + }, +} + +// spawnKnobs is the SpawnUser slice. Each zero keeps the matching +// boomer-worker flag. +type spawnKnobs struct { + TotalActors int `json:"total_actors"` + SpawnConcurrency int `json:"spawn_concurrency"` + ActorDeadline dynconfig.Seconds `json:"actor_deadline"` +} + +var spawnCodec = dynconfig.Typed[spawnKnobs]{ + Validate: func(k spawnKnobs) error { + if k.TotalActors < 0 { + return fmt.Errorf("total_actors cannot be negative: %d", k.TotalActors) + } + if k.SpawnConcurrency < 0 { + return fmt.Errorf("spawn_concurrency cannot be negative: %d", k.SpawnConcurrency) + } + if k.ActorDeadline < 0 { + return fmt.Errorf("actor_deadline cannot be negative: %v", k.ActorDeadline.Duration()) + } + return nil + }, +} diff --git a/internal/benchmarking/boomer/glutton/lifecycle.go b/internal/benchmarking/boomer/glutton/lifecycle.go index 53c3f76105..77239d04ad 100644 --- a/internal/benchmarking/boomer/glutton/lifecycle.go +++ b/internal/benchmarking/boomer/glutton/lifecycle.go @@ -76,6 +76,7 @@ func init() { Name: "glutton", LocustFile: "glutton.py", UserClass: userClass, + Config: gluttonCodec, Init: initPing, }) } @@ -181,7 +182,7 @@ func (r *taskRuntime) iterate() { // actor stays live for the full window. With the default zero window // the actor is hibernated right after the first ping. deadline := time.Now().Add(r.liveWait()) - maxPings := max(r.cfg.Dyn.Load().MaxPingsPerWake, 1) + maxPings := max(dynconfig.Get[gluttonKnobs](r.cfg.Dyn).MaxPingsPerWake, 1) actor.ping(ctx) for sent := 1; sent < maxPings; sent++ { gap := minPingGap + time.Duration(rand.Float64()*float64(maxPingGap-minPingGap)) @@ -261,24 +262,15 @@ func (r *taskRuntime) shutdown(ctx context.Context) { // dynamicWait is the gap between suspending one actor and resuming the // VU's next one, drawn uniformly from [MinWait, MaxWait]. func (r *taskRuntime) dynamicWait() time.Duration { - cfg := r.cfg.Dyn.Load() - return uniformWait(cfg.MinWait, cfg.MaxWait) + return dynconfig.Get[gluttonKnobs](r.cfg.Dyn).WaitTime.Draw() } // liveWait is how long an actor stays resumed between its first ping and // its suspend, drawn uniformly from [MinLive, MaxLive]. The default zero // window suspends right after the ping. func (r *taskRuntime) liveWait() time.Duration { - cfg := r.cfg.Dyn.Load() - return uniformWait(cfg.MinLive, cfg.MaxLive) -} - -// uniformWait draws from [lo, hi]; an inverted or empty range yields lo. -func uniformWait(lo, hi time.Duration) time.Duration { - if hi <= lo { - return lo - } - return lo + time.Duration(rand.Float64()*float64(hi-lo)) + knobs := dynconfig.Get[gluttonKnobs](r.cfg.Dyn) + return dynconfig.Uniform(knobs.MinLive.Duration(), knobs.MaxLive.Duration()) } // gluttonUser is one VU (boomer goroutine). It owns --actors-per-user actors @@ -459,7 +451,7 @@ func (u *gluttonActor) resume(ctx context.Context) error { // lifecycle mode selects: PauseActor keeps the snapshot on the node, while // SuspendActor writes it to durable storage. func (u *gluttonActor) hibernate(ctx context.Context) error { - if u.cfg.Dyn.Load().LifecycleMode == dynconfig.LifecycleModePause { + if dynconfig.Get[gluttonKnobs](u.cfg.Dyn).LifecycleMode == dynconfig.LifecycleModePause { return u.pause(ctx) } return u.suspend(ctx) @@ -600,7 +592,7 @@ func (u *gluttonActor) ensureRAMFilled(ctx context.Context) { if u.ramFilled { return } - target := u.cfg.Dyn.Load().MemTarget + target := dynconfig.Get[gluttonKnobs](u.cfg.Dyn).MemTarget if target == "" { u.ramFilled = true return @@ -631,8 +623,8 @@ func (u *gluttonActor) ensureCPULoad(ctx context.Context) { if u.cpuLoaded { return } - dyn := u.cfg.Dyn.Load() - if dyn.CPUCores == 0 { + knobs := dynconfig.Get[gluttonKnobs](u.cfg.Dyn) + if knobs.CPUCores == 0 { u.cpuLoaded = true return } @@ -642,8 +634,8 @@ func (u *gluttonActor) ensureCPULoad(ctx context.Context) { start := time.Now() err := u.postProto(ctx, useCPUPath, &gluttonpb.UseCPURequest{ - NumCores: int32(dyn.CPUCores), - DutyCycle: dyn.CPUDutyCycle, + NumCores: int32(knobs.CPUCores), + DutyCycle: knobs.CPUDutyCycle, }, &gluttonpb.UseCPUResponse{}) clientLatency := time.Since(start) boomerutil.LogSampledTrace(span, "GluttonUseCPU", clientLatency, boomerutil.SourceClient, err) @@ -665,7 +657,7 @@ func (u *gluttonActor) ensureCPULoad(ctx context.Context) { // iteration, only after the fill has succeeded, and reports as its own // GluttonChurnRAM stats row. func (u *gluttonActor) churnRAM(ctx context.Context) { - churn := u.cfg.Dyn.Load().MemChurn + churn := dynconfig.Get[gluttonKnobs](u.cfg.Dyn).MemChurn if churn == "" || !u.ramFilled { return } @@ -691,7 +683,7 @@ func (u *gluttonActor) churnRAM(ctx context.Context) { // memory; on an eagerly-restored actor it degenerates to a fast in-memory // scan, so the two restore modes are directly comparable. func (u *gluttonActor) readRAM(ctx context.Context) { - read := u.cfg.Dyn.Load().MemRead + read := dynconfig.Get[gluttonKnobs](u.cfg.Dyn).MemRead if read == "" || !u.ramFilled { return } diff --git a/internal/benchmarking/boomer/glutton/lifecycle_test.go b/internal/benchmarking/boomer/glutton/lifecycle_test.go index 96d0872780..6c12b6d537 100644 --- a/internal/benchmarking/boomer/glutton/lifecycle_test.go +++ b/internal/benchmarking/boomer/glutton/lifecycle_test.go @@ -33,8 +33,8 @@ func TestGluttonIterate_SuspendMode(t *testing.T) { cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModeSuspend, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModeSuspend}, }), }) @@ -56,8 +56,8 @@ func TestGluttonIterate_PauseMode(t *testing.T) { cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModePause, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModePause}, }), }) @@ -79,8 +79,8 @@ func TestGluttonShutdown_PauseModeRunningActor(t *testing.T) { cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModePause, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModePause}, }), }) @@ -110,8 +110,8 @@ func TestGluttonShutdown_DeleteSetsAnyState(t *testing.T) { cfg := newTestConfig(t, srv, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModeSuspend, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModeSuspend}, }), }) @@ -143,8 +143,8 @@ func newReplacementRuntime(t *testing.T, resumeErrs ...error) (*taskRuntime, *fa cfg := newTestConfig(t, &fake.Server{}, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModeSuspend, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModeSuspend}, }), }) return &taskRuntime{cfg: cfg}, fakeCtrl @@ -231,8 +231,8 @@ func TestGluttonIterate_RetriesStrandedHibernate(t *testing.T) { cfg := newTestConfig(t, &fake.Server{}, &userclass.Config{ APIStub: fakeCtrl, Atespace: "bench-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{ - LifecycleMode: dynconfig.LifecycleModeSuspend, + Dyn: dynconfig.Static(gluttonKnobs{ + Lifecycle: dynconfig.Lifecycle{LifecycleMode: dynconfig.LifecycleModeSuspend}, }), }) rt := &taskRuntime{cfg: cfg} diff --git a/internal/benchmarking/boomer/glutton/memfill_test.go b/internal/benchmarking/boomer/glutton/memfill_test.go index 00485ae8ce..4656dcb84b 100644 --- a/internal/benchmarking/boomer/glutton/memfill_test.go +++ b/internal/benchmarking/boomer/glutton/memfill_test.go @@ -24,9 +24,9 @@ import ( gluttonpb "github.com/agent-substrate/substrate/internal/proto/glutton" ) -func newTestGluttonActor(t *testing.T, srv *fake.Server, dyn dynconfig.Config) *gluttonActor { +func newTestGluttonActor(t *testing.T, srv *fake.Server, dyn gluttonKnobs) *gluttonActor { t.Helper() - cfg := newTestConfig(t, srv, &userclass.Config{Dyn: dynconfig.NewHolder(dyn)}) + cfg := newTestConfig(t, srv, &userclass.Config{Dyn: dynconfig.Static(dyn)}) return &gluttonActor{ cfg: cfg, actorName: "memactor", @@ -35,7 +35,7 @@ func newTestGluttonActor(t *testing.T, srv *fake.Server, dyn dynconfig.Config) * func TestEnsureRAMFilledRequestsTarget(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "2Gi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "2Gi"}) u.ensureRAMFilled(context.Background()) @@ -56,7 +56,7 @@ func TestEnsureRAMFilledRequestsTarget(t *testing.T) { func TestEnsureRAMFilledDisabledByDefault(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{}) + u := newTestGluttonActor(t, srv, gluttonKnobs{}) u.ensureRAMFilled(context.Background()) @@ -70,7 +70,7 @@ func TestEnsureRAMFilledDisabledByDefault(t *testing.T) { func TestChurnRAMOverwritesEachCycle(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Gi", MemChurn: "64Mi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Gi", MemChurn: "64Mi"}) ctx := context.Background() // Churn before fill is a no-op: there is nothing to overwrite yet. @@ -104,7 +104,7 @@ func TestChurnRAMOverwritesEachCycle(t *testing.T) { func TestChurnRAMDisabledByDefault(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Gi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Gi"}) ctx := context.Background() u.ensureRAMFilled(ctx) @@ -116,7 +116,7 @@ func TestChurnRAMDisabledByDefault(t *testing.T) { func TestReadRAMWalksAfterFill(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Gi", MemRead: "all"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Gi", MemRead: "all"}) ctx := context.Background() // Read before fill is a no-op: there is nothing to walk yet. @@ -144,7 +144,7 @@ func TestReadRAMWalksAfterFill(t *testing.T) { func TestReadRAMPassesSizeVerbatim(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Gi", MemRead: "512Mi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Gi", MemRead: "512Mi"}) ctx := context.Background() u.ensureRAMFilled(ctx) @@ -158,7 +158,7 @@ func TestReadRAMPassesSizeVerbatim(t *testing.T) { func TestReadRAMDisabledByDefault(t *testing.T) { srv := &fake.Server{} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Gi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Gi"}) ctx := context.Background() u.ensureRAMFilled(ctx) @@ -170,7 +170,7 @@ func TestReadRAMDisabledByDefault(t *testing.T) { func TestEnsureRAMFilledRetriesAfterFailure(t *testing.T) { srv := &fake.Server{Status: 503} - u := newTestGluttonActor(t, srv, dynconfig.Config{MemTarget: "1Mi"}) + u := newTestGluttonActor(t, srv, gluttonKnobs{MemTarget: "1Mi"}) u.ensureRAMFilled(context.Background()) if u.ramFilled { diff --git a/internal/benchmarking/boomer/glutton/spawn.go b/internal/benchmarking/boomer/glutton/spawn.go index d9df58df7e..82ef65e57c 100644 --- a/internal/benchmarking/boomer/glutton/spawn.go +++ b/internal/benchmarking/boomer/glutton/spawn.go @@ -30,6 +30,7 @@ import ( "github.com/agent-substrate/substrate/internal/ateinterceptors" "github.com/agent-substrate/substrate/internal/atenet" "github.com/agent-substrate/substrate/internal/benchmarking/boomer/boomerutil" + "github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig" bmetrics "github.com/agent-substrate/substrate/internal/benchmarking/boomer/metrics" "github.com/agent-substrate/substrate/internal/benchmarking/boomer/userclass" gluttonpb "github.com/agent-substrate/substrate/internal/proto/glutton" @@ -75,6 +76,7 @@ func init() { Name: "spawn", LocustFile: "spawn.py", UserClass: spawnUserClass, + Config: spawnCodec, Init: initSpawn, }) } @@ -161,17 +163,15 @@ func (r *spawnRuntime) runBatch(ctx context.Context) { spawnConcurrency := r.cfg.SpawnConcurrency deadline := r.cfg.ActorDeadline // Values set in the web UI form (dynconfig) override the flags. - if r.cfg.Dyn != nil { - dyn := r.cfg.Dyn.Load() - if dyn.TotalActors > 0 { - totalActors = dyn.TotalActors - } - if dyn.SpawnConcurrency > 0 { - spawnConcurrency = dyn.SpawnConcurrency - } - if dyn.ActorDeadline > 0 { - deadline = dyn.ActorDeadline - } + knobs := dynconfig.Get[spawnKnobs](r.cfg.Dyn) + if knobs.TotalActors > 0 { + totalActors = knobs.TotalActors + } + if knobs.SpawnConcurrency > 0 { + spawnConcurrency = knobs.SpawnConcurrency + } + if knobs.ActorDeadline > 0 { + deadline = knobs.ActorDeadline.Duration() } spawnConcurrency = min(spawnConcurrency, totalActors) if spawnConcurrency < 1 { diff --git a/internal/benchmarking/boomer/glutton/spawn_test.go b/internal/benchmarking/boomer/glutton/spawn_test.go index 0199d8109b..4c0fac212d 100644 --- a/internal/benchmarking/boomer/glutton/spawn_test.go +++ b/internal/benchmarking/boomer/glutton/spawn_test.go @@ -646,10 +646,10 @@ func TestSpawnRunBatch_DynConfigOverride(t *testing.T) { t.Run("overrides when set", func(t *testing.T) { fakeAPI := &fakeControlClient{} - dyn := dynconfig.NewHolder(dynconfig.Config{ + dyn := dynconfig.Static(spawnKnobs{ TotalActors: 3, SpawnConcurrency: 2, - ActorDeadline: 5 * time.Second, + ActorDeadline: dynconfig.Seconds(5 * time.Second), }) cfg := &userclass.Config{ APIStub: fakeAPI, @@ -684,7 +684,7 @@ func TestSpawnRunBatch_DynConfigOverride(t *testing.T) { t.Run("falls back to cfg when dynconfig is zero", func(t *testing.T) { fakeAPI := &fakeControlClient{} - dyn := dynconfig.NewHolder(dynconfig.Config{ + dyn := dynconfig.Static(spawnKnobs{ TotalActors: 0, // 0 = unset SpawnConcurrency: 0, ActorDeadline: 0, diff --git a/internal/benchmarking/boomer/glutton/wait_test.go b/internal/benchmarking/boomer/glutton/wait_test.go index de9313de40..d08a5a43e2 100644 --- a/internal/benchmarking/boomer/glutton/wait_test.go +++ b/internal/benchmarking/boomer/glutton/wait_test.go @@ -25,22 +25,22 @@ import ( func TestUniformWaitStaysInRange(t *testing.T) { lo, hi := 200*time.Millisecond, time.Second for i := 0; i < 1000; i++ { - got := uniformWait(lo, hi) + got := dynconfig.Uniform(lo, hi) if got < lo || got > hi { - t.Fatalf("uniformWait(%v, %v) = %v, outside range", lo, hi, got) + t.Fatalf("dynconfig.Uniform(%v, %v) = %v, outside range", lo, hi, got) } } } func TestUniformWaitDegenerateRanges(t *testing.T) { - if got := uniformWait(0, 0); got != 0 { - t.Errorf("uniformWait(0, 0) = %v, want 0", got) + if got := dynconfig.Uniform(0, 0); got != 0 { + t.Errorf("dynconfig.Uniform(0, 0) = %v, want 0", got) } - if got := uniformWait(3*time.Second, 3*time.Second); got != 3*time.Second { + if got := dynconfig.Uniform(3*time.Second, 3*time.Second); got != 3*time.Second { t.Errorf("equal bounds: got %v, want 3s", got) } // An inverted range yields the lower bound rather than a negative wait. - if got := uniformWait(5*time.Second, time.Second); got != 5*time.Second { + if got := dynconfig.Uniform(5*time.Second, time.Second); got != 5*time.Second { t.Errorf("inverted bounds: got %v, want 5s", got) } } @@ -48,9 +48,8 @@ func TestUniformWaitDegenerateRanges(t *testing.T) { // The wait window and the live window read different config fields, so a // run that sets only one of them must not leak into the other. func TestWaitAndLiveWindowsAreIndependent(t *testing.T) { - rt := &taskRuntime{cfg: &userclass.Config{Dyn: dynconfig.NewHolder(dynconfig.Config{ - MinWait: 200 * time.Millisecond, - MaxWait: time.Second, + rt := &taskRuntime{cfg: &userclass.Config{Dyn: dynconfig.Static(gluttonKnobs{ + WaitTime: dynconfig.WaitTime{MinWait: dynconfig.Seconds(200 * time.Millisecond), MaxWait: dynconfig.Seconds(time.Second)}, })}} for i := 0; i < 100; i++ { if got := rt.liveWait(); got != 0 { @@ -61,7 +60,7 @@ func TestWaitAndLiveWindowsAreIndependent(t *testing.T) { } } - rt.cfg.Dyn.Store(dynconfig.Config{MinLive: 9 * time.Second, MaxLive: 14 * time.Second}) + rt.cfg.Dyn = dynconfig.Static(gluttonKnobs{MinLive: dynconfig.Seconds(9 * time.Second), MaxLive: dynconfig.Seconds(14 * time.Second)}) for i := 0; i < 100; i++ { if got := rt.dynamicWait(); got != 0 { t.Fatalf("dynamicWait with zero wait window = %v, want 0", got) diff --git a/internal/benchmarking/boomer/sweperf/sweperf.go b/internal/benchmarking/boomer/sweperf/sweperf.go index 37dd467bb1..b4d6666964 100644 --- a/internal/benchmarking/boomer/sweperf/sweperf.go +++ b/internal/benchmarking/boomer/sweperf/sweperf.go @@ -28,7 +28,6 @@ import ( "fmt" "io" "log/slog" - "math/rand/v2" "net/http" "strings" "sync" @@ -37,6 +36,7 @@ import ( "github.com/agent-substrate/substrate/internal/ateinterceptors" "github.com/agent-substrate/substrate/internal/atenet" "github.com/agent-substrate/substrate/internal/benchmarking/boomer/boomerutil" + "github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig" bmetrics "github.com/agent-substrate/substrate/internal/benchmarking/boomer/metrics" "github.com/agent-substrate/substrate/internal/benchmarking/boomer/userclass" "github.com/agent-substrate/substrate/pkg/proto/ateapipb" @@ -84,10 +84,40 @@ func init() { Name: "sweperf", LocustFile: "sweperf.py", UserClass: sweperfUserClass, + Config: knobsCodec, Init: initSweperf, }) } +// knobs is the SweperfUser slice of the runtime config. Each zero falls +// back to the default above; the master populates the keys from the +// --sweperf-* locust flags (common/sweperf_config.py). +type knobs struct { + dynconfig.WaitTime + Template string `json:"sweperf_template"` + TotalSteps int `json:"sweperf_total_steps"` + NumCycles int `json:"sweperf_num_cycles"` + PollIntervalMs int `json:"sweperf_poll_interval_ms"` +} + +var knobsCodec = dynconfig.Typed[knobs]{ + Validate: func(k knobs) error { + if err := k.WaitTime.Validate(); err != nil { + return err + } + if k.TotalSteps < 0 { + return fmt.Errorf("sweperf_total_steps cannot be negative: %d", k.TotalSteps) + } + if k.NumCycles < 0 { + return fmt.Errorf("sweperf_num_cycles cannot be negative: %d", k.NumCycles) + } + if k.PollIntervalMs < 0 { + return fmt.Errorf("sweperf_poll_interval_ms cannot be negative: %d", k.PollIntervalMs) + } + return nil + }, +} + // chunk partitions totalSteps instructions into numCycles slices. type chunk struct { start int @@ -151,19 +181,19 @@ type sweperfRuntime struct { // the --sweperf-* locust flags (common/sweperf_config.py); an unset field // falls back to the built-in default below. func (r *sweperfRuntime) resolveConfig() (string, int, int) { - dyn := r.cfg.Dyn.Load() + dyn := dynconfig.Get[knobs](r.cfg.Dyn) - template := dyn.SweperfTemplate + template := dyn.Template if template == "" { template = defaultSweperfTemplate } - totalSteps := dyn.SweperfTotalSteps + totalSteps := dyn.TotalSteps if totalSteps <= 0 { totalSteps = defaultSweperfTotalSteps } - numCycles := dyn.SweperfNumCycles + numCycles := dyn.NumCycles if numCycles <= 0 { numCycles = defaultSweperfNumCycles } @@ -175,7 +205,7 @@ func (r *sweperfRuntime) resolveConfig() (string, int, int) { // (--sweperf-poll-interval-ms), or the default when unset. Read per job, so a // mid-run change applies to the next cycle. func pollInterval(cfg *userclass.Config) time.Duration { - if ms := cfg.Dyn.Load().SweperfPollIntervalMs; ms > 0 { + if ms := dynconfig.Get[knobs](cfg.Dyn).PollIntervalMs; ms > 0 { return time.Duration(ms) * time.Millisecond } return defaultSweperfPollInterval @@ -184,12 +214,7 @@ func pollInterval(cfg *userclass.Config) time.Duration { // dynamicWait is the think time between cycles: a uniform draw from // [MinWait, MaxWait), or MinWait when the range is empty. func (r *sweperfRuntime) dynamicWait() time.Duration { - cfg := r.cfg.Dyn.Load() - if cfg.MaxWait <= cfg.MinWait { - return cfg.MinWait - } - jitter := cfg.MaxWait - cfg.MinWait - return cfg.MinWait + time.Duration(rand.Float64()*float64(jitter)) + return dynconfig.Get[knobs](r.cfg.Dyn).WaitTime.Draw() } // iterate is the boomer task function, one call per goroutine per iteration. diff --git a/internal/benchmarking/boomer/sweperf/sweperf_test.go b/internal/benchmarking/boomer/sweperf/sweperf_test.go index 0993c52589..8e851c846a 100644 --- a/internal/benchmarking/boomer/sweperf/sweperf_test.go +++ b/internal/benchmarking/boomer/sweperf/sweperf_test.go @@ -190,7 +190,7 @@ func newTestConfig(t *testing.T, handler http.Handler) (*userclass.Config, *http HTTPClient: ts.Client(), RouterURL: ts.URL, Atespace: "benchmark-test", - Dyn: dynconfig.NewHolder(dynconfig.Config{}), + Dyn: dynconfig.Static(knobs{}), Tracer: otel.Tracer("test-sweperf"), } return cfg, ts, fakeCtrl @@ -283,7 +283,7 @@ func TestResolveConfig(t *testing.T) { t.Run("default fallback values", func(t *testing.T) { rt := &sweperfRuntime{ cfg: &userclass.Config{ - Dyn: dynconfig.NewHolder(dynconfig.Config{}), + Dyn: dynconfig.Static(knobs{}), }, } tmpl, steps, cycles := rt.resolveConfig() @@ -301,10 +301,10 @@ func TestResolveConfig(t *testing.T) { t.Run("dynamic config overrides", func(t *testing.T) { rt := &sweperfRuntime{ cfg: &userclass.Config{ - Dyn: dynconfig.NewHolder(dynconfig.Config{ - SweperfTemplate: "custom-template", - SweperfTotalSteps: 50, - SweperfNumCycles: 5, + Dyn: dynconfig.Static(knobs{ + Template: "custom-template", + TotalSteps: 50, + NumCycles: 5, }), }, } @@ -329,7 +329,7 @@ func TestPollInterval(t *testing.T) { {0, defaultSweperfPollInterval}, {250, 250 * time.Millisecond}, } { - cfg := &userclass.Config{Dyn: dynconfig.NewHolder(dynconfig.Config{SweperfPollIntervalMs: tt.ms})} + cfg := &userclass.Config{Dyn: dynconfig.Static(knobs{PollIntervalMs: tt.ms})} if got := pollInterval(cfg); got != tt.want { t.Errorf("pollInterval(%d) = %v, want %v", tt.ms, got, tt.want) } @@ -581,9 +581,8 @@ func TestDynamicWait(t *testing.T) { t.Run("returns MinWait when MaxWait <= MinWait", func(t *testing.T) { rt := &sweperfRuntime{ cfg: &userclass.Config{ - Dyn: dynconfig.NewHolder(dynconfig.Config{ - MinWait: 100 * time.Millisecond, - MaxWait: 50 * time.Millisecond, + Dyn: dynconfig.Static(knobs{ + WaitTime: dynconfig.WaitTime{MinWait: dynconfig.Seconds(100 * time.Millisecond), MaxWait: dynconfig.Seconds(50 * time.Millisecond)}, }), }, } @@ -597,9 +596,8 @@ func TestDynamicWait(t *testing.T) { maxW := 50 * time.Millisecond rt := &sweperfRuntime{ cfg: &userclass.Config{ - Dyn: dynconfig.NewHolder(dynconfig.Config{ - MinWait: minW, - MaxWait: maxW, + Dyn: dynconfig.Static(knobs{ + WaitTime: dynconfig.WaitTime{MinWait: dynconfig.Seconds(minW), MaxWait: dynconfig.Seconds(maxW)}, }), }, } @@ -620,9 +618,8 @@ func TestInitSweperfAndTaskFn(t *testing.T) { } }) cfg, _, _ := newTestConfig(t, handler) - cfg.Dyn = dynconfig.NewHolder(dynconfig.Config{ - MinWait: 1 * time.Millisecond, - MaxWait: 2 * time.Millisecond, + cfg.Dyn = dynconfig.Static(knobs{ + WaitTime: dynconfig.WaitTime{MinWait: dynconfig.Seconds(time.Millisecond), MaxWait: dynconfig.Seconds(2 * time.Millisecond)}, }) taskFn, shutdownFn := initSweperf(cfg) diff --git a/internal/benchmarking/boomer/userclass/config.go b/internal/benchmarking/boomer/userclass/config.go index 6e73314a6b..446522ddd8 100644 --- a/internal/benchmarking/boomer/userclass/config.go +++ b/internal/benchmarking/boomer/userclass/config.go @@ -35,9 +35,10 @@ type Config struct { // Atespace every actor this worker creates lives in. Required; caller // is responsible for having ensured it exists (see EnsureAtespace). Atespace string - // Dyn is the runtime-mutable config (wait-time bounds, trace - // probability). Required — every per-iteration read goes through it, - // so tests can mutate it without touching glutton internals. + // Dyn is the runtime-mutable config, built from the class's + // Entry.Config codec. Required: a class reads its typed knobs from it + // through dynconfig.Get on every iteration, so tests can swap the + // holder without touching class internals. Dyn *dynconfig.Holder // Tracer anchors sampled spans; falls back to the otel global if nil. Tracer trace.Tracer diff --git a/internal/benchmarking/boomer/userclass/registry.go b/internal/benchmarking/boomer/userclass/registry.go index f429c6a9b1..f1cd07b282 100644 --- a/internal/benchmarking/boomer/userclass/registry.go +++ b/internal/benchmarking/boomer/userclass/registry.go @@ -19,6 +19,8 @@ import ( "fmt" "slices" "sync" + + "github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig" ) // Entry declares one user class: the flag value that selects it, the Locust @@ -31,6 +33,12 @@ type Entry struct { // UserClass is the Python class name. It must equal boomer.Task.Name or // the master's spawn messages never match and no users start. UserClass string + // Config is the class's runtime config codec: a dynconfig.Typed over + // the struct it reads its knobs from, with its defaults and rules. The + // worker builds the holder from it, so a payload the class cannot run + // on is refused at fetch time rather than failing in the middle of a + // run. Nil when the class reads no runtime config. + Config dynconfig.Codec // Init builds the boomer task func and its shutdown hook. Init func(*Config) (task func(), shutdown func(context.Context)) }