diff --git a/CLAUDE.md b/CLAUDE.md index dcf9bcda..f6737894 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -277,6 +277,11 @@ actual_slot = finalized_slot + 1 + relative_index proving or proof decoding. `--prover-arena` opts into leanVM's bump arena, which recycles the prover's large buffers across proofs instead of re-faulting them, so its pages stay resident for the node's lifetime +- Every proof runs on one dedicated `leanvm-prover` thread (`prove` in + `ethlambda-crypto`). The arena pins a slab to each thread that ever drives a + proof, so proving from the calling thread (tokio's blocking pool, the actors) + grew RSS one proof peak per new thread until aggregators were OOM-killed. + `tests/arena_slabs.rs` guards this - `ethlambda keygen` generates genesis validator keys through the same `ValidatorSecretKey` the node loads them with, so a key set cannot be built against a different leanVM than the client reading it. Keys are only usable by diff --git a/Cargo.lock b/Cargo.lock index 6b75f7f4..6a65521b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1999,6 +1999,7 @@ dependencies = [ "postcard", "thiserror 2.0.18", "tracing", + "zk_alloc", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index a24fa802..a0cfcfd0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -87,6 +87,10 @@ clap = { version = "4.3", features = ["derive", "env"] } # the whole crypto stack arrives as one dependency at one revision. # Pinned to a `main` commit for reproducible builds; bump the rev to track main. leanvm = { git = "https://github.com/leanEthereum/leanVM.git", rev = "48a904208d682848dac0e18ef8b01ebfc40df9ad" } +# leanVM's prover arena, for tests that read its slab accounting (the facade does +# not re-export it). Keep the rev equal to leanvm's: only then is this the same +# crate instance, with the same statics, as the one leanvm proves on. +zk_alloc = { git = "https://github.com/leanEthereum/leanVM.git", rev = "48a904208d682848dac0e18ef8b01ebfc40df9ad" } # Secret-key (de)serialization for the leanVM xmss key format. postcard = { version = "1.1.3", features = ["alloc"] } diff --git a/crates/common/crypto/Cargo.toml b/crates/common/crypto/Cargo.toml index 5fcf5a40..c15944b1 100644 --- a/crates/common/crypto/Cargo.toml +++ b/crates/common/crypto/Cargo.toml @@ -23,3 +23,4 @@ shadow-integration = [] [dev-dependencies] hex.workspace = true +zk_alloc.workspace = true diff --git a/crates/common/crypto/src/lib.rs b/crates/common/crypto/src/lib.rs index 09fa0d43..aa4480a4 100644 --- a/crates/common/crypto/src/lib.rs +++ b/crates/common/crypto/src/lib.rs @@ -37,7 +37,9 @@ use ethlambda_types::{block::ByteList512KiB, primitives::H256}; use crate::signature::{ValidatorPublicKey, ValidatorSignature}; use leanvm::{ClaimSelection, EthereumProof, SignatureClaims, XmssClaimGroup, aggregate, xmss}; -use std::sync::{Mutex, MutexGuard}; +use std::panic::{self, AssertUnwindSafe}; +use std::sync::{OnceLock, mpsc}; +use std::thread; use thiserror::Error; use tracing::error; @@ -79,19 +81,59 @@ pub fn init_leanvm(use_arena: bool) { } } -/// Claims the exclusive right to prove. +/// A proving job queued for the prover thread. +type ProverJob = Box; + +/// Runs `job` on the process's one prover thread and blocks until it returns. /// /// leanVM allows one proof at a time per process; a second concurrent one panics. -/// Proving is legal only while the returned guard is alive, so take it immediately -/// before the prove call: decoding and argument conversion need no permit. +/// The single thread serializes them, and it is also what bounds the arena: +/// leanVM's arena hands each thread that proves its own slab and never takes it +/// back, and the thread driving a proof is the one whose slab fills. Proving from +/// whichever thread asked would pin one proof's peak per thread ever used, which +/// is how aggregators running `--prover-arena` grew until they were OOM-killed. +/// +/// Wrap only the prove call: decoding and argument conversion run fine on the +/// caller. A job must not call back into this function, since the prover thread +/// would then wait on itself. +/// +/// A panicking job is caught on the prover thread, so one bad proof cannot take +/// down every later one, and is re-raised on the caller. +fn prove(job: F) -> T +where + F: FnOnce() -> T + Send + 'static, + T: Send + 'static, +{ + let (reply_tx, reply_rx) = mpsc::sync_channel(1); + let task: ProverJob = Box::new(move || { + let outcome = panic::catch_unwind(AssertUnwindSafe(job)); + if outcome.is_err() { + error!("a proving job panicked; the prover thread carries on"); + } + // The caller waits on the reply until it arrives, so the send cannot fail. + let _ = reply_tx.send(outcome); + }); + prover_queue() + .send(task) + .expect("the prover thread never exits"); + reply_rx + .recv() + .expect("the prover thread replies to every job") + .unwrap_or_else(|payload| panic::resume_unwind(payload)) +} + +/// The prover thread's job queue, spawning the thread on first use. /// -/// The permit guards no data, so a poisoned lock is recovered rather than propagated: -/// one panicking prover must not brick every later proof. It is still an incident. -fn acquire_prover() -> MutexGuard<'static, ()> { - static PROVER_PERMIT: Mutex<()> = Mutex::new(()); - PROVER_PERMIT.lock().unwrap_or_else(|poisoned| { - error!("a previous proving job panicked while holding the permit; continuing"); - poisoned.into_inner() +/// The thread lives for the rest of the process. +fn prover_queue() -> &'static mpsc::Sender { + static QUEUE: OnceLock> = OnceLock::new(); + QUEUE.get_or_init(|| { + let (tx, rx) = mpsc::channel::(); + thread::Builder::new() + .name("leanvm-prover".to_string()) + .spawn(move || rx.into_iter().for_each(|job| job())) + .expect("failed to spawn the leanVM prover thread"); + tx }) } @@ -324,10 +366,8 @@ pub fn aggregate_signatures( let raw_xmss = raw_xmss_inputs(public_keys, signatures, message, slot); - let _permit = acquire_prover(); - - let proof = - aggregate(&[], raw_xmss, vec![], &[], None, LOG_INV_RATE).map_err(aggregation_failed)?; + let proof = prove(move || aggregate(&[], raw_xmss, vec![], &[], None, LOG_INV_RATE)) + .map_err(aggregation_failed)?; compress_to_byte_list(&proof) } @@ -375,10 +415,9 @@ pub fn aggregate_mixed( let children_native = decompress_children(children, message, slot)?; let raw_xmss = raw_xmss_inputs(raw_public_keys, raw_signatures, message, slot); - let _permit = acquire_prover(); - - let proof = aggregate(&children_native, raw_xmss, vec![], &[], None, LOG_INV_RATE) - .map_err(aggregation_failed)?; + let proof = + prove(move || aggregate(&children_native, raw_xmss, vec![], &[], None, LOG_INV_RATE)) + .map_err(aggregation_failed)?; compress_to_byte_list(&proof) } @@ -412,9 +451,7 @@ pub fn aggregate_proofs( let children_native = decompress_children(children, message, slot)?; - let _permit = acquire_prover(); - - let proof = aggregate(&children_native, vec![], vec![], &[], None, LOG_INV_RATE) + let proof = prove(move || aggregate(&children_native, vec![], vec![], &[], None, LOG_INV_RATE)) .map_err(aggregation_failed)?; compress_to_byte_list(&proof) @@ -500,9 +537,7 @@ pub fn merge_type_1s_into_type_2( }) .collect::>()?; - let _permit = acquire_prover(); - - let merged = aggregate(&type_1s_native, vec![], vec![], &[], None, LOG_INV_RATE) + let merged = prove(move || aggregate(&type_1s_native, vec![], vec![], &[], None, LOG_INV_RATE)) .map_err(aggregation_failed)?; compress_to_byte_list(&merged) @@ -574,17 +609,17 @@ pub fn split_type_2_by_message( xmss: vec![group], sphincs: Vec::new(), }; - // No blobs: ethlambda makes no LeanDA claim, so the selection publishes the - // one signature group and no DA roots. - let declare = ClaimSelection { - signatures: &kept, - da_commitments: &[], - }; - let _permit = acquire_prover(); - - let component = aggregate(&[type_2], vec![], vec![], &[], Some(declare), LOG_INV_RATE) - .map_err(aggregation_failed)?; + let component = prove(move || { + // No blobs: ethlambda makes no LeanDA claim, so the selection publishes the + // one signature group and no DA roots. + let declare = ClaimSelection { + signatures: &kept, + da_commitments: &[], + }; + aggregate(&[type_2], vec![], vec![], &[], Some(declare), LOG_INV_RATE) + }) + .map_err(aggregation_failed)?; compress_to_byte_list(&component) } @@ -639,11 +674,42 @@ mod tests { // (`OnceLock::get_or_init`). init_leanvm(false); init_leanvm(false); + } + + /// The arena bound rests on this: however many threads ask for a proof, one + /// thread runs them all, and it is none of the askers. + #[test] + fn prove_runs_every_job_on_one_thread() { + let askers: Vec<_> = (0..4) + .map(|_| { + thread::spawn(|| { + let asker = thread::current().id(); + let prover = prove(|| thread::current().id()); + (asker, prover) + }) + }) + .collect(); + let ids: Vec<_> = askers + .into_iter() + .map(|asker| asker.join().expect("asker thread")) + .collect(); + + let prover = ids[0].1; + assert!(ids.iter().all(|&(_, id)| id == prover)); + assert!(ids.iter().all(|&(asker, _)| asker != prover)); + assert_ne!(thread::current().id(), prover); + } + + /// A panicking job reaches its caller, and the prover thread outlives it. + #[test] + fn prove_reraises_a_panic_and_keeps_proving() { + let before = prove(|| thread::current().id()); + + let outcome = panic::catch_unwind(|| prove(|| panic!("bad proof"))); + let payload = outcome.expect_err("the job's panic reaches the caller"); + assert_eq!(payload.downcast_ref::<&str>(), Some(&"bad proof")); - // The permit is dropped between acquisitions: it is not reentrant, so holding - // both at once would deadlock. That also covers release-on-drop. - drop(acquire_prover()); - drop(acquire_prover()); + assert_eq!(prove(|| thread::current().id()), before); } /// The claim list a decode rebuilds has to match what was aggregated, and diff --git a/crates/common/crypto/tests/arena.rs b/crates/common/crypto/tests/arena.rs index a037ef9f..a4330cd1 100644 --- a/crates/common/crypto/tests/arena.rs +++ b/crates/common/crypto/tests/arena.rs @@ -4,33 +4,11 @@ //! the lib test binary: it would change the allocator under every other test. An //! integration test gets its own process. -use ethlambda_crypto::{ - aggregate_signatures, init_leanvm, - signature::{ValidatorPublicKey, ValidatorSignature}, - verify_aggregated_signature, -}; -use ethlambda_types::primitives::H256; -use leanvm::xmss::{self, Encode as _, key_gen_from_seed}; - -/// Mirrors the lib tests' helper: a small slot range keeps key generation fast. -fn keypair_and_signature( - seed: u64, - first_slot: u32, - signing_slot: u32, - message: &H256, -) -> (ValidatorPublicKey, ValidatorSignature) { - let mut seed_bytes = [0u8; 32]; - seed_bytes[..8].copy_from_slice(&seed.to_le_bytes()); - - let (sk, pk) = - key_gen_from_seed(seed_bytes, first_slot, first_slot + 63).expect("valid slot range"); - let sig = xmss::sign(&sk, &message.0, signing_slot).expect("sign"); +mod common; - ( - ValidatorPublicKey::from_bytes(&pk.as_ssz_bytes()).unwrap(), - ValidatorSignature::from_bytes(&sig.as_ssz_bytes()).unwrap(), - ) -} +use common::keypair_and_signature; +use ethlambda_crypto::{aggregate_signatures, init_leanvm, verify_aggregated_signature}; +use ethlambda_types::primitives::H256; #[test] #[ignore = "too slow"] diff --git a/crates/common/crypto/tests/arena_slabs.rs b/crates/common/crypto/tests/arena_slabs.rs new file mode 100644 index 00000000..cd683e1e --- /dev/null +++ b/crates/common/crypto/tests/arena_slabs.rs @@ -0,0 +1,49 @@ +//! The arena's memory bound under ethlambda's threading. +//! +//! leanVM's arena gives each thread that allocates during a proof its own slab and keeps +//! it for the life of the process; the slab that fills is the one of the thread driving +//! the proof. ethlambda asks for proofs from threads that come and go (tokio's blocking +//! pool, the actors), so proving on the asking thread would pin one slab per thread ever +//! used, which is how aggregators on `--prover-arena` were OOM-killed. +//! +//! Its own test binary, and so its own process: the arena engages process-wide and the +//! slab count is process-wide, so another test proving alongside would move it. + +mod common; + +use std::thread; + +use common::keypair_and_signature; +use ethlambda_crypto::{aggregate_signatures, init_leanvm}; +use ethlambda_types::primitives::H256; + +#[test] +#[ignore = "too slow"] +fn proving_from_fresh_threads_claims_no_new_slabs() { + init_leanvm(true); + + let message = H256::from([7u8; 32]); + let slot = 10u32; + let (pk, sig) = keypair_and_signature(1, 5, slot, &message); + let prove_from_a_fresh_thread = || { + let (pk, sig) = (pk.clone(), sig.clone()); + thread::spawn(move || aggregate_signatures(vec![pk], vec![sig], &message, slot)) + .join() + .expect("asking thread") + .expect("aggregation on the arena"); + }; + + prove_from_a_fresh_thread(); + let slabs = zk_alloc::stats().threads; + // Zero would mean this test reads another copy of the arena than leanvm's. + assert!(slabs > 0, "the first proof claimed no slab"); + + for _ in 0..4 { + prove_from_a_fresh_thread(); + } + assert_eq!( + zk_alloc::stats().threads, + slabs, + "proving from new threads claimed new slabs" + ); +} diff --git a/crates/common/crypto/tests/common/mod.rs b/crates/common/crypto/tests/common/mod.rs new file mode 100644 index 00000000..6dd16206 --- /dev/null +++ b/crates/common/crypto/tests/common/mod.rs @@ -0,0 +1,25 @@ +//! Helpers shared by the integration tests. + +use ethlambda_crypto::signature::{ValidatorPublicKey, ValidatorSignature}; +use ethlambda_types::primitives::H256; +use leanvm::xmss::{self, Encode as _, key_gen_from_seed}; + +/// Mirrors the lib tests' helper: a small slot range keeps key generation fast. +pub fn keypair_and_signature( + seed: u64, + first_slot: u32, + signing_slot: u32, + message: &H256, +) -> (ValidatorPublicKey, ValidatorSignature) { + let mut seed_bytes = [0u8; 32]; + seed_bytes[..8].copy_from_slice(&seed.to_le_bytes()); + + let (sk, pk) = + key_gen_from_seed(seed_bytes, first_slot, first_slot + 63).expect("valid slot range"); + let sig = xmss::sign(&sk, &message.0, signing_slot).expect("sign"); + + ( + ValidatorPublicKey::from_bytes(&pk.as_ssz_bytes()).unwrap(), + ValidatorSignature::from_bytes(&sig.as_ssz_bytes()).unwrap(), + ) +}