ledger_mempool.rs (22934B)
1 use std::{cmp::Ordering, collections::BTreeSet}; 2 3 use anyhow::{Context, Result, bail}; 4 use serde::Serialize; 5 6 use super::ledger_ops::{ 7 apply_transaction, compact_block_context, ensure_transaction_fits_empty_block, 8 transaction_has_missing_inputs, 9 }; 10 use super::transaction::{Transaction, transaction_inputs_spent_by}; 11 use super::{ 12 Ledger, MAX_ORPHAN_TRANSACTIONS, MAX_PENDING_POOL_BYTES, MAX_PENDING_TRANSACTIONS, 13 TransactionSubmitOutcome, TransactionV2, spend_inputs_with_lineage, 14 }; 15 16 impl Ledger { 17 pub fn submit_transaction(&mut self, transaction: Transaction) -> Result<bool> { 18 Ok(self.submit_transaction_with_outcome(transaction)?.added()) 19 } 20 21 /// Inserts a locally required block transaction ahead of the current 22 /// mempool, then rebuilds the pool around it. This is only used on a cloned 23 /// ledger while assembling a block; the node's public mempool is unchanged. 24 pub(crate) fn prioritize_transaction_for_block_building( 25 &mut self, 26 transaction: Transaction, 27 ) -> Result<bool> { 28 let mut displaced = std::mem::take(&mut self.pending); 29 displaced.append(&mut self.orphans); 30 self.pending_bytes = 0; 31 self.orphan_bytes = 0; 32 33 let added = self.submit_transaction(transaction)?; 34 if !added { 35 return Ok(false); 36 } 37 for candidate in displaced { 38 let _ = self.submit_transaction(candidate); 39 } 40 Ok(true) 41 } 42 43 /// Inserts a required v2 anchor ahead of the current v2 mempool and then 44 /// rebuilds the displaced entries in their original order. Callers use a 45 /// cloned ledger and only publish it after the priority transaction is 46 /// admitted, so a failed replacement cannot mutate live state. 47 pub(crate) fn prioritize_transaction_v2_for_block_building( 48 &mut self, 49 transaction: TransactionV2, 50 ) -> Result<bool> { 51 let displaced = std::mem::take(&mut self.pending_v2); 52 self.pending_v2_bytes = 0; 53 54 if !self.submit_transaction_v2(transaction)?.added() { 55 return Ok(false); 56 } 57 for candidate in displaced { 58 let _ = self.submit_transaction_v2(candidate); 59 } 60 Ok(true) 61 } 62 63 pub(crate) fn reserve_transaction_inputs(&mut self, transaction: &Transaction) -> Result<()> { 64 self.validate_new_transaction(transaction)?; 65 let mut utxos = self.utxos.clone(); 66 let mut utxo_lineage = self.utxo_lineage.clone(); 67 let mut lineage_values = self.lineage_values.clone(); 68 let mut lineage_owners = self.lineage_owners.clone(); 69 spend_inputs_with_lineage( 70 transaction, 71 &mut utxos, 72 &mut utxo_lineage, 73 &mut lineage_values, 74 &mut lineage_owners, 75 )?; 76 self.utxos = utxos; 77 self.utxo_lineage = utxo_lineage; 78 self.lineage_values = lineage_values; 79 self.lineage_owners = lineage_owners; 80 Ok(()) 81 } 82 83 pub fn submit_transaction_with_outcome( 84 &mut self, 85 transaction: Transaction, 86 ) -> Result<TransactionSubmitOutcome> { 87 if self.has_transaction(transaction.signature()) { 88 return Ok(TransactionSubmitOutcome::AlreadyKnown); 89 } 90 91 self.validate_transaction_terms(&transaction)?; 92 self.validate_transaction_anchor_for_pending(&transaction)?; 93 let signing_domain = self.transaction_signing_domain_for_pending(&transaction); 94 transaction.verify_signature(&signing_domain)?; 95 ensure_transaction_fits_empty_block( 96 compact_block_context(self), 97 &transaction, 98 self.launch_profile.max_block_bytes, 99 )?; 100 self.validate_mine_anchor_available(&transaction)?; 101 102 if transaction_inputs_spent_by(&transaction, &self.pending) { 103 return Ok(TransactionSubmitOutcome::ConflictsWithPending); 104 } 105 if transaction_inputs_spent_by(&transaction, &self.orphans) { 106 return Ok(TransactionSubmitOutcome::ConflictsWithPending); 107 } 108 if self.transaction_conflicts_with_pending_v2(&transaction) { 109 return Ok(TransactionSubmitOutcome::ConflictsWithPending); 110 } 111 112 let mut utxos = self.utxos_after_valid_pending()?; 113 if transaction_has_missing_inputs(&transaction, &utxos) { 114 if self.orphans.len() >= MAX_ORPHAN_TRANSACTIONS { 115 bail!("orphan transaction pool is full"); 116 } 117 let candidate_bytes = ensure_pending_pool_bytes( 118 "orphan transaction pool", 119 self.orphan_bytes, 120 &transaction, 121 MAX_PENDING_POOL_BYTES, 122 )?; 123 self.orphans.push(transaction); 124 self.orphan_bytes = self.orphan_bytes.saturating_add(candidate_bytes); 125 return Ok(TransactionSubmitOutcome::Added); 126 } 127 apply_transaction(&transaction, &mut utxos, &signing_domain)?; 128 let candidate_bytes = pending_pool_item_bytes(&transaction)?; 129 let replacement_backup = self.pending_room_required(candidate_bytes).then(|| { 130 ( 131 self.pending.clone(), 132 self.orphans.clone(), 133 self.pending_bytes, 134 self.orphan_bytes, 135 ) 136 }); 137 if let Err(error) = self.make_pending_room(&transaction, candidate_bytes) { 138 self.restore_pending_after_failed_replacement(replacement_backup); 139 return Err(error); 140 } 141 142 // An eviction may remove state the candidate depended on. Components 143 // containing candidate ancestors are protected, but revalidation keeps 144 // this admission step fail-closed if the graph is ever malformed. 145 let mut utxos = match self.utxos_after_valid_pending() { 146 Ok(utxos) => utxos, 147 Err(error) => { 148 self.restore_pending_after_failed_replacement(replacement_backup); 149 return Err(error); 150 } 151 }; 152 if transaction_has_missing_inputs(&transaction, &utxos) { 153 self.restore_pending_after_failed_replacement(replacement_backup); 154 bail!("mempool replacement removed a candidate dependency"); 155 } 156 if let Err(error) = apply_transaction(&transaction, &mut utxos, &signing_domain) { 157 self.restore_pending_after_failed_replacement(replacement_backup); 158 return Err(error); 159 } 160 self.pending.push(transaction); 161 self.pending_bytes = self.pending_bytes.saturating_add(candidate_bytes); 162 if let Err(error) = self.promote_orphan_transactions() { 163 if let Some((pending, orphans, pending_bytes, orphan_bytes)) = replacement_backup { 164 self.pending = pending; 165 self.orphans = orphans; 166 self.pending_bytes = pending_bytes; 167 self.orphan_bytes = orphan_bytes; 168 } 169 return Err(error); 170 } 171 Ok(TransactionSubmitOutcome::Added) 172 } 173 174 fn pending_room_required(&self, candidate_bytes: usize) -> bool { 175 self.pending.len() >= MAX_PENDING_TRANSACTIONS 176 || self 177 .pending_bytes 178 .checked_add(candidate_bytes) 179 .is_none_or(|bytes| bytes > MAX_PENDING_POOL_BYTES) 180 } 181 182 fn restore_pending_after_failed_replacement( 183 &mut self, 184 backup: Option<(Vec<Transaction>, Vec<Transaction>, usize, usize)>, 185 ) { 186 if let Some((pending, orphans, pending_bytes, orphan_bytes)) = backup { 187 self.pending = pending; 188 self.orphans = orphans; 189 self.pending_bytes = pending_bytes; 190 self.orphan_bytes = orphan_bytes; 191 } 192 } 193 194 fn make_pending_room(&mut self, candidate: &Transaction, candidate_bytes: usize) -> Result<()> { 195 if !self.pending_room_required(candidate_bytes) { 196 return Ok(()); 197 } 198 199 let evicted = 200 pending_eviction_plan(&self.pending, &self.orphans, candidate, candidate_bytes)?; 201 self.pending 202 .retain(|transaction| !evicted.contains(transaction.signature())); 203 self.orphans 204 .retain(|transaction| !evicted.contains(transaction.signature())); 205 self.refresh_pending_pool_byte_counters() 206 } 207 208 pub(super) fn refresh_pending_pool_byte_counters(&mut self) -> Result<()> { 209 self.pending_bytes = serialized_pool_len(&self.pending)?; 210 self.orphan_bytes = serialized_pool_len(&self.orphans)?; 211 Ok(()) 212 } 213 } 214 215 #[derive(Debug)] 216 struct EvictionPackage { 217 signatures: BTreeSet<String>, 218 fee: u128, 219 economic_bytes: u128, 220 pending_bytes: usize, 221 pending_count: usize, 222 tie_break: String, 223 } 224 225 fn pending_eviction_plan( 226 pending: &[Transaction], 227 orphans: &[Transaction], 228 candidate: &Transaction, 229 candidate_bytes: usize, 230 ) -> Result<BTreeSet<String>> { 231 if candidate_bytes > MAX_PENDING_POOL_BYTES { 232 bail!("mempool byte limit exceeded"); 233 } 234 235 let by_signature = pending 236 .iter() 237 .enumerate() 238 .map(|(index, transaction)| (transaction.signature().to_string(), index)) 239 .collect::<std::collections::BTreeMap<_, _>>(); 240 let mut edges = vec![Vec::new(); pending.len()]; 241 for (index, transaction) in pending.iter().enumerate() { 242 for input in transaction.inputs() { 243 if let Some(parent) = by_signature.get(&input.outpoint.txid).copied() { 244 edges[index].push(parent); 245 edges[parent].push(index); 246 } 247 } 248 } 249 250 let protected = candidate 251 .inputs() 252 .iter() 253 .filter_map(|input| by_signature.get(&input.outpoint.txid).copied()) 254 .collect::<BTreeSet<_>>(); 255 let mut visited = vec![false; pending.len()]; 256 let mut packages = Vec::new(); 257 for start in 0..pending.len() { 258 if visited[start] { 259 continue; 260 } 261 let mut stack = vec![start]; 262 let mut indices = Vec::new(); 263 let mut protects_candidate = false; 264 visited[start] = true; 265 while let Some(index) = stack.pop() { 266 indices.push(index); 267 protects_candidate |= protected.contains(&index); 268 for adjacent in &edges[index] { 269 if !visited[*adjacent] { 270 visited[*adjacent] = true; 271 stack.push(*adjacent); 272 } 273 } 274 } 275 if protects_candidate { 276 continue; 277 } 278 279 let signatures = indices 280 .iter() 281 .map(|index| pending[*index].signature().to_string()) 282 .collect::<BTreeSet<_>>(); 283 let fee = indices 284 .iter() 285 .map(|index| u128::from(pending[*index].fee())) 286 .sum(); 287 let economic_bytes = indices 288 .iter() 289 .map(|index| pending[*index].economic_size_bytes() as u128) 290 .sum(); 291 let pending_bytes = indices.iter().try_fold(0usize, |total, index| { 292 total 293 .checked_add(pending_pool_item_bytes(&pending[*index])?) 294 .context("pending eviction package byte size overflow") 295 })?; 296 let tie_break = signatures.iter().next().cloned().unwrap_or_default(); 297 packages.push(EvictionPackage { 298 signatures, 299 fee, 300 economic_bytes, 301 pending_bytes, 302 pending_count: indices.len(), 303 tie_break, 304 }); 305 } 306 307 packages.sort_by(compare_package_fee_rate); 308 let candidate_fee = u128::from(candidate.fee()); 309 let candidate_economic_bytes = candidate.economic_size_bytes() as u128; 310 let mut remaining_count = pending.len(); 311 let mut remaining_bytes = serialized_pool_len(pending)?; 312 let mut evicted = BTreeSet::new(); 313 for package in packages { 314 let count_fits = remaining_count < MAX_PENDING_TRANSACTIONS; 315 let bytes_fit = remaining_bytes 316 .checked_add(candidate_bytes) 317 .is_some_and(|bytes| bytes <= MAX_PENDING_POOL_BYTES); 318 if count_fits && bytes_fit { 319 extend_with_orphan_descendants(&mut evicted, orphans); 320 return Ok(evicted); 321 } 322 if candidate_fee.saturating_mul(package.economic_bytes) 323 <= package.fee.saturating_mul(candidate_economic_bytes) 324 { 325 break; 326 } 327 remaining_count = remaining_count.saturating_sub(package.pending_count); 328 remaining_bytes = remaining_bytes.saturating_sub(package.pending_bytes); 329 evicted.extend(package.signatures); 330 } 331 332 let count_fits = remaining_count < MAX_PENDING_TRANSACTIONS; 333 let bytes_fit = remaining_bytes 334 .checked_add(candidate_bytes) 335 .is_some_and(|bytes| bytes <= MAX_PENDING_POOL_BYTES); 336 if count_fits && bytes_fit { 337 extend_with_orphan_descendants(&mut evicted, orphans); 338 Ok(evicted) 339 } else { 340 bail!("mempool is full and candidate fee rate does not exceed an evictable package") 341 } 342 } 343 344 fn extend_with_orphan_descendants(evicted: &mut BTreeSet<String>, orphans: &[Transaction]) { 345 loop { 346 let mut changed = false; 347 for orphan in orphans { 348 if !evicted.contains(orphan.signature()) 349 && orphan 350 .inputs() 351 .iter() 352 .any(|input| evicted.contains(&input.outpoint.txid)) 353 { 354 changed |= evicted.insert(orphan.signature().to_string()); 355 } 356 } 357 if !changed { 358 return; 359 } 360 } 361 } 362 363 fn compare_package_fee_rate(left: &EvictionPackage, right: &EvictionPackage) -> Ordering { 364 left.fee 365 .saturating_mul(right.economic_bytes) 366 .cmp(&right.fee.saturating_mul(left.economic_bytes)) 367 .then_with(|| left.fee.cmp(&right.fee)) 368 .then_with(|| left.tie_break.cmp(&right.tie_break)) 369 } 370 371 fn ensure_pending_pool_bytes<T: Serialize>( 372 label: &str, 373 existing_bytes: usize, 374 candidate: &T, 375 max_bytes: usize, 376 ) -> Result<usize> { 377 let candidate_bytes = serialized_len(candidate)?; 378 let total_bytes = existing_bytes 379 .checked_add(candidate_bytes) 380 .context("pending pool byte size overflow")?; 381 if total_bytes > max_bytes { 382 bail!("{label} byte limit exceeded"); 383 } 384 Ok(candidate_bytes) 385 } 386 387 fn serialized_pool_len<T: Serialize>(items: &[T]) -> Result<usize> { 388 items.iter().try_fold(0usize, |total, item| { 389 total 390 .checked_add(serialized_len(item)?) 391 .context("pending pool byte size overflow") 392 }) 393 } 394 395 pub(super) fn pending_pool_item_bytes<T: Serialize>(item: &T) -> Result<usize> { 396 serialized_len(item) 397 } 398 399 fn serialized_len<T: Serialize>(item: &T) -> Result<usize> { 400 serde_json::to_vec(item) 401 .context("failed to serialize pending item for size check") 402 .map(|bytes| bytes.len()) 403 } 404 405 #[cfg(test)] 406 mod tests { 407 use std::collections::BTreeMap; 408 409 use super::*; 410 use crate::domain::transaction::{UnsignedTxInput, UnsignedUtxoTransaction}; 411 use crate::domain::{OutPoint, TxOutput, Wallet}; 412 413 fn dummy_mine(recipient: &str, signature_digit: char) -> Transaction { 414 Transaction::Mine { 415 recipient: recipient.to_string(), 416 anchor: "a".repeat(64), 417 salt: 1, 418 nonce: 1, 419 difficulty_bits: 10, 420 proof_header: None, 421 signature: signature_digit.to_string().repeat(64), 422 } 423 } 424 425 #[test] 426 fn pending_pool_byte_limit_accepts_the_boundary_and_rejects_one_byte_over() { 427 let item = "bounded-item"; 428 let item_bytes = serialized_len(&item).unwrap(); 429 430 assert_eq!( 431 ensure_pending_pool_bytes( 432 "test pool", 433 MAX_PENDING_POOL_BYTES - item_bytes, 434 &item, 435 MAX_PENDING_POOL_BYTES, 436 ) 437 .unwrap(), 438 item_bytes 439 ); 440 assert!( 441 ensure_pending_pool_bytes( 442 "test pool", 443 MAX_PENDING_POOL_BYTES - item_bytes + 1, 444 &item, 445 MAX_PENDING_POOL_BYTES, 446 ) 447 .unwrap_err() 448 .to_string() 449 .contains("byte limit exceeded") 450 ); 451 } 452 453 #[test] 454 fn full_mempool_accepts_a_higher_fee_independent_transaction() { 455 let alice = Wallet::from_seed("pending-limit-alice"); 456 let bob = Wallet::from_seed("pending-limit-bob"); 457 let mut ledger = Ledger::new( 458 BTreeMap::from([ 459 (alice.address().to_string(), 10_000_000), 460 (bob.address().to_string(), 10), 461 ]), 462 1, 463 ); 464 let candidate = ledger 465 .build_transfer(&alice, bob.address(), 1, 5_000_000) 466 .unwrap(); 467 ledger.pending = (0..MAX_PENDING_TRANSACTIONS) 468 .map(|index| { 469 let digit = char::from_digit((index % 15 + 1) as u32, 16).unwrap(); 470 let mut transaction = dummy_mine(bob.address(), digit); 471 if let Transaction::Mine { signature, .. } = &mut transaction { 472 *signature = format!("{index:064x}"); 473 } 474 transaction 475 }) 476 .collect(); 477 ledger.refresh_pending_pool_byte_counters().unwrap(); 478 479 assert_eq!(ledger.pending.len(), 10_000); 480 assert!(ledger.submit_transaction(candidate.clone()).unwrap()); 481 assert_eq!(ledger.pending.len(), 10_000); 482 assert!(ledger.has_transaction(candidate.signature())); 483 assert_eq!( 484 ledger.pending_bytes, 485 serialized_pool_len(&ledger.pending).unwrap() 486 ); 487 assert_eq!(ledger.orphan_bytes, 0); 488 } 489 490 #[test] 491 fn full_mempool_rejects_a_lower_fee_candidate_without_mutation() { 492 let alice = Wallet::from_seed("lower-fee-candidate-alice"); 493 let bob = Wallet::from_seed("lower-fee-candidate-bob"); 494 let mut ledger = Ledger::new( 495 BTreeMap::from([ 496 (alice.address().to_string(), 10), 497 (bob.address().to_string(), 10), 498 ]), 499 1, 500 ); 501 let candidate = ledger.build_transfer(&alice, bob.address(), 1, 1).unwrap(); 502 ledger.pending = (0..MAX_PENDING_TRANSACTIONS) 503 .map(|index| { 504 let mut transaction = dummy_mine(bob.address(), 'd'); 505 if let Transaction::Mine { signature, .. } = &mut transaction { 506 *signature = format!("{index:064x}"); 507 } 508 transaction 509 }) 510 .collect(); 511 ledger.refresh_pending_pool_byte_counters().unwrap(); 512 let before = ledger.pending.clone(); 513 let before_bytes = ledger.pending_bytes; 514 515 let error = ledger.submit_transaction(candidate).unwrap_err(); 516 517 assert!(error.to_string().contains("candidate fee rate")); 518 assert_eq!(ledger.pending, before); 519 assert_eq!(ledger.pending_bytes, before_bytes); 520 assert!(ledger.orphans.is_empty()); 521 assert_eq!(ledger.orphan_bytes, 0); 522 } 523 524 #[test] 525 fn eviction_package_includes_pending_and_orphan_descendants() { 526 let wallet = Wallet::from_seed("eviction-package-wallet"); 527 let parent_signature = "1".repeat(128); 528 let child_signature = "2".repeat(128); 529 let orphan_signature = "3".repeat(128); 530 let parent = synthetic_dependency_transaction( 531 wallet.address(), 532 "a".repeat(64), 533 parent_signature.clone(), 534 1, 535 ); 536 let child = synthetic_dependency_transaction( 537 wallet.address(), 538 parent_signature, 539 child_signature.clone(), 540 1, 541 ); 542 let orphan = synthetic_dependency_transaction( 543 wallet.address(), 544 child_signature, 545 orphan_signature.clone(), 546 u64::MAX, 547 ); 548 let candidate = synthetic_dependency_transaction( 549 wallet.address(), 550 "b".repeat(64), 551 "f".repeat(128), 552 10_000, 553 ); 554 let pending = vec![parent, child]; 555 let candidate_bytes = pending_pool_item_bytes(&candidate).unwrap(); 556 let mut padded = pending.clone(); 557 padded.extend( 558 (pending.len()..MAX_PENDING_TRANSACTIONS).map(|_| dummy_mine(wallet.address(), 'e')), 559 ); 560 561 let mut ledger = Ledger::new(BTreeMap::new(), 1); 562 ledger.pending = padded; 563 ledger.orphans = vec![orphan]; 564 ledger.refresh_pending_pool_byte_counters().unwrap(); 565 566 ledger 567 .make_pending_room(&candidate, candidate_bytes) 568 .unwrap(); 569 570 assert!(!ledger.has_transaction(&"1".repeat(128))); 571 assert!(!ledger.has_transaction(&"2".repeat(128))); 572 assert!(!ledger.has_transaction(&orphan_signature)); 573 assert_eq!( 574 ledger.pending_bytes, 575 serialized_pool_len(&ledger.pending).unwrap() 576 ); 577 assert_eq!( 578 ledger.orphan_bytes, 579 serialized_pool_len(&ledger.orphans).unwrap() 580 ); 581 } 582 583 #[test] 584 fn orphan_pool_rejects_transaction_after_1024_items() { 585 let wallet = Wallet::from_seed("orphan-limit-wallet"); 586 let mut ledger = Ledger::new(BTreeMap::from([(wallet.address().to_string(), 10)]), 1); 587 let orphan = UnsignedUtxoTransaction::Transfer { 588 inputs: vec![UnsignedTxInput { 589 outpoint: OutPoint { 590 txid: "b".repeat(64), 591 index: 0, 592 }, 593 owner: wallet.address().to_string(), 594 }], 595 outputs: vec![TxOutput { 596 address: wallet.address().to_string(), 597 amount: 1, 598 }], 599 fee: 1, 600 } 601 .sign(&wallet, &ledger.transaction_signing_domain()) 602 .unwrap(); 603 ledger.orphans = vec![dummy_mine(wallet.address(), 'e'); MAX_ORPHAN_TRANSACTIONS]; 604 605 assert_eq!(ledger.orphans.len(), 1_024); 606 assert!( 607 ledger 608 .submit_transaction(orphan) 609 .unwrap_err() 610 .to_string() 611 .contains("orphan transaction pool is full") 612 ); 613 } 614 615 fn synthetic_dependency_transaction( 616 owner: &str, 617 parent: String, 618 signature: String, 619 fee: u64, 620 ) -> Transaction { 621 Transaction::Transfer { 622 inputs: vec![super::super::TxInput { 623 outpoint: OutPoint { 624 txid: parent, 625 index: 0, 626 }, 627 owner: owner.to_string(), 628 signature: signature.clone(), 629 }], 630 outputs: vec![TxOutput { 631 address: owner.to_string(), 632 amount: 1, 633 }], 634 fee, 635 signature, 636 } 637 } 638 }