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)) }