Skip to content
Draft
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
150 changes: 141 additions & 9 deletions crates/anvil/src/eth/backend/beacon.rs
Original file line number Diff line number Diff line change
@@ -1,20 +1,41 @@
//! The immutable Beacon slot schedule of an execution fork.
//! Historical Beacon reads and the immutable slot schedule of an execution fork.

use alloy_rpc_types_beacon::genesis::{GenesisData, GenesisResponse};
use alloy_consensus::{Blob, EnvKzgSettings};
use alloy_eips::eip4844::kzg_to_versioned_hash;
use alloy_primitives::B256;
use alloy_rpc_types_beacon::{
genesis::{GenesisData, GenesisResponse},
sidecar::GetBlobsResponse,
};
use eyre::{Result, ensure, eyre};
use reqwest::{Client, StatusCode, Url, redirect::Policy};
use serde::de::DeserializeOwned;
use std::time::Duration;
use std::{fmt, time::Duration};

// JSON encodes each 128 KiB blob as hex. This permits over 100 blobs without accepting
// unbounded responses from an upstream provider.
const MAX_BLOB_RESPONSE_BYTES: usize = 32 * 1024 * 1024;
const MAX_METADATA_RESPONSE_BYTES: usize = 1024 * 1024;

/// Beacon metadata belonging to an execution fork.
#[derive(Clone, Debug)]
/// Beacon connectivity and metadata belonging to an execution fork.
#[derive(Clone)]
pub struct ForkBeacon {
url: Url,
client: Client,
pub genesis: GenesisData,
pub seconds_per_slot: u64,
}

impl fmt::Debug for ForkBeacon {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
// Endpoint paths and query parameters may contain credentials.
f.debug_struct("ForkBeacon")
.field("genesis", &self.genesis)
.field("seconds_per_slot", &self.seconds_per_slot)
.finish_non_exhaustive()
}
}

