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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,7 @@ The pick is work-conserving: a connection with a free permit, searched from a ro
`--r2-max-concurrent-requests` stays the cap across all of them, split evenly and rounded up (rounding down would leave some connection at zero permits and wedge every GET routed to it), so raising the connection count alone spreads the same concurrency thinner instead of raising the ceiling; the per-connection share is what must stay at or below the stream limit, and the fetcher warns at startup when it exceeds it.
The count is published as `debug_trace_r2_connections` / `r2_connections`, and is rejected by name at zero, on a non-numeric or blank value, on the S3 target (HTTP/1.1 already opens a socket per in-flight GET there), and above the cap it divides — more connections than permits would leave some of them permanently idle.
It travels as text and is parsed after clap, so a blank env line — what a templated env file renders for a variable a role does not set — is named rather than aborting startup through clap's unnamed value error, and stays inert on the validator under `--witness-source rpc`, where every `--r2-*` flag is deliberately unread.
On the validator the shared semaphore is `--witness-max-concurrent-requests`, which sizes the RPC gateway too, so the two consumers trade off against each other.
The validator splits the two the same way: `--r2-max-concurrent-requests` caps R2 GETs while `--witness-max-concurrent-requests` sizes only the RPC witness path, so a budget written for one service cannot silently become the other's. They were one flag until the split, and carrying the old spelling into `--witness-source r2` is refused at startup by name rather than left to drop the cap — that mode has no RPC fallback, so an uncapped fetcher aims its whole in-flight window at the bucket.
`stateless-common`'s shared JSON-RPC client pins `http1_only`: `stateless-r2` enables reqwest's `http2` feature and Cargo unifies it workspace-wide, which would otherwise move the multi-MB witness RPC payloads onto one non-adaptive h2 connection per host.
Client-side routing, budgets, and fallback match the S3 target, but edge behavior is zone configuration: **a cache rule making these objects cacheable must set 404s to bypass cache**, or a pre-upload frontier miss gets pinned for the negative-cache TTL (stalling the validator's tip-following in its fallback-less R2 mode) and a cached 404 can false-fire the below-band `kind="missing"` bucket-integrity alarm.
The bucket is the same store the public gateway reads and can lead the generator at the frontier (uploader and generator RPC server publish from different files), so frontier hits are real; the frontier band is a small near-tip window (`R2_FRONTIER_WINDOW`, 32 blocks of uploader-lag grace on either side of the local tip — deliberately far narrower than the 4096-block routing window, so a stale catching-up tip cannot silence holes above it), hits there are labeled `witness_r2_frontier` (vs `witness_r2` past the band), the speculative frontier probe runs on an eighth of the remaining stage (vs half for blocks R2 must hold, so degraded R2 cannot burn half of every near-tip request's budget), and a `missing` classifies by band: in-band is the expected probe-ahead outcome (excluded from the alarm), below-band feeds `debug_trace_r2_witness_errors_total{kind="missing"}` (the bucket-integrity alarm, still covering recent-but-below-tip holes), and above-band — only reachable behind a stale catching-up tip — lands on its own `kind="missing_above_tip"` series, visible without flooding the alarm on every catch-up.
Expand Down
1 change: 0 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 0 additions & 1 deletion bin/debug-trace-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@ exclude.workspace = true
alloy-consensus.workspace = true
alloy-evm.workspace = true
alloy-genesis.workspace = true
alloy-op-evm.workspace = true
alloy-primitives.workspace = true
alloy-rpc-types-eth.workspace = true
alloy-rpc-types-trace.workspace = true
Expand Down
156 changes: 144 additions & 12 deletions bin/stateless-validator/src/app.rs
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,7 @@ pub struct CommandLineArgs {
/// when the first saturates, so this is the only way past the edge's per-connection stream
/// limit — and the only way one dropped connection stops taking every in-flight GET with
/// it, which matters here because R2 mode has no RPC fallback.
/// `--witness-max-concurrent-requests` is still the cap across all of them, split evenly
/// `--r2-max-concurrent-requests` is still the cap across all of them, split evenly
/// and rounded up, so raising this alone spreads the same concurrency thinner rather than
/// raising the ceiling; a count larger than that cap is rejected, since the surplus
/// connections could never be filled.
Expand Down Expand Up @@ -213,16 +213,25 @@ pub struct CommandLineArgs {
#[clap(long, env = "STATELESS_VALIDATOR_DATA_MAX_CONCURRENT_REQUESTS")]
pub data_max_concurrent_requests: Option<usize>,

/// Maximum concurrent in-flight witness fetches, independent of the data cap. Omit for
/// unlimited. Applies to both RPC witness calls and, with `--witness-source r2`, R2 GETs.
///
/// Against `--r2-custom-domain` this is also what bounds the GETs multiplexed onto the
/// HTTP/2 connection, so keep it at or below the edge's per-connection stream limit
/// (Cloudflare's is 100): above it the surplus queues inside the connection instead, where
/// the wait is unobservable and still counts against the per-attempt timeout.
/// Maximum concurrent in-flight RPC witness fetches, independent of the data cap. Omit
/// for unlimited. Applies to `--witness-source rpc` only; R2 GETs are capped by
/// `--r2-max-concurrent-requests`.
#[clap(long, env = "STATELESS_VALIDATOR_WITNESS_MAX_CONCURRENT_REQUESTS")]
pub witness_max_concurrent_requests: Option<usize>,

/// Maximum concurrent in-flight R2 witness GETs. Omit for unlimited. Deliberately
/// separate from `--witness-max-concurrent-requests`: that one sizes what we ask of the
/// RPC gateway, while R2 is a different service that tolerates far higher parallelism,
/// and under `--witness-source r2` the RPC witness path is not used at all.
///
/// Against `--r2-custom-domain` this is what bounds the GETs multiplexed onto each HTTP/2
/// connection, so keep the per-connection share (this value divided by
/// `--r2-connections`) at or below the edge's per-connection stream limit (Cloudflare's
/// is 100): above it the surplus queues inside the connection instead, where the wait is
/// unobservable and still counts against the per-attempt timeout.
#[clap(long, env = "STATELESS_VALIDATOR_R2_MAX_CONCURRENT_REQUESTS")]
pub r2_max_concurrent_requests: Option<usize>,
Comment thread
flyq marked this conversation as resolved.
Comment thread
flyq marked this conversation as resolved.

/// Fetcher caught-up poll interval (milliseconds). Also rate-limits `eth_blockNumber`.
/// Lower values reduce tip-following lag at the cost of more RPC traffic when caught up.
#[clap(long, env = "STATELESS_VALIDATOR_POLL_INTERVAL_MS")]
Expand Down Expand Up @@ -459,6 +468,20 @@ fn build_r2_client(
timeouts: stateless_r2::fetch::FetchTimeouts,
retry: BackoffPolicy,
) -> Result<R2WitnessClient> {
// `--witness-max-concurrent-requests` capped R2 GETs too before the caps were split.
// Refuse the pre-split spelling by name rather than leave R2 uncapped: this mode has no
// RPC fallback, so an uncapped fetcher aims its whole in-flight window at the bucket.
// Both spellings together stay legal — one env template can feed rpc-mode and r2-mode
// roles alike, each mode reading only its own cap — so only old-spelling-alone is refused.
// The message names the env spelling too: the deployments this guard exists for configure
// through env files, where the flag spelling alone costs a name-translation round trip.
if args.witness_max_concurrent_requests.is_some() && args.r2_max_concurrent_requests.is_none() {
return Err(eyre::eyre!(
"--witness-max-concurrent-requests no longer caps R2 GETs under --witness-source \
r2 (it now sizes only the RPC witness path): set --r2-max-concurrent-requests \
(env STATELESS_VALIDATOR_R2_MAX_CONCURRENT_REQUESTS) instead"
));
}
// Every coherence rule lives in the shared validator, so the reads below rest on an
// invariant that was actually checked: no empty values, exactly one target, and an Access
// pair that is either whole or absent.
Expand All @@ -484,14 +507,15 @@ fn build_r2_client(
access,
timeouts,
retry,
args.witness_max_concurrent_requests,
args.r2_max_concurrent_requests,
connections,
)?;
metrics::record_r2_connections(client.connections());
info!(
domain = %client.origin(),
cf_access,
connections = client.connections(),
max_concurrent_requests = ?client.max_concurrent_requests(),
"Witness source: R2 (custom domain)"
);
client
Expand All @@ -505,11 +529,12 @@ fn build_r2_client(
args.r2_secret_access_key.as_ref().expect("S3 target").as_ref().to_string(),
timeouts,
retry,
args.witness_max_concurrent_requests,
args.r2_max_concurrent_requests,
)?;
info!(
endpoint = %client.origin(),
bucket = args.r2_bucket.as_deref().unwrap_or_default(),
max_concurrent_requests = ?client.max_concurrent_requests(),
"Witness source: R2 (direct S3)"
);
client
Expand Down Expand Up @@ -540,8 +565,8 @@ fn r2_flags(args: &CommandLineArgs) -> R2Flags<'_> {
),
connections: R2Flag::new("--r2-connections", args.r2_connections.as_deref()),
max_concurrent_requests: R2CountFlag::new(
"--witness-max-concurrent-requests",
args.witness_max_concurrent_requests,
"--r2-max-concurrent-requests",
args.r2_max_concurrent_requests,
),
// Empty on purpose. The orphan-tuning rule exists for a binary that validates R2 flags
// on every startup; here they are only read under `--witness-source r2`, where a target
Expand All @@ -550,3 +575,110 @@ fn r2_flags(args: &CommandLineArgs) -> R2Flags<'_> {
tuning: &[],
}
}

#[cfg(test)]
mod tests {
use super::*;

/// A loopback custom-domain target, the default shape for rules that are target-agnostic.
const CUSTOM_DOMAIN_TARGET: &[&str] = &["--r2-custom-domain", "http://127.0.0.1:9000"];

/// The signed S3 target: the endpoint plus its credential quad.
const S3_TARGET: &[&str] = &[
"--r2-endpoint",
"https://acc.r2.cloudflarestorage.com",
"--r2-bucket",
"witness",
"--r2-access-key-id",
"key-id",
"--r2-secret-access-key",
"secret",
];

/// Argv for `--witness-source r2` with the given target flags, so [`build_r2_client`] —
/// the seam every R2-mode rule is gated behind — runs the rules from the path production
/// takes.
fn parse_r2_with_target(target: &[&str], extra: &[&str]) -> CommandLineArgs {
let argv = [
"stateless-validator",
"--data-dir",
"/tmp/x",
"--rpc-endpoint",
"http://rpc",
"--witness-source",
"r2",
];
CommandLineArgs::try_parse_from(argv.iter().chain(target).chain(extra)).expect("parses")
}

fn parse_r2(extra: &[&str]) -> CommandLineArgs {
parse_r2_with_target(CUSTOM_DOMAIN_TARGET, extra)
}

fn build(args: &CommandLineArgs) -> Result<R2WitnessClient> {
let timeouts = stateless_r2::fetch::FetchTimeouts {
per_attempt: Duration::from_secs(1),
connect: Duration::from_secs(1),
};
let retry =
BackoffPolicy { initial: Duration::from_millis(1), max: Duration::from_millis(1) };
build_r2_client(args, timeouts, retry)
}

/// Carrying the pre-split spelling of the R2 concurrency cap into `--witness-source r2`
/// must fail by name rather than leave R2 uncapped: that mode has no RPC fallback, so an
/// uncapped fetcher aims its whole in-flight window at the bucket. Outside r2 mode the
/// rule is unreachable by construction — `build_r2_client` is only called from the
/// `WitnessSource::R2` arm, the same call-site gating as every other R2 rule.
#[test]
fn r2_mode_refuses_the_pre_split_concurrency_spelling() {
let _guard = stateless_test_utils::env::env_lock();

let stale = parse_r2(&["--witness-max-concurrent-requests", "48"]);
let msg =
build(&stale).expect_err("the old spelling must be refused in r2 mode").to_string();
assert!(msg.contains("--witness-max-concurrent-requests"), "{msg}");
assert!(msg.contains("--r2-max-concurrent-requests"), "{msg}");
// The env spelling too: the deployments this guard exists for configure through env
// files, and the flag spelling alone would cost a name-translation round trip.
assert!(msg.contains("STATELESS_VALIDATOR_R2_MAX_CONCURRENT_REQUESTS"), "{msg}");

// Migrated: the new spelling alone is accepted.
build(&parse_r2(&["--r2-max-concurrent-requests", "48"]))
.expect("migrated spelling builds");

// Both set is accepted — the RPC cap is simply unread in this mode — and each
// spelling lands on its own field.
let both = parse_r2(&[
"--witness-max-concurrent-requests",
"16",
"--r2-max-concurrent-requests",
"48",
]);
assert_eq!(both.witness_max_concurrent_requests, Some(16));
assert_eq!(both.r2_max_concurrent_requests, Some(48));
build(&both).expect("both caps set builds");
}

/// The migration guard fires on the flags alone, so its test above would still pass with
/// the constructors wired to the old field. This is the assertion that observes which cap
/// actually reaches the client — with both spellings set it must be the R2 one, not the
/// RPC one — on both target arms, plus the uncapped default, so a revert of either arm's
/// wiring fails here by value.
#[test]
fn the_r2_cap_not_the_rpc_one_reaches_the_client() {
let _guard = stateless_test_utils::env::env_lock();

for target in [CUSTOM_DOMAIN_TARGET, S3_TARGET] {
let both = parse_r2_with_target(
target,
&["--witness-max-concurrent-requests", "16", "--r2-max-concurrent-requests", "48"],
);
let capped = build(&both).expect("both caps set builds");
assert_eq!(capped.max_concurrent_requests(), Some(48), "{target:?}");

let uncapped = build(&parse_r2_with_target(target, &[])).expect("no caps builds");
assert_eq!(uncapped.max_concurrent_requests(), None, "{target:?}");
}
}
}
2 changes: 1 addition & 1 deletion bin/stateless-validator/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -274,7 +274,7 @@ pub fn record_r2_negotiated_version(version: &'static str) {
/// How many HTTP/2 connections the custom-domain target spreads its GETs over.
///
/// A plain value rather than an info label: it is the divisor for the per-connection stream
/// budget, so a dashboard reads it against `--witness-max-concurrent-requests` and against the
/// budget, so a dashboard reads it against `--r2-max-concurrent-requests` and against the
/// edge's limit rather than grouping by it. Published only for the custom-domain target, where
/// one client is one connection and the count is a real property of the transport.
pub fn record_r2_connections(connections: usize) {
Expand Down
15 changes: 12 additions & 3 deletions bin/stateless-validator/src/r2_witness.rs
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,9 @@ impl R2WitnessError {
#[derive(Debug)]
pub struct R2WitnessClient {
fetcher: R2ObjectFetcher,
/// The configured in-flight GET cap, retained here because the fetcher decomposes it into
/// per-connection permits and cannot report the configured value back.
max_concurrent_requests: Option<usize>,
}

/// The fetcher's pacing view of a `BackoffPolicy` — the adapter-layer conversion that keeps
Expand All @@ -128,6 +131,12 @@ impl R2WitnessClient {
self.fetcher.connections()
}

/// The configured cap on in-flight GETs (`None` = unlimited; see [`Self::new`] for the
/// exact semantics), for startup logging.
pub fn max_concurrent_requests(&self) -> Option<usize> {
self.max_concurrent_requests
}

/// Builds a client from an R2 endpoint origin, bucket, and bucket-scoped S3 credentials.
///
/// `timeouts` bounds each individual GET (end-to-end and connect). `retry_backoff` paces the
Expand All @@ -136,7 +145,7 @@ impl R2WitnessClient {
/// `--rpc-initial-backoff-ms` / `--rpc-max-backoff-ms`, so one pair of flags tunes both
/// paths. `max_concurrent_requests` caps the number of GETs in flight at once (`None` =
/// unlimited, `Some(0)` clamps to 1 — same semantics as the RPC witness semaphore; in R2
/// mode this client is the only enforcement of `--witness-max-concurrent-requests`). Fails
/// mode this client is the only enforcement of `--r2-max-concurrent-requests`). Fails
/// if the endpoint is not a bare `scheme://host[:port]` origin or the HTTP client cannot be
/// built.
pub fn new(
Expand All @@ -158,7 +167,7 @@ impl R2WitnessClient {
max_concurrent_requests,
)
.map_err(|e| eyre::eyre!(e))?;
Ok(Self { fetcher })
Ok(Self { fetcher, max_concurrent_requests })
}

/// Builds a client that fetches unsigned through a Cloudflare custom domain fronting the
Expand All @@ -182,7 +191,7 @@ impl R2WitnessClient {
)
.map(|fetcher| fetcher.on_version_observed(metrics::record_r2_negotiated_version))
.map_err(|e| eyre::eyre!(e))?;
Ok(Self { fetcher })
Ok(Self { fetcher, max_concurrent_requests })
}

/// Fetches and decodes the witness for `(number, hash)` from R2.
Expand Down
9 changes: 9 additions & 0 deletions bin/stateless-validator/tests/integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,15 @@ fn witness_max_concurrent_requests_flag_and_env() {
);
}

#[test]
fn r2_max_concurrent_requests_flag_and_env() {
assert_optional_numeric_flag::<usize>(
"--r2-max-concurrent-requests",
"STATELESS_VALIDATOR_R2_MAX_CONCURRENT_REQUESTS",
|a| a.r2_max_concurrent_requests,
);
}

#[test]
fn tip_buffer_flag_and_env() {
assert_optional_numeric_flag::<u64>("--tip-buffer", "STATELESS_VALIDATOR_TIP_BUFFER", |a| {
Expand Down
Loading