1use bitcoin::block::Header;
27use bitcoin::hash_types::{BlockHash, Txid};
28
29use bitcoin::secp256k1::PublicKey;
30
31use crate::chain;
32use crate::chain::chaininterface::{BroadcasterInterface, FeeEstimator};
33#[cfg(peer_storage)]
34use crate::chain::channelmonitor::write_chanmon_internal;
35use crate::chain::channelmonitor::{
36 Balance, ChannelMonitor, ChannelMonitorUpdate, MonitorEvent, TransactionOutputs,
37 WithChannelMonitor,
38};
39use crate::chain::transaction::{OutPoint, TransactionData};
40use crate::chain::{BestBlock, ChannelMonitorUpdateStatus, Filter, WatchedOutput};
41use crate::events::{self, Event, EventHandler, ReplayEvent};
42use crate::ln::channel_state::ChannelDetails;
43#[cfg(peer_storage)]
44use crate::ln::msgs::PeerStorage;
45use crate::ln::msgs::{BaseMessageHandler, Init, MessageSendEvent, SendOnlyMessageHandler};
46#[cfg(peer_storage)]
47use crate::ln::our_peer_storage::{DecryptedOurPeerStorage, PeerStorageMonitorHolder};
48use crate::ln::types::ChannelId;
49use crate::prelude::*;
50use crate::sign::ecdsa::EcdsaChannelSigner;
51use crate::sign::{EntropySource, PeerStorageKey, SignerProvider};
52use crate::sync::{Mutex, MutexGuard, RwLock, RwLockReadGuard};
53use crate::types::features::{InitFeatures, NodeFeatures};
54use crate::util::async_poll::{MaybeSend, MaybeSync};
55use crate::util::errors::APIError;
56use crate::util::logger::{Logger, WithContext};
57use crate::util::native_async::FutureSpawner;
58use crate::util::persist::{KVStore, MonitorName, MonitorUpdatingPersisterAsync};
59#[cfg(peer_storage)]
60use crate::util::ser::{VecWriter, Writeable};
61use crate::util::wakers::{Future, Notifier};
62
63use alloc::sync::Arc;
64#[cfg(peer_storage)]
65use core::iter::Cycle;
66use core::ops::Deref;
67use core::sync::atomic::{AtomicUsize, Ordering};
68
69pub trait Persist<ChannelSigner: EcdsaChannelSigner> {
125 fn persist_new_channel(
144 &self, monitor_name: MonitorName, monitor: &ChannelMonitor<ChannelSigner>,
145 ) -> ChannelMonitorUpdateStatus;
146
147 fn update_persisted_channel(
185 &self, monitor_name: MonitorName, monitor_update: Option<&ChannelMonitorUpdate>,
186 monitor: &ChannelMonitor<ChannelSigner>,
187 ) -> ChannelMonitorUpdateStatus;
188 fn archive_persisted_channel(&self, monitor_name: MonitorName);
200
201 #[doc(hidden)]
208 fn get_and_clear_completed_updates(&self) -> Vec<(ChannelId, u64)> {
209 Vec::new()
210 }
211}
212
213struct MonitorHolder<ChannelSigner: EcdsaChannelSigner> {
214 monitor: ChannelMonitor<ChannelSigner>,
215 pending_monitor_updates: Mutex<Vec<u64>>,
230}
231
232impl<ChannelSigner: EcdsaChannelSigner> MonitorHolder<ChannelSigner> {
233 fn has_pending_updates(&self, pending_monitor_updates_lock: &MutexGuard<Vec<u64>>) -> bool {
234 !pending_monitor_updates_lock.is_empty()
235 }
236}
237
238pub struct LockedChannelMonitor<'a, ChannelSigner: EcdsaChannelSigner> {
243 lock: RwLockReadGuard<'a, HashMap<ChannelId, MonitorHolder<ChannelSigner>>>,
244 channel_id: ChannelId,
245}
246
247impl<ChannelSigner: EcdsaChannelSigner> Deref for LockedChannelMonitor<'_, ChannelSigner> {
248 type Target = ChannelMonitor<ChannelSigner>;
249 fn deref(&self) -> &ChannelMonitor<ChannelSigner> {
250 &self.lock.get(&self.channel_id).expect("Checked at construction").monitor
251 }
252}
253
254pub struct AsyncPersister<
259 K: Deref + MaybeSend + MaybeSync + 'static,
260 S: FutureSpawner,
261 L: Deref + MaybeSend + MaybeSync + 'static,
262 ES: Deref + MaybeSend + MaybeSync + 'static,
263 SP: Deref + MaybeSend + MaybeSync + 'static,
264 BI: Deref + MaybeSend + MaybeSync + 'static,
265 FE: Deref + MaybeSend + MaybeSync + 'static,
266> where
267 K::Target: KVStore + MaybeSync,
268 L::Target: Logger,
269 ES::Target: EntropySource + Sized,
270 SP::Target: SignerProvider + Sized,
271 BI::Target: BroadcasterInterface,
272 FE::Target: FeeEstimator,
273{
274 persister: MonitorUpdatingPersisterAsync<K, S, L, ES, SP, BI, FE>,
275 event_notifier: Arc<Notifier>,
276}
277
278impl<
279 K: Deref + MaybeSend + MaybeSync + 'static,
280 S: FutureSpawner,
281 L: Deref + MaybeSend + MaybeSync + 'static,
282 ES: Deref + MaybeSend + MaybeSync + 'static,
283 SP: Deref + MaybeSend + MaybeSync + 'static,
284 BI: Deref + MaybeSend + MaybeSync + 'static,
285 FE: Deref + MaybeSend + MaybeSync + 'static,
286 > Deref for AsyncPersister<K, S, L, ES, SP, BI, FE>
287where
288 K::Target: KVStore + MaybeSync,
289 L::Target: Logger,
290 ES::Target: EntropySource + Sized,
291 SP::Target: SignerProvider + Sized,
292 BI::Target: BroadcasterInterface,
293 FE::Target: FeeEstimator,
294{
295 type Target = Self;
296 fn deref(&self) -> &Self {
297 self
298 }
299}
300
301impl<
302 K: Deref + MaybeSend + MaybeSync + 'static,
303 S: FutureSpawner,
304 L: Deref + MaybeSend + MaybeSync + 'static,
305 ES: Deref + MaybeSend + MaybeSync + 'static,
306 SP: Deref + MaybeSend + MaybeSync + 'static,
307 BI: Deref + MaybeSend + MaybeSync + 'static,
308 FE: Deref + MaybeSend + MaybeSync + 'static,
309 > Persist<<SP::Target as SignerProvider>::EcdsaSigner> for AsyncPersister<K, S, L, ES, SP, BI, FE>
310where
311 K::Target: KVStore + MaybeSync,
312 L::Target: Logger,
313 ES::Target: EntropySource + Sized,
314 SP::Target: SignerProvider + Sized,
315 BI::Target: BroadcasterInterface,
316 FE::Target: FeeEstimator,
317 <SP::Target as SignerProvider>::EcdsaSigner: MaybeSend + 'static,
318{
319 fn persist_new_channel(
320 &self, monitor_name: MonitorName,
321 monitor: &ChannelMonitor<<SP::Target as SignerProvider>::EcdsaSigner>,
322 ) -> ChannelMonitorUpdateStatus {
323 let notifier = Arc::clone(&self.event_notifier);
324 self.persister.spawn_async_persist_new_channel(monitor_name, monitor, notifier);
325 ChannelMonitorUpdateStatus::InProgress
326 }
327
328 fn update_persisted_channel(
329 &self, monitor_name: MonitorName, monitor_update: Option<&ChannelMonitorUpdate>,
330 monitor: &ChannelMonitor<<SP::Target as SignerProvider>::EcdsaSigner>,
331 ) -> ChannelMonitorUpdateStatus {
332 let notifier = Arc::clone(&self.event_notifier);
333 self.persister.spawn_async_update_channel(monitor_name, monitor_update, monitor, notifier);
334 ChannelMonitorUpdateStatus::InProgress
335 }
336
337 fn archive_persisted_channel(&self, monitor_name: MonitorName) {
338 self.persister.spawn_async_archive_persisted_channel(monitor_name);
339 }
340
341 fn get_and_clear_completed_updates(&self) -> Vec<(ChannelId, u64)> {
342 self.persister.get_and_clear_completed_updates()
343 }
344}
345
346pub struct ChainMonitor<
363 ChannelSigner: EcdsaChannelSigner,
364 C: Deref,
365 T: Deref,
366 F: Deref,
367 L: Deref,
368 P: Deref,
369 ES: Deref,
370> where
371 C::Target: chain::Filter,
372 T::Target: BroadcasterInterface,
373 F::Target: FeeEstimator,
374 L::Target: Logger,
375 P::Target: Persist<ChannelSigner>,
376 ES::Target: EntropySource,
377{
378 monitors: RwLock<HashMap<ChannelId, MonitorHolder<ChannelSigner>>>,
379 chain_source: Option<C>,
380 broadcaster: T,
381 logger: L,
382 fee_estimator: F,
383 persister: P,
384 _entropy_source: ES,
385 pending_monitor_events: Mutex<Vec<(OutPoint, ChannelId, Vec<MonitorEvent>, PublicKey)>>,
388 highest_chain_height: AtomicUsize,
390
391 event_notifier: Arc<Notifier>,
394
395 pending_send_only_events: Mutex<Vec<MessageSendEvent>>,
397
398 #[cfg(peer_storage)]
399 our_peerstorage_encryption_key: PeerStorageKey,
400}
401
402impl<
403 K: Deref + MaybeSend + MaybeSync + 'static,
404 S: FutureSpawner,
405 SP: Deref + MaybeSend + MaybeSync + 'static,
406 C: Deref,
407 T: Deref + MaybeSend + MaybeSync + 'static,
408 F: Deref + MaybeSend + MaybeSync + 'static,
409 L: Deref + MaybeSend + MaybeSync + 'static,
410 ES: Deref + MaybeSend + MaybeSync + 'static,
411 >
412 ChainMonitor<
413 <SP::Target as SignerProvider>::EcdsaSigner,
414 C,
415 T,
416 F,
417 L,
418 AsyncPersister<K, S, L, ES, SP, T, F>,
419 ES,
420 > where
421 K::Target: KVStore + MaybeSync,
422 SP::Target: SignerProvider + Sized,
423 C::Target: chain::Filter,
424 T::Target: BroadcasterInterface,
425 F::Target: FeeEstimator,
426 L::Target: Logger,
427 ES::Target: EntropySource + Sized,
428 <SP::Target as SignerProvider>::EcdsaSigner: MaybeSend + 'static,
429{
430 pub fn new_async_beta(
439 chain_source: Option<C>, broadcaster: T, logger: L, feeest: F,
440 persister: MonitorUpdatingPersisterAsync<K, S, L, ES, SP, T, F>, _entropy_source: ES,
441 _our_peerstorage_encryption_key: PeerStorageKey,
442 ) -> Self {
443 let event_notifier = Arc::new(Notifier::new());
444 Self {
445 monitors: RwLock::new(new_hash_map()),
446 chain_source,
447 broadcaster,
448 logger,
449 fee_estimator: feeest,
450 _entropy_source,
451 pending_monitor_events: Mutex::new(Vec::new()),
452 highest_chain_height: AtomicUsize::new(0),
453 event_notifier: Arc::clone(&event_notifier),
454 persister: AsyncPersister { persister, event_notifier },
455 pending_send_only_events: Mutex::new(Vec::new()),
456 #[cfg(peer_storage)]
457 our_peerstorage_encryption_key: _our_peerstorage_encryption_key,
458 }
459 }
460}
461
462impl<
463 ChannelSigner: EcdsaChannelSigner,
464 C: Deref,
465 T: Deref,
466 F: Deref,
467 L: Deref,
468 P: Deref,
469 ES: Deref,
470 > ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
471where
472 C::Target: chain::Filter,
473 T::Target: BroadcasterInterface,
474 F::Target: FeeEstimator,
475 L::Target: Logger,
476 P::Target: Persist<ChannelSigner>,
477 ES::Target: EntropySource,
478{
479 fn process_chain_data<FN>(
491 &self, header: &Header, best_height: Option<u32>, txdata: &TransactionData, process: FN,
492 ) where
493 FN: Fn(&ChannelMonitor<ChannelSigner>, &TransactionData) -> Vec<TransactionOutputs>,
494 {
495 let err_str = "ChannelMonitor[Update] persistence failed unrecoverably. This indicates we cannot continue normal operation and must shut down.";
496 let channel_ids = hash_set_from_iter(self.monitors.read().unwrap().keys().cloned());
497 let channel_count = channel_ids.len();
498 for channel_id in channel_ids.iter() {
499 let monitor_lock = self.monitors.read().unwrap();
500 if let Some(monitor_state) = monitor_lock.get(channel_id) {
501 let update_res = self.update_monitor_with_chain_data(
502 header,
503 best_height,
504 txdata,
505 &process,
506 channel_id,
507 &monitor_state,
508 channel_count,
509 );
510 if update_res.is_err() {
511 core::mem::drop(monitor_lock);
514 let _poison = self.monitors.write().unwrap();
515 log_error!(self.logger, "{}", err_str);
516 panic!("{}", err_str);
517 }
518 }
519 }
520
521 let monitor_states = self.monitors.write().unwrap();
523 for (channel_id, monitor_state) in monitor_states.iter() {
524 if !channel_ids.contains(channel_id) {
525 let update_res = self.update_monitor_with_chain_data(
526 header,
527 best_height,
528 txdata,
529 &process,
530 channel_id,
531 &monitor_state,
532 channel_count,
533 );
534 if update_res.is_err() {
535 log_error!(self.logger, "{}", err_str);
536 panic!("{}", err_str);
537 }
538 }
539 }
540
541 if let Some(height) = best_height {
542 let old_height = self.highest_chain_height.load(Ordering::Acquire);
545 let new_height = height as usize;
546 if new_height > old_height {
547 self.highest_chain_height.store(new_height, Ordering::Release);
548 }
549 }
550 }
551
552 fn update_monitor_with_chain_data<FN>(
553 &self, header: &Header, best_height: Option<u32>, txdata: &TransactionData, process: FN,
554 channel_id: &ChannelId, monitor_state: &MonitorHolder<ChannelSigner>, channel_count: usize,
555 ) -> Result<(), ()>
556 where
557 FN: Fn(&ChannelMonitor<ChannelSigner>, &TransactionData) -> Vec<TransactionOutputs>,
558 {
559 let monitor = &monitor_state.monitor;
560 let logger = WithChannelMonitor::from(&self.logger, &monitor, None);
561
562 let mut txn_outputs = process(monitor, txdata);
563
564 let get_partition_key = |channel_id: &ChannelId| {
565 let channel_id_bytes = channel_id.0;
566 let channel_id_u32 = u32::from_be_bytes([
567 channel_id_bytes[0],
568 channel_id_bytes[1],
569 channel_id_bytes[2],
570 channel_id_bytes[3],
571 ]);
572 best_height.map(|height| channel_id_u32.wrapping_add(height))
573 };
574
575 let partition_factor = if channel_count < 15 {
576 5
577 } else {
578 50 };
580
581 let has_pending_claims = monitor_state.monitor.has_pending_claims();
582 if has_pending_claims
583 || get_partition_key(channel_id).map(|key| key % partition_factor == 0).unwrap_or(false)
584 {
585 log_trace!(
586 logger,
587 "Syncing Channel Monitor for channel {}",
588 log_funding_info!(monitor)
589 );
590 let _pending_monitor_updates = monitor_state.pending_monitor_updates.lock().unwrap();
595 match self.persister.update_persisted_channel(monitor.persistence_key(), None, monitor)
596 {
597 ChannelMonitorUpdateStatus::Completed => log_trace!(
598 logger,
599 "Finished syncing Channel Monitor for channel {} for block-data",
600 log_funding_info!(monitor)
601 ),
602 ChannelMonitorUpdateStatus::InProgress => {
603 log_trace!(
604 logger,
605 "Channel Monitor sync for channel {} in progress.",
606 log_funding_info!(monitor)
607 );
608 },
609 ChannelMonitorUpdateStatus::UnrecoverableError => {
610 return Err(());
611 },
612 }
613 }
614
615 if let Some(ref chain_source) = self.chain_source {
618 let block_hash = header.block_hash();
619 for (txid, mut outputs) in txn_outputs.drain(..) {
620 for (idx, output) in outputs.drain(..) {
621 let output = WatchedOutput {
623 block_hash: Some(block_hash),
624 outpoint: OutPoint { txid, index: idx as u16 },
625 script_pubkey: output.script_pubkey,
626 };
627 log_trace!(
628 logger,
629 "Adding monitoring for spends of outpoint {} to the filter",
630 output.outpoint
631 );
632 chain_source.register_output(output);
633 }
634 }
635 }
636 Ok(())
637 }
638
639 pub fn new(
659 chain_source: Option<C>, broadcaster: T, logger: L, feeest: F, persister: P,
660 _entropy_source: ES, _our_peerstorage_encryption_key: PeerStorageKey,
661 ) -> Self {
662 Self {
663 monitors: RwLock::new(new_hash_map()),
664 chain_source,
665 broadcaster,
666 logger,
667 fee_estimator: feeest,
668 persister,
669 _entropy_source,
670 pending_monitor_events: Mutex::new(Vec::new()),
671 highest_chain_height: AtomicUsize::new(0),
672 event_notifier: Arc::new(Notifier::new()),
673 pending_send_only_events: Mutex::new(Vec::new()),
674 #[cfg(peer_storage)]
675 our_peerstorage_encryption_key: _our_peerstorage_encryption_key,
676 }
677 }
678
679 pub fn get_claimable_balances(&self, ignored_channels: &[&ChannelDetails]) -> Vec<Balance> {
688 let mut ret = Vec::new();
689 let monitor_states = self.monitors.read().unwrap();
690 for (_, monitor_state) in monitor_states.iter().filter(|(channel_id, _)| {
691 for chan in ignored_channels {
692 if chan.channel_id == **channel_id {
693 return false;
694 }
695 }
696 true
697 }) {
698 ret.append(&mut monitor_state.monitor.get_claimable_balances());
699 }
700 ret
701 }
702
703 pub fn get_monitor(
709 &self, channel_id: ChannelId,
710 ) -> Result<LockedChannelMonitor<'_, ChannelSigner>, ()> {
711 let lock = self.monitors.read().unwrap();
712 if lock.get(&channel_id).is_some() {
713 Ok(LockedChannelMonitor { lock, channel_id })
714 } else {
715 Err(())
716 }
717 }
718
719 pub fn list_monitors(&self) -> Vec<ChannelId> {
724 self.monitors.read().unwrap().keys().copied().collect()
725 }
726
727 #[cfg(not(c_bindings))]
728 pub fn list_pending_monitor_updates(&self) -> HashMap<ChannelId, Vec<u64>> {
733 hash_map_from_iter(self.monitors.read().unwrap().iter().map(|(channel_id, holder)| {
734 (*channel_id, holder.pending_monitor_updates.lock().unwrap().clone())
735 }))
736 }
737
738 #[cfg(c_bindings)]
739 pub fn list_pending_monitor_updates(&self) -> Vec<(ChannelId, Vec<u64>)> {
744 let monitors = self.monitors.read().unwrap();
745 monitors
746 .iter()
747 .map(|(channel_id, holder)| {
748 (*channel_id, holder.pending_monitor_updates.lock().unwrap().clone())
749 })
750 .collect()
751 }
752
753 #[cfg(any(test, feature = "_test_utils"))]
754 pub fn remove_monitor(&self, channel_id: &ChannelId) -> ChannelMonitor<ChannelSigner> {
755 self.monitors.write().unwrap().remove(channel_id).unwrap().monitor
756 }
757
758 pub fn channel_monitor_updated(
779 &self, channel_id: ChannelId, completed_update_id: u64,
780 ) -> Result<(), APIError> {
781 let monitors = self.monitors.read().unwrap();
782 let monitor_data = if let Some(mon) = monitors.get(&channel_id) {
783 mon
784 } else {
785 return Err(APIError::APIMisuseError {
786 err: format!("No ChannelMonitor matching channel ID {} found", channel_id),
787 });
788 };
789 let mut pending_monitor_updates = monitor_data.pending_monitor_updates.lock().unwrap();
790 pending_monitor_updates.retain(|update_id| *update_id != completed_update_id);
791
792 let monitor_is_pending_updates = monitor_data.has_pending_updates(&pending_monitor_updates);
795 log_debug!(
796 self.logger,
797 "Completed off-chain monitor update {} for channel with channel ID {}, {}",
798 completed_update_id,
799 channel_id,
800 if monitor_is_pending_updates {
801 "still have pending off-chain updates"
802 } else {
803 "all off-chain updates complete, returning a MonitorEvent"
804 }
805 );
806 if monitor_is_pending_updates {
807 return Ok(());
810 }
811 let funding_txo = monitor_data.monitor.get_funding_txo();
812 self.pending_monitor_events.lock().unwrap().push((
813 funding_txo,
814 channel_id,
815 vec![MonitorEvent::Completed {
816 funding_txo,
817 channel_id,
818 monitor_update_id: monitor_data.monitor.get_latest_update_id(),
819 }],
820 monitor_data.monitor.get_counterparty_node_id(),
821 ));
822
823 self.event_notifier.notify();
824 Ok(())
825 }
826
827 #[cfg(any(test, fuzzing))]
831 pub fn force_channel_monitor_updated(&self, channel_id: ChannelId, monitor_update_id: u64) {
832 let monitors = self.monitors.read().unwrap();
833 let monitor_state = monitors.get(&channel_id).unwrap();
834 let monitor = &monitor_state.monitor;
835 let mut pending_monitor_updates = monitor_state.pending_monitor_updates.lock().unwrap();
836 pending_monitor_updates.retain(|update_id| *update_id > monitor_update_id);
837 let funding_txo = monitor.get_funding_txo();
838 self.pending_monitor_events.lock().unwrap().push((
839 funding_txo,
840 channel_id,
841 vec![MonitorEvent::Completed { funding_txo, channel_id, monitor_update_id }],
842 monitor.get_counterparty_node_id(),
843 ));
844 self.event_notifier.notify();
845 }
846
847 #[cfg(any(test, feature = "_test_utils"))]
848 pub fn get_and_clear_pending_events(&self) -> Vec<events::Event> {
849 use crate::events::EventsProvider;
850 let events = core::cell::RefCell::new(Vec::new());
851 let event_handler = |event: events::Event| Ok(events.borrow_mut().push(event));
852 self.process_pending_events(&event_handler);
853 events.into_inner()
854 }
855
856 pub async fn process_pending_events_async<
863 Future: core::future::Future<Output = Result<(), ReplayEvent>>,
864 H: Fn(Event) -> Future,
865 >(
866 &self, handler: H,
867 ) {
868 let mons_to_process = self.monitors.read().unwrap().keys().cloned().collect::<Vec<_>>();
871 for channel_id in mons_to_process {
872 let mut ev;
873 match super::channelmonitor::process_events_body!(
874 self.monitors.read().unwrap().get(&channel_id).map(|m| &m.monitor),
875 self.logger,
876 ev,
877 handler(ev).await
878 ) {
879 Ok(()) => {},
880 Err(ReplayEvent()) => {
881 self.event_notifier.notify();
882 },
883 }
884 }
885 }
886
887 pub fn get_update_future(&self) -> Future {
896 self.event_notifier.get_future()
897 }
898
899 pub fn rebroadcast_pending_claims(&self) {
905 let monitors = self.monitors.read().unwrap();
906 for (_, monitor_holder) in &*monitors {
907 monitor_holder.monitor.rebroadcast_pending_claims(
908 &*self.broadcaster,
909 &*self.fee_estimator,
910 &self.logger,
911 )
912 }
913 }
914
915 pub fn signer_unblocked(&self, monitor_opt: Option<ChannelId>) {
920 let monitors = self.monitors.read().unwrap();
921 if let Some(channel_id) = monitor_opt {
922 if let Some(monitor_holder) = monitors.get(&channel_id) {
923 monitor_holder.monitor.signer_unblocked(
924 &*self.broadcaster,
925 &*self.fee_estimator,
926 &self.logger,
927 )
928 }
929 } else {
930 for (_, monitor_holder) in &*monitors {
931 monitor_holder.monitor.signer_unblocked(
932 &*self.broadcaster,
933 &*self.fee_estimator,
934 &self.logger,
935 )
936 }
937 }
938 }
939
940 pub fn archive_fully_resolved_channel_monitors(&self) {
950 let mut have_monitors_to_prune = false;
951 for monitor_holder in self.monitors.read().unwrap().values() {
952 let logger = WithChannelMonitor::from(&self.logger, &monitor_holder.monitor, None);
953 let (is_fully_resolved, needs_persistence) =
954 monitor_holder.monitor.check_and_update_full_resolution_status(&logger);
955 if is_fully_resolved {
956 have_monitors_to_prune = true;
957 }
958 if needs_persistence {
959 self.persister.update_persisted_channel(
960 monitor_holder.monitor.persistence_key(),
961 None,
962 &monitor_holder.monitor,
963 );
964 }
965 }
966 if have_monitors_to_prune {
967 let mut monitors = self.monitors.write().unwrap();
968 monitors.retain(|channel_id, monitor_holder| {
969 let logger = WithChannelMonitor::from(&self.logger, &monitor_holder.monitor, None);
970 let (is_fully_resolved, _) =
971 monitor_holder.monitor.check_and_update_full_resolution_status(&logger);
972 if is_fully_resolved {
973 log_info!(
974 logger,
975 "Archiving fully resolved ChannelMonitor for channel ID {}",
976 channel_id
977 );
978 self.persister
979 .archive_persisted_channel(monitor_holder.monitor.persistence_key());
980 false
981 } else {
982 true
983 }
984 });
985 }
986 }
987
988 #[cfg(peer_storage)]
991 fn all_counterparty_node_ids(&self) -> HashSet<PublicKey> {
992 let mon = self.monitors.read().unwrap();
993 mon.values().map(|monitor| monitor.monitor.get_counterparty_node_id()).collect()
994 }
995
996 #[cfg(peer_storage)]
997 fn send_peer_storage(&self, their_node_id: PublicKey) {
998 let mut monitors_list: Vec<PeerStorageMonitorHolder> = Vec::new();
999 let random_bytes = self._entropy_source.get_secure_random_bytes();
1000
1001 const MAX_PEER_STORAGE_SIZE: usize = 65531;
1002 const USIZE_LEN: usize = core::mem::size_of::<usize>();
1003 let mut random_bytes_cycle_iter = random_bytes.iter().cycle();
1004
1005 let mut current_size = 0;
1006 let monitors_lock = self.monitors.read().unwrap();
1007 let mut channel_ids = monitors_lock.keys().copied().collect();
1008
1009 fn next_random_id(
1010 channel_ids: &mut Vec<ChannelId>,
1011 random_bytes_cycle_iter: &mut Cycle<core::slice::Iter<u8>>,
1012 ) -> Option<ChannelId> {
1013 if channel_ids.is_empty() {
1014 return None;
1015 }
1016 let random_idx = {
1017 let mut usize_bytes = [0u8; USIZE_LEN];
1018 usize_bytes.iter_mut().for_each(|b| {
1019 *b = *random_bytes_cycle_iter.next().expect("A cycle never ends")
1020 });
1021 random_bytes_cycle_iter.next().expect("A cycle never ends");
1023 usize::from_le_bytes(usize_bytes) % channel_ids.len()
1024 };
1025 Some(channel_ids.swap_remove(random_idx))
1026 }
1027
1028 while let Some(channel_id) = next_random_id(&mut channel_ids, &mut random_bytes_cycle_iter)
1029 {
1030 let monitor_holder = if let Some(monitor_holder) = monitors_lock.get(&channel_id) {
1031 monitor_holder
1032 } else {
1033 debug_assert!(
1034 false,
1035 "Tried to access non-existing monitor, this should never happen"
1036 );
1037 break;
1038 };
1039
1040 let mut serialized_channel = VecWriter(Vec::new());
1041 let min_seen_secret = monitor_holder.monitor.get_min_seen_secret();
1042 let counterparty_node_id = monitor_holder.monitor.get_counterparty_node_id();
1043 {
1044 let inner_lock = monitor_holder.monitor.inner.lock().unwrap();
1045
1046 write_chanmon_internal(&inner_lock, true, &mut serialized_channel)
1047 .expect("can not write Channel Monitor for peer storage message");
1048 }
1049 let peer_storage_monitor = PeerStorageMonitorHolder {
1050 channel_id,
1051 min_seen_secret,
1052 counterparty_node_id,
1053 monitor_bytes: serialized_channel.0,
1054 };
1055
1056 let serialized_length = peer_storage_monitor.serialized_length();
1057
1058 if current_size + serialized_length > MAX_PEER_STORAGE_SIZE {
1059 continue;
1060 } else {
1061 current_size += serialized_length;
1062 monitors_list.push(peer_storage_monitor);
1063 }
1064 }
1065
1066 let serialised_channels = monitors_list.encode();
1067 let our_peer_storage = DecryptedOurPeerStorage::new(serialised_channels);
1068 let cipher = our_peer_storage.encrypt(&self.our_peerstorage_encryption_key, &random_bytes);
1069
1070 log_debug!(self.logger, "Sending Peer Storage to {}", log_pubkey!(their_node_id));
1071 let send_peer_storage_event = MessageSendEvent::SendPeerStorage {
1072 node_id: their_node_id,
1073 msg: PeerStorage { data: cipher.into_vec() },
1074 };
1075
1076 self.pending_send_only_events.lock().unwrap().push(send_peer_storage_event)
1077 }
1078
1079 pub fn load_existing_monitor(
1098 &self, channel_id: ChannelId, monitor: ChannelMonitor<ChannelSigner>,
1099 ) -> Result<ChannelMonitorUpdateStatus, ()> {
1100 if !monitor.written_by_0_1_or_later() {
1101 return chain::Watch::watch_channel(self, channel_id, monitor);
1102 }
1103
1104 let logger = WithChannelMonitor::from(&self.logger, &monitor, None);
1105 let mut monitors = self.monitors.write().unwrap();
1106 let entry = match monitors.entry(channel_id) {
1107 hash_map::Entry::Occupied(_) => {
1108 log_error!(logger, "Failed to add new channel data: channel monitor for given channel ID is already present");
1109 return Err(());
1110 },
1111 hash_map::Entry::Vacant(e) => e,
1112 };
1113 log_trace!(
1114 logger,
1115 "Loaded existing ChannelMonitor for channel {}",
1116 log_funding_info!(monitor)
1117 );
1118 if let Some(ref chain_source) = self.chain_source {
1119 monitor.load_outputs_to_watch(chain_source, &self.logger);
1120 }
1121 entry.insert(MonitorHolder { monitor, pending_monitor_updates: Mutex::new(Vec::new()) });
1122
1123 Ok(ChannelMonitorUpdateStatus::Completed)
1124 }
1125}
1126
1127impl<
1128 ChannelSigner: EcdsaChannelSigner,
1129 C: Deref,
1130 T: Deref,
1131 F: Deref,
1132 L: Deref,
1133 P: Deref,
1134 ES: Deref,
1135 > BaseMessageHandler for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1136where
1137 C::Target: chain::Filter,
1138 T::Target: BroadcasterInterface,
1139 F::Target: FeeEstimator,
1140 L::Target: Logger,
1141 P::Target: Persist<ChannelSigner>,
1142 ES::Target: EntropySource,
1143{
1144 fn get_and_clear_pending_msg_events(&self) -> Vec<MessageSendEvent> {
1145 let mut pending_events = self.pending_send_only_events.lock().unwrap();
1146 core::mem::take(&mut *pending_events)
1147 }
1148
1149 fn peer_disconnected(&self, _their_node_id: PublicKey) {}
1150
1151 fn provided_node_features(&self) -> NodeFeatures {
1152 NodeFeatures::empty()
1153 }
1154
1155 fn provided_init_features(&self, _their_node_id: PublicKey) -> InitFeatures {
1156 InitFeatures::empty()
1157 }
1158
1159 fn peer_connected(
1160 &self, _their_node_id: PublicKey, _msg: &Init, _inbound: bool,
1161 ) -> Result<(), ()> {
1162 Ok(())
1163 }
1164}
1165
1166impl<
1167 ChannelSigner: EcdsaChannelSigner,
1168 C: Deref,
1169 T: Deref,
1170 F: Deref,
1171 L: Deref,
1172 P: Deref,
1173 ES: Deref,
1174 > SendOnlyMessageHandler for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1175where
1176 C::Target: chain::Filter,
1177 T::Target: BroadcasterInterface,
1178 F::Target: FeeEstimator,
1179 L::Target: Logger,
1180 P::Target: Persist<ChannelSigner>,
1181 ES::Target: EntropySource,
1182{
1183}
1184
1185impl<
1186 ChannelSigner: EcdsaChannelSigner,
1187 C: Deref,
1188 T: Deref,
1189 F: Deref,
1190 L: Deref,
1191 P: Deref,
1192 ES: Deref,
1193 > chain::Listen for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1194where
1195 C::Target: chain::Filter,
1196 T::Target: BroadcasterInterface,
1197 F::Target: FeeEstimator,
1198 L::Target: Logger,
1199 P::Target: Persist<ChannelSigner>,
1200 ES::Target: EntropySource,
1201{
1202 fn filtered_block_connected(&self, header: &Header, txdata: &TransactionData, height: u32) {
1203 log_debug!(
1204 self.logger,
1205 "New best block {} at height {} provided via block_connected",
1206 header.block_hash(),
1207 height
1208 );
1209 self.process_chain_data(header, Some(height), &txdata, |monitor, txdata| {
1210 monitor.block_connected(
1211 header,
1212 txdata,
1213 height,
1214 &*self.broadcaster,
1215 &*self.fee_estimator,
1216 &self.logger,
1217 )
1218 });
1219
1220 #[cfg(peer_storage)]
1221 for node_id in self.all_counterparty_node_ids() {
1223 self.send_peer_storage(node_id);
1224 }
1225
1226 self.event_notifier.notify();
1228 }
1229
1230 fn blocks_disconnected(&self, fork_point: BestBlock) {
1231 let monitor_states = self.monitors.read().unwrap();
1232 log_debug!(
1233 self.logger,
1234 "Block(s) removed to height {} via blocks_disconnected. New best block is {}",
1235 fork_point.height,
1236 fork_point.block_hash,
1237 );
1238 for monitor_state in monitor_states.values() {
1239 monitor_state.monitor.blocks_disconnected(
1240 fork_point,
1241 &*self.broadcaster,
1242 &*self.fee_estimator,
1243 &self.logger,
1244 );
1245 }
1246 }
1247}
1248
1249impl<
1250 ChannelSigner: EcdsaChannelSigner,
1251 C: Deref,
1252 T: Deref,
1253 F: Deref,
1254 L: Deref,
1255 P: Deref,
1256 ES: Deref,
1257 > chain::Confirm for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1258where
1259 C::Target: chain::Filter,
1260 T::Target: BroadcasterInterface,
1261 F::Target: FeeEstimator,
1262 L::Target: Logger,
1263 P::Target: Persist<ChannelSigner>,
1264 ES::Target: EntropySource,
1265{
1266 fn transactions_confirmed(&self, header: &Header, txdata: &TransactionData, height: u32) {
1267 log_debug!(
1268 self.logger,
1269 "{} provided transactions confirmed at height {} in block {}",
1270 txdata.len(),
1271 height,
1272 header.block_hash()
1273 );
1274 self.process_chain_data(header, None, txdata, |monitor, txdata| {
1275 monitor.transactions_confirmed(
1276 header,
1277 txdata,
1278 height,
1279 &*self.broadcaster,
1280 &*self.fee_estimator,
1281 &self.logger,
1282 )
1283 });
1284 self.event_notifier.notify();
1286 }
1287
1288 fn transaction_unconfirmed(&self, txid: &Txid) {
1289 log_debug!(self.logger, "Transaction {} reorganized out of chain", txid);
1290 let monitor_states = self.monitors.read().unwrap();
1291 for monitor_state in monitor_states.values() {
1292 monitor_state.monitor.transaction_unconfirmed(
1293 txid,
1294 &*self.broadcaster,
1295 &*self.fee_estimator,
1296 &self.logger,
1297 );
1298 }
1299 }
1300
1301 fn best_block_updated(&self, header: &Header, height: u32) {
1302 log_debug!(
1303 self.logger,
1304 "New best block {} at height {} provided via best_block_updated",
1305 header.block_hash(),
1306 height
1307 );
1308 self.process_chain_data(header, Some(height), &[], |monitor, txdata| {
1309 debug_assert!(txdata.is_empty());
1312 monitor.best_block_updated(
1313 header,
1314 height,
1315 &*self.broadcaster,
1316 &*self.fee_estimator,
1317 &self.logger,
1318 )
1319 });
1320
1321 #[cfg(peer_storage)]
1322 for node_id in self.all_counterparty_node_ids() {
1324 self.send_peer_storage(node_id);
1325 }
1326
1327 self.event_notifier.notify();
1329 }
1330
1331 fn get_relevant_txids(&self) -> Vec<(Txid, u32, Option<BlockHash>)> {
1332 let mut txids = Vec::new();
1333 let monitor_states = self.monitors.read().unwrap();
1334 for monitor_state in monitor_states.values() {
1335 txids.append(&mut monitor_state.monitor.get_relevant_txids());
1336 }
1337
1338 txids.sort_unstable_by(|a, b| a.0.cmp(&b.0).then(b.1.cmp(&a.1)));
1339 txids.dedup_by_key(|(txid, _, _)| *txid);
1340 txids
1341 }
1342}
1343
1344impl<
1345 ChannelSigner: EcdsaChannelSigner,
1346 C: Deref,
1347 T: Deref,
1348 F: Deref,
1349 L: Deref,
1350 P: Deref,
1351 ES: Deref,
1352 > chain::Watch<ChannelSigner> for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1353where
1354 C::Target: chain::Filter,
1355 T::Target: BroadcasterInterface,
1356 F::Target: FeeEstimator,
1357 L::Target: Logger,
1358 P::Target: Persist<ChannelSigner>,
1359 ES::Target: EntropySource,
1360{
1361 fn watch_channel(
1362 &self, channel_id: ChannelId, monitor: ChannelMonitor<ChannelSigner>,
1363 ) -> Result<ChannelMonitorUpdateStatus, ()> {
1364 let logger = WithChannelMonitor::from(&self.logger, &monitor, None);
1365 let mut monitors = self.monitors.write().unwrap();
1366 let entry = match monitors.entry(channel_id) {
1367 hash_map::Entry::Occupied(_) => {
1368 log_error!(logger, "Failed to add new channel data: channel monitor for given channel ID is already present");
1369 return Err(());
1370 },
1371 hash_map::Entry::Vacant(e) => e,
1372 };
1373 log_trace!(logger, "Got new ChannelMonitor for channel {}", log_funding_info!(monitor));
1374 let update_id = monitor.get_latest_update_id();
1375 let mut pending_monitor_updates = Vec::new();
1376 let persist_res = self.persister.persist_new_channel(monitor.persistence_key(), &monitor);
1377 match persist_res {
1378 ChannelMonitorUpdateStatus::InProgress => {
1379 log_info!(
1380 logger,
1381 "Persistence of new ChannelMonitor for channel {} in progress",
1382 log_funding_info!(monitor)
1383 );
1384 pending_monitor_updates.push(update_id);
1385 },
1386 ChannelMonitorUpdateStatus::Completed => {
1387 log_info!(
1388 logger,
1389 "Persistence of new ChannelMonitor for channel {} completed",
1390 log_funding_info!(monitor)
1391 );
1392 },
1393 ChannelMonitorUpdateStatus::UnrecoverableError => {
1394 let err_str = "ChannelMonitor[Update] persistence failed unrecoverably. This indicates we cannot continue normal operation and must shut down.";
1395 log_error!(logger, "{}", err_str);
1396 panic!("{}", err_str);
1397 },
1398 }
1399 if let Some(ref chain_source) = self.chain_source {
1400 monitor.load_outputs_to_watch(chain_source, &self.logger);
1401 }
1402 entry.insert(MonitorHolder {
1403 monitor,
1404 pending_monitor_updates: Mutex::new(pending_monitor_updates),
1405 });
1406 Ok(persist_res)
1407 }
1408
1409 fn update_channel(
1410 &self, channel_id: ChannelId, update: &ChannelMonitorUpdate,
1411 ) -> ChannelMonitorUpdateStatus {
1412 debug_assert_eq!(update.channel_id.unwrap(), channel_id);
1415 let monitors = self.monitors.read().unwrap();
1417 match monitors.get(&channel_id) {
1418 None => {
1419 let logger = WithContext::from(&self.logger, None, Some(channel_id), None);
1420 log_error!(logger, "Failed to update channel monitor: no such monitor registered");
1421
1422 #[cfg(debug_assertions)]
1426 panic!("ChannelManager generated a channel update for a channel that was not yet registered!");
1427 #[cfg(not(debug_assertions))]
1428 ChannelMonitorUpdateStatus::InProgress
1429 },
1430 Some(monitor_state) => {
1431 let monitor = &monitor_state.monitor;
1432 let logger = WithChannelMonitor::from(&self.logger, &monitor, None);
1433 log_trace!(
1434 logger,
1435 "Updating ChannelMonitor to id {} for channel {}",
1436 update.update_id,
1437 log_funding_info!(monitor)
1438 );
1439
1440 let mut pending_monitor_updates =
1444 monitor_state.pending_monitor_updates.lock().unwrap();
1445 let update_res = monitor.update_monitor(
1446 update,
1447 &self.broadcaster,
1448 &self.fee_estimator,
1449 &self.logger,
1450 );
1451
1452 let update_id = update.update_id;
1453 let persist_res = if update_res.is_err() {
1454 log_warn!(logger, "Failed to update ChannelMonitor for channel {}. Going ahead and persisting the entire ChannelMonitor", log_funding_info!(monitor));
1460 self.persister.update_persisted_channel(
1461 monitor.persistence_key(),
1462 None,
1463 monitor,
1464 )
1465 } else {
1466 self.persister.update_persisted_channel(
1467 monitor.persistence_key(),
1468 Some(update),
1469 monitor,
1470 )
1471 };
1472 match persist_res {
1473 ChannelMonitorUpdateStatus::InProgress => {
1474 if update_res.is_ok() {
1477 pending_monitor_updates.push(update_id);
1478 }
1479 log_debug!(
1480 logger,
1481 "Persistence of ChannelMonitorUpdate id {:?} for channel {} in progress",
1482 update_id,
1483 log_funding_info!(monitor)
1484 );
1485 },
1486 ChannelMonitorUpdateStatus::Completed => {
1487 log_debug!(
1488 logger,
1489 "Persistence of ChannelMonitorUpdate id {:?} for channel {} completed",
1490 update_id,
1491 log_funding_info!(monitor)
1492 );
1493 },
1494 ChannelMonitorUpdateStatus::UnrecoverableError => {
1495 core::mem::drop(pending_monitor_updates);
1498 core::mem::drop(monitors);
1499 let _poison = self.monitors.write().unwrap();
1500 let err_str = "ChannelMonitor[Update] persistence failed unrecoverably. This indicates we cannot continue normal operation and must shut down.";
1501 log_error!(logger, "{}", err_str);
1502 panic!("{}", err_str);
1503 },
1504 }
1505
1506 if let Some(ref chain_source) = self.chain_source {
1508 for (funding_outpoint, funding_script) in
1509 update.internal_renegotiated_funding_data()
1510 {
1511 log_trace!(
1512 logger,
1513 "Registering renegotiated funding outpoint {} with the filter to monitor confirmations and spends",
1514 funding_outpoint
1515 );
1516 chain_source.register_tx(&funding_outpoint.txid, &funding_script);
1517 chain_source.register_output(WatchedOutput {
1518 block_hash: None,
1519 outpoint: funding_outpoint,
1520 script_pubkey: funding_script,
1521 });
1522 }
1523 }
1524
1525 if update_res.is_err() {
1526 ChannelMonitorUpdateStatus::InProgress
1527 } else {
1528 persist_res
1529 }
1530 },
1531 }
1532 }
1533
1534 fn release_pending_monitor_events(
1535 &self,
1536 ) -> Vec<(OutPoint, ChannelId, Vec<MonitorEvent>, PublicKey)> {
1537 for (channel_id, update_id) in self.persister.get_and_clear_completed_updates() {
1538 let _ = self.channel_monitor_updated(channel_id, update_id);
1539 }
1540 let mut pending_monitor_events = self.pending_monitor_events.lock().unwrap().split_off(0);
1541 let monitors = self.monitors.read().unwrap();
1542 for monitor_state in monitors.values() {
1543 let monitor_events = {
1549 let pending_updates = monitor_state.pending_monitor_updates.lock().unwrap();
1550 if monitor_state.has_pending_updates(&pending_updates) {
1551 monitor_state.monitor.get_and_clear_pending_non_htlc_fail_events()
1552 } else {
1553 monitor_state.monitor.get_and_clear_pending_monitor_events()
1554 }
1555 };
1556 if monitor_events.len() > 0 {
1557 let monitor_funding_txo = monitor_state.monitor.get_funding_txo();
1558 let monitor_channel_id = monitor_state.monitor.channel_id();
1559 let counterparty_node_id = monitor_state.monitor.get_counterparty_node_id();
1560 pending_monitor_events.push((
1561 monitor_funding_txo,
1562 monitor_channel_id,
1563 monitor_events,
1564 counterparty_node_id,
1565 ));
1566 }
1567 }
1568 pending_monitor_events
1569 }
1570}
1571
1572impl<
1573 ChannelSigner: EcdsaChannelSigner,
1574 C: Deref,
1575 T: Deref,
1576 F: Deref,
1577 L: Deref,
1578 P: Deref,
1579 ES: Deref,
1580 > events::EventsProvider for ChainMonitor<ChannelSigner, C, T, F, L, P, ES>
1581where
1582 C::Target: chain::Filter,
1583 T::Target: BroadcasterInterface,
1584 F::Target: FeeEstimator,
1585 L::Target: Logger,
1586 P::Target: Persist<ChannelSigner>,
1587 ES::Target: EntropySource,
1588{
1589 fn process_pending_events<H: Deref>(&self, handler: H)
1603 where
1604 H::Target: EventHandler,
1605 {
1606 for monitor_state in self.monitors.read().unwrap().values() {
1607 match monitor_state.monitor.process_pending_events(&handler, &self.logger) {
1608 Ok(()) => {},
1609 Err(ReplayEvent()) => {
1610 self.event_notifier.notify();
1611 },
1612 }
1613 }
1614 }
1615}
1616
1617#[cfg(test)]
1618mod tests {
1619 use crate::chain::channelmonitor::ANTI_REORG_DELAY;
1620 use crate::chain::{ChannelMonitorUpdateStatus, Watch};
1621 use crate::events::{ClosureReason, Event};
1622 use crate::ln::functional_test_utils::*;
1623 use crate::ln::msgs::{BaseMessageHandler, ChannelMessageHandler, MessageSendEvent};
1624 use crate::{check_added_monitors, check_closed_event};
1625 use crate::{expect_payment_path_successful, get_event_msg};
1626 use crate::{get_htlc_update_msgs, get_revoke_commit_msgs};
1627
1628 const CHAINSYNC_MONITOR_PARTITION_FACTOR: u32 = 5;
1629
1630 #[test]
1631 fn test_async_ooo_offchain_updates() {
1632 let chanmon_cfgs = create_chanmon_cfgs(2);
1636 let node_cfgs = create_node_cfgs(2, &chanmon_cfgs);
1637 let node_chanmgrs = create_node_chanmgrs(2, &node_cfgs, &[None, None]);
1638 let nodes = create_network(2, &node_cfgs, &node_chanmgrs);
1639 let channel_id = create_announced_chan_between_nodes(&nodes, 0, 1).2;
1640
1641 let node_a_id = nodes[0].node.get_our_node_id();
1642 let node_b_id = nodes[1].node.get_our_node_id();
1643
1644 let (payment_preimage_1, payment_hash_1, ..) =
1646 route_payment(&nodes[0], &[&nodes[1]], 1_000_000);
1647 let (payment_preimage_2, payment_hash_2, ..) =
1648 route_payment(&nodes[0], &[&nodes[1]], 1_000_000);
1649
1650 chanmon_cfgs[1].persister.offchain_monitor_updates.lock().unwrap().clear();
1651 chanmon_cfgs[1].persister.set_update_ret(ChannelMonitorUpdateStatus::InProgress);
1652 chanmon_cfgs[1].persister.set_update_ret(ChannelMonitorUpdateStatus::InProgress);
1653
1654 nodes[1].node.claim_funds(payment_preimage_1);
1655 check_added_monitors!(nodes[1], 1);
1656 nodes[1].node.claim_funds(payment_preimage_2);
1657 check_added_monitors!(nodes[1], 1);
1658
1659 let persistences =
1660 chanmon_cfgs[1].persister.offchain_monitor_updates.lock().unwrap().clone();
1661 assert_eq!(persistences.len(), 1);
1662 let (_, updates) = persistences.iter().next().unwrap();
1663 assert_eq!(updates.len(), 2);
1664
1665 let mut update_iter = updates.iter();
1668 let next_update = update_iter.next().unwrap().clone();
1669 let node_b_mon = &nodes[1].chain_monitor.chain_monitor;
1670
1671 let pending_updates = node_b_mon.list_pending_monitor_updates();
1673 #[cfg(not(c_bindings))]
1674 let pending_chan_updates = pending_updates.get(&channel_id).unwrap();
1675 #[cfg(c_bindings)]
1676 let pending_chan_updates =
1677 &pending_updates.iter().find(|(chan_id, _)| *chan_id == channel_id).unwrap().1;
1678 assert!(pending_chan_updates.contains(&next_update));
1679
1680 node_b_mon.channel_monitor_updated(channel_id, next_update.clone()).unwrap();
1681
1682 let pending_updates = node_b_mon.list_pending_monitor_updates();
1684 #[cfg(not(c_bindings))]
1685 let pending_chan_updates = pending_updates.get(&channel_id).unwrap();
1686 #[cfg(c_bindings)]
1687 let pending_chan_updates =
1688 &pending_updates.iter().find(|(chan_id, _)| *chan_id == channel_id).unwrap().1;
1689 assert!(!pending_chan_updates.contains(&next_update));
1690
1691 assert!(nodes[1].chain_monitor.release_pending_monitor_events().is_empty());
1692 assert!(nodes[1].node.get_and_clear_pending_msg_events().is_empty());
1693 assert!(nodes[1].node.get_and_clear_pending_events().is_empty());
1694
1695 let next_update = update_iter.next().unwrap().clone();
1696 node_b_mon.channel_monitor_updated(channel_id, next_update).unwrap();
1697
1698 let claim_events = nodes[1].node.get_and_clear_pending_events();
1699 assert_eq!(claim_events.len(), 2);
1700 match claim_events[0] {
1701 Event::PaymentClaimed { ref payment_hash, amount_msat: 1_000_000, .. } => {
1702 assert_eq!(payment_hash_1, *payment_hash);
1703 },
1704 _ => panic!("Unexpected event"),
1705 }
1706 match claim_events[1] {
1707 Event::PaymentClaimed { ref payment_hash, amount_msat: 1_000_000, .. } => {
1708 assert_eq!(payment_hash_2, *payment_hash);
1709 },
1710 _ => panic!("Unexpected event"),
1711 }
1712
1713 let mut updates = get_htlc_update_msgs!(nodes[1], node_a_id);
1717 nodes[0].node.handle_update_fulfill_htlc(node_b_id, updates.update_fulfill_htlcs.remove(0));
1718 expect_payment_sent(&nodes[0], payment_preimage_1, None, false, false);
1719 nodes[0].node.handle_commitment_signed_batch_test(node_b_id, &updates.commitment_signed);
1720 check_added_monitors!(nodes[0], 1);
1721 let (as_first_raa, as_first_update) = get_revoke_commit_msgs!(nodes[0], node_b_id);
1722
1723 nodes[1].node.handle_revoke_and_ack(node_a_id, &as_first_raa);
1724 check_added_monitors!(nodes[1], 1);
1725 let mut bs_2nd_updates = get_htlc_update_msgs!(nodes[1], node_a_id);
1726 nodes[1].node.handle_commitment_signed_batch_test(node_a_id, &as_first_update);
1727 check_added_monitors!(nodes[1], 1);
1728 let bs_first_raa = get_event_msg!(nodes[1], MessageSendEvent::SendRevokeAndACK, node_a_id);
1729
1730 nodes[0]
1731 .node
1732 .handle_update_fulfill_htlc(node_b_id, bs_2nd_updates.update_fulfill_htlcs.remove(0));
1733 expect_payment_sent(&nodes[0], payment_preimage_2, None, false, false);
1734 nodes[0]
1735 .node
1736 .handle_commitment_signed_batch_test(node_b_id, &bs_2nd_updates.commitment_signed);
1737 check_added_monitors!(nodes[0], 1);
1738 nodes[0].node.handle_revoke_and_ack(node_b_id, &bs_first_raa);
1739 expect_payment_path_successful!(nodes[0]);
1740 check_added_monitors!(nodes[0], 1);
1741 let (as_second_raa, as_second_update) = get_revoke_commit_msgs!(nodes[0], node_b_id);
1742
1743 nodes[1].node.handle_revoke_and_ack(node_a_id, &as_second_raa);
1744 check_added_monitors!(nodes[1], 1);
1745 nodes[1].node.handle_commitment_signed_batch_test(node_a_id, &as_second_update);
1746 check_added_monitors!(nodes[1], 1);
1747 let bs_second_raa = get_event_msg!(nodes[1], MessageSendEvent::SendRevokeAndACK, node_a_id);
1748
1749 nodes[0].node.handle_revoke_and_ack(node_b_id, &bs_second_raa);
1750 expect_payment_path_successful!(nodes[0]);
1751 check_added_monitors!(nodes[0], 1);
1752 }
1753
1754 #[test]
1755 fn test_chainsync_triggers_distributed_monitor_persistence() {
1756 let chanmon_cfgs = create_chanmon_cfgs(3);
1757 let node_cfgs = create_node_cfgs(3, &chanmon_cfgs);
1758 let node_chanmgrs = create_node_chanmgrs(3, &node_cfgs, &[None, None, None]);
1759 let nodes = create_network(3, &node_cfgs, &node_chanmgrs);
1760
1761 let node_a_id = nodes[0].node.get_our_node_id();
1762 let node_c_id = nodes[2].node.get_our_node_id();
1763
1764 *nodes[0].connect_style.borrow_mut() = ConnectStyle::FullBlockViaListen;
1767 *nodes[1].connect_style.borrow_mut() = ConnectStyle::FullBlockViaListen;
1768 *nodes[2].connect_style.borrow_mut() = ConnectStyle::FullBlockViaListen;
1769
1770 let _channel_1 = create_announced_chan_between_nodes(&nodes, 0, 1).2;
1771 let channel_2 =
1772 create_announced_chan_between_nodes_with_value(&nodes, 0, 2, 1_000_000, 0).2;
1773
1774 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1775 chanmon_cfgs[1].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1776 chanmon_cfgs[2].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1777
1778 connect_blocks(&nodes[0], CHAINSYNC_MONITOR_PARTITION_FACTOR * 2);
1779 connect_blocks(&nodes[1], CHAINSYNC_MONITOR_PARTITION_FACTOR * 2);
1780 connect_blocks(&nodes[2], CHAINSYNC_MONITOR_PARTITION_FACTOR * 2);
1781
1782 assert_eq!(
1785 2 * 2,
1786 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1787 );
1788 assert_eq!(
1789 2,
1790 chanmon_cfgs[1].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1791 );
1792 assert_eq!(
1793 2,
1794 chanmon_cfgs[2].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1795 );
1796
1797 let message = "Channel force-closed".to_owned();
1800 nodes[0]
1801 .node
1802 .force_close_broadcasting_latest_txn(&channel_2, &node_c_id, message.clone())
1803 .unwrap();
1804 let closure_reason =
1805 ClosureReason::HolderForceClosed { broadcasted_latest_txn: Some(true), message };
1806 check_closed_event!(&nodes[0], 1, closure_reason, false, [node_c_id], 1000000);
1807 check_closed_broadcast(&nodes[0], 1, true);
1808 let close_tx = nodes[0].tx_broadcaster.txn_broadcasted.lock().unwrap().split_off(0);
1809 assert_eq!(close_tx.len(), 1);
1810
1811 mine_transaction(&nodes[2], &close_tx[0]);
1812 check_closed_broadcast(&nodes[2], 1, true);
1813 check_added_monitors(&nodes[2], 1);
1814 let closure_reason = ClosureReason::CommitmentTxConfirmed;
1815 check_closed_event!(&nodes[2], 1, closure_reason, false, [node_a_id], 1000000);
1816
1817 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1818 chanmon_cfgs[2].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1819
1820 connect_blocks(&nodes[0], CHAINSYNC_MONITOR_PARTITION_FACTOR);
1825 connect_blocks(&nodes[2], CHAINSYNC_MONITOR_PARTITION_FACTOR);
1826
1827 assert_eq!(
1830 (CHAINSYNC_MONITOR_PARTITION_FACTOR + 1) as usize,
1831 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1832 );
1833 assert_eq!(
1835 1,
1836 chanmon_cfgs[2].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1837 );
1838
1839 mine_transaction(&nodes[0], &close_tx[0]);
1841 connect_blocks(&nodes[0], ANTI_REORG_DELAY - 1);
1842 check_added_monitors(&nodes[0], 1);
1843 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().clear();
1844
1845 connect_blocks(&nodes[0], CHAINSYNC_MONITOR_PARTITION_FACTOR);
1848 assert_eq!(
1849 2,
1850 chanmon_cfgs[0].persister.chain_sync_monitor_persistences.lock().unwrap().len()
1851 );
1852 }
1853
1854 #[test]
1855 #[cfg(feature = "std")]
1856 fn update_during_chainsync_poisons_channel() {
1857 let chanmon_cfgs = create_chanmon_cfgs(2);
1858 let node_cfgs = create_node_cfgs(2, &chanmon_cfgs);
1859 let node_chanmgrs = create_node_chanmgrs(2, &node_cfgs, &[None, None]);
1860 let nodes = create_network(2, &node_cfgs, &node_chanmgrs);
1861 create_announced_chan_between_nodes(&nodes, 0, 1);
1862 *nodes[0].connect_style.borrow_mut() = ConnectStyle::FullBlockViaListen;
1863
1864 chanmon_cfgs[0].persister.set_update_ret(ChannelMonitorUpdateStatus::UnrecoverableError);
1865
1866 assert!(std::panic::catch_unwind(|| {
1867 connect_blocks(&nodes[0], CHAINSYNC_MONITOR_PARTITION_FACTOR);
1871 })
1872 .is_err());
1873 assert!(std::panic::catch_unwind(|| {
1874 core::mem::drop(nodes);
1876 })
1877 .is_err());
1878 }
1879}