Skip to content
Draft
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
1 change: 1 addition & 0 deletions crates/blockifier/src/blockifier.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
pub mod block;
pub mod concurrent_transaction_executor;
pub mod config;
pub mod native_execution_thread;
pub mod stateful_validator;
pub mod transaction_executor;
#[cfg(test)]
Expand Down
152 changes: 152 additions & 0 deletions crates/blockifier/src/blockifier/native_execution_thread.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
//! Stack-sized execution threads whose thread-local Native caches survive sequential batches.

use std::cell::Cell;
use std::io;
use std::thread::{self, JoinHandle};

thread_local! {
// Only threads allocated by this module may advertise their stack size.
static NATIVE_STACK_SIZE: Cell<usize> = const { Cell::new(0) };
}

/// Starts a long-lived execution loop with a stack suitable for Native execution.
///
/// Use the same stack size as `TransactionExecutorConfig::stack_size`. Sequential executors
/// on this thread can reuse it, including across block transitions. Ordinary callers retain
/// the scoped-thread fallback; a larger requested stack also falls back rather than trusting
/// an undersized worker. This does not change the concurrent executor's worker pool.
pub fn spawn_native_execution_thread<F, T>(
name: String,
stack_size: usize,
run: F,
) -> io::Result<JoinHandle<T>>
where
F: FnOnce() -> T + Send + 'static,
T: Send + 'static,
{
thread::Builder::new().name(name).stack_size(stack_size).spawn(move || {
NATIVE_STACK_SIZE.set(stack_size);
run()
})
}

#[cfg(feature = "cairo_native")]
pub(crate) fn on_native_stack<F, T>(stack_size: usize, run: F) -> T
where
F: FnOnce() -> T + Send,
T: Send,
{
if NATIVE_STACK_SIZE.get() >= stack_size && NATIVE_STACK_SIZE.get() != 0 {
return run();
}
// Preserve the existing protection for RPC/validation and other unregistered callers.
// This temporary thread is deliberately not registered as a reusable execution loop.
thread::scope(|scope| {
thread::Builder::new()
.stack_size(stack_size)
.spawn_scoped(scope, run)
.expect("Failed to spawn Native execution thread")
.join()
.expect("Failed to join Native execution thread")
})
}

#[cfg(all(test, feature = "cairo_native"))]
mod tests {
use super::*;

const STACK: usize = 2 * 1024 * 1024;
thread_local! {
static BATCHES: Cell<usize> = const { Cell::new(0) };
}

#[test]
fn registered_worker_reuses_thread_and_locals_across_batches() {
spawn_native_execution_thread("native-test".into(), STACK, || {
let worker = thread::current().id();
for expected in 1..=3 {
on_native_stack(STACK, || {
assert_eq!(thread::current().id(), worker);
BATCHES.set(BATCHES.get() + 1);
assert_eq!(BATCHES.get(), expected);
});
}
assert_eq!(BATCHES.get(), 3);
})
.unwrap()
.join()
.unwrap();
}

#[test]
fn unregistered_caller_keeps_scoped_thread_fallback() {
let caller = thread::current().id();
let mut result = 0;
for _ in 0..2 {
on_native_stack(STACK, || {
assert_ne!(thread::current().id(), caller);
assert_eq!(BATCHES.get(), 0);
BATCHES.set(1);
result += 1;
});
}
assert_eq!(result, 2);
assert_eq!(NATIVE_STACK_SIZE.get(), 0);
}

#[test]
fn larger_stack_request_falls_back_without_losing_worker_locals() {
spawn_native_execution_thread("native-small".into(), STACK, || {
let worker = thread::current().id();
BATCHES.set(7);
on_native_stack(STACK * 2, || {
assert_ne!(thread::current().id(), worker);
assert_eq!(BATCHES.get(), 0);
});
on_native_stack(STACK, || assert_eq!(BATCHES.get(), 7));
})
.unwrap()
.join()
.unwrap();
}

#[test]
fn registration_is_not_inherited_by_unregistered_child_threads() {
spawn_native_execution_thread("native-parent".into(), STACK, || {
thread::spawn(|| {
assert_eq!(NATIVE_STACK_SIZE.get(), 0);
let child = thread::current().id();
on_native_stack(STACK, || assert_ne!(thread::current().id(), child));
})
.join()
.unwrap();
})
.unwrap()
.join()
.unwrap();
}

#[test]
fn worker_panic_reaches_joiner_and_replacement_starts_cold() {
let failed = spawn_native_execution_thread("native-panic".into(), STACK, || {
on_native_stack(STACK, || {
BATCHES.set(9);
panic!("execution failed");
});
})
.unwrap()
.join();
assert!(failed.is_err());
spawn_native_execution_thread("native-replacement".into(), STACK, || {
on_native_stack(STACK, || assert_eq!(BATCHES.get(), 0));
})
.unwrap()
.join()
.unwrap();
}

#[test]
fn fallback_panic_reaches_caller() {
assert!(std::panic::catch_unwind(|| on_native_stack(STACK, || panic!("failed"))).is_err());
}
}
23 changes: 5 additions & 18 deletions crates/blockifier/src/blockifier/transaction_executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -330,24 +330,11 @@ impl<S: StateReader + Send> TransactionExecutor<S> {
// TODO(meshi): find a way to access the contract class manager config from transaction
// executor.
let txs = txs.to_vec();
std::thread::scope(|s| {
std::thread::Builder::new()
// when running Cairo natively, the real stack is used and could get overflowed
// (unlike the VM where the stack is simulated in the heap as a memory segment).
//
// We pre-allocate the stack here, and not during Native execution (not trivial), so it
// needs to be big enough ahead.
// However, making it very big is wasteful (especially with multi-threading).
// So, the stack size should support calls with a reasonable gas limit, for extremely deep
// recursions to reach out-of-gas before hitting the bottom of the recursion.
//
// The gas upper bound is MAX_POSSIBLE_SIERRA_GAS, and sequencers must not raise it without
// adjusting the stack size.
.stack_size(self.config.stack_size)
.spawn_scoped(s, || self.execute_txs_sequentially_inner(&txs, execution_deadline))
.expect("Failed to spawn thread")
.join()
.expect("Failed to join thread.")
// Native uses the real stack. Reuse only a registered execution worker with at
// least the configured stack; other callers still get a scoped thread. The stack
// must cover MAX_POSSIBLE_SIERRA_GAS, including deep recursion before out-of-gas.
super::native_execution_thread::on_native_stack(self.config.stack_size, || {
self.execute_txs_sequentially_inner(&txs, execution_deadline)
})
}
}
Expand Down
Loading