Skip to content

Commit bc9e76b

Browse files
sorenbsclaude
andcommitted
multi-generator ladder: K gens on distinct source IPs (edge per-IP connection cap) — gen slices via BENCH_CERT_SUB_OFFSET, one writer (BENCH_CERT_WRITERS), subscriber-only gens wait for creates; ladder WC_GENS x WC_GEN_REGIONS round-robin with per-k services, URLs and harvest
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 20277f1 commit bc9e76b

2 files changed

Lines changed: 80 additions & 23 deletions

File tree

bench/awsbench/src/main.rs

Lines changed: 33 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1279,6 +1279,13 @@ async fn run_cert(args: &Args, stats: Arc<Stats>) -> anyhow::Result<()> {
12791279
let tenants: usize = wide_env("BENCH_CERT_TENANTS", 10_000);
12801280
let sub_tenants: usize = wide_env("BENCH_CERT_SUB_TENANTS", 1_000);
12811281
let subs_n: usize = wide_env("BENCH_CERT_SUBS_N", 0);
1282+
// Multi-generator rungs (the edge caps concurrent connections per
1283+
// SOURCE IP at ~1.2-2k): K gens on distinct hosts each park a
1284+
// disjoint slice of subscribers; exactly one of them writes.
1285+
let sub_offset: usize = wide_env("BENCH_CERT_SUB_OFFSET", 0);
1286+
let writers: bool = std::env::var("BENCH_CERT_WRITERS")
1287+
.map(|v| v != "0")
1288+
.unwrap_or(true);
12821289
let active: usize = wide_env("BENCH_CERT_ACTIVE", 100);
12831290
let fanout_active: usize = wide_env("BENCH_CERT_FANOUT_ACTIVE", 10);
12841291
let window_ms: u64 = wide_env("BENCH_CERT_WINDOW_MS", 5_000);
@@ -1406,9 +1413,28 @@ async fn run_cert(args: &Args, stats: Arc<Stats>) -> anyhow::Result<()> {
14061413
}
14071414
};
14081415