impl ForkBeacon {
/// Captures Beacon metadata before the fork starts serving requests.
pub async fn connect(endpoint: &str, timeout: Duration, fork_timestamp: u64) -> Result<Self> {
Expand All @@ -25,6 +46,7 @@ impl ForkBeacon {
&client,
&url,
"eth/v1/beacon/genesis",
&[],
MAX_METADATA_RESPONSE_BYTES,
)
.await?
Expand All @@ -34,6 +56,7 @@ impl ForkBeacon {
&client,
&url,
"eth/v1/config/spec",
&[],
MAX_METADATA_RESPONSE_BYTES,
)
.await?
Expand All @@ -50,18 +73,53 @@ impl ForkBeacon {
elapsed.is_multiple_of(seconds_per_slot),
"execution fork timestamp is not on the Beacon slot grid"
);
Ok(Self { genesis, seconds_per_slot })
Ok(Self { url, client, genesis, seconds_per_slot })
}

/// Converts a slot to its exact execution timestamp without wrapping on user input.
pub fn timestamp(&self, slot: u64) -> Option<u64> {
slot.checked_mul(self.seconds_per_slot)?.checked_add(self.genesis.genesis_time)
}

/// Reads historical blobs. The caller must enforce the inclusive execution fork boundary.
pub async fn blobs(&self, slot: u64, hashes: &[B256]) -> Result<Option<Vec<Blob>>> {
let Some(response) = Self::get::<GetBlobsResponse>(
&self.client,
&self.url,
&format!("eth/v1/beacon/blobs/{slot}"),
hashes,
MAX_BLOB_RESPONSE_BYTES,
)
.await?
else {
return Ok(None);
};
let mut blobs = response.data;
if !hashes.is_empty() {
// Providers may ignore the filter. Apply it to the actual blob content, not
// to untrusted upstream labels; Base also validates the requested hashes.
let hashes = hashes.to_vec();
blobs = tokio::task::spawn_blocking(move || -> Result<Vec<Blob>> {
let settings = EnvKzgSettings::Default;
let mut filtered = Vec::new();
for blob in blobs {
let commitment = settings.get().blob_to_kzg_commitment(&blob.0.into())?;
if hashes.contains(&kzg_to_versioned_hash(commitment.as_slice())) {
filtered.push(blob);
}
}
Ok(filtered)
})
.await??;
}
Ok(Some(blobs))
}

async fn get<T: DeserializeOwned>(
client: &Client,
endpoint: &Url,
path: &str,
hashes: &[B256],
max_bytes: usize,
) -> Result<Option<T>> {
// Append to the endpoint's path and keep its query; either may carry credentials.
Expand All @@ -70,7 +128,11 @@ impl ForkBeacon {
.map_err(|_| eyre!("invalid fork Beacon URL"))?
.pop_if_empty()
.extend(path.split('/'));
let request = client.get(url).header("Accept", "application/json");
let mut request = client.get(url).header("Accept", "application/json");
if !hashes.is_empty() {
let hashes = hashes.iter().map(ToString::to_string).collect::<Vec<_>>().join(",");
request = request.query(&[("versioned_hashes", hashes)]);
}
let response = request.send().await.map_err(reqwest::Error::without_url)?;
if response.status() == StatusCode::NOT_FOUND {
return Ok(None);
Expand All @@ -89,13 +151,15 @@ impl ForkBeacon {
#[cfg(test)]
mod tests {
use super::*;
use alloy_primitives::{B256, FixedBytes};
use alloy_consensus::{BlobTransactionSidecar, SidecarBuilder, SimpleCoder};
use alloy_primitives::FixedBytes;
use axum::{
Json, Router,
extract::RawQuery,
extract::{Path, RawQuery},
response::{IntoResponse, Redirect},
routing::get,
};
use std::sync::{Arc, Mutex};
use tokio::net::TcpListener;

const TIMEOUT: Duration = Duration::from_secs(1);
Expand Down Expand Up @@ -227,4 +291,72 @@ mod tests {
.await;
assert!(connect_err(&slow, FORK).await.contains("timed out"));
}

fn sidecar(data: &[u8]) -> BlobTransactionSidecar {
SidecarBuilder::<SimpleCoder>::from_slice(data).build().unwrap()
}

/// Connects through a credential-bearing endpoint whose slot 20 serves `blobs` regardless of
/// filters. Slot 16 is oversized, slot 18 fails, other slots are missing. Records blob queries.
async fn blob_upstream(blobs: Vec<Blob>) -> (ForkBeacon, Arc<Mutex<Vec<String>>>) {
let queries = Arc::new(Mutex::new(Vec::new()));
let recorded = Arc::clone(&queries);
let handler = move |Path(slot): Path<u64>, RawQuery(query): RawQuery| {
recorded.lock().unwrap().push(query.unwrap_or_default());
let data = blobs.clone();
async move {
match slot {
16 => vec![b' '; MAX_BLOB_RESPONSE_BYTES + 1].into_response(),
18 => StatusCode::INTERNAL_SERVER_ERROR.into_response(),
20 => Json(GetBlobsResponse {
execution_optimistic: false,
finalized: true,
data,
})
.into_response(),
_ => StatusCode::NOT_FOUND.into_response(),
}
}
};
let router = metadata("/secret-key", "12")
.route("/secret-key/eth/v1/beacon/blobs/{slot}", get(handler));
let url = serve(router).await;
let beacon = ForkBeacon::connect(&format!("{url}/secret-key?token=secret"), TIMEOUT, FORK)
.await
.unwrap();
(beacon, queries)
}

#[tokio::test]
async fn blobs_are_filtered_by_actual_content() {
let [first, second, absent] = [b"first".as_slice(), b"second", b"absent"].map(sidecar);
let hash = |sidecar: &BlobTransactionSidecar| sidecar.versioned_hashes().next().unwrap();
let (beacon, queries) = blob_upstream(vec![first.blobs[0], second.blobs[0]]).await;
assert!(!format!("{beacon:?}").contains("secret"));

let all = beacon.blobs(20, &[]).await.unwrap().unwrap();
assert!(all == [first.blobs[0], second.blobs[0]], "unfiltered blobs");
let filtered = beacon.blobs(20, &[hash(&second)]).await.unwrap().unwrap();
assert!(filtered == [second.blobs[0]], "filtered blobs");
assert!(beacon.blobs(20, &[hash(&absent)]).await.unwrap().unwrap().is_empty());
assert_eq!(
*queries.lock().unwrap(),
[
"token=secret".to_string(),
format!("token=secret&versioned_hashes={}", hash(&second)),
format!("token=secret&versioned_hashes={}", hash(&absent)),
]
);
}

#[tokio::test]
async fn blobs_distinguish_missing_from_upstream_failures() {
let (beacon, _) = blob_upstream(Vec::new()).await;
assert!(beacon.blobs(19, &[]).await.unwrap().is_none());
assert!(beacon.blobs(20, &[]).await.unwrap().unwrap().is_empty());
for (slot, error) in [(16, "exceeds size limit"), (18, "500")] {
let err = format!("{:?}", beacon.blobs(slot, &[]).await.unwrap_err());
assert!(err.contains(error) && !err.contains("secret"), "slot {slot}: {err}");
}
}
}
Loading