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
258 changes: 223 additions & 35 deletions client/rpc-v2/src/privacy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,10 +45,8 @@ use sp_crypto_hashing::{blake2_128, twox_128};
use sp_runtime::traits::Block as BlockT;
use std::{
marker::PhantomData,
sync::{
atomic::{AtomicUsize, Ordering},
Arc,
},
sync::{Arc, Condvar, Mutex, MutexGuard},
time::{Duration, Instant},
};

// ============================================================================
Expand Down Expand Up @@ -159,10 +157,24 @@ pub trait PrivacyApi {
// RPC server
// ============================================================================

/// Merkle proofs computed at once, at most. A sealed tree's path rebuilds
/// pruned siblings from its leaves, so a proof costs real CPU; past this, a
/// request is refused as busy rather than queued behind the others.
const MAX_CONCURRENT_PROOFS: usize = 4;
/// Merkle proofs a request may wait behind before it is refused as busy.
///
/// Each waiting request holds a thread of the RPC's blocking pool (512 by
/// default), so this stays well below it. Full, the queue drains in a fraction
/// of a second; more than this at once is abuse, not load.
const MAX_QUEUED_PROOFS: usize = 128;

/// The longest a request waits for a slot. A full queue drains far sooner, so
/// running out of time means the node itself is saturated.
const MAX_PROOF_WAIT: Duration = Duration::from_secs(2);

/// Merkle proofs computed at once on a machine of `cores` cores: half of them,
/// at least one. A proof is pure CPU (Poseidon over a sealed tree's leaves), and
/// the other half stays free for importing and authoring blocks, GRANDPA,
/// networking and the rest of the RPC.
fn proof_slots(cores: usize) -> usize {
(cores / 2).max(1)
}

/// Privacy RPC server for the Orbinum shielded pool.
pub struct PrivacyRpc<C, B, BE> {
Expand All @@ -175,43 +187,93 @@ impl<C, B, BE> PrivacyRpc<C, B, BE> {
pub fn new(client: Arc<C>) -> Self {
Self {
client,
proofs: ProofGate::new(MAX_CONCURRENT_PROOFS),
proofs: ProofGate::new(
proof_slots(std::thread::available_parallelism().map_or(1, |n| n.get())),
MAX_QUEUED_PROOFS,
MAX_PROOF_WAIT,
),
_ph: PhantomData,
}
}
}

/// Caps how many proofs run at once. A permit is held for the duration of one
/// request and released when dropped.
/// Caps how many proofs run at once, queueing the rest.
///
/// A request takes a free slot, or waits for one behind at most `max_queued`
/// others and for at most `max_wait`; it is refused only when the queue is full
/// or the wait runs out. Both proof RPCs are `blocking`, so each waits on its own
/// thread. A slot is freed when its permit drops, even if the proof panics.
struct ProofGate {
in_flight: AtomicUsize,
max: usize,
state: Mutex<GateState>,
freed: Condvar,
slots: usize,
max_queued: usize,
max_wait: Duration,
}

#[derive(Default)]
struct GateState {
in_flight: usize,
waiting: usize,
}

struct ProofPermit<'a>(&'a AtomicUsize);
struct ProofPermit<'a>(&'a ProofGate);

impl ProofGate {
fn new(max: usize) -> Self {
fn new(slots: usize, max_queued: usize, max_wait: Duration) -> Self {
Self {
in_flight: AtomicUsize::new(0),
max,
state: Mutex::default(),
freed: Condvar::new(),
slots: slots.max(1),
max_queued,
max_wait,
}
}

/// A permit, or `None` when `max` proofs are already running.
/// The state, recovered if a holder panicked: two counters cannot be left
/// inconsistent by a panic between lock and unlock.
fn lock(&self) -> MutexGuard<'_, GateState> {
self.state.lock().unwrap_or_else(|e| e.into_inner())
}

/// A permit, waiting for a slot if needed; `None` when the queue is full or
/// no slot frees up within `max_wait`.
fn enter(&self) -> Option<ProofPermit<'_>> {
self.in_flight
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |n| {
(n < self.max).then_some(n + 1)
})
.ok()
.map(|_| ProofPermit(&self.in_flight))
let mut state = self.lock();
if state.in_flight < self.slots {
state.in_flight += 1;
return Some(ProofPermit(self));
}
if state.waiting >= self.max_queued {
return None;
}
state.waiting += 1;
let deadline = Instant::now() + self.max_wait;
let got_slot = loop {
if state.in_flight < self.slots {
break true;
}
let Some(left) = deadline.checked_duration_since(Instant::now()) else {
break false;
};
state = self
.freed
.wait_timeout(state, left)
.unwrap_or_else(|e| e.into_inner())
.0;
};
state.waiting -= 1;
if got_slot {
state.in_flight += 1;
}
got_slot.then(|| ProofPermit(self))
}
}

impl Drop for ProofPermit<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
self.0.lock().in_flight -= 1;
self.0.freed.notify_one();
}
}

Expand Down Expand Up @@ -243,10 +305,11 @@ fn pool_not_initialized() -> ErrorObject<'static> {
)
}