1409-
// Phase 1: create every stream the run can touch.
1410-
eprintln!("CERT: creating {tenants} streams (conc {setup_conc})");
1416+
// Phase 1: create every stream the run can touch — the WRITER gen
1417+
// only; subscriber-only gens wait until the writer's creates land.
14111418
let t_create = Instant::now();
1419+
if !writers {
1420+
let t0 = sub_offset % sub_tenants.max(1);
1421+
let probe = format!("{base}/v1/streams/{}", stream_of(t0));
1422+
eprintln!("CERT: subscriber-only gen (offset {sub_offset}): waiting for {probe}");
1423+
for _ in 0..300 {
1424+
let r = http
1425+
.get(&probe)
1426+
.timeout(Duration::from_secs(15))
1427+
.header("authorization", format!("Bearer {}", tokens[t0]))
1428+
.header("prisma-encryption-key", key.clone())
1429+
.send()
1430+
.await;
1431+
if matches!(&r, Ok(resp) if resp.status().is_success()) {
1432+
break;
1433+
}
1434+
tokio::time::sleep(Duration::from_secs(2)).await;
1435+
}
1436+
} else {
1437+
eprintln!("CERT: creating {tenants} streams (conc {setup_conc})");
14121438
for chunk in (0..tenants).collect::<Vec<_>>().chunks(10_000) {
14131439
let fails: usize = futures_util::stream::iter(chunk.iter().map(|&t| {
14141440
let http = http.clone();
@@ -1439,6 +1465,7 @@ async fn run_cert(args: &Args, stats: Arc<Stats>) -> anyhow::Result<()> {
14391465
.await;
14401466
anyhow::ensure!(fails == 0, "{fails} creates failed");
14411467
}
1468+
}
14421469
let create_ms = t_create.elapsed().as_millis();
14431470
eprintln!("CERT: creates done in {create_ms}ms");
14441471

@@ -1452,7 +1479,7 @@ async fn run_cert(args: &Args, stats: Arc<Stats>) -> anyhow::Result<()> {
14521479
let lag_hist = WideHist::new();
14531480
let stop = Arc::new(AtomicU64::new(0));
14541481
for j in 0..subs_n {
1455-
let t = j % sub_tenants;
1482+
let t = (j + sub_offset) % sub_tenants;
14561483
let http = http.clone();
14571484
let base = base.clone();
14581485
let name = stream_of(t);
@@ -1616,6 +1643,9 @@ async fn run_cert(args: &Args, stats: Arc<Stats>) -> anyhow::Result<()> {
16161643
}
16171644
};
16181645
tokio::spawn(async move {
1646+
if !writers {
1647+
return; // subscriber-only gen: no rotation writer
1648+
}
16191649
let mut window: u64 = 0;
16201650
while Instant::now() < deadline && stop.load(Ordering::Relaxed) == 0 {
16211651
// Deterministic active set for this window.

bench/soak/wc-ladder.sh

Lines changed: 47 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,12 @@ R=${WC_REGION:-eu-central-1}
2222
# run their gen OUT-of-region (WC_GEN_REGION) until the platform fixes
2323
# it. Server stays in R.
2424
GR=${WC_GEN_REGION:-$R}
25+
# K generators on distinct source IPs (the edge caps concurrent
26+
# connections PER SOURCE IP at ~1.2-2k): gen k parks subscribers
27+
# [k*slice, (k+1)*slice) and only gen 0 writes. Regions round-robin
28+
# from WC_GEN_REGIONS so the gens really are distinct hosts/IPs.
29+
GENS=${WC_GENS:-1}
30+
GEN_REGIONS=(${WC_GEN_REGIONS:-$GR})
2531
export SOAK_RUN_ID=${SOAK_RUN_ID:-"wc-$(date -u +%Y%m%dT%H%M%SZ)"}
2632
export BIN_TAG=$SOAK_RUN_ID
2733
OUT="$S/results/$SOAK_RUN_ID"; mkdir -p "$OUT"
@@ -87,7 +93,13 @@ for path, key in [(sys.argv[1], sys.argv[2]), (sys.argv[3], sys.argv[4])]:
8793
print(f"uploaded {key} ({len(data)} bytes)")
8894
PY
8995

90-
role_region() { [ "$1" = gen ] && echo "$GR" || echo "$R"; }
96+
role_region() {
97+
case "$1" in
98+
gen) echo "${GEN_REGIONS[0]}";;
99+
gen-*) local k=${1#gen-}; echo "${GEN_REGIONS[$((k % ${#GEN_REGIONS[@]}))]}";;
100+
*) echo "$R";;
101+
esac
102+
}
91103
role_project() { cat "$S/proj-$(role_region "$1").txt"; }
92104
svc_id() {
93105
local ROLE=$1 DR=$(role_region "$1") DP=$(role_project "$1")
@@ -108,7 +120,7 @@ svc_id() {
108120
}
109121
deploy() {
110122
local ROLE=$1; shift
111-
local DIR="$S/app-$ROLE-$R"
123+
local DIR="$S/app-${ROLE%%-*}-$R" # gen-k roles share the gen app dir
112124
local DR=$(role_region "$ROLE") DP=$(role_project "$ROLE")
113125
local SVC=$(svc_id "$ROLE"); local SVCARG=()
114126
[ -n "$SVC" ] && SVCARG=(--service "$SVC")
@@ -180,25 +192,35 @@ verify_server_live
180192
curl -sf --max-time 20 -H "authorization: Bearer $AUTH" \
181193
"$(cat "$S/url-server-$R.txt")/v1/debug/load" > "$OUT/server-before-$STAGE.json" || true
182194

183-
echo "== gen: cert stage $STAGE"
184-
deploy gen \
185-
--env AWSBENCH_S3_KEY="bin/awsbench-$BIN_TAG-x64" \
186-
--env S3_ENDPOINT=$BINEP --env S3_BUCKET=$BINBKT --env S3_REGION=auto \
187-
--env S3_ACCESS_KEY_ID="$BINID" --env S3_SECRET_ACCESS_KEY="$BINSEC" \
188-
--env BENCH_SYSTEM=prisma --env BENCH_SHAPE=cert \
189-
--env BENCH_TARGET="$(cat "$S/url-server-$R.txt")" \
190-
--env TOKENS_S3_KEY="wc/$SOAK_RUN_ID/tokens.json" \
191-
--env BENCH_CERT_TENANTS="$TENANTS" --env BENCH_CERT_SUB_TENANTS="$SUB_TENANTS" \
192-
--env BENCH_CERT_SUBS_N="$SUBS_N" --env BENCH_CERT_ACTIVE="$ACTIVE" \
193-
--env BENCH_CERT_FANOUT_ACTIVE="$FANOUT" --env BENCH_CERT_WINDOW_MS=5000 \
194-
--env BENCH_CERT_CONNECT_CONC=${WC_CONNECT_CONC:-48} \
195-
--env BENCH_CERT_WPS="$WPS" --env BENCH_CERT_SECS="$SECS" \
196-
--env BENCH_WIDE_SETUP_CONC=64 --env BENCH_RECORD_BYTES=1024 --env BENCH_BATCH=1 \
197-
--env BENCH_HOLD=1 --env BENCH_START_GATED=false \
198-
--env AUTH_TOKEN="$AUTH" --env STREAM_KEY="$KEY" \
199-
--env BENCH_STREAM="wc$STAGE-" --env KEEP_AWAKE=1
195+
echo "== gen: cert stage $STAGE (gens=$GENS regions=${GEN_REGIONS[*]})"
196+
SLICE=$(( (SUBS_N + GENS - 1) / GENS ))
197+
for k in $(seq 0 $((GENS - 1))); do
198+
ROLE=gen; [ "$k" -gt 0 ] && ROLE="gen-$k"
199+
OFFSET=$(( k * SLICE )); N=$SLICE
200+
[ $(( OFFSET + N )) -gt "$SUBS_N" ] && N=$(( SUBS_N - OFFSET ))
201+
[ "$N" -lt 0 ] && N=0
202+
W=1; [ "$k" -gt 0 ] && W=0
203+
echo " gen $k ($ROLE @ $(role_region "$ROLE")): subs $N from $OFFSET, writers=$W"
204+
deploy "$ROLE" \
205+
--env AWSBENCH_S3_KEY="bin/awsbench-$BIN_TAG-x64" \
206+
--env S3_ENDPOINT=$BINEP --env S3_BUCKET=$BINBKT --env S3_REGION=auto \
207+
--env S3_ACCESS_KEY_ID="$BINID" --env S3_SECRET_ACCESS_KEY="$BINSEC" \
208+
--env BENCH_SYSTEM=prisma --env BENCH_SHAPE=cert \
209+
--env BENCH_TARGET="$(cat "$S/url-server-$R.txt")" \
210+
--env TOKENS_S3_KEY="wc/$SOAK_RUN_ID/tokens.json" \
211+
--env BENCH_CERT_TENANTS="$TENANTS" --env BENCH_CERT_SUB_TENANTS="$SUB_TENANTS" \
212+
--env BENCH_CERT_SUBS_N="$N" --env BENCH_CERT_SUB_OFFSET="$OFFSET" --env BENCH_CERT_WRITERS="$W" \
213+
--env BENCH_CERT_ACTIVE="$ACTIVE" \
214+
--env BENCH_CERT_FANOUT_ACTIVE="$FANOUT" --env BENCH_CERT_WINDOW_MS=5000 \
215+
--env BENCH_CERT_CONNECT_CONC=${WC_CONNECT_CONC:-48} \
216+
--env BENCH_CERT_WPS="$WPS" --env BENCH_CERT_SECS="$SECS" \
217+
--env BENCH_WIDE_SETUP_CONC=64 --env BENCH_RECORD_BYTES=1024 --env BENCH_BATCH=1 \
218+
--env BENCH_HOLD=1 --env BENCH_START_GATED=false \
219+
--env AUTH_TOKEN="$AUTH" --env STREAM_KEY="$KEY" \
220+
--env BENCH_STREAM="wc$STAGE-" --env KEEP_AWAKE=1
221+
done
200222

201-
GURL=$(cat "$S/url-gen-$GR.txt")
223+
GURL=$(cat "$S/url-gen-$(role_region gen).txt") # the writer gen decides cert_done
202224
SURL=$(cat "$S/url-server-$R.txt")
203225
DEADLINE=$(( $(date +%s) + SECS + 1500 ))
204226
# RSS timeline: correlate memory peaks with shed windows.
@@ -215,6 +237,11 @@ while [ "$(date +%s)" -lt "$DEADLINE" ]; do
215237
BODY=$(curl -sf --max-time 15 "$GURL/" || true)
216238
if [ -n "$BODY" ] && echo "$BODY" | grep -q '"cert_done"'; then
217239
echo "$BODY" > "$OUT/stage-$STAGE.json"
240+
# Every subscriber-only gen's stats too (subsLive adds up across them).
241+
for k in $(seq 1 $((GENS - 1))); do
242+
U=$(cat "$S/url-gen-$k-$(role_region "gen-$k").txt" 2>/dev/null) || continue
243+
curl -sf --max-time 15 "$U/" > "$OUT/stage-$STAGE-gen$k.json" || true
244+
done
218245
curl -sf --max-time 20 -H "authorization: Bearer $AUTH" \
219246
"$(cat "$S/url-server-$R.txt")/v1/debug/load" > "$OUT/server-after-$STAGE.json" || true
220247
echo "WC_STAGE_DONE $STAGE"

0 commit comments

Comments
 (0)