mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
benchmarking: agent-session workload — a scripted coding agent that suspends while the LLM thinks (#1934)
## What this is We keep saying Substrate's sweet spot is "agents that are idle most of the time" — but none of our benchmarks actually behave like one. Glutton hammers one resource at a time, and sweperf needs an external image with a replay trace. This adds a workload that acts like the thing we're building for: a coding agent working through a task. Each locust user is one session. The actor gets told to do the kind of things a coding agent does — clone a repo, install deps, build, hit a failing test, fix it, write tests, refactor, package — twenty steps, each one costing the sandbox the CPU, memory, disk, and network the real action would. Between steps the "LLM is thinking," so the driver suspends the actor, and the next step's first request wakes it back up through the router. That parked wake (`WakeFirstTouch` in the stats) is the number this whole benchmark exists to measure. The part I care most about: the entire workload is one table in `internal/benchmarking/boomer/agentsession/script.go`. Every step says in plain English what the agent is doing and what it costs. If you want to know what step 6 does to the sandbox, you read step 6. If you want a different workload, you edit the table — tests will catch you if you write a step that reads a file nothing wrote, or blow the actor's memory budget. To act the steps out, glutton grew two RPCs: `BurnCPU` (compute-bound work) and `Ingest` (bytes that actually cross the network before hitting disk, so a "git clone" is a real download, not a local write). ## How it went when we ran it Validated on a fresh 2-node GKE cluster, micro-VM first, then gVisor on the same hardware. Smoke runs were clean on both classes (gVisor: 1340 requests over 7 full laps, zero failures, wake p50 1.4s; micro-VM: wake p50 2.1s). At 20 concurrent sessions with realistic think times, gVisor held 2.6% failures with wake p50 1.3s. Pushing past the knee (~7 concurrently-active sessions on 8 vCPUs) was also useful: it reproduced the ateom-socket-vanishing failure from #1133 and left six actors permanently wedged in DELETING — a live repro of #1665. Two things the first live run taught us are already folded in: RAM refills are in-place so repeat laps don't transiently double the guest heap (512Mi micro-VM actors OOM'd without this — use 1Gi), and the driver replaces an actor after three failed steps in a row, because a CRASHED actor never comes back on its own. ## Future changes Right now the script is compiled in — changing what steps do means editing the table and rebuilding the image. That's deliberate for this PR (one reviewable, test-guarded source of truth), but the follow-up we've agreed on is to make the script runtime-configurable: named script variants selectable per run first, then accepting a full script as a file so operators can define workloads without touching Go. That lands as its own PR once this one is in. --------- Co-authored-by: Aditya Shantanu <aditya-shantanu@users.noreply.github.com>
This commit is contained in:
co-authored by
Aditya Shantanu
parent
f0e3ef64e7
commit
bc3cbc519d
@@ -143,6 +143,80 @@ The liveness check at session start already has the actor running, so the first
|
||||
resume is a no-op. Its successful `ResumeActor`, `ResumeActor_rtt` and `ResumeToFirstExec`
|
||||
samples are not recorded; failures still are.
|
||||
|
||||
### Agent-Session Benchmark
|
||||
|
||||
The agent-session benchmark (`--user-class agentsession`) emulates a fleet of
|
||||
coding agents on Substrate. Each locust user is one session: an actor driven
|
||||
through a scripted 20-step coding task ("clone a repo, build it, fix a test,
|
||||
refactor, package it"), where every step costs the sandbox the CPU, memory,
|
||||
disk, and network a real coding agent's action would. Between steps the agent
|
||||
is "waiting for the LLM to think": the driver **suspends the actor** for the
|
||||
step's think time, and the next step's first request **wakes it through the
|
||||
atenet router** (request parking). Mostly-idle sessions plus fast wake is
|
||||
exactly the oversubscription story this measures.
|
||||
|
||||
The entire workload is the declarative script in
|
||||
[`internal/benchmarking/boomer/agentsession/script.go`](../internal/benchmarking/boomer/agentsession/script.go)
|
||||
— one table entry per step, naming what the agent is doing and the resource
|
||||
ops that act it out. To change the workload, edit the table. Each session is
|
||||
an actor from the stock `glutton` template; steps are sequences of glutton
|
||||
RPCs:
|
||||
|
||||
| Step | The agent is… | Sandbox effect |
|
||||
|---|---|---|
|
||||
| 01_read_task | reading the task prompt | fill 32Mi RAM (agent context) |
|
||||
| 02_clone_repo | `git clone` | 16Mi arrives over the network → disk; 0.5s CPU |
|
||||
| 03_explore_tree | listing/grepping the tree | disk read (digest); 0.2s CPU |
|
||||
| 04_read_key_files | opening files into context | disk read shipped back out; 8Mi RAM churn |
|
||||
| 05_install_deps | `pip install` / `go mod download` | 32Mi over the network → disk; 1.5s CPU ×2 |
|
||||
| 06_first_build | first full build | fill 64Mi RAM; 3s CPU ×2; 24Mi disk write |
|
||||
| 07_run_unit_tests | running the test suite (one fails) | disk read; 2.5s CPU ×2 |
|
||||
| 08_reason_about_failure | tracing the bug (long LLM turn, 8s think) | full RAM page-walk after the wake |
|
||||
| 09_edit_source | applying the fix | 64Ki patch over the network; 4Mi disk write |
|
||||
| 10_incremental_build | rebuilding changed packages | 1.2s CPU ×2; 8Mi disk write |
|
||||
| 11_rerun_failed_test | re-running the failing test | 0.8s CPU |
|
||||
| 12_write_new_tests | authoring regression tests (6s think) | 128Ki over the network; 2Mi disk write |
|
||||
| 13_run_new_tests | running the new tests | disk read; 1s CPU |
|
||||
| 14_full_test_suite | full-suite regression run | 4s CPU ×2; disk read; 8Mi RAM churn |
|
||||
| 15_lint_format | lint + format pass | disk read; 0.9s CPU |
|
||||
| 16_refactor | multi-file refactor (8s think) | 512Ki over the network; 6Mi disk write; 0.6s CPU |
|
||||
| 17_rebuild | full rebuild | 32Mi RAM churn; 2s CPU ×2; 16Mi disk write |
|
||||
| 18_final_test_suite | final full-suite run | 3.5s CPU ×2 |
|
||||
| 19_package_artifact | building the release package | disk read; 1s CPU; 24Mi disk write |
|
||||
| 20_commit_and_summarize | committing + summarizing | 256Ki disk write; 24Mi shipped back out; RAM walk |
|
||||
|
||||
The script needs bigger actors than the 256Mi default: it holds ~96Mi of
|
||||
RAM arrays + ~110Mi of tmpfs files, and the observed guest peak with
|
||||
allocator transients is ~320Mi (512Mi OOMs). Deploy the workloads with
|
||||
`--actor-memory 1Gi`. `TestSessionBudgets` bounds the script-declared bytes
|
||||
only, so script edits that grow the working set fail the test and force this
|
||||
guidance to be revisited.
|
||||
|
||||
```sh
|
||||
./benchmarking/deploy_locust.sh --deploy --sandbox-class gvisor --actor-memory 1Gi
|
||||
./benchmarking/locust/deploy.sh --deploy --user-class agentsession
|
||||
```
|
||||
|
||||
#### Agent-Session Configuration Knobs
|
||||
|
||||
* `--agentsession-think-scale` — multiplier on every think gap; 0.5 makes the
|
||||
fleet twice as chatty, 4.0 models slow reasoning models (default 1.0). Each
|
||||
gap gets ±20% jitter so sessions don't move in lockstep.
|
||||
* `--resume-mode implicit|explicit` — implicit (default) lets the parked
|
||||
first request wake the actor; explicit issues ResumeActor before traffic.
|
||||
* `--lifecycle-mode suspend|pause` — durable suspend (default) or node-local
|
||||
pause between steps.
|
||||
|
||||
#### Agent-Session Reported Metrics
|
||||
|
||||
* `WakeFirstTouch`: latency of a dedicated ping sent before each step's ops —
|
||||
in implicit mode that ping is what triggers the parked wake, so this row
|
||||
**is** the user-visible wake latency, unpolluted by the step's own work.
|
||||
* `Step_<name>` (e.g. `Step_06_first_build`): wall time of that step's ops,
|
||||
think gap excluded.
|
||||
* `SuspendActor` / `ResumeActor` / `CreateActor` / `DeleteActor`: control-plane
|
||||
lifecycle latencies.
|
||||
|
||||
### Viewing Traces
|
||||
You must have enabled otel tracing for your cluster to view traces.
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
# 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.
|
||||
|
||||
"""Coding-agent-session benchmark runtime flags."""
|
||||
|
||||
from locust import events
|
||||
from locust.argument_parser import LocustArgumentParser
|
||||
|
||||
|
||||
@events.init_command_line_parser.add_listener
|
||||
def add_agentsession_arguments(parser: LocustArgumentParser) -> None:
|
||||
parser.add_argument(
|
||||
"--agentsession-think-scale",
|
||||
type=float,
|
||||
default=1.0,
|
||||
env_var="LOCUST_AGENTSESSION_THINK_SCALE",
|
||||
help=(
|
||||
"Multiplier on the script's per-step LLM think times: the gap the "
|
||||
"actor spends suspended between steps. 0.5 halves every gap, 2.0 "
|
||||
"doubles them (default: 1.0)"
|
||||
),
|
||||
include_in_web_ui=True,
|
||||
)
|
||||
@@ -67,6 +67,7 @@ _FLAGS = {
|
||||
"--sweperf-total-steps": int,
|
||||
"--sweperf-num-cycles": int,
|
||||
"--sweperf-poll-interval-ms": int,
|
||||
"--agentsession-think-scale": float,
|
||||
}
|
||||
|
||||
|
||||
@@ -142,6 +143,7 @@ def init_boomer_config() -> None:
|
||||
from locust.argument_parser import LocustArgumentParser
|
||||
from locust.env import Environment
|
||||
|
||||
from common.agentsession_config import add_agentsession_arguments # noqa: F401
|
||||
from common.durdir_config import add_durdir_arguments
|
||||
from common.lifecycle_mode import add_lifecycle_mode_arguments
|
||||
from common.memload_config import add_memload_arguments
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
# 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.
|
||||
|
||||
"""Stub AgentSessionUser declaration.
|
||||
|
||||
The real load implementation lives in the boomer-Go worker at
|
||||
internal/benchmarking/boomer/agentsession/ (the scripted coding-agent
|
||||
session in script.go); this Python class is declared only so the master
|
||||
recognizes the name and attributes boomer's stats rows to it. The master
|
||||
loads this stub file, selected by ${BENCHMARK_USER_CLASS} in
|
||||
locust/manifests/locust.yaml. The Python worker container sets
|
||||
LOCUST_NO_AGENTSESSION_USER=1 to skip loading this file, leaving boomer as
|
||||
the sole owner of AgentSessionUser load.
|
||||
"""
|
||||
|
||||
import os
|
||||
|
||||
if os.environ.get("LOCUST_NO_AGENTSESSION_USER") != "1":
|
||||
from locust import User, task
|
||||
|
||||
from common.boomer_config import init_boomer_config
|
||||
|
||||
# Master serves /boomer-config so the boomer-worker workers can fetch
|
||||
# runtime flag values (think-time scale, resume mode) the
|
||||
# operator set in the web UI form. No-op on workers without a web UI.
|
||||
init_boomer_config()
|
||||
|
||||
class AgentSessionUser(User):
|
||||
host = "api.ate-system.svc.cluster.local:443"
|
||||
|
||||
@task
|
||||
def noop(self) -> None:
|
||||
# Unreached under normal operation: the Python worker container
|
||||
# does not load this file (LOCUST_NO_AGENTSESSION_USER=1). Body
|
||||
# is required because locust validates that every User has at
|
||||
# least one @task method.
|
||||
pass
|
||||
@@ -36,6 +36,7 @@ import (
|
||||
"github.com/myzhan/boomer"
|
||||
|
||||
// Register user classes via init():
|
||||
_ "github.com/agent-substrate/substrate/internal/benchmarking/boomer/agentsession"
|
||||
_ "github.com/agent-substrate/substrate/internal/benchmarking/boomer/glutton"
|
||||
_ "github.com/agent-substrate/substrate/internal/benchmarking/boomer/sweperf"
|
||||
)
|
||||
|
||||
@@ -0,0 +1,573 @@
|
||||
// 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 agentsession implements the AgentSessionUser locust test: each VU
|
||||
// is one coding-agent session replaying the scripted steps in script.go
|
||||
// against a glutton actor. After every step the actor is suspended for the
|
||||
// step's "LLM thinking" gap; by default the next step's first request wakes
|
||||
// it through the atenet router (request parking), so the benchmark measures
|
||||
// the user-visible wake latency Substrate's oversubscription story rests on.
|
||||
//
|
||||
// Per-step timings land in locust stats as Step_<name>. The first request
|
||||
// after each suspend is additionally recorded as WakeFirstTouch — in
|
||||
// implicit resume mode that row IS the parking wake latency distribution.
|
||||
package agentsession
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
mathrand "math/rand/v2"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"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/internal/benchmarking/glutton"
|
||||
gluttonpb "github.com/agent-substrate/substrate/internal/proto/glutton"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
"github.com/google/uuid"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.opentelemetry.io/otel"
|
||||
"go.opentelemetry.io/otel/propagation"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/metadata"
|
||||
"google.golang.org/grpc/status"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
const (
|
||||
agentSessionUserClass = "AgentSessionUser"
|
||||
// templateNS is the atespace holding the benchmark ActorTemplates.
|
||||
templateNS = "benchmark-workloads"
|
||||
// templateName is the ActorTemplate instantiated per session: the stock
|
||||
// glutton actor. The script needs it deployed with --actor-memory 1Gi
|
||||
// (see benchmarking/README.md).
|
||||
templateName = "glutton"
|
||||
|
||||
// controlRPCTimeout bounds every control-plane RPC. The router HTTP
|
||||
// client already times out at 30s; without a deadline here a
|
||||
// server-side hang in SuspendActor or DeleteActor would park the VU
|
||||
// goroutine for the rest of the run.
|
||||
controlRPCTimeout = 60 * time.Second
|
||||
|
||||
// maxConsecutiveStepFailures is how many ReplaceIfPersistent failures in
|
||||
// a row a session tolerates before its actor is deleted and recreated.
|
||||
maxConsecutiveStepFailures = 3
|
||||
)
|
||||
|
||||
// burnRate is the per-goroutine sha256 iteration rate BurnCPU reports.
|
||||
// BurnCPU runs for a fixed wall-clock duration, so Step_* latency stays flat
|
||||
// when the sandbox is CPU-starved; this rate is what drops.
|
||||
var burnRate = prometheus.NewHistogram(prometheus.HistogramOpts{
|
||||
Name: "agentsession_burn_cpu_iterations_per_second",
|
||||
Help: "BurnCPU sha256 iterations per second per goroutine; drops under CPU contention while Step_* latency stays flat.",
|
||||
Buckets: prometheus.ExponentialBuckets(256, 2, 16),
|
||||
})
|
||||
|
||||
func init() {
|
||||
prometheus.MustRegister(burnRate)
|
||||
userclass.Add(userclass.Entry{
|
||||
Name: "agentsession",
|
||||
LocustFile: "agentsession.py",
|
||||
UserClass: agentSessionUserClass,
|
||||
Init: initAgentSession,
|
||||
})
|
||||
}
|
||||
|
||||
func initAgentSession(cfg *userclass.Config) (taskFn func(), shutdown func(context.Context)) {
|
||||
if cfg.Tracer == nil {
|
||||
cfg.Tracer = otel.Tracer("substrate-boomer/agentsession")
|
||||
}
|
||||
rt := &runtime{cfg: cfg, steps: Session()}
|
||||
rt.ingestBuf = makeIngestBuf(rt.steps)
|
||||
return rt.iterate, rt.shutdown
|
||||
}
|
||||
|
||||
// makeIngestBuf pre-generates one random payload sized to the script's
|
||||
// largest ingest. Generating bytes per op would count client-side
|
||||
// crypto/rand time inside the step timings and allocate tens of MiB per VU
|
||||
// per op; glutton only writes the payload, so every session can safely
|
||||
// slice the same buffer.
|
||||
func makeIngestBuf(steps []Step) []byte {
|
||||
var maxBytes int64
|
||||
for _, s := range steps {
|
||||
for _, o := range s.Ops {
|
||||
if o.kind == opIngest && o.bytes > maxBytes {
|
||||
maxBytes = o.bytes
|
||||
}
|
||||
}
|
||||
}
|
||||
buf := make([]byte, maxBytes)
|
||||
if _, err := rand.Read(buf); err != nil {
|
||||
// Payload content is irrelevant to the benchmark; zeros still move
|
||||
// the same bytes over the wire.
|
||||
slog.Warn("failed to randomize ingest payload; using zeros", slog.String("err", err.Error()))
|
||||
}
|
||||
return buf
|
||||
}
|
||||
|
||||
// runtime is the per-worker state shared by every boomer goroutine. Each
|
||||
// goroutine keeps its own session in users, keyed by goroutine ID, because
|
||||
// boomer offers no per-VU context.
|
||||
type runtime struct {
|
||||
cfg *userclass.Config
|
||||
steps []Step
|
||||
ingestBuf []byte // shared read-only ingest payload, sized to the largest ingest op
|
||||
users sync.Map // goroutineID -> *sessionUser
|
||||
}
|
||||
|
||||
// think returns the suspended "LLM thinking" gap before step s, which is the
|
||||
// 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
|
||||
if scale <= 0 {
|
||||
scale = 1.0
|
||||
}
|
||||
jitter := 0.8 + 0.4*mathrand.Float64()
|
||||
return time.Duration(float64(s.Think) * scale * jitter)
|
||||
}
|
||||
|
||||
// iterate is the boomer task function: one scripted step per call. The VU's
|
||||
// session is created lazily on first use and wraps back to step 1 after the
|
||||
// script completes, so one actor keeps producing suspend/resume cycles for
|
||||
// the whole run.
|
||||
func (r *runtime) iterate() {
|
||||
gid := boomerutil.GoroutineID()
|
||||
val, loaded := r.users.Load(gid)
|
||||
if !loaded {
|
||||
u, err := r.startUser(context.Background())
|
||||
if err != nil {
|
||||
slog.Warn("agentsession start failed; goroutine will retry next iter",
|
||||
slog.String("err", err.Error()))
|
||||
time.Sleep(2 * time.Second)
|
||||
return
|
||||
}
|
||||
val, _ = r.users.LoadOrStore(gid, u)
|
||||
}
|
||||
u := val.(*sessionUser)
|
||||
|
||||
step := r.steps[u.stepIndex]
|
||||
|
||||
// The LLM is "thinking": the actor stays suspended for the gap.
|
||||
time.Sleep(r.think(step))
|
||||
|
||||
if u.runStep(context.Background(), step) {
|
||||
u.stepIndex = (u.stepIndex + 1) % len(r.steps)
|
||||
if u.stepIndex == 0 {
|
||||
slog.Info("agent session completed the full script; starting over",
|
||||
slog.String("actor", u.actorName))
|
||||
}
|
||||
}
|
||||
|
||||
if u.broken {
|
||||
slog.Warn("agent session wedged; deleting its actor and starting a fresh session",
|
||||
slog.String("actor", u.actorName),
|
||||
slog.Int("consecutive_failures", u.consecutiveFailures))
|
||||
u.suspendAndDelete(context.Background())
|
||||
r.users.Delete(gid)
|
||||
}
|
||||
}
|
||||
|
||||
// startUser creates the session's actor, waits for its sandbox to serve,
|
||||
// then suspends it so the first scripted step begins — like every later
|
||||
// step — with a wake from suspension.
|
||||
func (r *runtime) startUser(ctx context.Context) (*sessionUser, error) {
|
||||
u := &sessionUser{
|
||||
cfg: r.cfg,
|
||||
actorName: "agent-" + uuid.NewString(),
|
||||
ingestBuf: r.ingestBuf,
|
||||
}
|
||||
slog.Info("Creating agent session", slog.String("actor", u.actorName))
|
||||
bmetrics.UpdateUsers(agentSessionUserClass, 1)
|
||||
|
||||
if err := u.ensureAtespace(ctx); err != nil {
|
||||
bmetrics.UpdateUsers(agentSessionUserClass, -1)
|
||||
return nil, fmt.Errorf("ensureAtespace: %w", err)
|
||||
}
|
||||
if err := u.create(ctx); err != nil {
|
||||
bmetrics.UpdateUsers(agentSessionUserClass, -1)
|
||||
return nil, fmt.Errorf("createActor: %w", err)
|
||||
}
|
||||
if err := u.waitServing(ctx); err != nil {
|
||||
// suspendAndDelete decrements the user gauge; no extra decrement here.
|
||||
u.suspendAndDelete(ctx)
|
||||
return nil, fmt.Errorf("waitServing: %w", err)
|
||||
}
|
||||
// Park the fresh actor; step 01 wakes it like any other step. A failed
|
||||
// park is re-driven at the top of the first step.
|
||||
u.hibernate(ctx)
|
||||
return u, nil
|
||||
}
|
||||
|
||||
func (r *runtime) shutdown(ctx context.Context) {
|
||||
r.users.Range(func(_, val any) bool {
|
||||
val.(*sessionUser).suspendAndDelete(ctx)
|
||||
return true
|
||||
})
|
||||
}
|
||||
|
||||
// sessionUser is one coding-agent session: a single actor plus its progress
|
||||
// through the script. Owned by one boomer goroutine; no locking needed.
|
||||
type sessionUser struct {
|
||||
cfg *userclass.Config
|
||||
actorName string
|
||||
stepIndex int
|
||||
cleanedUp bool
|
||||
// ingestBuf is the runtime's shared read-only ingest payload.
|
||||
ingestBuf []byte
|
||||
// hibernatePending is set by a failed Pause/Suspend: the actor is
|
||||
// stranded RUNNING or SUSPENDING. Waking it from there would misreport
|
||||
// WakeFirstTouch, so runStep finishes the hibernate first.
|
||||
hibernatePending bool
|
||||
// consecutiveFailures counts ReplaceIfPersistent failures since the
|
||||
// last successful step; see noteFailure.
|
||||
consecutiveFailures int
|
||||
// broken is set once the actor should be replaced: a CRASHED or
|
||||
// otherwise stuck actor never recovers on its own, and without
|
||||
// replacement it would wedge its VU for the rest of the run.
|
||||
broken bool
|
||||
}
|
||||
|
||||
// noteFailure classifies a failed step or lifecycle call and marks the actor
|
||||
// broken when a replacement is warranted. Cluster-wide errors (no capacity,
|
||||
// ate-api-server restarting) are not the actor's fault and do not count.
|
||||
func (u *sessionUser) noteFailure(err error) {
|
||||
switch boomerutil.ClassifyLifecycleFailure(err) {
|
||||
case boomerutil.ReplaceNow:
|
||||
u.broken = true
|
||||
case boomerutil.RetryLater:
|
||||
default:
|
||||
u.consecutiveFailures++
|
||||
if u.consecutiveFailures >= maxConsecutiveStepFailures {
|
||||
u.broken = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// noteSuccess clears the failure count: the actor just completed a whole
|
||||
// step, so the earlier failures were transient after all.
|
||||
func (u *sessionUser) noteSuccess() {
|
||||
u.consecutiveFailures = 0
|
||||
}
|
||||
|
||||
func (u *sessionUser) ref() *ateapipb.ObjectRef {
|
||||
return &ateapipb.ObjectRef{Atespace: u.cfg.Atespace, Name: u.actorName}
|
||||
}
|
||||
|
||||
// runStep wakes the actor, executes the step's ops in order, and suspends
|
||||
// it again, reporting whether everything succeeded. The wake is measured by
|
||||
// a dedicated ping sent ahead of the step's ops and recorded as
|
||||
// WakeFirstTouch; in implicit resume mode that ping is also what triggers
|
||||
// the parked wake, so the row measures the wake alone rather than the wake
|
||||
// plus whatever heavy op happens to lead the step.
|
||||
func (u *sessionUser) runStep(ctx context.Context, step Step) bool {
|
||||
// A failed hibernate left the actor awake. In implicit mode the wake
|
||||
// ping would return in a few ms and be booked as a WakeFirstTouch
|
||||
// success; in explicit mode ResumeActor would fail FailedPrecondition.
|
||||
// SuspendActor and PauseActor are re-entrant, so finish the hibernate
|
||||
// and pick the step up next iteration.
|
||||
if u.hibernatePending {
|
||||
u.hibernate(ctx)
|
||||
return false
|
||||
}
|
||||
|
||||
if u.cfg.Dyn.Load().ResumeMode == dynconfig.ResumeModeExplicit {
|
||||
if err := u.resume(ctx); err != nil {
|
||||
u.noteFailure(err)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
wakeStart := time.Now()
|
||||
if err := u.execOp(ctx, ping()); err != nil {
|
||||
u.noteFailure(err)
|
||||
bmetrics.RecordFailure("http", "WakeFirstTouch", agentSessionUserClass, time.Since(wakeStart), err.Error())
|
||||
slog.Warn("agent session wake failed",
|
||||
slog.String("actor", u.actorName),
|
||||
slog.String("step", step.Name),
|
||||
slog.String("err", err.Error()))
|
||||
u.hibernate(ctx)
|
||||
return false
|
||||
}
|
||||
bmetrics.RecordSuccess("http", "WakeFirstTouch", agentSessionUserClass, time.Since(wakeStart), 0)
|
||||
|
||||
metricName := "Step_" + step.Name
|
||||
stepStart := time.Now()
|
||||
for i, o := range step.Ops {
|
||||
if err := u.execOp(ctx, o); err != nil {
|
||||
u.noteFailure(err)
|
||||
bmetrics.RecordFailure("http", metricName, agentSessionUserClass, time.Since(stepStart), err.Error())
|
||||
slog.Warn("agent session step failed",
|
||||
slog.String("actor", u.actorName),
|
||||
slog.String("step", step.Name),
|
||||
slog.Int("op", i),
|
||||
slog.String("err", err.Error()))
|
||||
u.hibernate(ctx)
|
||||
return false
|
||||
}
|
||||
}
|
||||
u.noteSuccess()
|
||||
bmetrics.RecordSuccess("http", metricName, agentSessionUserClass, time.Since(stepStart), 0)
|
||||
|
||||
u.hibernate(ctx)
|
||||
return true
|
||||
}
|
||||
|
||||
// execOp performs one scripted op as a glutton RPC through the router.
|
||||
func (u *sessionUser) execOp(ctx context.Context, o op) error {
|
||||
switch o.kind {
|
||||
case opIngest:
|
||||
payload := u.ingestBuf
|
||||
if int64(len(payload)) < o.bytes {
|
||||
// Fallback for callers (tests) that built the user by hand.
|
||||
payload = make([]byte, o.bytes)
|
||||
}
|
||||
return u.postProto(ctx, glutton.IngestRoute,
|
||||
&gluttonpb.IngestRequest{Key: o.key, Payload: payload[:o.bytes]},
|
||||
&gluttonpb.IngestResponse{})
|
||||
case opBurnCPU:
|
||||
resp := &gluttonpb.BurnCPUResponse{}
|
||||
if err := u.postProto(ctx, glutton.BurnCPURoute,
|
||||
&gluttonpb.BurnCPURequest{DurationMs: o.millis, Parallelism: o.parallel}, resp); err != nil {
|
||||
return err
|
||||
}
|
||||
observeBurnRate(o, resp.GetIterations())
|
||||
return nil
|
||||
case opWriteDisk:
|
||||
return u.postProto(ctx, glutton.WriteDiskRoute,
|
||||
&gluttonpb.WriteDiskRequest{Key: o.key, Size: int32(o.bytes), WriteMode: gluttonpb.WriteMode_WRITE_MODE_TRUNCATE},
|
||||
&gluttonpb.WriteDiskResponse{})
|
||||
case opReadDiskDigest:
|
||||
return u.postProto(ctx, glutton.ReadDiskRoute,
|
||||
&gluttonpb.ReadDiskRequest{Key: o.key, ReadMode: gluttonpb.ReadMode_READ_MODE_DIGEST_ONLY},
|
||||
&gluttonpb.ReadDiskResponse{})
|
||||
case opReadDiskData:
|
||||
return u.postProto(ctx, glutton.ReadDiskRoute,
|
||||
&gluttonpb.ReadDiskRequest{Key: o.key, ReadMode: gluttonpb.ReadMode_READ_MODE_DATA},
|
||||
&gluttonpb.ReadDiskResponse{})
|
||||
case opFillRAM:
|
||||
// OVERWRITE grows the array on first touch and re-randomizes it in
|
||||
// place on later laps; TRUNCATE would reallocate and transiently
|
||||
// double the guest heap, which OOMs a tightly-sized sandbox.
|
||||
return u.postProto(ctx, glutton.WriteRAMRoute,
|
||||
&gluttonpb.WriteRAMRequest{Key: o.key, Size: fmt.Sprintf("%d", o.bytes), WriteMode: gluttonpb.WriteMode_WRITE_MODE_OVERWRITE},
|
||||
&gluttonpb.WriteRAMResponse{})
|
||||
case opChurnRAM:
|
||||
return u.postProto(ctx, glutton.WriteRAMRoute,
|
||||
&gluttonpb.WriteRAMRequest{Key: o.key, Size: fmt.Sprintf("%d", o.bytes), WriteMode: gluttonpb.WriteMode_WRITE_MODE_OVERWRITE_ROTATE},
|
||||
&gluttonpb.WriteRAMResponse{})
|
||||
case opWalkRAM:
|
||||
return u.postProto(ctx, glutton.ReadRAMRoute,
|
||||
&gluttonpb.ReadRAMRequest{Key: o.key},
|
||||
&gluttonpb.ReadRAMResponse{})
|
||||
case opPing:
|
||||
return u.postProto(ctx, glutton.PingRoute,
|
||||
&gluttonpb.PingRequest{Message: "think"},
|
||||
&gluttonpb.PingResponse{})
|
||||
default:
|
||||
return fmt.Errorf("unknown op kind %d", o.kind)
|
||||
}
|
||||
}
|
||||
|
||||
// observeBurnRate records a burn's iterations per goroutine-second.
|
||||
func observeBurnRate(o op, iterations int64) {
|
||||
if rate, ok := burnRatePerGoroutine(o, iterations); ok {
|
||||
burnRate.Observe(rate)
|
||||
}
|
||||
}
|
||||
|
||||
// burnRatePerGoroutine normalizes a BurnCPU result by the burn's goroutine
|
||||
// count and wall-clock seconds. A zero-duration burn has no rate.
|
||||
func burnRatePerGoroutine(o op, iterations int64) (float64, bool) {
|
||||
if o.millis <= 0 {
|
||||
return 0, false
|
||||
}
|
||||
goroutines := max(int64(o.parallel), 1)
|
||||
return float64(iterations) / float64(goroutines) / (float64(o.millis) / 1000), true
|
||||
}
|
||||
|
||||
// postProto POSTs req to the actor's route through the router and
|
||||
// unmarshals the reply into resp. HTTP status >= 400 is an error carrying
|
||||
// the response body.
|
||||
func (u *sessionUser) postProto(ctx context.Context, route string, req, resp proto.Message) error {
|
||||
body, err := proto.Marshal(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, u.cfg.RouterURL+route, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
httpReq.Header.Set(atenet.TargetActorHeader, u.cfg.Atespace+"/"+u.actorName)
|
||||
httpReq.Header.Set("Content-Type", "application/x-protobuf")
|
||||
otel.GetTextMapPropagator().Inject(ctx, propagation.HeaderCarrier(httpReq.Header))
|
||||
|
||||
httpResp, err := u.cfg.HTTPClient.Do(httpReq)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer httpResp.Body.Close()
|
||||
respBody, err := io.ReadAll(httpResp.Body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if httpResp.StatusCode >= 400 {
|
||||
return fmt.Errorf("%s: HTTP %d: %s", route, httpResp.StatusCode, strings.TrimSpace(string(respBody)))
|
||||
}
|
||||
return proto.Unmarshal(respBody, resp)
|
||||
}
|
||||
|
||||
// ensureAtespace creates the configured atespace, treating AlreadyExists as
|
||||
// success so concurrent sessions racing the first creation all proceed.
|
||||
func (u *sessionUser) ensureAtespace(ctx context.Context) error {
|
||||
return u.tracedCall(ctx, "CreateAtespace", func(callCtx context.Context, tr *metadata.MD) error {
|
||||
_, err := u.cfg.APIStub.CreateAtespace(callCtx, &ateapipb.CreateAtespaceRequest{
|
||||
Atespace: &ateapipb.Atespace{
|
||||
Metadata: &ateapipb.ResourceMetadata{Name: u.cfg.Atespace},
|
||||
},
|
||||
}, grpc.Trailer(tr))
|
||||
if s, ok := status.FromError(err); ok && s.Code() == codes.AlreadyExists {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
func (u *sessionUser) create(ctx context.Context) error {
|
||||
return u.tracedCall(ctx, "CreateActor", func(callCtx context.Context, tr *metadata.MD) error {
|
||||
_, err := u.cfg.APIStub.CreateActor(callCtx, &ateapipb.CreateActorRequest{
|
||||
Actor: &ateapipb.Actor{
|
||||
Metadata: &ateapipb.ResourceMetadata{Atespace: u.cfg.Atespace, Name: u.actorName},
|
||||
ActorTemplate: &ateapipb.ObjectRef{Atespace: templateNS, Name: templateName},
|
||||
},
|
||||
}, grpc.Trailer(tr))
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
// waitServing pings the actor through the router until glutton answers. The
|
||||
// first request is also what wakes a newly created actor, so this doubles
|
||||
// as the initial activation.
|
||||
func (u *sessionUser) waitServing(ctx context.Context) error {
|
||||
const maxRetries = 30
|
||||
const retryInterval = 2 * time.Second
|
||||
var lastErr error
|
||||
for attempt := 0; attempt < maxRetries; attempt++ {
|
||||
lastErr = u.postProto(ctx, glutton.PingRoute,
|
||||
&gluttonpb.PingRequest{Message: "hello"}, &gluttonpb.PingResponse{})
|
||||
if lastErr == nil {
|
||||
return nil
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(retryInterval):
|
||||
}
|
||||
}
|
||||
return fmt.Errorf("actor %s never served: %w", u.actorName, lastErr)
|
||||
}
|
||||
|
||||
// resume issues ResumeActor, retrying concurrent-update conflicts inside
|
||||
// the traced call so the reported latency spans every attempt.
|
||||
func (u *sessionUser) resume(ctx context.Context) error {
|
||||
err := u.tracedCall(ctx, "ResumeActor", func(callCtx context.Context, tr *metadata.MD) error {
|
||||
return boomerutil.RetryOnConflict(callCtx, func() error {
|
||||
_, err := u.cfg.APIStub.ResumeActor(callCtx, &ateapipb.ResumeActorRequest{Actor: u.ref()}, grpc.Trailer(tr))
|
||||
return err
|
||||
})
|
||||
})
|
||||
if err != nil {
|
||||
if boomerutil.IsCrashed(err) {
|
||||
bmetrics.RecordFailure("actor", "CrashCount", agentSessionUserClass, 0, "actor entered ACTOR_STATE_CRASHED")
|
||||
}
|
||||
slog.Error("ResumeActor failed", slog.String("actor", u.actorName), slog.String("err", err.Error()))
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// hibernate suspends (or pauses, per LifecycleMode) the actor at the end of
|
||||
// a step. A failure leaves hibernatePending set so the next step re-drives
|
||||
// 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 {
|
||||
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
|
||||
})
|
||||
} else {
|
||||
err = u.tracedCall(ctx, "SuspendActor", func(callCtx context.Context, tr *metadata.MD) error {
|
||||
_, err := u.cfg.APIStub.SuspendActor(callCtx, &ateapipb.SuspendActorRequest{Actor: u.ref()}, grpc.Trailer(tr))
|
||||
return err
|
||||
})
|
||||
}
|
||||
u.hibernatePending = err != nil
|
||||
if err != nil {
|
||||
u.noteFailure(err)
|
||||
}
|
||||
}
|
||||
|
||||
// suspendAndDelete releases the actor and its worker. Suspend first: a
|
||||
// worker only frees an actor it can account for, and an actor still awake
|
||||
// at delete risks being left CRASHED. Safe to call more than once.
|
||||
func (u *sessionUser) suspendAndDelete(ctx context.Context) {
|
||||
if u.cleanedUp {
|
||||
return
|
||||
}
|
||||
suspendCtx, cancel := context.WithTimeout(ctx, controlRPCTimeout)
|
||||
_, _ = u.cfg.APIStub.SuspendActor(suspendCtx, &ateapipb.SuspendActorRequest{Actor: u.ref()})
|
||||
cancel()
|
||||
_ = u.tracedCall(ctx, "DeleteActor", func(callCtx context.Context, tr *metadata.MD) error {
|
||||
_, err := u.cfg.APIStub.DeleteActor(callCtx, &ateapipb.DeleteActorRequest{Actor: u.ref(), AnyState: true}, grpc.Trailer(tr))
|
||||
return err
|
||||
})
|
||||
bmetrics.UpdateUsers(agentSessionUserClass, -1)
|
||||
u.cleanedUp = true
|
||||
}
|
||||
|
||||
// tracedCall runs one control-plane RPC under a span and a deadline, and
|
||||
// records it as a locust stats row of the same name, preferring the
|
||||
// server-measured elapsed time from the response trailer.
|
||||
func (u *sessionUser) tracedCall(ctx context.Context, name string, do func(context.Context, *metadata.MD) error) error {
|
||||
ctx, cancel := context.WithTimeout(ctx, controlRPCTimeout)
|
||||
defer cancel()
|
||||
ctx, span := u.cfg.Tracer.Start(ctx, name)
|
||||
defer span.End()
|
||||
|
||||
start := time.Now()
|
||||
var tr metadata.MD
|
||||
err := do(ctx, &tr)
|
||||
latency := time.Since(start)
|
||||
|
||||
serverLatency, source := boomerutil.ElapsedFromMD(tr, ateinterceptors.ServerElapsedTrailer, latency)
|
||||
boomerutil.LogSampledTrace(span, name, serverLatency, source, err)
|
||||
if err != nil {
|
||||
bmetrics.RecordFailure("grpc", name, agentSessionUserClass, latency, err.Error())
|
||||
return err
|
||||
}
|
||||
bmetrics.RecordSuccess("grpc", name, agentSessionUserClass, latency, 0)
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,426 @@
|
||||
// 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 agentsession
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/boomer/dynconfig"
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/boomer/userclass"
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/glutton/fake"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
"go.opentelemetry.io/otel"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// TestSessionScriptIsWellFormed pins the script's invariants: unique step
|
||||
// names, positive think times, and no op that consumes a sandbox object
|
||||
// before an earlier step created it. A broken ordering would fail at run
|
||||
// time with NotFound from glutton; this catches it at test time.
|
||||
func TestSessionScriptIsWellFormed(t *testing.T) {
|
||||
steps := Session()
|
||||
if len(steps) != 20 {
|
||||
t.Fatalf("Session has %d steps, want 20", len(steps))
|
||||
}
|
||||
|
||||
seen := map[string]bool{}
|
||||
ramFilled := map[string]bool{}
|
||||
diskWritten := map[string]bool{}
|
||||
|
||||
for i, s := range steps {
|
||||
if s.Name == "" || s.Agent == "" {
|
||||
t.Errorf("step %d: Name and Agent must be set", i)
|
||||
}
|
||||
if seen[s.Name] {
|
||||
t.Errorf("step %q: duplicate name", s.Name)
|
||||
}
|
||||
seen[s.Name] = true
|
||||
if s.Think <= 0 {
|
||||
t.Errorf("step %q: think time must be positive", s.Name)
|
||||
}
|
||||
if len(s.Ops) == 0 {
|
||||
t.Errorf("step %q: has no ops", s.Name)
|
||||
}
|
||||
for _, o := range s.Ops {
|
||||
switch o.kind {
|
||||
case opFillRAM:
|
||||
ramFilled[o.key] = true
|
||||
case opChurnRAM, opWalkRAM:
|
||||
if !ramFilled[o.key] {
|
||||
t.Errorf("step %q: %s RAM op before any fill of %q", s.Name, opName(o.kind), o.key)
|
||||
}
|
||||
case opIngest, opWriteDisk:
|
||||
diskWritten[o.key] = true
|
||||
case opReadDiskDigest, opReadDiskData:
|
||||
if !diskWritten[o.key] {
|
||||
t.Errorf("step %q: read of %q before any write", s.Name, o.key)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func opName(k opKind) string {
|
||||
switch k {
|
||||
case opChurnRAM:
|
||||
return "churn"
|
||||
case opWalkRAM:
|
||||
return "walk"
|
||||
default:
|
||||
return "op"
|
||||
}
|
||||
}
|
||||
|
||||
// TestSessionBudgets bounds the bytes the script itself declares: resident
|
||||
// RAM (largest fill per key) and disk (largest object per key). It does NOT
|
||||
// model the guest's real peak — kernel, kata-agent, and allocator transients
|
||||
// sit on top — so the budgets are deliberately far below the 1Gi the
|
||||
// template calls for. A script edit that outgrows them must come with a
|
||||
// fresh look at the actor memory guidance.
|
||||
func TestSessionBudgets(t *testing.T) {
|
||||
const (
|
||||
ramBudget = 128 << 20 // bytes
|
||||
diskBudget = 256 << 20
|
||||
)
|
||||
ramMax := map[string]int64{}
|
||||
diskMax := map[string]int64{}
|
||||
for _, s := range Session() {
|
||||
for _, o := range s.Ops {
|
||||
switch o.kind {
|
||||
case opFillRAM:
|
||||
if o.bytes > ramMax[o.key] {
|
||||
ramMax[o.key] = o.bytes
|
||||
}
|
||||
case opIngest, opWriteDisk:
|
||||
if o.bytes > diskMax[o.key] {
|
||||
diskMax[o.key] = o.bytes
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
var ramTotal, diskTotal int64
|
||||
for _, v := range ramMax {
|
||||
ramTotal += v
|
||||
}
|
||||
for _, v := range diskMax {
|
||||
diskTotal += v
|
||||
}
|
||||
if ramTotal > ramBudget {
|
||||
t.Errorf("script fills %d bytes of RAM, budget %d", ramTotal, ramBudget)
|
||||
}
|
||||
if diskTotal > diskBudget {
|
||||
t.Errorf("script writes %d bytes of disk, budget %d", diskTotal, diskBudget)
|
||||
}
|
||||
}
|
||||
|
||||
// TestExecOpAgainstFake replays every scripted op against the fake glutton
|
||||
// server, proving each op marshals a request the actor-side routes accept.
|
||||
func TestExecOpAgainstFake(t *testing.T) {
|
||||
fakeSrv := &fake.Server{Data: []byte("filecontents")}
|
||||
ts := fakeSrv.Start(t)
|
||||
|
||||
u := &sessionUser{
|
||||
cfg: &userclass.Config{
|
||||
HTTPClient: http.DefaultClient,
|
||||
RouterURL: ts.URL,
|
||||
Atespace: "benchmark",
|
||||
Dyn: dynconfig.NewHolder(dynconfig.Config{}),
|
||||
},
|
||||
actorName: "agent-test",
|
||||
}
|
||||
|
||||
var opCount int
|
||||
for _, s := range Session() {
|
||||
for i, o := range s.Ops {
|
||||
// Cap ingest payloads in the unit test: transport shape is what
|
||||
// matters here, not moving tens of MiB through httptest.
|
||||
if o.kind == opIngest && o.bytes > 1<<10 {
|
||||
o.bytes = 1 << 10
|
||||
}
|
||||
if o.kind == opBurnCPU {
|
||||
o.millis = 1
|
||||
}
|
||||
if err := u.execOp(context.Background(), o); err != nil {
|
||||
t.Fatalf("step %q op %d: %v", s.Name, i, err)
|
||||
}
|
||||
opCount++
|
||||
}
|
||||
}
|
||||
if got := len(fakeSrv.RecordedPaths()); got != opCount {
|
||||
t.Errorf("fake served %d requests, want %d", got, opCount)
|
||||
}
|
||||
for _, n := range fakeSrv.RecordedIngestSizes() {
|
||||
if n > 1<<10 {
|
||||
t.Errorf("ingest payload of %d bytes reached the server, cap is %d", n, 1<<10)
|
||||
}
|
||||
}
|
||||
for _, ms := range fakeSrv.RecordedBurnMillis() {
|
||||
if ms != 1 {
|
||||
t.Errorf("burn of %dms reached the server, override is 1ms", ms)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestThinkScaling checks the think-time multiplier and its jitter bounds.
|
||||
func TestThinkScaling(t *testing.T) {
|
||||
r := &runtime{
|
||||
cfg: &userclass.Config{
|
||||
Dyn: dynconfig.NewHolder(dynconfig.Config{AgentSessionThinkScale: 0.5}),
|
||||
},
|
||||
}
|
||||
s := Step{Think: 10 * time.Second}
|
||||
for i := 0; i < 100; i++ {
|
||||
got := r.think(s)
|
||||
if got < 4*time.Second || got > 6*time.Second {
|
||||
t.Fatalf("think = %v, want within [4s, 6s] (scale 0.5, jitter ±20%%)", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Zero scale reads as 1.0.
|
||||
r.cfg.Dyn.Store(dynconfig.Config{})
|
||||
for i := 0; i < 100; i++ {
|
||||
got := r.think(s)
|
||||
if got < 8*time.Second || got > 12*time.Second {
|
||||
t.Fatalf("think = %v with unset scale, want within [8s, 12s]", got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// fakeControlClient records control-plane calls and fails ResumeActor /
|
||||
// SuspendActor with queued errors, in order, until each queue drains.
|
||||
type fakeControlClient struct {
|
||||
ateapipb.ControlClient
|
||||
mu sync.Mutex
|
||||
calls []string
|
||||
resumeErrs []error
|
||||
suspendErrs []error
|
||||
// sawDeadline is set when a call arrived with a context deadline.
|
||||
sawDeadline bool
|
||||
}
|
||||
|
||||
func nextErr(errs *[]error) error {
|
||||
if len(*errs) == 0 {
|
||||
return nil
|
||||
}
|
||||
err := (*errs)[0]
|
||||
*errs = (*errs)[1:]
|
||||
return err
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) record(ctx context.Context, name string) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.calls = append(f.calls, name)
|
||||
if _, ok := ctx.Deadline(); ok {
|
||||
f.sawDeadline = true
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) ResumeActor(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
f.record(ctx, "ResumeActor")
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
if err := nextErr(&f.resumeErrs); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &ateapipb.ResumeActorResponse{}, nil
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) SuspendActor(ctx context.Context, in *ateapipb.SuspendActorRequest, opts ...grpc.CallOption) (*ateapipb.SuspendActorResponse, error) {
|
||||
f.record(ctx, "SuspendActor")
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
if err := nextErr(&f.suspendErrs); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &ateapipb.SuspendActorResponse{}, nil
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) PauseActor(ctx context.Context, in *ateapipb.PauseActorRequest, opts ...grpc.CallOption) (*ateapipb.PauseActorResponse, error) {
|
||||
f.record(ctx, "PauseActor")
|
||||
return &ateapipb.PauseActorResponse{}, nil
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) DeleteActor(ctx context.Context, in *ateapipb.DeleteActorRequest, opts ...grpc.CallOption) (*ateapipb.Actor, error) {
|
||||
f.record(ctx, "DeleteActor")
|
||||
return &ateapipb.Actor{}, nil
|
||||
}
|
||||
|
||||
func (f *fakeControlClient) recordedCalls() []string {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]string(nil), f.calls...)
|
||||
}
|
||||
|
||||
func countCalls(calls []string, name string) int {
|
||||
n := 0
|
||||
for _, c := range calls {
|
||||
if c == name {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func newTestUser(t *testing.T, srv *fake.Server, ctl *fakeControlClient, dyn dynconfig.Config) *sessionUser {
|
||||
t.Helper()
|
||||
ts := srv.Start(t)
|
||||
return &sessionUser{
|
||||
cfg: &userclass.Config{
|
||||
APIStub: ctl,
|
||||
HTTPClient: ts.Client(),
|
||||
RouterURL: ts.URL,
|
||||
Atespace: "benchmark",
|
||||
Dyn: dynconfig.NewHolder(dyn),
|
||||
Tracer: otel.Tracer("test"),
|
||||
},
|
||||
actorName: "agent-test",
|
||||
}
|
||||
}
|
||||
|
||||
var pingStep = Step{Name: "01_test", Agent: "testing", Think: time.Second, Ops: []op{ping()}}
|
||||
|
||||
// A retryable suspend failure strands the actor RUNNING. The next step must
|
||||
// finish the suspend rather than wake an already-awake actor, which would
|
||||
// book a few-ms WakeFirstTouch success.
|
||||
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{})
|
||||
|
||||
if !u.runStep(context.Background(), pingStep) {
|
||||
t.Fatal("runStep = false; the step's ops succeeded and must count")
|
||||
}
|
||||
if !u.hibernatePending || u.broken {
|
||||
t.Fatalf("after failed suspend: hibernatePending=%v broken=%v, want true/false", u.hibernatePending, u.broken)
|
||||
}
|
||||
served := len(srv.RecordedPaths())
|
||||
|
||||
if u.runStep(context.Background(), pingStep) {
|
||||
t.Error("runStep = true while re-driving the hibernate; the step must not advance")
|
||||
}
|
||||
calls := ctl.recordedCalls()
|
||||
if got := calls[len(calls)-1]; got != "SuspendActor" {
|
||||
t.Errorf("last call = %q, want SuspendActor; calls = %v", got, calls)
|
||||
}
|
||||
if got := len(srv.RecordedPaths()); got != served {
|
||||
t.Errorf("router requests = %d, want %d (no wake ping against a stranded actor)", got, served)
|
||||
}
|
||||
if u.hibernatePending || u.broken {
|
||||
t.Errorf("after re-driven suspend: hibernatePending=%v broken=%v, want false/false", u.hibernatePending, u.broken)
|
||||
}
|
||||
|
||||
if !u.runStep(context.Background(), pingStep) {
|
||||
t.Error("runStep = false once the suspend cleared")
|
||||
}
|
||||
}
|
||||
|
||||
// ateapi reports a CRASHED actor on ResumeActor; it never recovers, so the
|
||||
// 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})
|
||||
|
||||
if u.runStep(context.Background(), pingStep) {
|
||||
t.Fatal("runStep = true with a crashed actor")
|
||||
}
|
||||
if !u.broken {
|
||||
t.Error("broken = false after a crashed verdict; want immediate replacement")
|
||||
}
|
||||
}
|
||||
|
||||
// In implicit mode the wake is an HTTP request, so the gRPC verdict arrives
|
||||
// from the hibernate that follows the failed step: FailedPrecondition means
|
||||
// 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{})
|
||||
|
||||
if u.runStep(context.Background(), pingStep) {
|
||||
t.Fatal("runStep = true with a failing router")
|
||||
}
|
||||
if !u.broken {
|
||||
t.Error("broken = false after FailedPrecondition on suspend; want immediate replacement")
|
||||
}
|
||||
}
|
||||
|
||||
// A saturated pool reports ResourceExhausted to every VU at once; a
|
||||
// replacement would need the same capacity, so the actor must be kept.
|
||||
func TestRunStep_KeepsActorThroughCapacityShortage(t *testing.T) {
|
||||
const rounds = maxConsecutiveStepFailures + 2
|
||||
errs := make([]error, rounds)
|
||||
for i := range errs {
|
||||
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})
|
||||
|
||||
for range rounds {
|
||||
if u.runStep(context.Background(), pingStep) {
|
||||
t.Fatal("runStep = true while resume is failing")
|
||||
}
|
||||
}
|
||||
if u.broken || u.consecutiveFailures != 0 {
|
||||
t.Errorf("broken=%v consecutiveFailures=%d after capacity errors, want false/0", u.broken, u.consecutiveFailures)
|
||||
}
|
||||
}
|
||||
|
||||
// HTTP failures carry no gRPC code, so they count toward the threshold: an
|
||||
// actor whose sandbox is dead but whose record looks healthy is replaced
|
||||
// after maxConsecutiveStepFailures steps.
|
||||
func TestRunStep_ReplacesActorAfterRepeatedStepFailures(t *testing.T) {
|
||||
ctl := &fakeControlClient{}
|
||||
u := newTestUser(t, &fake.Server{Status: 502}, ctl, dynconfig.Config{})
|
||||
|
||||
for i := 1; i <= maxConsecutiveStepFailures; i++ {
|
||||
u.runStep(context.Background(), pingStep)
|
||||
if want := i == maxConsecutiveStepFailures; u.broken != want {
|
||||
t.Fatalf("after %d failed steps: broken = %v, want %v", i, u.broken, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestControlRPCsCarryADeadline(t *testing.T) {
|
||||
ctl := &fakeControlClient{}
|
||||
u := newTestUser(t, &fake.Server{}, ctl, dynconfig.Config{ResumeMode: dynconfig.ResumeModeExplicit})
|
||||
|
||||
u.runStep(context.Background(), pingStep)
|
||||
u.suspendAndDelete(context.Background())
|
||||
|
||||
if !ctl.sawDeadline {
|
||||
t.Error("control-plane calls arrived without a context deadline")
|
||||
}
|
||||
for _, name := range []string{"ResumeActor", "SuspendActor", "DeleteActor"} {
|
||||
if got := countCalls(ctl.recordedCalls(), name); got == 0 {
|
||||
t.Errorf("%s was never called; calls = %v", name, ctl.recordedCalls())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestBurnRatePerGoroutine(t *testing.T) {
|
||||
rate, ok := burnRatePerGoroutine(burn(500, 2), 10_000)
|
||||
if !ok || rate != 10_000 {
|
||||
t.Errorf("burnRatePerGoroutine(500ms x2, 10k) = %v, %v; want 10000, true", rate, ok)
|
||||
}
|
||||
if _, ok := burnRatePerGoroutine(burn(0, 1), 5); ok {
|
||||
t.Error("a zero-duration burn must not report a rate")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,314 @@
|
||||
// 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 agentsession
|
||||
|
||||
import "time"
|
||||
|
||||
// This file is the whole benchmark, spelled out: one coding-agent session as
|
||||
// a script of steps. Each step names what the agent is doing in plain
|
||||
// English and lists the resource operations the real action would cost the
|
||||
// sandbox. Nothing is hidden in the driver — to understand or change the
|
||||
// workload, edit this table.
|
||||
//
|
||||
// Between steps the agent is "waiting for the LLM to think": the driver
|
||||
// suspends the actor, sleeps the step's think time, and lets the NEXT step's
|
||||
// first request wake the actor through the atenet router (request parking).
|
||||
// That idle gap is where Substrate earns its keep, so the think times are
|
||||
// first-class script data, not driver noise.
|
||||
|
||||
// op is one resource effect inside a step, executed as a single glutton RPC
|
||||
// through the router. Build them with the constructors below so the script
|
||||
// stays declarative and greppable.
|
||||
type op struct {
|
||||
kind opKind
|
||||
key string // file or RAM-array name inside the sandbox
|
||||
bytes int64 // payload / file / RAM size
|
||||
millis int64 // CPU burn wall-clock
|
||||
parallel int32 // CPU burn goroutines
|
||||
}
|
||||
|
||||
type opKind int
|
||||
|
||||
const (
|
||||
// opIngest pushes bytes from the driver through the router into the
|
||||
// actor, which writes them to disk: real network ingress + disk write,
|
||||
// the shape of a download.
|
||||
opIngest opKind = iota
|
||||
// opBurnCPU spins the sandbox's CPU for a wall-clock duration.
|
||||
opBurnCPU
|
||||
// opWriteDisk writes locally generated random bytes to a sandbox file.
|
||||
opWriteDisk
|
||||
// opReadDiskDigest reads and sha256-hashes a sandbox file without
|
||||
// shipping the bytes back: disk read I/O only.
|
||||
opReadDiskDigest
|
||||
// opReadDiskData reads a sandbox file AND returns its bytes to the
|
||||
// driver: disk read + network egress through the router response.
|
||||
opReadDiskData
|
||||
// opFillRAM allocates a resident RAM array of random bytes.
|
||||
opFillRAM
|
||||
// opChurnRAM re-randomizes part of an existing RAM array in place,
|
||||
// dirtying pages so the next suspend snapshot has fresh content.
|
||||
opChurnRAM
|
||||
// opWalkRAM touches one byte per page of a RAM array, forcing every
|
||||
// page resident — after a resume this measures demand-paging cost.
|
||||
opWalkRAM
|
||||
// opPing is a minimal round-trip through the router.
|
||||
opPing
|
||||
)
|
||||
|
||||
func ingest(key string, bytes int64) op { return op{kind: opIngest, key: key, bytes: bytes} }
|
||||
func burn(millis int64, parallel int32) op {
|
||||
return op{kind: opBurnCPU, millis: millis, parallel: parallel}
|
||||
}
|
||||
func writeDisk(key string, bytes int64) op { return op{kind: opWriteDisk, key: key, bytes: bytes} }
|
||||
func readDigest(key string) op { return op{kind: opReadDiskDigest, key: key} }
|
||||
func readData(key string) op { return op{kind: opReadDiskData, key: key} }
|
||||
func fillRAM(key string, bytes int64) op { return op{kind: opFillRAM, key: key, bytes: bytes} }
|
||||
func churnRAM(key string, bytes int64) op { return op{kind: opChurnRAM, key: key, bytes: bytes} }
|
||||
func walkRAM(key string) op { return op{kind: opWalkRAM, key: key} }
|
||||
func ping() op { return op{kind: opPing} }
|
||||
|
||||
// Step is one agent action: what a coding agent would be doing, the think
|
||||
// time that precedes it (the LLM producing this step), and the resource
|
||||
// operations acting it out.
|
||||
type Step struct {
|
||||
// Name keys the step's locust stats row: Step_<Name>.
|
||||
Name string
|
||||
// Agent says what the coding agent is doing, for humans.
|
||||
Agent string
|
||||
// Think is how long the LLM "thinks" before this step. The actor is
|
||||
// suspended for this gap (scaled by --agentsession-think-scale).
|
||||
Think time.Duration
|
||||
// Ops are the resource effects, executed in order.
|
||||
Ops []op
|
||||
}
|
||||
|
||||
const (
|
||||
_ = iota
|
||||
kib = int64(1) << (10 * iota)
|
||||
mib
|
||||
)
|
||||
|
||||
// Sandbox object names, so the script reads like a filesystem.
|
||||
const (
|
||||
repoBlob = "repo_tarball" // the cloned repository
|
||||
depsBlob = "deps_cache" // downloaded dependency archives
|
||||
buildOut = "build_artifacts" // compiler output
|
||||
patchFile = "patch" // edits arriving from the LLM
|
||||
packageBlob = "release_package" // the final packaged artifact
|
||||
contextRAM = "agent_context" // the agent process's resident working set
|
||||
compilerRAM = "compiler_ws" // build-time working set
|
||||
testDataBlob = "test_fixtures" // fixtures written by the test steps
|
||||
)
|
||||
|
||||
// Session is the default 20-step coding-agent session. The script itself
|
||||
// declares ≈96Mi of resident RAM (contextRAM 32Mi + compilerRAM 64Mi) and
|
||||
// ≈110Mi of files; the guest peak on top of that (kernel, kata-agent, Go
|
||||
// allocator transients) has been observed around 320Mi, which is why the
|
||||
// template calls for 1Gi actors. TestSessionBudgets bounds only the
|
||||
// script-declared bytes, so edits that grow the working set fail the test
|
||||
// and force the memory guidance to be revisited.
|
||||
//
|
||||
// The narrative: the agent is told to fetch a repository, get it building,
|
||||
// fix a failing test, extend the test suite, refactor, and package the
|
||||
// result — with an LLM round-trip (seconds of suspension) before every step.
|
||||
func Session() []Step {
|
||||
return []Step{
|
||||
{
|
||||
Name: "01_read_task",
|
||||
Agent: "Boots, reads the task prompt, loads its context window",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
fillRAM(contextRAM, 32*mib), // the agent's resident working set for the whole session
|
||||
ping(),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "02_clone_repo",
|
||||
Agent: "git clone of the target repository",
|
||||
Think: 3 * time.Second,
|
||||
Ops: []op{
|
||||
ingest(repoBlob, 16*mib), // tarball arrives over the network, lands on disk
|
||||
burn(500, 1), // checkout: decompress + write tree
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "03_explore_tree",
|
||||
Agent: "Lists directories, greps for entry points",
|
||||
Think: 4 * time.Second,
|
||||
Ops: []op{
|
||||
readDigest(repoBlob), // walk the repo bytes on disk
|
||||
burn(200, 1), // grep-ish scanning
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "04_read_key_files",
|
||||
Agent: "Opens the files it plans to change; they enter the LLM context",
|
||||
Think: 3 * time.Second,
|
||||
Ops: []op{
|
||||
readData(repoBlob), // file contents also travel back out to the 'LLM'
|
||||
churnRAM(contextRAM, 8*mib), // context window grows/changes
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "05_install_deps",
|
||||
Agent: "Installs dependencies (pip install / go mod download)",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
ingest(depsBlob, 32*mib), // packages arrive over the network
|
||||
burn(1500, 2), // unpack, byte-compile, link
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "06_first_build",
|
||||
Agent: "First full build of the project",
|
||||
Think: 3 * time.Second,
|
||||
Ops: []op{
|
||||
fillRAM(compilerRAM, 64*mib), // compiler working set
|
||||
burn(3000, 2), // the compile itself
|
||||
writeDisk(buildOut, 24*mib), // object files and binaries
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "07_run_unit_tests",
|
||||
Agent: "Runs the existing unit tests; one fails",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
readDigest(buildOut), // load test binaries
|
||||
burn(2500, 2), // the test run
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "08_reason_about_failure",
|
||||
Agent: "Re-reads the failing test and traces the bug (mostly thinking)",
|
||||
Think: 8 * time.Second, // the long LLM analysis turn
|
||||
Ops: []op{
|
||||
walkRAM(contextRAM), // re-touch the whole context after the long suspend
|
||||
ping(),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "09_edit_source",
|
||||
Agent: "Applies the LLM's fix to the source files",
|
||||
Think: 5 * time.Second,
|
||||
Ops: []op{
|
||||
ingest(patchFile, 64*kib), // the patch arrives from the LLM
|
||||
writeDisk(repoBlob, 4*mib), // rewrite the touched sources
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "10_incremental_build",
|
||||
Agent: "Rebuilds just the changed packages",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
burn(1200, 2),
|
||||
writeDisk(buildOut, 8*mib),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "11_rerun_failed_test",
|
||||
Agent: "Re-runs the previously failing test; it passes",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
burn(800, 1),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "12_write_new_tests",
|
||||
Agent: "Writes regression tests for the fix",
|
||||
Think: 6 * time.Second, // LLM authors the test file
|
||||
Ops: []op{
|
||||
ingest(patchFile, 128*kib),
|
||||
writeDisk(testDataBlob, 2*mib), // fixtures the new tests need
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "13_run_new_tests",
|
||||
Agent: "Runs the new tests in isolation",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
readDigest(testDataBlob),
|
||||
burn(1000, 1),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "14_full_test_suite",
|
||||
Agent: "Runs the entire test suite to check for regressions",
|
||||
Think: 3 * time.Second,
|
||||
Ops: []op{
|
||||
burn(4000, 2), // the longest compute step in the session
|
||||
readDigest(buildOut),
|
||||
churnRAM(contextRAM, 8*mib), // test output enters the context
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "15_lint_format",
|
||||
Agent: "Runs the linter and formatter over the tree",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
readDigest(repoBlob),
|
||||
burn(900, 1),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "16_refactor",
|
||||
Agent: "Applies the LLM's multi-file cleanup refactor",
|
||||
Think: 8 * time.Second, // LLM plans the refactor
|
||||
Ops: []op{
|
||||
ingest(patchFile, 512*kib),
|
||||
writeDisk(repoBlob, 6*mib),
|
||||
burn(600, 1),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "17_rebuild",
|
||||
Agent: "Full rebuild after the refactor",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
churnRAM(compilerRAM, 32*mib), // compiler working set changes with the new tree
|
||||
burn(2000, 2),
|
||||
writeDisk(buildOut, 16*mib),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "18_final_test_suite",
|
||||
Agent: "Final full-suite run before shipping",
|
||||
Think: 2 * time.Second,
|
||||
Ops: []op{
|
||||
burn(3500, 2),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "19_package_artifact",
|
||||
Agent: "Builds the release package / container image",
|
||||
Think: 3 * time.Second,
|
||||
Ops: []op{
|
||||
readDigest(buildOut),
|
||||
burn(1000, 1),
|
||||
writeDisk(packageBlob, 24*mib),
|
||||
},
|
||||
},
|
||||
{
|
||||
Name: "20_commit_and_summarize",
|
||||
Agent: "Commits, writes the summary, hands the session back",
|
||||
Think: 5 * time.Second, // LLM writes the summary
|
||||
Ops: []op{
|
||||
writeDisk(repoBlob, 256*kib), // the commit
|
||||
readData(packageBlob), // artifact ships back out over the network
|
||||
walkRAM(contextRAM), // final context read
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,136 @@
|
||||
// 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 boomerutil
|
||||
|
||||
import (
|
||||
"context"
|
||||
"math/rand/v2"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// FailureAction is what a user class does about a failed lifecycle RPC: is
|
||||
// this actor wedged, or is the cluster busy? Only the first is worth a
|
||||
// delete + create.
|
||||
type FailureAction int
|
||||
|
||||
const (
|
||||
// RetryLater: a replacement would hit the same error, so keep the actor.
|
||||
RetryLater FailureAction = iota
|
||||
// ReplaceNow: this actor can never make progress again.
|
||||
ReplaceNow
|
||||
// ReplaceIfPersistent: counts toward the caller's consecutive-failure
|
||||
// threshold.
|
||||
ReplaceIfPersistent
|
||||
)
|
||||
|
||||
const (
|
||||
// ConcurrentUpdateMsg is the message ateapi returns with codes.Aborted
|
||||
// when two callers race an actor update (see
|
||||
// cmd/ateapi/internal/controlapi/workflow_resume.go). It's transient —
|
||||
// the loser retries and one of them wins.
|
||||
ConcurrentUpdateMsg = "concurrent update conflict, please retry"
|
||||
|
||||
// ConflictRetryAttempts bounds RetryOnConflict. The first retry is
|
||||
// immediate (conflicts often clear the instant the racing writer
|
||||
// commits); later gaps are conflictRetryBackoff plus uniform jitter in
|
||||
// [0, conflictRetryJitter) so the losers of one race do not retry in
|
||||
// lockstep and lose the next one the same way.
|
||||
ConflictRetryAttempts = 5
|
||||
// ConflictRetryBackoff is the gap before the second and later retries.
|
||||
ConflictRetryBackoff = 50 * time.Millisecond
|
||||
conflictRetryJitter = 5 * time.Millisecond
|
||||
)
|
||||
|
||||
// ClassifyLifecycleFailure maps an error from ResumeActor / SuspendActor /
|
||||
// PauseActor onto what to do about it. Unrecognized codes are recoverable
|
||||
// until proven otherwise: a wrapped atelet error arrives as Unknown.
|
||||
func ClassifyLifecycleFailure(err error) FailureAction {
|
||||
s, ok := status.FromError(err)
|
||||
if !ok {
|
||||
return ReplaceIfPersistent
|
||||
}
|
||||
switch s.Code() {
|
||||
case codes.NotFound, codes.DataLoss:
|
||||
// The actor, or the snapshot it would resume from, is gone.
|
||||
return ReplaceNow
|
||||
case codes.FailedPrecondition:
|
||||
// A state this operation has no edge out of — the left-in-SUSPENDING
|
||||
// case, or a CRASHED actor. Nothing a client can call moves it on.
|
||||
return ReplaceNow
|
||||
case codes.Aborted:
|
||||
if IsCrashed(err) {
|
||||
return ReplaceNow
|
||||
}
|
||||
// Concurrent update conflict. The caller already spent its retry
|
||||
// budget, but losing every race in one burst is still a race.
|
||||
return ReplaceIfPersistent
|
||||
case codes.ResourceExhausted, codes.Unavailable, codes.DeadlineExceeded, codes.Canceled:
|
||||
// Cluster-wide and load-dependent ("no free workers available", a
|
||||
// restarting ate-api-server): every VU sees these at once, and the
|
||||
// replacement needs the capacity the original was denied.
|
||||
return RetryLater
|
||||
case codes.InvalidArgument, codes.PermissionDenied, codes.Unauthenticated, codes.Unimplemented:
|
||||
// A misconfigured run. The replacement is created the same way and
|
||||
// fails the same way.
|
||||
return RetryLater
|
||||
default:
|
||||
return ReplaceIfPersistent
|
||||
}
|
||||
}
|
||||
|
||||
// IsCrashed reports whether err is ateapi's ResumeActor verdict on an actor
|
||||
// in ACTOR_STATE_CRASHED: codes.Aborted with "crashed" in the message.
|
||||
// ateapi never rehabilitates one.
|
||||
func IsCrashed(err error) bool {
|
||||
s, ok := status.FromError(err)
|
||||
return ok && s.Code() == codes.Aborted && strings.Contains(s.Message(), "crashed")
|
||||
}
|
||||
|
||||
// IsConcurrentUpdateConflict identifies the transient racy-update error
|
||||
// ateapi's workflow_*.go returns as codes.Aborted with the retry-me message.
|
||||
// Kept distinct from IsCrashed because the two look the same at the code
|
||||
// level and mean opposite things.
|
||||
func IsConcurrentUpdateConflict(err error) bool {
|
||||
s, ok := status.FromError(err)
|
||||
return ok && s.Code() == codes.Aborted && strings.Contains(s.Message(), ConcurrentUpdateMsg)
|
||||
}
|
||||
|
||||
// RetryOnConflict runs call up to ConflictRetryAttempts times while it keeps
|
||||
// returning a concurrent-update conflict, and returns any other error
|
||||
// unchanged. A canceled ctx ends the loop between attempts.
|
||||
func RetryOnConflict(ctx context.Context, call func() error) error {
|
||||
var backoff time.Duration // 0 → first retry runs immediately
|
||||
var lastErr error
|
||||
for range ConflictRetryAttempts {
|
||||
lastErr = call()
|
||||
if lastErr == nil || !IsConcurrentUpdateConflict(lastErr) {
|
||||
return lastErr
|
||||
}
|
||||
if backoff > 0 {
|
||||
jitter := time.Duration(rand.Float64() * float64(conflictRetryJitter))
|
||||
select {
|
||||
case <-time.After(backoff + jitter):
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
}
|
||||
backoff = ConflictRetryBackoff
|
||||
}
|
||||
return lastErr
|
||||
}
|
||||
@@ -0,0 +1,142 @@
|
||||
// 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 boomerutil
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
func conflictErr() error {
|
||||
return status.Error(codes.Aborted, ConcurrentUpdateMsg)
|
||||
}
|
||||
|
||||
func TestClassifyLifecycleFailure(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
err error
|
||||
want FailureAction
|
||||
}{
|
||||
{"actor gone", status.Error(codes.NotFound, "Actor not found"), ReplaceNow},
|
||||
{"snapshot unreadable", status.Error(codes.DataLoss, "external snapshot"), ReplaceNow},
|
||||
{"stuck state", status.Error(codes.FailedPrecondition, "MarkSuspending prerequisite not met"), ReplaceNow},
|
||||
{"crashed", status.Error(codes.Aborted, "actor bench/sb-1 crashed"), ReplaceNow},
|
||||
{"no capacity", status.Error(codes.ResourceExhausted, "no free workers available"), RetryLater},
|
||||
{"api server down", status.Error(codes.Unavailable, "connection refused"), RetryLater},
|
||||
{"timeout", status.Error(codes.DeadlineExceeded, "context deadline exceeded"), RetryLater},
|
||||
{"canceled", status.Error(codes.Canceled, "context canceled"), RetryLater},
|
||||
{"misconfigured run", status.Error(codes.InvalidArgument, "bad template"), RetryLater},
|
||||
{"unauthorized", status.Error(codes.PermissionDenied, "denied"), RetryLater},
|
||||
{"update conflict", conflictErr(), ReplaceIfPersistent},
|
||||
{"wrapped atelet error", status.Error(codes.Unknown, "while checkpointing workload"), ReplaceIfPersistent},
|
||||
{"internal", status.Error(codes.Internal, "boom"), ReplaceIfPersistent},
|
||||
{"not a status", errors.New("plain error"), ReplaceIfPersistent},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if got := ClassifyLifecycleFailure(tc.err); got != tc.want {
|
||||
t.Errorf("ClassifyLifecycleFailure(%v) = %v, want %v", tc.err, got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsCrashedVsConflict(t *testing.T) {
|
||||
crashed := status.Error(codes.Aborted, "actor bench/sb-1 crashed")
|
||||
if !IsCrashed(crashed) || IsConcurrentUpdateConflict(crashed) {
|
||||
t.Error("crashed verdict misclassified")
|
||||
}
|
||||
if IsCrashed(conflictErr()) || !IsConcurrentUpdateConflict(conflictErr()) {
|
||||
t.Error("update conflict misclassified")
|
||||
}
|
||||
}
|
||||
|
||||
// queue returns a call that pops errs in order and succeeds once drained,
|
||||
// counting the attempts. Safe for concurrent callers.
|
||||
func queue(errs ...error) (call func() error, attempts *atomic.Int64) {
|
||||
var mu sync.Mutex
|
||||
attempts = new(atomic.Int64)
|
||||
return func() error {
|
||||
attempts.Add(1)
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if len(errs) == 0 {
|
||||
return nil
|
||||
}
|
||||
err := errs[0]
|
||||
errs = errs[1:]
|
||||
return err
|
||||
}, attempts
|
||||
}
|
||||
|
||||
func TestRetryOnConflictRetriesThenSucceeds(t *testing.T) {
|
||||
errs := []error{conflictErr(), conflictErr()}
|
||||
call, attempts := queue(errs...)
|
||||
start := time.Now()
|
||||
if err := RetryOnConflict(context.Background(), call); err != nil {
|
||||
t.Fatalf("RetryOnConflict = %v, want nil after conflicts clear", err)
|
||||
}
|
||||
if want := int64(len(errs) + 1); attempts.Load() != want {
|
||||
t.Errorf("attempts = %d, want %d (every conflict, then success)", attempts.Load(), want)
|
||||
}
|
||||
if elapsed := time.Since(start); elapsed < ConflictRetryBackoff {
|
||||
t.Errorf("elapsed = %v, want >= %v (second retry must back off)", elapsed, ConflictRetryBackoff)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryOnConflictPassesOtherErrorsThrough(t *testing.T) {
|
||||
want := status.Error(codes.Unavailable, "down")
|
||||
call, attempts := queue(want)
|
||||
if err := RetryOnConflict(context.Background(), call); !errors.Is(err, want) {
|
||||
t.Fatalf("RetryOnConflict = %v, want %v", err, want)
|
||||
}
|
||||
if attempts.Load() != 1 {
|
||||
t.Errorf("attempts = %d, want 1 (no retry for non-conflict errors)", attempts.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryOnConflictGivesUp(t *testing.T) {
|
||||
errs := make([]error, ConflictRetryAttempts+1)
|
||||
for i := range errs {
|
||||
errs[i] = conflictErr()
|
||||
}
|
||||
call, attempts := queue(errs...)
|
||||
if err := RetryOnConflict(context.Background(), call); !IsConcurrentUpdateConflict(err) {
|
||||
t.Fatalf("RetryOnConflict = %v, want the last conflict", err)
|
||||
}
|
||||
if attempts.Load() != ConflictRetryAttempts {
|
||||
t.Errorf("attempts = %d, want %d", attempts.Load(), ConflictRetryAttempts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryOnConflictHonorsContext(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
call, attempts := queue(conflictErr(), conflictErr(), conflictErr())
|
||||
if err := RetryOnConflict(ctx, call); !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("RetryOnConflict = %v, want context.Canceled", err)
|
||||
}
|
||||
// The first retry is immediate; the second sees the canceled context
|
||||
// instead of sleeping.
|
||||
if attempts.Load() != 2 {
|
||||
t.Errorf("attempts = %d, want 2 after cancellation", attempts.Load())
|
||||
}
|
||||
}
|
||||
@@ -68,6 +68,8 @@ type Config struct {
|
||||
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
|
||||
|
||||
AgentSessionThinkScale float64 // multiplier on the script's per-step think times; 0 reads as 1.0
|
||||
}
|
||||
|
||||
// Holder lets readers Load() the current Config and writers Store() a new
|
||||
@@ -116,6 +118,8 @@ type payload struct {
|
||||
SweperfTotalSteps *float64 `json:"sweperf_total_steps"`
|
||||
SweperfNumCycles *float64 `json:"sweperf_num_cycles"`
|
||||
SweperfPollIntervalMs *float64 `json:"sweperf_poll_interval_ms"`
|
||||
|
||||
AgentSessionThinkScale *float64 `json:"agentsession_think_scale"`
|
||||
}
|
||||
|
||||
// Parse decodes a JSON blob (typically from a CLI flag) and merges its
|
||||
@@ -209,6 +213,9 @@ func (c Config) Validate() error {
|
||||
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)
|
||||
}
|
||||
// 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.
|
||||
@@ -278,6 +285,9 @@ func (p payload) merge(current Config) Config {
|
||||
if p.SweperfPollIntervalMs != nil {
|
||||
out.SweperfPollIntervalMs = int(*p.SweperfPollIntervalMs)
|
||||
}
|
||||
if p.AgentSessionThinkScale != nil {
|
||||
out.AgentSessionThinkScale = *p.AgentSessionThinkScale
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
@@ -354,6 +364,7 @@ func StartPoll(
|
||||
slog.Int("sweperf_total_steps", next.SweperfTotalSteps),
|
||||
slog.Int("sweperf_num_cycles", next.SweperfNumCycles),
|
||||
slog.Int("sweperf_poll_interval_ms", next.SweperfPollIntervalMs),
|
||||
slog.Float64("agentsession_think_scale", next.AgentSessionThinkScale),
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -401,6 +412,7 @@ func SubscribeSpawn(url string, holder *Holder, sampler ProbabilityUpdater, fetc
|
||||
slog.Int("sweperf_total_steps", next.SweperfTotalSteps),
|
||||
slog.Int("sweperf_num_cycles", next.SweperfNumCycles),
|
||||
slog.Int("sweperf_poll_interval_ms", next.SweperfPollIntervalMs),
|
||||
slog.Float64("agentsession_think_scale", next.AgentSessionThinkScale),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -58,20 +58,7 @@ const (
|
||||
memLoadKey = "memload"
|
||||
memReadAll = "all"
|
||||
|
||||
// ateapi returns Aborted with this message when two callers race an
|
||||
// actor update (see cmd/ateapi/internal/controlapi/workflow_resume.go).
|
||||
// It's transient — the loser retries and one of them wins.
|
||||
concurrentUpdateMsg = "concurrent update conflict, please retry"
|
||||
// Retry budget for ResumeActor concurrent-update conflicts. First retry
|
||||
// is immediate (no initial backoff — conflicts often clear the instant
|
||||
// the racing writer commits); subsequent gaps are 50ms plus uniform
|
||||
// jitter in [0, resumeBackoffJitter) so the losers of one race do not
|
||||
// retry in lockstep and lose the next one the same way.
|
||||
resumeMaxAttempts = 5
|
||||
resumeMaxBackoff = 50 * time.Millisecond
|
||||
resumeBackoffJitter = 5 * time.Millisecond
|
||||
|
||||
// Consecutive replaceIfPersistent failures it takes to replace an actor.
|
||||
// Consecutive ReplaceIfPersistent failures it takes to replace an actor.
|
||||
// One blip is not evidence of a wedged actor; three in a row, each
|
||||
// spaced by the wait window, is.
|
||||
maxConsecutiveFailures = 3
|
||||
@@ -361,63 +348,13 @@ type gluttonActor struct {
|
||||
consecutiveFailures int
|
||||
}
|
||||
|
||||
// failureAction is what iterate() does about a failed lifecycle RPC: is this
|
||||
// actor wedged, or is the cluster busy? Only the first is worth a
|
||||
// delete + create.
|
||||
type failureAction int
|
||||
|
||||
const (
|
||||
// retryLater: a replacement would hit the same error, so keep the actor.
|
||||
retryLater failureAction = iota
|
||||
// replaceNow: this actor can never make progress again.
|
||||
replaceNow
|
||||
// replaceIfPersistent: counts toward maxConsecutiveFailures.
|
||||
replaceIfPersistent
|
||||
)
|
||||
|
||||
// classifyLifecycleFailure maps an error from ResumeActor / SuspendActor /
|
||||
// PauseActor onto what to do about it. Unrecognized codes are recoverable
|
||||
// until proven otherwise: a wrapped atelet error arrives as Unknown.
|
||||
func classifyLifecycleFailure(err error) failureAction {
|
||||
s, ok := status.FromError(err)
|
||||
if !ok {
|
||||
return replaceIfPersistent
|
||||
}
|
||||
switch s.Code() {
|
||||
case codes.NotFound, codes.DataLoss:
|
||||
// The actor, or the snapshot it would resume from, is gone.
|
||||
return replaceNow
|
||||
case codes.FailedPrecondition:
|
||||
// A state this operation has no edge out of — the left-in-SUSPENDING
|
||||
// case. Nothing glutton can call moves the actor on.
|
||||
return replaceNow
|
||||
case codes.Aborted:
|
||||
if strings.Contains(s.Message(), "crashed") {
|
||||
return replaceNow
|
||||
}
|
||||
// Concurrent update conflict. resume() already spent its retry
|
||||
// budget, but losing every race in one burst is still a race.
|
||||
return replaceIfPersistent
|
||||
case codes.ResourceExhausted, codes.Unavailable, codes.DeadlineExceeded, codes.Canceled:
|
||||
// Cluster-wide and load-dependent ("no free workers available", a
|
||||
// restarting ate-api-server): every VU sees these at once, and the
|
||||
// replacement needs the capacity the original was denied.
|
||||
return retryLater
|
||||
case codes.InvalidArgument, codes.PermissionDenied, codes.Unauthenticated, codes.Unimplemented:
|
||||
// A misconfigured run. The replacement is created the same way and
|
||||
// fails the same way.
|
||||
return retryLater
|
||||
}
|
||||
return replaceIfPersistent
|
||||
}
|
||||
|
||||
// noteFailure records a failed lifecycle RPC against the actor and reports
|
||||
// whether the VU should replace it.
|
||||
func (u *gluttonActor) noteFailure(err error) bool {
|
||||
switch classifyLifecycleFailure(err) {
|
||||
case replaceNow:
|
||||
switch boomerutil.ClassifyLifecycleFailure(err) {
|
||||
case boomerutil.ReplaceNow:
|
||||
return true
|
||||
case retryLater:
|
||||
case boomerutil.RetryLater:
|
||||
return false
|
||||
}
|
||||
u.consecutiveFailures++
|
||||
@@ -487,36 +424,18 @@ func (u *gluttonActor) resume(ctx context.Context) error {
|
||||
// the ateapi contract. Kept inside the tracedCall closure so the
|
||||
// reported latency spans every attempt and the span carries the
|
||||
// last attempt's server trailer, same as any other single-shot RPC.
|
||||
var backoff time.Duration // 0 → first retry runs immediately
|
||||
var lastErr error
|
||||
for range resumeMaxAttempts {
|
||||
_, lastErr = u.cfg.APIStub.ResumeActor(callCtx, &ateapipb.ResumeActorRequest{
|
||||
return boomerutil.RetryOnConflict(callCtx, func() error {
|
||||
_, err := u.cfg.APIStub.ResumeActor(callCtx, &ateapipb.ResumeActorRequest{
|
||||
Actor: u.ref(),
|
||||
}, grpc.Trailer(tr))
|
||||
if lastErr == nil {
|
||||
return nil
|
||||
}
|
||||
if !isConcurrentUpdateConflict(lastErr) {
|
||||
return lastErr
|
||||
}
|
||||
if backoff > 0 {
|
||||
jitter := time.Duration(rand.Float64() * float64(resumeBackoffJitter))
|
||||
select {
|
||||
case <-time.After(backoff + jitter):
|
||||
case <-callCtx.Done():
|
||||
return callCtx.Err()
|
||||
}
|
||||
}
|
||||
backoff = resumeMaxBackoff
|
||||
}
|
||||
return lastErr
|
||||
return err
|
||||
})
|
||||
})
|
||||
if err != nil {
|
||||
// ateapi reports a crashed actor as codes.Aborted with "crashed" in
|
||||
// the message (see workflow_resume.go). Mark the user so iterate()
|
||||
// stops touching it, and surface a CrashCount tick so operators can
|
||||
// see the crash total in the locust stats table.
|
||||
if s, ok := status.FromError(err); ok && s.Code() == codes.Aborted && strings.Contains(s.Message(), "crashed") {
|
||||
// Mark a crashed actor so iterate() stops touching it, and surface a
|
||||
// CrashCount tick so operators can see the crash total in the
|
||||
// locust stats table.
|
||||
if boomerutil.IsCrashed(err) {
|
||||
u.crashed = true
|
||||
bmetrics.RecordFailure("actor", "CrashCount", userClass, 0, "actor entered ACTOR_STATE_CRASHED")
|
||||
slog.Warn("glutton actor crashed; will stop sending requests",
|
||||
@@ -531,15 +450,6 @@ func (u *gluttonActor) resume(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// isConcurrentUpdateConflict identifies the transient racy-update error
|
||||
// ateapi's workflow_*.go returns as codes.Aborted with the retry-me message.
|
||||
// Kept distinct from the "crashed" Aborted check in resume() because the two
|
||||
// look the same at the code level and mean opposite things.
|
||||
func isConcurrentUpdateConflict(err error) bool {
|
||||
s, ok := status.FromError(err)
|
||||
return ok && s.Code() == codes.Aborted && strings.Contains(s.Message(), concurrentUpdateMsg)
|
||||
}
|
||||
|
||||
// hibernate takes the actor off its worker by whichever operation the
|
||||
// lifecycle mode selects: PauseActor keeps the snapshot on the node, while
|
||||
// SuspendActor writes it to durable storage.
|
||||
|
||||
@@ -16,7 +16,6 @@ package glutton
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"slices"
|
||||
"testing"
|
||||
|
||||
@@ -203,7 +202,7 @@ func TestGluttonIterate_KeepsActorOnTransientResumeFailure(t *testing.T) {
|
||||
// An exhausted conflict-retry budget only replaces the actor once it has
|
||||
// repeated maxConsecutiveFailures times.
|
||||
func TestGluttonIterate_ReplacesActorAfterRepeatedConflicts(t *testing.T) {
|
||||
errs := make([]error, resumeMaxAttempts*maxConsecutiveFailures)
|
||||
errs := make([]error, boomerutil.ConflictRetryAttempts*maxConsecutiveFailures)
|
||||
for i := range errs {
|
||||
errs[i] = conflictErr()
|
||||
}
|
||||
@@ -260,35 +259,6 @@ func TestGluttonIterate_RetriesStrandedHibernate(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestClassifyLifecycleFailure(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
err error
|
||||
want failureAction
|
||||
}{
|
||||
{"actor gone", status.Error(codes.NotFound, "Actor not found"), replaceNow},
|
||||
{"snapshot unreadable", status.Error(codes.DataLoss, "external snapshot"), replaceNow},
|
||||
{"stuck state", status.Error(codes.FailedPrecondition, "MarkSuspending prerequisite not met"), replaceNow},
|
||||
{"crashed", status.Error(codes.Aborted, "actor bench/sb-1 crashed"), replaceNow},
|
||||
{"no capacity", status.Error(codes.ResourceExhausted, "no free workers available"), retryLater},
|
||||
{"api server down", status.Error(codes.Unavailable, "connection refused"), retryLater},
|
||||
{"timeout", status.Error(codes.DeadlineExceeded, "context deadline exceeded"), retryLater},
|
||||
{"canceled", status.Error(codes.Canceled, "context canceled"), retryLater},
|
||||
{"misconfigured run", status.Error(codes.InvalidArgument, "bad template"), retryLater},
|
||||
{"unauthorized", status.Error(codes.PermissionDenied, "denied"), retryLater},
|
||||
{"update conflict", conflictErr(), replaceIfPersistent},
|
||||
{"wrapped atelet error", status.Error(codes.Unknown, "while checkpointing workload"), replaceIfPersistent},
|
||||
{"internal", status.Error(codes.Internal, "boom"), replaceIfPersistent},
|
||||
{"not a status", errors.New("plain error"), replaceIfPersistent},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if got := classifyLifecycleFailure(tc.err); got != tc.want {
|
||||
t.Errorf("classifyLifecycleFailure(%v) = %v, want %v", tc.err, got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNoteFailureResetsOnSuccess(t *testing.T) {
|
||||
a := &gluttonActor{actorName: "sb-1"}
|
||||
for range maxConsecutiveFailures - 1 {
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/boomer/boomerutil"
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/boomer/userclass"
|
||||
"github.com/agent-substrate/substrate/internal/benchmarking/glutton/fake"
|
||||
"google.golang.org/grpc/codes"
|
||||
@@ -27,7 +28,7 @@ import (
|
||||
)
|
||||
|
||||
func conflictErr() error {
|
||||
return status.Error(codes.Aborted, concurrentUpdateMsg)
|
||||
return status.Error(codes.Aborted, boomerutil.ConcurrentUpdateMsg)
|
||||
}
|
||||
|
||||
func newResumeTestActor(t *testing.T, resumeErrs ...error) (*gluttonActor, *fakeControlClient) {
|
||||
@@ -54,8 +55,8 @@ func TestResumeRetriesConcurrentUpdateConflict(t *testing.T) {
|
||||
if got := resumeCalls(fakeCtrl); got != 3 {
|
||||
t.Errorf("ResumeActor calls = %d, want 3 (two conflicts, then success)", got)
|
||||
}
|
||||
if elapsed < resumeMaxBackoff {
|
||||
t.Errorf("elapsed = %v, want >= %v (second retry must back off)", elapsed, resumeMaxBackoff)
|
||||
if elapsed < boomerutil.ConflictRetryBackoff {
|
||||
t.Errorf("elapsed = %v, want >= %v (second retry must back off)", elapsed, boomerutil.ConflictRetryBackoff)
|
||||
}
|
||||
if u.crashed {
|
||||
t.Error("crashed = true after a transient conflict")
|
||||
@@ -74,7 +75,7 @@ func TestResumeDoesNotRetryOtherErrors(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestResumeGivesUpAfterMaxAttempts(t *testing.T) {
|
||||
errs := make([]error, resumeMaxAttempts+1)
|
||||
errs := make([]error, boomerutil.ConflictRetryAttempts+1)
|
||||
for i := range errs {
|
||||
errs[i] = conflictErr()
|
||||
}
|
||||
@@ -83,8 +84,8 @@ func TestResumeGivesUpAfterMaxAttempts(t *testing.T) {
|
||||
if err := u.resume(context.Background()); err == nil {
|
||||
t.Fatal("resume = nil, want an error when every attempt conflicts")
|
||||
}
|
||||
if got := resumeCalls(fakeCtrl); got != resumeMaxAttempts {
|
||||
t.Errorf("ResumeActor calls = %d, want %d", got, resumeMaxAttempts)
|
||||
if got := resumeCalls(fakeCtrl); got != boomerutil.ConflictRetryAttempts {
|
||||
t.Errorf("ResumeActor calls = %d, want %d", got, boomerutil.ConflictRetryAttempts)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -37,6 +37,9 @@ const (
|
||||
ReadDiskRoute = glutton.ReadDiskRoute
|
||||
WriteRAMRoute = glutton.WriteRAMRoute
|
||||
ReadRAMRoute = glutton.ReadRAMRoute
|
||||
BurnCPURoute = glutton.BurnCPURoute
|
||||
IngestRoute = glutton.IngestRoute
|
||||
PingRoute = glutton.PingRoute
|
||||
)
|
||||
|
||||
// Server is an httptest-backed stand-in for a glutton actor holding one file.
|
||||
@@ -65,6 +68,8 @@ type Server struct {
|
||||
ramWriteSizes []string
|
||||
ramWriteModes []gluttonpb.WriteMode
|
||||
ramReadSizes []string
|
||||
burnMillis []int64
|
||||
ingestSizes []int64
|
||||
}
|
||||
|
||||
func (s *Server) reportedDigest() []byte {
|
||||
@@ -231,7 +236,75 @@ func (s *Server) serve(w http.ResponseWriter, r *http.Request) {
|
||||
resp, _ := proto.Marshal(&gluttonpb.ReadRAMResponse{Size: int64(len(s.Data))})
|
||||
_, _ = w.Write(resp)
|
||||
|
||||
case BurnCPURoute:
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
var req gluttonpb.BurnCPURequest
|
||||
if err := proto.Unmarshal(body, &req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.burnMillis = append(s.burnMillis, req.GetDurationMs())
|
||||
s.mu.Unlock()
|
||||
|
||||
resp, _ := proto.Marshal(&gluttonpb.BurnCPUResponse{Iterations: 1})
|
||||
_, _ = w.Write(resp)
|
||||
|
||||
case IngestRoute:
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
var req gluttonpb.IngestRequest
|
||||
if err := proto.Unmarshal(body, &req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.ingestSizes = append(s.ingestSizes, int64(len(req.GetPayload())))
|
||||
s.mu.Unlock()
|
||||
|
||||
digest := sha256.Sum256(req.GetPayload())
|
||||
resp, _ := proto.Marshal(&gluttonpb.IngestResponse{
|
||||
Size: int64(len(req.GetPayload())),
|
||||
Sha256: digest[:],
|
||||
})
|
||||
_, _ = w.Write(resp)
|
||||
|
||||
case PingRoute:
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
var req gluttonpb.PingRequest
|
||||
if err := proto.Unmarshal(body, &req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
resp, _ := proto.Marshal(&gluttonpb.PingResponse{Message: req.GetMessage()})
|
||||
_, _ = w.Write(resp)
|
||||
|
||||
default:
|
||||
http.NotFound(w, r)
|
||||
}
|
||||
}
|
||||
|
||||
// RecordedBurnMillis returns each /burncpu request's duration_ms.
|
||||
func (s *Server) RecordedBurnMillis() []int64 {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return append([]int64(nil), s.burnMillis...)
|
||||
}
|
||||
|
||||
// RecordedIngestSizes returns each /ingest request's payload length.
|
||||
func (s *Server) RecordedIngestSizes() []int64 {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return append([]int64(nil), s.ingestSizes...)
|
||||
}
|
||||
|
||||
@@ -41,4 +41,6 @@ const (
|
||||
ReadDiskRoute = "/readdisk"
|
||||
WriteRAMRoute = "/writeram"
|
||||
ReadRAMRoute = "/readram"
|
||||
BurnCPURoute = "/burncpu"
|
||||
IngestRoute = "/ingest"
|
||||
)
|
||||
|
||||
@@ -102,6 +102,8 @@ func newMux(svc *Service) *http.ServeMux {
|
||||
mux.HandleFunc(ReadDiskRoute, protoRoute("ReadDisk", svc.ReadDisk))
|
||||
mux.HandleFunc(WriteRAMRoute, protoRoute("WriteRAM", svc.WriteRAM))
|
||||
mux.HandleFunc(ReadRAMRoute, protoRoute("ReadRAM", svc.ReadRAM))
|
||||
mux.HandleFunc(BurnCPURoute, protoRoute("BurnCPU", svc.BurnCPU))
|
||||
mux.HandleFunc(IngestRoute, protoRoute("Ingest", svc.Ingest))
|
||||
return mux
|
||||
}
|
||||
|
||||
|
||||
@@ -25,7 +25,10 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
goruntime "runtime"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
"google.golang.org/grpc/codes"
|
||||
@@ -334,6 +337,83 @@ func (s *Service) Ping(ctx context.Context, req *gluttonpb.PingRequest) (*glutto
|
||||
return &gluttonpb.PingResponse{Message: req.GetMessage()}, nil
|
||||
}
|
||||
|
||||
// burnBlockBytes sizes the buffer each burner goroutine re-hashes; large
|
||||
// enough that hashing dominates loop overhead, small enough to stay in cache
|
||||
// so the burn is compute-bound rather than memory-bound.
|
||||
const burnBlockBytes = 64 << 10 // 64 KiB
|
||||
|
||||
// BurnCPU spins the requested number of goroutines in a sha256 loop until
|
||||
// the wall-clock duration elapses. The iteration count in the response keeps
|
||||
// the work observable. Deliberately unsynchronized with s.mu: burning must
|
||||
// not block the other RPCs.
|
||||
func (s *Service) BurnCPU(ctx context.Context, req *gluttonpb.BurnCPURequest) (*gluttonpb.BurnCPUResponse, error) {
|
||||
if req.GetDurationMs() < 0 {
|
||||
return nil, status.Error(codes.InvalidArgument, "duration_ms must be non-negative")
|
||||
}
|
||||
parallelism := int(req.GetParallelism())
|
||||
if parallelism < 1 {
|
||||
parallelism = 1
|
||||
}
|
||||
// Cap the goroutine count: past a few per CPU the burn gains nothing,
|
||||
// and an uncapped value lets a single request spawn without bound.
|
||||
if maxPar := goruntime.NumCPU() * 4; parallelism > maxPar {
|
||||
parallelism = maxPar
|
||||
}
|
||||
deadline := time.Now().Add(time.Duration(req.GetDurationMs()) * time.Millisecond)
|
||||
|
||||
var total atomic.Int64
|
||||
var wg sync.WaitGroup
|
||||
for i := 0; i < parallelism; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
block := make([]byte, burnBlockBytes)
|
||||
var n int64
|
||||
for time.Now().Before(deadline) && ctx.Err() == nil {
|
||||
sum := sha256.Sum256(block)
|
||||
copy(block, sum[:])
|
||||
n++
|
||||
}
|
||||
total.Add(n)
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
return &gluttonpb.BurnCPUResponse{Iterations: total.Load()}, nil
|
||||
}
|
||||
|
||||
// Ingest writes the caller-supplied payload to a file under the data dir.
|
||||
// The payload crossed the network to get here, which is the point: WriteDisk
|
||||
// generates its bytes locally, Ingest models a download arriving through the
|
||||
// actor's ingress path before it hits disk.
|
||||
func (s *Service) Ingest(ctx context.Context, req *gluttonpb.IngestRequest) (*gluttonpb.IngestResponse, error) {
|
||||
if !diskKeyRE.MatchString(req.GetKey()) {
|
||||
return nil, status.Errorf(codes.InvalidArgument, "key %q must match %s", req.GetKey(), diskKeyRE)
|
||||
}
|
||||
|
||||
path := filepath.Join(s.dataDir, req.GetKey())
|
||||
flag := os.O_WRONLY | os.O_CREATE | os.O_TRUNC
|
||||
if req.GetAppend() {
|
||||
flag = os.O_WRONLY | os.O_CREATE | os.O_APPEND
|
||||
}
|
||||
f, err := os.OpenFile(path, flag, 0o600)
|
||||
if err != nil {
|
||||
return nil, status.Errorf(codes.Internal, "open %s: %v", path, err)
|
||||
}
|
||||
defer f.Close()
|
||||
|
||||
if _, err := f.Write(req.GetPayload()); err != nil {
|
||||
return nil, status.Errorf(codes.Internal, "write %s: %v", path, err)
|
||||
}
|
||||
info, err := f.Stat()
|
||||
if err != nil {
|
||||
return nil, status.Errorf(codes.Internal, "stat %s: %v", path, err)
|
||||
}
|
||||
|
||||
digest := sha256.Sum256(req.GetPayload())
|
||||
s.diskWriteBytes.Add(ctx, int64(len(req.GetPayload())))
|
||||
return &gluttonpb.IngestResponse{Size: info.Size(), Sha256: digest[:]}, nil
|
||||
}
|
||||
|
||||
func randomBytes(n int) ([]byte, error) {
|
||||
buf := make([]byte, n)
|
||||
if _, err := rand.Read(buf); err != nil {
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
// 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 (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
|
||||
gluttonpb "github.com/agent-substrate/substrate/internal/proto/glutton"
|
||||
)
|
||||
|
||||
// newTestService builds a Service rooted in a per-test temp dir.
|
||||
func newTestService(t *testing.T) *Service {
|
||||
t.Helper()
|
||||
svc, err := New(t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
t.Cleanup(svc.Close)
|
||||
return svc
|
||||
}
|
||||
|
||||
func TestBurnCPU(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
|
||||
start := time.Now()
|
||||
resp, err := svc.BurnCPU(context.Background(), &gluttonpb.BurnCPURequest{
|
||||
DurationMs: 50,
|
||||
Parallelism: 2,
|
||||
})
|
||||
elapsed := time.Since(start)
|
||||
if err != nil {
|
||||
t.Fatalf("BurnCPU: %v", err)
|
||||
}
|
||||
if resp.GetIterations() <= 0 {
|
||||
t.Errorf("BurnCPU iterations = %d, want > 0", resp.GetIterations())
|
||||
}
|
||||
if elapsed < 50*time.Millisecond {
|
||||
t.Errorf("BurnCPU returned after %v, want >= 50ms", elapsed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBurnCPU_ZeroDurationReturnsImmediately(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
resp, err := svc.BurnCPU(context.Background(), &gluttonpb.BurnCPURequest{})
|
||||
if err != nil {
|
||||
t.Fatalf("BurnCPU: %v", err)
|
||||
}
|
||||
if resp.GetIterations() != 0 {
|
||||
t.Errorf("BurnCPU iterations = %d, want 0 for zero duration", resp.GetIterations())
|
||||
}
|
||||
}
|
||||
|
||||
func TestBurnCPU_NegativeDurationRejected(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
_, err := svc.BurnCPU(context.Background(), &gluttonpb.BurnCPURequest{DurationMs: -1})
|
||||
if status.Code(err) != codes.InvalidArgument {
|
||||
t.Errorf("BurnCPU(-1ms) code = %v, want InvalidArgument", status.Code(err))
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngest(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
payload := []byte("pretend this is a repository tarball")
|
||||
|
||||
resp, err := svc.Ingest(context.Background(), &gluttonpb.IngestRequest{
|
||||
Key: "clone",
|
||||
Payload: payload,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Ingest: %v", err)
|
||||
}
|
||||
if resp.GetSize() != int64(len(payload)) {
|
||||
t.Errorf("Ingest size = %d, want %d", resp.GetSize(), len(payload))
|
||||
}
|
||||
want := sha256.Sum256(payload)
|
||||
if !bytes.Equal(resp.GetSha256(), want[:]) {
|
||||
t.Errorf("Ingest sha256 mismatch")
|
||||
}
|
||||
got, err := os.ReadFile(filepath.Join(svc.dataDir, "clone"))
|
||||
if err != nil {
|
||||
t.Fatalf("read ingested file: %v", err)
|
||||
}
|
||||
if !bytes.Equal(got, payload) {
|
||||
t.Errorf("ingested file contents differ from payload")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngest_AppendGrowsFile(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
ctx := context.Background()
|
||||
|
||||
if _, err := svc.Ingest(ctx, &gluttonpb.IngestRequest{Key: "deps", Payload: []byte("aaaa")}); err != nil {
|
||||
t.Fatalf("first Ingest: %v", err)
|
||||
}
|
||||
resp, err := svc.Ingest(ctx, &gluttonpb.IngestRequest{Key: "deps", Payload: []byte("bb"), Append: true})
|
||||
if err != nil {
|
||||
t.Fatalf("append Ingest: %v", err)
|
||||
}
|
||||
if resp.GetSize() != 6 {
|
||||
t.Errorf("appended file size = %d, want 6", resp.GetSize())
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngest_RejectsPathEscape(t *testing.T) {
|
||||
svc := newTestService(t)
|
||||
_, err := svc.Ingest(context.Background(), &gluttonpb.IngestRequest{
|
||||
Key: "../escape",
|
||||
Payload: []byte("x"),
|
||||
})
|
||||
if status.Code(err) != codes.InvalidArgument {
|
||||
t.Errorf("Ingest(../escape) code = %v, want InvalidArgument", status.Code(err))
|
||||
}
|
||||
}
|
||||
@@ -881,6 +881,223 @@ func (x *Peer) GetDelayMs() int32 {
|
||||
return 0
|
||||
}
|
||||
|
||||
type BurnCPURequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
// Wall-clock time to keep burning, in milliseconds.
|
||||
DurationMs int64 `protobuf:"varint,1,opt,name=duration_ms,json=durationMs,proto3" json:"duration_ms,omitempty"`
|
||||
// Number of goroutines burning concurrently; values < 1 read as 1.
|
||||
Parallelism int32 `protobuf:"varint,2,opt,name=parallelism,proto3" json:"parallelism,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *BurnCPURequest) Reset() {
|
||||
*x = BurnCPURequest{}
|
||||
mi := &file_glutton_proto_msgTypes[15]
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
|
||||
func (x *BurnCPURequest) String() string {
|
||||
return protoimpl.X.MessageStringOf(x)
|
||||
}
|
||||
|
||||
func (*BurnCPURequest) ProtoMessage() {}
|
||||
|
||||
func (x *BurnCPURequest) ProtoReflect() protoreflect.Message {
|
||||
mi := &file_glutton_proto_msgTypes[15]
|
||||
if x != nil {
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
if ms.LoadMessageInfo() == nil {
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
return ms
|
||||
}
|
||||
return mi.MessageOf(x)
|
||||
}
|
||||
|
||||
// Deprecated: Use BurnCPURequest.ProtoReflect.Descriptor instead.
|
||||
func (*BurnCPURequest) Descriptor() ([]byte, []int) {
|
||||
return file_glutton_proto_rawDescGZIP(), []int{15}
|
||||
}
|
||||
|
||||
func (x *BurnCPURequest) GetDurationMs() int64 {
|
||||
if x != nil {
|
||||
return x.DurationMs
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *BurnCPURequest) GetParallelism() int32 {
|
||||
if x != nil {
|
||||
return x.Parallelism
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
type BurnCPUResponse struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
// Total hash iterations completed across all goroutines, so the burn
|
||||
// is observable and cannot be optimized away.
|
||||
Iterations int64 `protobuf:"varint,1,opt,name=iterations,proto3" json:"iterations,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *BurnCPUResponse) Reset() {
|
||||
*x = BurnCPUResponse{}
|
||||
mi := &file_glutton_proto_msgTypes[16]
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
|
||||
func (x *BurnCPUResponse) String() string {
|
||||
return protoimpl.X.MessageStringOf(x)
|
||||
}
|
||||
|
||||
func (*BurnCPUResponse) ProtoMessage() {}
|
||||
|
||||
func (x *BurnCPUResponse) ProtoReflect() protoreflect.Message {
|
||||
mi := &file_glutton_proto_msgTypes[16]
|
||||
if x != nil {
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
if ms.LoadMessageInfo() == nil {
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
return ms
|
||||
}
|
||||
return mi.MessageOf(x)
|
||||
}
|
||||
|
||||
// Deprecated: Use BurnCPUResponse.ProtoReflect.Descriptor instead.
|
||||
func (*BurnCPUResponse) Descriptor() ([]byte, []int) {
|
||||
return file_glutton_proto_rawDescGZIP(), []int{16}
|
||||
}
|
||||
|
||||
func (x *BurnCPUResponse) GetIterations() int64 {
|
||||
if x != nil {
|
||||
return x.Iterations
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
type IngestRequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
// name of the file the payload is written to
|
||||
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"`
|
||||
// the bytes to write, carried in the request body over the network
|
||||
Payload []byte `protobuf:"bytes,2,opt,name=payload,proto3" json:"payload,omitempty"`
|
||||
// append to the file instead of truncating it
|
||||
Append bool `protobuf:"varint,3,opt,name=append,proto3" json:"append,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *IngestRequest) Reset() {
|
||||
*x = IngestRequest{}
|
||||
mi := &file_glutton_proto_msgTypes[17]
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
|
||||
func (x *IngestRequest) String() string {
|
||||
return protoimpl.X.MessageStringOf(x)
|
||||
}
|
||||
|
||||
func (*IngestRequest) ProtoMessage() {}
|
||||
|
||||
func (x *IngestRequest) ProtoReflect() protoreflect.Message {
|
||||
mi := &file_glutton_proto_msgTypes[17]
|
||||
if x != nil {
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
if ms.LoadMessageInfo() == nil {
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
return ms
|
||||
}
|
||||
return mi.MessageOf(x)
|
||||
}
|
||||
|
||||
// Deprecated: Use IngestRequest.ProtoReflect.Descriptor instead.
|
||||
func (*IngestRequest) Descriptor() ([]byte, []int) {
|
||||
return file_glutton_proto_rawDescGZIP(), []int{17}
|
||||
}
|
||||
|
||||
func (x *IngestRequest) GetKey() string {
|
||||
if x != nil {
|
||||
return x.Key
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *IngestRequest) GetPayload() []byte {
|
||||
if x != nil {
|
||||
return x.Payload
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (x *IngestRequest) GetAppend() bool {
|
||||
if x != nil {
|
||||
return x.Append
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
type IngestResponse struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
// size of the file after the write
|
||||
Size int64 `protobuf:"varint,1,opt,name=size,proto3" json:"size,omitempty"`
|
||||
// sha256 of the payload as received, so the caller can verify transport
|
||||
Sha256 []byte `protobuf:"bytes,2,opt,name=sha256,proto3" json:"sha256,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *IngestResponse) Reset() {
|
||||
*x = IngestResponse{}
|
||||
mi := &file_glutton_proto_msgTypes[18]
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
|
||||
func (x *IngestResponse) String() string {
|
||||
return protoimpl.X.MessageStringOf(x)
|
||||
}
|
||||
|
||||
func (*IngestResponse) ProtoMessage() {}
|
||||
|
||||
func (x *IngestResponse) ProtoReflect() protoreflect.Message {
|
||||
mi := &file_glutton_proto_msgTypes[18]
|
||||
if x != nil {
|
||||
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
|
||||
if ms.LoadMessageInfo() == nil {
|
||||
ms.StoreMessageInfo(mi)
|
||||
}
|
||||
return ms
|
||||
}
|
||||
return mi.MessageOf(x)
|
||||
}
|
||||
|
||||
// Deprecated: Use IngestResponse.ProtoReflect.Descriptor instead.
|
||||
func (*IngestResponse) Descriptor() ([]byte, []int) {
|
||||
return file_glutton_proto_rawDescGZIP(), []int{18}
|
||||
}
|
||||
|
||||
func (x *IngestResponse) GetSize() int64 {
|
||||
if x != nil {
|
||||
return x.Size
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *IngestResponse) GetSha256() []byte {
|
||||
if x != nil {
|
||||
return x.Sha256
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
var File_glutton_proto protoreflect.FileDescriptor
|
||||
|
||||
const file_glutton_proto_rawDesc = "" +
|
||||
@@ -925,14 +1142,29 @@ const file_glutton_proto_rawDesc = "" +
|
||||
"\x0eGossipResponse\"5\n" +
|
||||
"\x04Peer\x12\x12\n" +
|
||||
"\x04host\x18\x01 \x01(\tR\x04host\x12\x19\n" +
|
||||
"\bdelay_ms\x18\x02 \x01(\x05R\adelayMs*_\n" +
|
||||
"\bdelay_ms\x18\x02 \x01(\x05R\adelayMs\"S\n" +
|
||||
"\x0eBurnCPURequest\x12\x1f\n" +
|
||||
"\vduration_ms\x18\x01 \x01(\x03R\n" +
|
||||
"durationMs\x12 \n" +
|
||||
"\vparallelism\x18\x02 \x01(\x05R\vparallelism\"1\n" +
|
||||
"\x0fBurnCPUResponse\x12\x1e\n" +
|
||||
"\n" +
|
||||
"iterations\x18\x01 \x01(\x03R\n" +
|
||||
"iterations\"S\n" +
|
||||
"\rIngestRequest\x12\x10\n" +
|
||||
"\x03key\x18\x01 \x01(\tR\x03key\x12\x18\n" +
|
||||
"\apayload\x18\x02 \x01(\fR\apayload\x12\x16\n" +
|
||||
"\x06append\x18\x03 \x01(\bR\x06append\"<\n" +
|
||||
"\x0eIngestResponse\x12\x12\n" +
|
||||
"\x04size\x18\x01 \x01(\x03R\x04size\x12\x16\n" +
|
||||
"\x06sha256\x18\x02 \x01(\fR\x06sha256*_\n" +
|
||||
"\tWriteMode\x12\x17\n" +
|
||||
"\x13WRITE_MODE_TRUNCATE\x10\x00\x12\x18\n" +
|
||||
"\x14WRITE_MODE_OVERWRITE\x10\x01\x12\x1f\n" +
|
||||
"\x1bWRITE_MODE_OVERWRITE_ROTATE\x10\x02*9\n" +
|
||||
"\bReadMode\x12\x12\n" +
|
||||
"\x0eREAD_MODE_DATA\x10\x00\x12\x19\n" +
|
||||
"\x15READ_MODE_DIGEST_ONLY\x10\x012\xc6\x03\n" +
|
||||
"\x15READ_MODE_DIGEST_ONLY\x10\x012\xc3\x04\n" +
|
||||
"\aGlutton\x12A\n" +
|
||||
"\bWriteRAM\x12\x18.glutton.WriteRAMRequest\x1a\x19.glutton.WriteRAMResponse\"\x00\x12>\n" +
|
||||
"\aReadRAM\x12\x17.glutton.ReadRAMRequest\x1a\x18.glutton.ReadRAMResponse\"\x00\x12D\n" +
|
||||
@@ -940,7 +1172,9 @@ const file_glutton_proto_rawDesc = "" +
|
||||
"\bReadDisk\x12\x18.glutton.ReadDiskRequest\x1a\x19.glutton.ReadDiskResponse\"\x00\x12;\n" +
|
||||
"\x06OpenFD\x12\x16.glutton.OpenFDRequest\x1a\x17.glutton.OpenFDResponse\"\x00\x125\n" +
|
||||
"\x04Ping\x12\x14.glutton.PingRequest\x1a\x15.glutton.PingResponse\"\x00\x12;\n" +
|
||||
"\x06Gossip\x12\x16.glutton.GossipRequest\x1a\x17.glutton.GossipResponse\"\x00B=Z;github.com/agent-substrate/substrate/internal/proto/gluttonb\x06proto3"
|
||||
"\x06Gossip\x12\x16.glutton.GossipRequest\x1a\x17.glutton.GossipResponse\"\x00\x12>\n" +
|
||||
"\aBurnCPU\x12\x17.glutton.BurnCPURequest\x1a\x18.glutton.BurnCPUResponse\"\x00\x12;\n" +
|
||||
"\x06Ingest\x12\x16.glutton.IngestRequest\x1a\x17.glutton.IngestResponse\"\x00B=Z;github.com/agent-substrate/substrate/internal/proto/gluttonb\x06proto3"
|
||||
|
||||
var (
|
||||
file_glutton_proto_rawDescOnce sync.Once
|
||||
@@ -955,7 +1189,7 @@ func file_glutton_proto_rawDescGZIP() []byte {
|
||||
}
|
||||
|
||||
var file_glutton_proto_enumTypes = make([]protoimpl.EnumInfo, 2)
|
||||
var file_glutton_proto_msgTypes = make([]protoimpl.MessageInfo, 15)
|
||||
var file_glutton_proto_msgTypes = make([]protoimpl.MessageInfo, 19)
|
||||
var file_glutton_proto_goTypes = []any{
|
||||
(WriteMode)(0), // 0: glutton.WriteMode
|
||||
(ReadMode)(0), // 1: glutton.ReadMode
|
||||
@@ -974,6 +1208,10 @@ var file_glutton_proto_goTypes = []any{
|
||||
(*GossipRequest)(nil), // 14: glutton.GossipRequest
|
||||
(*GossipResponse)(nil), // 15: glutton.GossipResponse
|
||||
(*Peer)(nil), // 16: glutton.Peer
|
||||
(*BurnCPURequest)(nil), // 17: glutton.BurnCPURequest
|
||||
(*BurnCPUResponse)(nil), // 18: glutton.BurnCPUResponse
|
||||
(*IngestRequest)(nil), // 19: glutton.IngestRequest
|
||||
(*IngestResponse)(nil), // 20: glutton.IngestResponse
|
||||
}
|
||||
var file_glutton_proto_depIdxs = []int32{
|
||||
0, // 0: glutton.WriteRAMRequest.write_mode:type_name -> glutton.WriteMode
|
||||
@@ -987,15 +1225,19 @@ var file_glutton_proto_depIdxs = []int32{
|
||||
10, // 8: glutton.Glutton.OpenFD:input_type -> glutton.OpenFDRequest
|
||||
12, // 9: glutton.Glutton.Ping:input_type -> glutton.PingRequest
|
||||
14, // 10: glutton.Glutton.Gossip:input_type -> glutton.GossipRequest
|
||||
3, // 11: glutton.Glutton.WriteRAM:output_type -> glutton.WriteRAMResponse
|
||||
5, // 12: glutton.Glutton.ReadRAM:output_type -> glutton.ReadRAMResponse
|
||||
7, // 13: glutton.Glutton.WriteDisk:output_type -> glutton.WriteDiskResponse
|
||||
9, // 14: glutton.Glutton.ReadDisk:output_type -> glutton.ReadDiskResponse
|
||||
11, // 15: glutton.Glutton.OpenFD:output_type -> glutton.OpenFDResponse
|
||||
13, // 16: glutton.Glutton.Ping:output_type -> glutton.PingResponse
|
||||
15, // 17: glutton.Glutton.Gossip:output_type -> glutton.GossipResponse
|
||||
11, // [11:18] is the sub-list for method output_type
|
||||
4, // [4:11] is the sub-list for method input_type
|
||||
17, // 11: glutton.Glutton.BurnCPU:input_type -> glutton.BurnCPURequest
|
||||
19, // 12: glutton.Glutton.Ingest:input_type -> glutton.IngestRequest
|
||||
3, // 13: glutton.Glutton.WriteRAM:output_type -> glutton.WriteRAMResponse
|
||||
5, // 14: glutton.Glutton.ReadRAM:output_type -> glutton.ReadRAMResponse
|
||||
7, // 15: glutton.Glutton.WriteDisk:output_type -> glutton.WriteDiskResponse
|
||||
9, // 16: glutton.Glutton.ReadDisk:output_type -> glutton.ReadDiskResponse
|
||||
11, // 17: glutton.Glutton.OpenFD:output_type -> glutton.OpenFDResponse
|
||||
13, // 18: glutton.Glutton.Ping:output_type -> glutton.PingResponse
|
||||
15, // 19: glutton.Glutton.Gossip:output_type -> glutton.GossipResponse
|
||||
18, // 20: glutton.Glutton.BurnCPU:output_type -> glutton.BurnCPUResponse
|
||||
20, // 21: glutton.Glutton.Ingest:output_type -> glutton.IngestResponse
|
||||
13, // [13:22] is the sub-list for method output_type
|
||||
4, // [4:13] is the sub-list for method input_type
|
||||
4, // [4:4] is the sub-list for extension type_name
|
||||
4, // [4:4] is the sub-list for extension extendee
|
||||
0, // [0:4] is the sub-list for field type_name
|
||||
@@ -1012,7 +1254,7 @@ func file_glutton_proto_init() {
|
||||
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
|
||||
RawDescriptor: unsafe.Slice(unsafe.StringData(file_glutton_proto_rawDesc), len(file_glutton_proto_rawDesc)),
|
||||
NumEnums: 2,
|
||||
NumMessages: 15,
|
||||
NumMessages: 19,
|
||||
NumExtensions: 0,
|
||||
NumServices: 1,
|
||||
},
|
||||
|
||||
@@ -50,6 +50,17 @@ service Glutton {
|
||||
// Tells the glutton to send network traffic to a peer glutton.
|
||||
// messages will be sent on regular intervals separated by delay_ms.
|
||||
rpc Gossip(GossipRequest) returns (GossipResponse) {}
|
||||
|
||||
// Tells glutton to spin CPU for a wall-clock duration, hashing in a
|
||||
// tight loop across the requested number of goroutines. Models
|
||||
// compute-bound work such as compiling or running a test suite.
|
||||
rpc BurnCPU(BurnCPURequest) returns (BurnCPUResponse) {}
|
||||
|
||||
// Carries caller-supplied bytes over the wire into the glutton, which
|
||||
// writes them to a file under its data dir. Unlike WriteDisk, the bytes
|
||||
// cross the network, so the request models a download (git clone,
|
||||
// dependency fetch) arriving through the actor's ingress path.
|
||||
rpc Ingest(IngestRequest) returns (IngestResponse) {}
|
||||
}
|
||||
|
||||
enum WriteMode {
|
||||
@@ -169,3 +180,36 @@ message Peer {
|
||||
string host = 1;
|
||||
int32 delay_ms = 2;
|
||||
}
|
||||
|
||||
message BurnCPURequest {
|
||||
// Wall-clock time to keep burning, in milliseconds.
|
||||
int64 duration_ms = 1;
|
||||
|
||||
// Number of goroutines burning concurrently; values < 1 read as 1.
|
||||
int32 parallelism = 2;
|
||||
}
|
||||
|
||||
message BurnCPUResponse {
|
||||
// Total hash iterations completed across all goroutines, so the burn
|
||||
// is observable and cannot be optimized away.
|
||||
int64 iterations = 1;
|
||||
}
|
||||
|
||||
message IngestRequest {
|
||||
// name of the file the payload is written to
|
||||
string key = 1;
|
||||
|
||||
// the bytes to write, carried in the request body over the network
|
||||
bytes payload = 2;
|
||||
|
||||
// append to the file instead of truncating it
|
||||
bool append = 3;
|
||||
}
|
||||
|
||||
message IngestResponse {
|
||||
// size of the file after the write
|
||||
int64 size = 1;
|
||||
|
||||
// sha256 of the payload as received, so the caller can verify transport
|
||||
bytes sha256 = 2;
|
||||
}
|
||||
|
||||
@@ -40,6 +40,8 @@ const (
|
||||
Glutton_OpenFD_FullMethodName = "/glutton.Glutton/OpenFD"
|
||||
Glutton_Ping_FullMethodName = "/glutton.Glutton/Ping"
|
||||
Glutton_Gossip_FullMethodName = "/glutton.Glutton/Gossip"
|
||||
Glutton_BurnCPU_FullMethodName = "/glutton.Glutton/BurnCPU"
|
||||
Glutton_Ingest_FullMethodName = "/glutton.Glutton/Ingest"
|
||||
)
|
||||
|
||||
// GluttonClient is the client API for Glutton service.
|
||||
@@ -72,6 +74,15 @@ type GluttonClient interface {
|
||||
// Tells the glutton to send network traffic to a peer glutton.
|
||||
// messages will be sent on regular intervals separated by delay_ms.
|
||||
Gossip(ctx context.Context, in *GossipRequest, opts ...grpc.CallOption) (*GossipResponse, error)
|
||||
// Tells glutton to spin CPU for a wall-clock duration, hashing in a
|
||||
// tight loop across the requested number of goroutines. Models
|
||||
// compute-bound work such as compiling or running a test suite.
|
||||
BurnCPU(ctx context.Context, in *BurnCPURequest, opts ...grpc.CallOption) (*BurnCPUResponse, error)
|
||||
// Carries caller-supplied bytes over the wire into the glutton, which
|
||||
// writes them to a file under its data dir. Unlike WriteDisk, the bytes
|
||||
// cross the network, so the request models a download (git clone,
|
||||
// dependency fetch) arriving through the actor's ingress path.
|
||||
Ingest(ctx context.Context, in *IngestRequest, opts ...grpc.CallOption) (*IngestResponse, error)
|
||||
}
|
||||
|
||||
type gluttonClient struct {
|
||||
@@ -152,6 +163,26 @@ func (c *gluttonClient) Gossip(ctx context.Context, in *GossipRequest, opts ...g
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (c *gluttonClient) BurnCPU(ctx context.Context, in *BurnCPURequest, opts ...grpc.CallOption) (*BurnCPUResponse, error) {
|
||||
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
|
||||
out := new(BurnCPUResponse)
|
||||
err := c.cc.Invoke(ctx, Glutton_BurnCPU_FullMethodName, in, out, cOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (c *gluttonClient) Ingest(ctx context.Context, in *IngestRequest, opts ...grpc.CallOption) (*IngestResponse, error) {
|
||||
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
|
||||
out := new(IngestResponse)
|
||||
err := c.cc.Invoke(ctx, Glutton_Ingest_FullMethodName, in, out, cOpts...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// GluttonServer is the server API for Glutton service.
|
||||
// All implementations must embed UnimplementedGluttonServer
|
||||
// for forward compatibility.
|
||||
@@ -182,6 +213,15 @@ type GluttonServer interface {
|
||||
// Tells the glutton to send network traffic to a peer glutton.
|
||||
// messages will be sent on regular intervals separated by delay_ms.
|
||||
Gossip(context.Context, *GossipRequest) (*GossipResponse, error)
|
||||
// Tells glutton to spin CPU for a wall-clock duration, hashing in a
|
||||
// tight loop across the requested number of goroutines. Models
|
||||
// compute-bound work such as compiling or running a test suite.
|
||||
BurnCPU(context.Context, *BurnCPURequest) (*BurnCPUResponse, error)
|
||||
// Carries caller-supplied bytes over the wire into the glutton, which
|
||||
// writes them to a file under its data dir. Unlike WriteDisk, the bytes
|
||||
// cross the network, so the request models a download (git clone,
|
||||
// dependency fetch) arriving through the actor's ingress path.
|
||||
Ingest(context.Context, *IngestRequest) (*IngestResponse, error)
|
||||
mustEmbedUnimplementedGluttonServer()
|
||||
}
|
||||
|
||||
@@ -213,6 +253,12 @@ func (UnimplementedGluttonServer) Ping(context.Context, *PingRequest) (*PingResp
|
||||
func (UnimplementedGluttonServer) Gossip(context.Context, *GossipRequest) (*GossipResponse, error) {
|
||||
return nil, status.Error(codes.Unimplemented, "method Gossip not implemented")
|
||||
}
|
||||
func (UnimplementedGluttonServer) BurnCPU(context.Context, *BurnCPURequest) (*BurnCPUResponse, error) {
|
||||
return nil, status.Error(codes.Unimplemented, "method BurnCPU not implemented")
|
||||
}
|
||||
func (UnimplementedGluttonServer) Ingest(context.Context, *IngestRequest) (*IngestResponse, error) {
|
||||
return nil, status.Error(codes.Unimplemented, "method Ingest not implemented")
|
||||
}
|
||||
func (UnimplementedGluttonServer) mustEmbedUnimplementedGluttonServer() {}
|
||||
func (UnimplementedGluttonServer) testEmbeddedByValue() {}
|
||||
|
||||
@@ -360,6 +406,42 @@ func _Glutton_Gossip_Handler(srv interface{}, ctx context.Context, dec func(inte
|
||||
return interceptor(ctx, in, info, handler)
|
||||
}
|
||||
|
||||
func _Glutton_BurnCPU_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
|
||||
in := new(BurnCPURequest)
|
||||
if err := dec(in); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if interceptor == nil {
|
||||
return srv.(GluttonServer).BurnCPU(ctx, in)
|
||||
}
|
||||
info := &grpc.UnaryServerInfo{
|
||||
Server: srv,
|
||||
FullMethod: Glutton_BurnCPU_FullMethodName,
|
||||
}
|
||||
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
|
||||
return srv.(GluttonServer).BurnCPU(ctx, req.(*BurnCPURequest))
|
||||
}
|
||||
return interceptor(ctx, in, info, handler)
|
||||
}
|
||||
|
||||
func _Glutton_Ingest_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
|
||||
in := new(IngestRequest)
|
||||
if err := dec(in); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if interceptor == nil {
|
||||
return srv.(GluttonServer).Ingest(ctx, in)
|
||||
}
|
||||
info := &grpc.UnaryServerInfo{
|
||||
Server: srv,
|
||||
FullMethod: Glutton_Ingest_FullMethodName,
|
||||
}
|
||||
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
|
||||
return srv.(GluttonServer).Ingest(ctx, req.(*IngestRequest))
|
||||
}
|
||||
return interceptor(ctx, in, info, handler)
|
||||
}
|
||||
|
||||
// Glutton_ServiceDesc is the grpc.ServiceDesc for Glutton service.
|
||||
// It's only intended for direct use with grpc.RegisterService,
|
||||
// and not to be introspected or modified (even as a copy)
|
||||
@@ -395,6 +477,14 @@ var Glutton_ServiceDesc = grpc.ServiceDesc{
|
||||
MethodName: "Gossip",
|
||||
Handler: _Glutton_Gossip_Handler,
|
||||
},
|
||||
{
|
||||
MethodName: "BurnCPU",
|
||||
Handler: _Glutton_BurnCPU_Handler,
|
||||
},
|
||||
{
|
||||
MethodName: "Ingest",
|
||||
Handler: _Glutton_Ingest_Handler,
|
||||
},
|
||||
},
|
||||
Streams: []grpc.StreamDesc{},
|
||||
Metadata: "glutton.proto",
|
||||
|
||||
Reference in New Issue
Block a user