Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion benchmarking/locust/common/boomer_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
31 changes: 15 additions & 16 deletions cmd/benchmarking/boomer-worker/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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(), "|")))
Expand Down Expand Up @@ -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()))
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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",
Expand Down
44 changes: 35 additions & 9 deletions internal/benchmarking/boomer/agentsession/agentsession.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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
}
Expand All @@ -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
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
50 changes: 25 additions & 25 deletions internal/benchmarking/boomer/agentsession/agentsession_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
}
Expand Down Expand Up @@ -256,15 +256,15 @@ 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")
}
if rt.loaded != nil {
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)
Expand All @@ -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")
}
Expand All @@ -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)
Expand Down Expand Up @@ -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)})
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)})
Expand Down Expand Up @@ -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 {
Expand All @@ -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 {
Expand All @@ -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)
Expand All @@ -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}
Expand All @@ -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 {
Expand Down Expand Up @@ -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{
Expand All @@ -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",
Expand All @@ -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")
Expand Down Expand Up @@ -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")
Expand All @@ -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")
Expand All @@ -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) {
Expand All @@ -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")
Expand All @@ -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")
Expand All @@ -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)
Expand All @@ -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())
Expand Down
Loading
Loading