/// The proof queue is full, or no slot freed up in time.
fn busy() -> ErrorObject<'static> {
ErrorObject::owned(
ErrorCode::ServerIsBusy.code(),
"Too many Merkle proofs in flight, try again later",
"Merkle proof queue is full, try again later",
None::<()>,
)
}
Expand Down Expand Up @@ -501,16 +564,141 @@ mod tests {
mod proof_gate {
use super::*;

use std::{
sync::atomic::{AtomicUsize, Ordering},
thread,
};

fn counters(gate: &ProofGate) -> (usize, usize) {
let s = gate.lock();
(s.in_flight, s.waiting)
}

#[test]
fn half_the_cores_at_least_one() {
for (cores, slots) in [
(0, 1),
(1, 1),
(2, 1),
(3, 1),
(4, 2),
(8, 4),
(16, 8),
(64, 32),
] {
assert_eq!(proof_slots(cores), slots, "{cores} cores");
}
}

#[test]
fn a_free_slot_is_taken_without_waiting() {
let gate = ProofGate::new(2, 0, Duration::ZERO);
let a = gate.enter().expect("first slot");
let b = gate.enter().expect("second slot");
assert_eq!(counters(&gate), (2, 0));
drop((a, b));
assert_eq!(counters(&gate), (0, 0));
}

/// With no queue, a request past the slots is refused at once.
#[test]
fn a_full_queue_refuses_at_once() {
let gate = ProofGate::new(1, 0, Duration::from_secs(60));
let _held = gate.enter().unwrap();
let start = Instant::now();
assert!(gate.enter().is_none());
assert!(
start.elapsed() < Duration::from_millis(100),
"it must not wait"
);
}

#[test]
fn a_waiter_takes_the_slot_a_permit_frees() {
let gate = Arc::new(ProofGate::new(1, 4, Duration::from_secs(10)));
let held = gate.enter().unwrap();
let waiter = {
let gate = gate.clone();
thread::spawn(move || gate.enter().map(drop).is_some())
};
while counters(&gate).1 == 0 {
thread::yield_now();
}
drop(held);
assert!(waiter.join().unwrap(), "the waiter gets the freed slot");
assert_eq!(counters(&gate), (0, 0));
}

#[test]
fn a_wait_past_the_deadline_is_refused() {
let gate = ProofGate::new(1, 4, Duration::from_millis(50));
let _held = gate.enter().unwrap();
let start = Instant::now();
assert!(gate.enter().is_none());
assert!(start.elapsed() >= Duration::from_millis(50));
assert_eq!(counters(&gate), (1, 0), "the waiter leaves the queue");
}

/// A burst far larger than the slots, within the queue: every request is
/// served, never more than `slots` at once.
#[test]
fn a_burst_within_the_queue_is_all_served() {
let gate = Arc::new(ProofGate::new(3, 64, Duration::from_secs(10)));
let running = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let served: usize = (0..64)
.map(|_| {
let (gate, running, peak) = (gate.clone(), running.clone(), peak.clone());
thread::spawn(move || {
let Some(_permit) = gate.enter() else {
return 0;
};
let now = running.fetch_add(1, Ordering::SeqCst) + 1;
peak.fetch_max(now, Ordering::SeqCst);
thread::sleep(Duration::from_millis(2));
running.fetch_sub(1, Ordering::SeqCst);
1
})
})
.collect::<Vec<_>>()
.into_iter()
.map(|h| h.join().unwrap())
.sum();
assert_eq!(served, 64);
assert_eq!(peak.load(Ordering::SeqCst), 3);
assert_eq!(counters(&gate), (0, 0));
}

/// Past slots + queue, the excess is refused and the rest served.
#[test]
fn refuses_past_the_cap_and_frees_a_slot_on_drop() {
let gate = ProofGate::new(2);
let a = gate.enter().expect("first permit");
let b = gate.enter().expect("second permit");
assert!(gate.enter().is_none(), "a third proof must be refused");
drop(a);
let c = gate.enter().expect("a slot frees when a permit drops");
drop((b, c));
assert_eq!(gate.in_flight.load(Ordering::Acquire), 0);
fn a_burst_past_the_queue_refuses_only_the_excess() {
let gate = Arc::new(ProofGate::new(2, 8, Duration::from_secs(10)));
let held: Vec<_> = (0..2).map(|_| gate.enter().unwrap()).collect();
let waiters: Vec<_> = (0..8)
.map(|_| {
let gate = gate.clone();
thread::spawn(move || gate.enter().map(drop).is_some())
})
.collect();
while counters(&gate).1 < 8 {
thread::yield_now();
}
assert!(gate.enter().is_none(), "the queue is full");
drop(held);
assert!(waiters.into_iter().all(|w| w.join().unwrap()));
assert_eq!(counters(&gate), (0, 0));
}

#[test]
fn a_panicking_proof_frees_its_slot() {
let gate = Arc::new(ProofGate::new(1, 0, Duration::ZERO));
let g = gate.clone();
let _ = thread::spawn(move || {
let _permit = g.enter().unwrap();
panic!("proof failed");
})
.join();
assert!(gate.enter().is_some(), "the slot came back");
}

#[test]
Expand Down
7 changes: 4 additions & 3 deletions template/runtime/RUNTIME_VERSIONS.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,9 +92,10 @@ active.**
public RPC, so that cost was what one request could make a node pay. Lowering
the cut is safe whenever it lands: a node pruned under the old cut is rebuilt
on demand.
- **Node:** the `privacy_getMerkleProof*` RPCs run on blocking threads, at most 4
at once; the excess is refused as busy (`-32009`) instead of stalling every
other RPC.
- **Node:** the `privacy_getMerkleProof*` RPCs run on blocking threads, half the
node's cores at once; the rest wait in a queue of 128 for up to 2 s. Only past
that (a full queue or a saturated node) is a request refused as busy
(`-32009`), so a burst no longer stalls every other RPC or gets refused.
- **zk-verifier:** a spend circuit never takes a base-layout key (no memo
binding), at any version: not on registration, activation or unretiring, and a
base key already in storage verifies nothing.
Expand Down
Loading