@@ -5,8 +5,9 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH};
55
66use bip157:: chain:: ChainState ;
77use bip157:: {
8- chain:: BlockHeaderChanges , Builder as KyotoBuilder , Client , Event , HashCheckpoint , Header ,
9- IndexedBlock , Info , Node as KyotoNode , Package , Requester , TrustedPeer , Warning ,
8+ chain:: BlockHeaderChanges , Builder as KyotoBuilder , Client , Event as KyotoEvent ,
9+ HashCheckpoint , Header , IndexedBlock , Info , Node as KyotoNode , Package , Requester , TrustedPeer ,
10+ Warning ,
1011} ;
1112use bitcoin:: { BlockHash , FeeRate , Network , Script , ScriptBuf , Transaction , Txid } ;
1213use electrum_client:: { Client as ElectrumClient , ConfigBuilder as ElectrumConfigBuilder } ;
@@ -576,14 +577,14 @@ impl CbfChainSource {
576577 }
577578
578579 async fn process_kyoto_events (
579- logger : Arc < Logger > , mut event_rx : mpsc:: UnboundedReceiver < Event > ,
580+ logger : Arc < Logger > , mut event_rx : mpsc:: UnboundedReceiver < KyotoEvent > ,
580581 registered_scripts : Arc < Mutex < HashSet < ScriptBuf > > > ,
581582 cbf_runtime_status : Arc < Mutex < CbfRuntimeStatus > > , ops_tx : mpsc:: UnboundedSender < ChainOp > ,
582583 onchain_wallet : std:: sync:: Weak < Wallet > , sync_state_tx : watch:: Sender < CbfSyncState > ,
583584 ) {
584585 while let Some ( event) = event_rx. recv ( ) . await {
585586 match event {
586- Event :: IndexedFilter ( indexed_filter) => {
587+ KyotoEvent :: IndexedFilter ( indexed_filter) => {
587588 let Some ( onchain_wallet) = onchain_wallet. upgrade ( ) else {
588589 log_debug ! ( logger, "Onchain wallet dropped; stopping CBF event processing" ) ;
589590 break ;
@@ -671,61 +672,30 @@ impl CbfChainSource {
671672 } ;
672673 ChainOp :: ConnectFull { block }
673674 } else {
674- let height = indexed_filter. height ( ) ;
675- //TODO we need to recheck that a particular height has not been
676- //reorganized, and we retrieve indeed the same block header that we
677- //received `IndexedFilter` event of.
678- match requester. get_header ( height) . await {
679- Ok ( Some ( indexed_header) ) => {
680- if indexed_header. block_hash ( ) != block_hash {
681- log_debug ! (
682- logger,
683- "Filter for {} reorged; skipping" ,
684- block_hash
685- ) ;
686- continue ;
687- }
688- ChainOp :: ConnectFiltered {
689- header : indexed_header. header ,
690- height : indexed_header. height ,
691- }
692- } ,
693- Ok ( None ) => {
694- log_error ! ( logger, "No header at height {}" , height, ) ;
695- let _ = ops_tx. send ( ChainOp :: Failed { error : Error :: TxSyncFailed } ) ;
696- break ;
697- } ,
698- Err ( e) => {
699- log_error ! (
700- logger,
701- "Failed to fetch header at height {}: {:?}" ,
702- height,
703- e,
704- ) ;
705- let _ = ops_tx. send ( ChainOp :: Failed { error : Error :: TxSyncFailed } ) ;
706- break ;
707- } ,
675+ ChainOp :: ConnectFiltered {
676+ header : indexed_filter. header ( ) ,
677+ height : indexed_filter. height ( ) ,
708678 }
709679 } ;
710680 if let Err ( e) = ops_tx. send ( chop) {
711681 log_debug ! ( logger, "ops_rx gone: {}" , e) ;
712682 }
713683 } ,
714- Event :: FiltersSynced ( sync_update) => {
684+ KyotoEvent :: FiltersSynced ( sync_update) => {
715685 //Because application of blocks is async, the fact that kyoto synced up to the
716686 //tip does NOT mean that we caught everything up, that's why we send a ChainOp,
717687 //only processing of which means we processed all blocks up to the tip.
718688 log_info ! ( logger, "Kyoto synced up to the tip {}" , sync_update. tip( ) . height) ;
719689 let _ = ops_tx. send ( ChainOp :: Synced { tip_height : sync_update. tip ( ) . height } ) ;
720690 } ,
721- Event :: ChainUpdate ( BlockHeaderChanges :: Connected ( indexed_header) ) => {
691+ KyotoEvent :: ChainUpdate ( BlockHeaderChanges :: Connected ( indexed_header) ) => {
722692 log_debug ! (
723693 logger,
724694 "Kyoto connected header at height {}" ,
725695 indexed_header. height
726696 ) ;
727697 } ,
728- Event :: ChainUpdate ( BlockHeaderChanges :: Reorganized {
698+ KyotoEvent :: ChainUpdate ( BlockHeaderChanges :: Reorganized {
729699 reorganized,
730700 accepted : _,
731701 } ) => {
@@ -738,7 +708,7 @@ impl CbfChainSource {
738708 let _ = ops_tx. send ( ChainOp :: Disconnect { fork_point } ) ;
739709 }
740710 } ,
741- Event :: ChainUpdate ( BlockHeaderChanges :: ForkAdded ( fork) ) => {
711+ KyotoEvent :: ChainUpdate ( BlockHeaderChanges :: ForkAdded ( fork) ) => {
742712 log_debug ! ( logger, "Kyoto added fork header at height {}" , fork. height) ;
743713 } ,
744714 }
0 commit comments