1use std::{collections::BTreeSet, sync::Arc};
6
7use either::Either;
8use fastcrypto::traits::{AggregateAuthenticator, ToFromBytes};
9use futures::pin_mut;
10use iota_metrics::monitored_scope;
11use iota_sdk_types::{
12 CertificateDigest, SenderSignedDataDigest, SenderSignedTransaction, crypto::Intent,
13};
14use iota_types::{
15 base_types::AuthorityName,
16 committee::Committee,
17 crypto::{AggregateAuthorityPublicKey, AuthoritySignInfoTrait, VerificationObligation},
18 error::{IotaError, IotaResult},
19 message_envelope::Message,
20 messages_checkpoint::{CheckpointSummaryExt, SignedCheckpointSummary},
21 messages_consensus::{AuthorityCapabilitiesDigest, SignedAuthorityCapabilitiesV1},
22 signature::VerifyParams,
23 signature_verification::{VerifiedDigestCache, verify_sender_signed_data_message_signatures},
24 transaction::{CertifiedTransaction, VerifiedCertificate},
25};
26use itertools::{Itertools as _, izip};
27use parking_lot::{Mutex, MutexGuard};
28use prometheus_filtered::{IntCounter, Registry, register_int_counter_with_registry};
29use tap::TapFallible;
30use tokio::{
31 runtime::Handle,
32 sync::oneshot,
33 time::{Duration, timeout},
34};
35use tracing::{Instrument, instrument, trace_span};
36const BATCH_TIMEOUT_MS: Duration = Duration::from_millis(10);
39
40const MAX_BATCH_SIZE: usize = 8;
46
47type Sender = oneshot::Sender<IotaResult<VerifiedCertificate>>;
48
49struct CertBuffer {
50 certs: Vec<CertifiedTransaction>,
51 senders: Vec<Sender>,
52 id: u64,
53}
54
55impl CertBuffer {
56 fn new(capacity: usize) -> Self {
57 Self {
58 certs: Vec::with_capacity(capacity),
59 senders: Vec::with_capacity(capacity),
60 id: 0,
61 }
62 }
63
64 fn take_and_replace(mut guard: MutexGuard<'_, Self>) -> Self {
67 let this = &mut *guard;
68 let mut new = CertBuffer::new(this.capacity());
69 new.id = this.id + 1;
70 std::mem::swap(&mut new, this);
71 new
72 }
73
74 fn capacity(&self) -> usize {
75 debug_assert_eq!(self.certs.capacity(), self.senders.capacity());
76 self.certs.capacity()
77 }
78
79 fn len(&self) -> usize {
80 debug_assert_eq!(self.certs.len(), self.senders.len());
81 self.certs.len()
82 }
83
84 fn push(&mut self, tx: Sender, cert: CertifiedTransaction) {
85 self.senders.push(tx);
86 self.certs.push(cert);
87 }
88}
89
90pub struct SignatureVerifier {
95 committee: Arc<Committee>,
96 non_committee_validators: BTreeSet<AuthorityName>,
97
98 certificate_cache: VerifiedDigestCache<CertificateDigest>,
99 signed_data_cache: VerifiedDigestCache<SenderSignedDataDigest>,
100 authority_capability_cache: VerifiedDigestCache<AuthorityCapabilitiesDigest>,
101
102 verify_params: VerifyParams,
104
105 queue: Mutex<CertBuffer>,
106 pub metrics: Arc<SignatureVerifierMetrics>,
107}
108
109impl SignatureVerifier {
110 pub fn new_with_batch_size(
111 committee: Arc<Committee>,
112 non_committee_validators: BTreeSet<AuthorityName>,
113 batch_size: usize,
114 metrics: Arc<SignatureVerifierMetrics>,
115 accept_passkey_in_multisig: bool,
116 additional_multisig_checks: bool,
117 ) -> Self {
118 Self {
119 committee,
120 non_committee_validators,
121 certificate_cache: VerifiedDigestCache::new(
122 metrics.certificate_signatures_cache_hits.clone(),
123 metrics.certificate_signatures_cache_misses.clone(),
124 metrics.certificate_signatures_cache_evictions.clone(),
125 ),
126 signed_data_cache: VerifiedDigestCache::new(
127 metrics.signed_data_cache_hits.clone(),
128 metrics.signed_data_cache_misses.clone(),
129 metrics.signed_data_cache_evictions.clone(),
130 ),
131 authority_capability_cache: VerifiedDigestCache::new(
132 metrics.authority_capabilities_cache_hits.clone(),
133 metrics.authority_capabilities_cache_misses.clone(),
134 metrics.authority_capabilities_cache_evictions.clone(),
135 ),
136 queue: Mutex::new(CertBuffer::new(batch_size)),
137 metrics,
138 verify_params: VerifyParams::new(
139 accept_passkey_in_multisig,
140 additional_multisig_checks,
141 ),
142 }
143 }
144
145 pub fn new(
146 committee: Arc<Committee>,
147 non_committee_validators: BTreeSet<AuthorityName>,
148 metrics: Arc<SignatureVerifierMetrics>,
149 accept_passkey_in_multisig: bool,
150 additional_multisig_checks: bool,
151 ) -> Self {
152 Self::new_with_batch_size(
153 committee,
154 non_committee_validators,
155 MAX_BATCH_SIZE,
156 metrics,
157 accept_passkey_in_multisig,
158 additional_multisig_checks,
159 )
160 }
161
162 #[instrument(level = "trace", skip_all)]
164 pub fn verify_certs_and_checkpoints(
165 &self,
166 certs: Vec<&CertifiedTransaction>,
167 checkpoints: Vec<&SignedCheckpointSummary>,
168 authority_capabilities: Vec<&SignedAuthorityCapabilitiesV1>,
169 ) -> IotaResult {
170 for cert in &certs {
173 self.verify_tx(cert.data())?;
174 }
175
176 for cap in &authority_capabilities {
179 self.verify_authority_capabilities(cap)?;
180 }
181
182 batch_verify_all_certificates_and_checkpoints(&self.committee, &certs, &checkpoints)?;
183 Ok(())
184 }
185
186 pub async fn verify_cert(&self, cert: CertifiedTransaction) -> IotaResult<VerifiedCertificate> {
188 let cert_digest = cert.certificate_digest();
189 if self.certificate_cache.is_cached(&cert_digest) {
190 return Ok(VerifiedCertificate::new_unchecked(cert));
191 }
192 self.verify_tx(cert.data())?;
193 self.verify_cert_skip_cache(cert)
194 .await
195 .tap_ok(|_| self.certificate_cache.cache_digest(cert_digest))
196 }
197
198 pub async fn multi_verify_certs(
199 &self,
200 certs: Vec<CertifiedTransaction>,
201 ) -> Vec<IotaResult<VerifiedCertificate>> {
202 let mut futures = Vec::with_capacity(certs.len());
205 for cert in certs {
206 futures.push(self.verify_cert(cert));
207 }
208 futures::future::join_all(futures).await
209 }
210
211 pub async fn verify_cert_skip_cache(
213 &self,
214 cert: CertifiedTransaction,
215 ) -> IotaResult<VerifiedCertificate> {
216 if cert.auth_sig().epoch != self.committee.epoch() {
219 return Err(IotaError::WrongEpoch {
220 expected_epoch: self.committee.epoch(),
221 actual_epoch: cert.auth_sig().epoch,
222 });
223 }
224
225 self.verify_cert_inner(cert).await
226 }
227
228 async fn verify_cert_inner(
229 &self,
230 cert: CertifiedTransaction,
231 ) -> IotaResult<VerifiedCertificate> {
232 let (tx, rx) = oneshot::channel();
237 pin_mut!(rx);
238
239 let prev_id_or_buffer = {
240 let mut queue = self.queue.lock();
241 queue.push(tx, cert);
242 if queue.len() == queue.capacity() {
243 Either::Right(CertBuffer::take_and_replace(queue))
244 } else {
245 Either::Left(queue.id)
246 }
247 };
248 let prev_id = match prev_id_or_buffer {
249 Either::Left(prev_id) => prev_id,
250 Either::Right(buffer) => {
251 self.metrics.full_batches.inc();
252 self.process_queue(buffer)
253 .instrument(trace_span!("SignatureVerifier::process_queue"))
254 .await;
255 return rx.try_recv().unwrap();
257 }
258 };
259
260 if let Ok(res) = timeout(BATCH_TIMEOUT_MS, &mut rx).await {
261 return res.unwrap();
263 }
264 self.metrics.timeouts.inc();
265
266 let buffer = {
267 let queue = self.queue.lock();
268 if prev_id == queue.id {
270 debug_assert_ne!(queue.len(), queue.capacity());
271 Some(CertBuffer::take_and_replace(queue))
272 } else {
273 None
274 }
275 };
276
277 if let Some(buffer) = buffer {
278 self.metrics.partial_batches.inc();
279 self.process_queue(buffer).await;
280 return rx.try_recv().unwrap();
282 }
283
284 rx.await.unwrap()
287 }
288
289 async fn process_queue(&self, buffer: CertBuffer) {
290 let committee = self.committee.clone();
291 let metrics = self.metrics.clone();
292 Handle::current()
293 .spawn_blocking(move || Self::process_queue_sync(committee, metrics, buffer))
294 .await
295 .expect("Spawn blocking should not fail");
296 }
297
298 #[instrument(level = "trace", skip_all)]
299 fn process_queue_sync(
300 committee: Arc<Committee>,
301 metrics: Arc<SignatureVerifierMetrics>,
302 buffer: CertBuffer,
303 ) {
304 let _scope = monitored_scope("BatchCertificateVerifier::process_queue");
305
306 let results = batch_verify_certificates(&committee, &buffer.certs.iter().collect_vec());
307 izip!(
308 results.into_iter(),
309 buffer.certs.into_iter(),
310 buffer.senders.into_iter(),
311 )
312 .for_each(|(result, cert, tx)| {
313 tx.send(match result {
314 Ok(()) => {
315 metrics.total_verified_certs.inc();
316 Ok(VerifiedCertificate::new_unchecked(cert))
317 }
318 Err(e) => {
319 metrics.total_failed_certs.inc();
320 Err(e)
321 }
322 })
323 .ok();
324 });
325 }
326
327 #[instrument(level = "trace", skip_all, fields(tx_digest = ?signed_tx.digest()))]
328 pub fn verify_tx(&self, signed_tx: &SenderSignedTransaction) -> IotaResult {
329 self.signed_data_cache.is_verified(
330 signed_tx.full_message_digest(),
331 || verify_sender_signed_data_message_signatures(signed_tx, &self.verify_params),
332 || Ok(()),
333 )
334 }
335
336 #[instrument(level = "trace", skip_all)]
337 pub fn verify_authority_capabilities(
338 &self,
339 signed_authority_capabilities: &SignedAuthorityCapabilitiesV1,
340 ) -> IotaResult {
341 let epoch = self.committee.epoch();
342 self.authority_capability_cache.is_verified(
343 signed_authority_capabilities.cache_digest(epoch),
344 || {
345 let authority_name = signed_authority_capabilities.data().authority;
347 if !self.non_committee_validators.contains(&authority_name) {
348 return Err(IotaError::IncorrectSigner {
349 error: "Signer must be part of non-committee active validators".to_string(),
350 });
351 }
352
353 let mut obligation = VerificationObligation::default();
355 let idx = obligation.add_message(
356 signed_authority_capabilities.data(),
357 epoch, Intent::iota_app(signed_authority_capabilities.scope()),
360 );
361
362 let authority_key = AggregateAuthorityPublicKey::from_bytes(
364 authority_name.as_bytes(),
365 )
366 .map_err(|_| IotaError::IncorrectSigner {
367 error: "Invalid authority public key bytes".to_string(),
368 })?;
369 obligation
370 .public_keys
371 .get_mut(idx)
372 .ok_or(IotaError::InvalidAuthenticator)?
373 .push(&authority_key);
374
375 obligation
376 .signatures
377 .get_mut(idx)
378 .ok_or(IotaError::InvalidAuthenticator)?
379 .add_signature(signed_authority_capabilities.auth_sig().clone())
380 .map_err(|_| IotaError::InvalidSignature {
381 error: "Failed to add authority signature to obligation".to_string(),
382 })?;
383
384 obligation.verify_all()
385 },
386 || Ok(()),
387 )
388 }
389
390 pub fn clear_signature_cache(&self) {
391 self.certificate_cache.clear();
392 self.authority_capability_cache.clear();
393 self.signed_data_cache.clear();
394 }
395}
396
397pub struct SignatureVerifierMetrics {
398 pub certificate_signatures_cache_hits: IntCounter,
399 pub certificate_signatures_cache_misses: IntCounter,
400 pub certificate_signatures_cache_evictions: IntCounter,
401 pub signed_data_cache_hits: IntCounter,
402 pub signed_data_cache_misses: IntCounter,
403 pub signed_data_cache_evictions: IntCounter,
404 pub authority_capabilities_cache_hits: IntCounter,
405 pub authority_capabilities_cache_misses: IntCounter,
406 pub authority_capabilities_cache_evictions: IntCounter,
407 timeouts: IntCounter,
408 full_batches: IntCounter,
409 partial_batches: IntCounter,
410 total_verified_certs: IntCounter,
411 total_failed_certs: IntCounter,
412}
413
414impl SignatureVerifierMetrics {
415 pub fn new(registry: &Registry) -> Arc<Self> {
416 Arc::new(Self {
417 certificate_signatures_cache_hits: register_int_counter_with_registry!(
418 "certificate_signatures_cache_hits",
419 "Number of certificates which were known to be verified because of signature cache.",
420 registry
421 )
422 .unwrap(),
423 certificate_signatures_cache_misses: register_int_counter_with_registry!(
424 "certificate_signatures_cache_misses",
425 "Number of certificates which missed the signature cache",
426 registry
427 )
428 .unwrap(),
429 certificate_signatures_cache_evictions: register_int_counter_with_registry!(
430 "certificate_signatures_cache_evictions",
431 "Number of times we evict a pre-existing key were known to be verified because of signature cache.",
432 registry
433 )
434 .unwrap(),
435 signed_data_cache_hits: register_int_counter_with_registry!(
436 "signed_data_cache_hits",
437 "Number of signed data which were known to be verified because of signature cache.",
438 registry
439 )
440 .unwrap(),
441 signed_data_cache_misses: register_int_counter_with_registry!(
442 "signed_data_cache_misses",
443 "Number of signed data which missed the signature cache.",
444 registry
445 )
446 .unwrap(),
447 signed_data_cache_evictions: register_int_counter_with_registry!(
448 "signed_data_cache_evictions",
449 "Number of times we evict a pre-existing signed data were known to be verified because of signature cache.",
450 registry
451 )
452 .unwrap(),
453 authority_capabilities_cache_hits: register_int_counter_with_registry!(
454 "authority_capabilities_cache_hits",
455 "Number of authority capabilities which were known to be verified because of capabilities cache.",
456 registry
457 )
458 .unwrap(),
459 authority_capabilities_cache_misses: register_int_counter_with_registry!(
460 "authority_capabilities_cache_misses",
461 "Number of authority capabilities which missed the capabilities cache.",
462 registry
463 )
464 .unwrap(),
465 authority_capabilities_cache_evictions: register_int_counter_with_registry!(
466 "authority_capabilities_cache_evictions",
467 "Number of times we evict a pre-existing authority capabilities that were known to be verified.",
468 registry
469 )
470 .unwrap(),
471 timeouts: register_int_counter_with_registry!(
472 "async_batch_verifier_timeouts",
473 "Number of times batch verifier times out and verifies a partial batch",
474 registry
475 )
476 .unwrap(),
477 full_batches: register_int_counter_with_registry!(
478 "async_batch_verifier_full_batches",
479 "Number of times batch verifier verifies a full batch",
480 registry
481 )
482 .unwrap(),
483 partial_batches: register_int_counter_with_registry!(
484 "async_batch_verifier_partial_batches",
485 "Number of times batch verifier verifies a partial batch",
486 registry
487 )
488 .unwrap(),
489 total_verified_certs: register_int_counter_with_registry!(
490 "async_batch_verifier_total_verified_certs",
491 "Total number of certs batch verifier has verified",
492 registry
493 )
494 .unwrap(),
495 total_failed_certs: register_int_counter_with_registry!(
496 "async_batch_verifier_total_failed_certs",
497 "Total number of certs batch verifier has rejected",
498 registry
499 )
500 .unwrap(),
501 })
502 }
503}
504
505#[instrument(level = "trace", skip_all)]
507pub fn batch_verify_all_certificates_and_checkpoints(
508 committee: &Committee,
509 certs: &[&CertifiedTransaction],
510 checkpoints: &[&SignedCheckpointSummary],
511) -> IotaResult {
512 for ckpt in checkpoints {
515 ckpt.data().verify_epoch(committee.epoch())?;
516 }
517
518 batch_verify(committee, certs, checkpoints)
519}
520
521#[instrument(level = "trace", skip_all)]
527pub fn batch_verify_certificates(
528 committee: &Committee,
529 certs: &[&CertifiedTransaction],
530) -> Vec<IotaResult> {
531 match batch_verify(committee, certs, &[]) {
532 Ok(_) => vec![Ok(()); certs.len()],
533
534 Err(_) if certs.len() > 1 => certs
536 .iter()
537 .map(|c| c.verify_committee_sigs_only(committee))
538 .collect(),
539
540 Err(e) => vec![Err(e)],
541 }
542}
543
544fn batch_verify(
545 committee: &Committee,
546 certs: &[&CertifiedTransaction],
547 checkpoints: &[&SignedCheckpointSummary],
548) -> IotaResult {
549 let mut obligation = VerificationObligation::default();
550
551 for cert in certs {
552 let idx = obligation.add_message(cert.data(), cert.epoch(), Intent::iota_app(cert.scope()));
553 cert.auth_sig()
554 .add_to_verification_obligation(committee, &mut obligation, idx)?;
555 }
556
557 for ckpt in checkpoints {
558 let idx = obligation.add_message(ckpt.data(), ckpt.epoch(), Intent::iota_app(ckpt.scope()));
559 ckpt.auth_sig()
560 .add_to_verification_obligation(committee, &mut obligation, idx)?;
561 }
562
563 obligation.verify_all()
564}