From 05e0fa397e4618d6956aea58eec1396bbd0addbc Mon Sep 17 00:00:00 2001 From: Youngjoon Lee <5462944+youngjoon-lee@users.noreply.github.com> Date: Thu, 27 Aug 2026 14:50:57 +0900 Subject: [PATCH] refactor: Direction -> Outcome --- services/sdp/src/intent.rs | 225 ++++++++++++++++++++----------------- services/sdp/src/lib.rs | 64 ++++++----- 2 files changed, 160 insertions(+), 129 deletions(-) diff --git a/services/sdp/src/intent.rs b/services/sdp/src/intent.rs index 1c73695bf..17793a341 100644 --- a/services/sdp/src/intent.rs +++ b/services/sdp/src/intent.rs @@ -5,7 +5,6 @@ use lb_core::header::HeaderId; use lb_ledger::{IntentStatus, LedgerState}; use overwatch::DynError; use serde::{Deserialize, Serialize}; -use tracing::debug; /// Tracks the status of a submitted intent on the ledger. /// @@ -54,7 +53,7 @@ impl IntentTracker { impl IntentTracker where - Intent: lb_ledger::Intent + Clone, + Intent: lb_ledger::Intent + Clone, Provider: LedgerStateProvider, { /// Feeds a new tip to the tracker. A tip the tracker has already seen is @@ -76,17 +75,18 @@ where /// or is not applicable, because a chain reorg can change its status. /// Instead, it keeps tracking the intent for up to /// [`Config::max_status_checks`]. - pub async fn handle_tip( - mut self, - tip: HeaderId, - ) -> Result, (Error, Option)> { + pub async fn handle_tip(&mut self, tip: HeaderId) -> Result, Error> { + if self.status_checks >= self.config.max_status_checks.get() { + return Ok(Outcome::Exhausted); + } + if std::mem::replace(&mut self.last_tip, tip) == tip { - return Ok(Direction::Track(self)); // tip unchanged + return Ok(Outcome::WaitingforMoreTipChanges); // tip unchanged } self.tip_changes = self.tip_changes.saturating_add(1); if self.tip_changes < self.config.status_check_interval_in_tip_changes.get() { - return Ok(Direction::Track(self)); // not time to check the status yet + return Ok(Outcome::WaitingforMoreTipChanges); // not time to check the status yet } self.tip_changes = 0; @@ -94,45 +94,29 @@ where let ledger = match self.ledger_state_provider.get(tip).await { Ok(Some(ledger)) => ledger, - Ok(None) => return Err((Error::LedgerStateNotFound(tip), self.check_deadline())), - Err(e) => return Err((Error::LedgerStateProvider(Box::new(e)), self.decheck_deadline())), + Ok(None) => return Err(Error::LedgerStateNotFound(tip)), + Err(e) => { + return Err(Error::LedgerStateProvider(Box::new(e))); + } }; match self.intent.status(&ledger) { - Ok(IntentStatus::NotApplied) => { - let intent = self.intent.clone(); - match self.check_deadline() { - Some(tracker) => Ok(Direction::ResubmitAndTrack { intent, tracker }), - None => Ok(Direction::ResubmitAndDone(intent)), - } - } - Ok(IntentStatus::Applied) => self.check_deadline().map_or_else( - || Ok(Direction::Done), - |tracker| Ok(Direction::Track(tracker)), - ), - Err(err) => { - debug!(%err, "failed to check status of intent"); - self.check_deadline().map_or_else( - || Ok(Direction::Done), - |tracker| Ok(Direction::Track(tracker)), - ) - } + Ok(status) => Ok(Outcome::StatusChecked { + intent: self.intent.clone(), + status, + }), + Err(e) => Err(Error::StatusCheckFailed(Box::new(e))), } } - - fn check_deadline(self) -> Option { - (self.status_checks < self.config.max_status_checks.get()).then_some(self) - } } -pub enum Direction { - ResubmitAndTrack { +pub enum Outcome { + StatusChecked { intent: Intent, - tracker: IntentTracker, + status: IntentStatus, }, - ResubmitAndDone(Intent), - Track(IntentTracker), - Done, + WaitingforMoreTipChanges, + Exhausted, } #[async_trait] @@ -148,6 +132,8 @@ pub enum Error { LedgerStateNotFound(HeaderId), #[error("ledger state provider error: {0}")] LedgerStateProvider(DynError), + #[error("intent status check failed: {0}")] + StatusCheckFailed(DynError), } #[cfg(test)] @@ -171,10 +157,11 @@ mod tests { #[tokio::test] async fn unchanged_tip_never_trigger_status_check() { - let tracker = tracker(IntentStatus::NotApplied, 2, 3); + let mut tracker = tracker(IntentStatus::NotApplied, 2, 3); let tip = tracker.last_tip; - let tracker = expect_track(tracker.handle_tip(tip).await); + let out = tracker.handle_tip(tip).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); assert_eq!(tracker.last_tip, tip); assert_eq!(tracker.tip_changes, 0); assert_eq!(tracker.status_checks, 0); @@ -182,91 +169,134 @@ mod tests { #[tokio::test] async fn status_check_is_triggered_after_interval() { - let tracker = tracker(IntentStatus::Applied, 2, 3); + let mut tracker = tracker(IntentStatus::Applied, 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); assert_eq!(tracker.status_checks, 0); assert_eq!(tracker.tip_changes, 1); - let tracker = expect_track(tracker.handle_tip(tip(2)).await); + + let out = tracker.handle_tip(tip(2)).await.unwrap(); + expect_applied(&out); assert_eq!(tracker.status_checks, 1); assert_eq!(tracker.tip_changes, 0); // was reset - let tracker = expect_track(tracker.handle_tip(tip(3)).await); + let out = tracker.handle_tip(tip(3)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); assert_eq!(tracker.status_checks, 1); assert_eq!(tracker.tip_changes, 1); - let tracker = expect_track(tracker.handle_tip(tip(4)).await); + + let out = tracker.handle_tip(tip(4)).await.unwrap(); + expect_applied(&out); assert_eq!(tracker.status_checks, 2); assert_eq!(tracker.tip_changes, 0); // was reset - let tracker = expect_track(tracker.handle_tip(tip(5)).await); + let out = tracker.handle_tip(tip(5)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); assert_eq!(tracker.status_checks, 2); assert_eq!(tracker.tip_changes, 1); - expect_done(tracker.handle_tip(tip(6)).await); + + let out = tracker.handle_tip(tip(6)).await.unwrap(); + expect_applied(&out); + assert_eq!(tracker.status_checks, 3); + assert_eq!(tracker.tip_changes, 0); // was reset + + let out = tracker.handle_tip(tip(7)).await.unwrap(); + assert!(matches!(out, Outcome::Exhausted)); } #[tokio::test] async fn resubmit_if_intent_not_applied() { - let tracker = tracker(IntentStatus::NotApplied, 2, 3); + let mut tracker = tracker(IntentStatus::NotApplied, 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); - let tracker = expect_resubmit_and_track(tracker.handle_tip(tip(2)).await); - let tracker = expect_track(tracker.handle_tip(tip(3)).await); - let tracker = expect_resubmit_and_track(tracker.handle_tip(tip(4)).await); - let tracker = expect_track(tracker.handle_tip(tip(5)).await); - expect_resubmit_and_done(tracker.handle_tip(tip(6)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + + let out = tracker.handle_tip(tip(2)).await.unwrap(); + expect_not_applied(&out); + + let out = tracker.handle_tip(tip(3)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + + let out = tracker.handle_tip(tip(4)).await.unwrap(); + expect_not_applied(&out); + + let out = tracker.handle_tip(tip(5)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + + let out = tracker.handle_tip(tip(6)).await.unwrap(); + expect_not_applied(&out); + + let out = tracker.handle_tip(tip(7)).await.unwrap(); + assert!(matches!(out, Outcome::Exhausted)); } #[tokio::test] async fn resubmit_after_applied_intent_reverted() { let status = Arc::new(Mutex::new(Some(IntentStatus::Applied))); - let tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); + let mut tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); - let tracker = expect_track(tracker.handle_tip(tip(2)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + let out = tracker.handle_tip(tip(2)).await.unwrap(); + expect_applied(&out); // e.g., a reorg reverted the applied intent. *status.lock().unwrap() = Some(IntentStatus::NotApplied); - let tracker = expect_track(tracker.handle_tip(tip(3)).await); - expect_resubmit_and_track(tracker.handle_tip(tip(4)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + let out = tracker.handle_tip(tip(2)).await.unwrap(); + expect_not_applied(&out); } #[tokio::test] async fn keep_tracking_if_status_check_failed() { let status = Arc::new(Mutex::new(None)); - let tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); + let mut tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); - let tracker = expect_track(tracker.handle_tip(tip(2)).await); - let tracker = expect_track(tracker.handle_tip(tip(3)).await); - let tracker = expect_track(tracker.handle_tip(tip(4)).await); - let tracker = expect_track(tracker.handle_tip(tip(5)).await); - expect_done(tracker.handle_tip(tip(6)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + tracker.handle_tip(tip(2)).await.err().unwrap(); + + let out = tracker.handle_tip(tip(3)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + tracker.handle_tip(tip(4)).await.err().unwrap(); + + let out = tracker.handle_tip(tip(5)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + tracker.handle_tip(tip(6)).await.err().unwrap(); + + let out = tracker.handle_tip(tip(7)).await.unwrap(); + assert!(matches!(out, Outcome::Exhausted)); } #[tokio::test] async fn resubmit_after_status_check_becomes_available_again() { let status = Arc::new(Mutex::new(None)); - let tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); + let mut tracker = tracker_with(MockIntent(Arc::clone(&status)), 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); - let tracker = expect_track(tracker.handle_tip(tip(2)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + tracker.handle_tip(tip(2)).await.err().unwrap(); // e.g., the intent became applicable after a chain reorg. *status.lock().unwrap() = Some(IntentStatus::NotApplied); - let tracker = expect_track(tracker.handle_tip(tip(3)).await); - expect_resubmit_and_track(tracker.handle_tip(tip(4)).await); + let out = tracker.handle_tip(tip(3)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); + let out = tracker.handle_tip(tip(4)).await.unwrap(); + expect_not_applied(&out); } #[tokio::test] async fn status_check_counts_error() { - let tracker = tracker(IntentStatus::Applied, 2, 3); + let mut tracker = tracker(IntentStatus::Applied, 2, 3); - let tracker = expect_track(tracker.handle_tip(tip(1)).await); + let out = tracker.handle_tip(tip(1)).await.unwrap(); + assert!(matches!(out, Outcome::WaitingforMoreTipChanges)); assert_eq!(tracker.status_checks, 0); let unknown_tip = tip(99); - let Err((Error::LedgerStateNotFound(_), tracker)) = tracker.handle_tip(unknown_tip).await - else { + let Err(Error::LedgerStateNotFound(_)) = tracker.handle_tip(unknown_tip).await else { panic!("expected error"); }; assert_eq!(tracker.status_checks, 1); @@ -274,7 +304,6 @@ mod tests { } type MockTracker = IntentTracker; - type MockResult = Result, (Error, MockTracker)>; fn tracker(status: IntentStatus, interval: u64, max: u64) -> MockTracker { tracker_with( @@ -293,36 +322,24 @@ mod tests { ) } - fn expect_track(result: MockResult) -> MockTracker { - match result { - Ok(Direction::Track(tracker)) => tracker, - Ok(_) => panic!("unexpected direction"), - Err((e, _)) => panic!("unexpected error: {e}"), - } + fn expect_applied(out: &Outcome) { + assert!(matches!( + out, + Outcome::StatusChecked { + status: IntentStatus::Applied, + .. + } + )); } - fn expect_done(result: MockResult) { - match result { - Ok(Direction::Done) => (), - Ok(_) => panic!("unexpected direction"), - Err((e, _)) => panic!("unexpected error: {e}"), - } - } - - fn expect_resubmit_and_track(result: MockResult) -> MockTracker { - match result { - Ok(Direction::ResubmitAndTrack { tracker, .. }) => tracker, - Ok(_) => panic!("unexpected direction"), - Err((e, _)) => panic!("unexpected error: {e}"), - } - } - - fn expect_resubmit_and_done(result: MockResult) { - match result { - Ok(Direction::ResubmitAndDone(_)) => (), - Ok(_) => panic!("unexpected direction"), - Err((e, _)) => panic!("unexpected error: {e}"), - } + fn expect_not_applied(out: &Outcome) { + assert!(matches!( + out, + Outcome::StatusChecked { + status: IntentStatus::NotApplied, + .. + } + )); } fn config(interval: u64, max: u64) -> Config { diff --git a/services/sdp/src/lib.rs b/services/sdp/src/lib.rs index 6d1280bb3..f98aaccd3 100644 --- a/services/sdp/src/lib.rs +++ b/services/sdp/src/lib.rs @@ -21,7 +21,7 @@ use lb_core::{ sdp::{ActiveMessage, ActivityMetadata, DeclarationId, DeclarationMessage, WithdrawMessage}, }; use lb_key_management_system_keys::keys::ZkPublicKey; -use lb_ledger::LedgerState; +use lb_ledger::{IntentStatus, LedgerState}; use lb_services_utils::overwatch::{RecoveryData, RecoveryOperator, StorageRecoverySettings}; use overwatch::{ DynError, OpaqueServiceResourcesHandle, @@ -34,7 +34,7 @@ use tracing::{debug, error, trace}; pub use crate::{api::SdpServiceApi, intent::Config as ActiveMessageTrackerConfig}; use crate::{ - intent::{Direction, IntentTracker}, + intent::IntentTracker, mempool::SdpMempoolAdapter, state::{SdpState, SdpStateStorage}, wallet::{SdpWalletAdapter, SdpWalletConfig}, @@ -288,63 +288,77 @@ where mempool_adapter: &MempoolAdapter, chain_api: &CryptarchiaServiceApi, ) { - let Some(tracker) = self.active_message_tracker.take() else { + let Some(tracker) = self.active_message_tracker.as_mut() else { trace!("no active message tracker exists"); return; }; match tracker.handle_tip(event.tip).await { - Ok(direction) => { - self.handle_active_message_tracker_direction( - direction, + Ok(outcome) => { + self.handle_active_message_tracker_outcome( + outcome, wallet_adapter, mempool_adapter, chain_api, ) .await; } - Err((err, tracker)) => { + Err(err) => { error!(%err, "active message tracker failed to handle tip"); - self.active_message_tracker = tracker; } } } - async fn handle_active_message_tracker_direction( + async fn handle_active_message_tracker_outcome( &mut self, - direction: Direction>, + outcome: intent::Outcome, wallet_adapter: &WalletAdapter, mempool_adapter: &MempoolAdapter, chain_api: &CryptarchiaServiceApi, ) { - match direction { - Direction::ResubmitAndTrack { intent, tracker } => { - self.active_message_tracker = Some(tracker); - self.submit_activity( - intent.declaration_id, - intent.metadata, + match outcome { + intent::Outcome::StatusChecked { intent, status } => { + self.handle_active_message_status( + intent, + status, wallet_adapter, mempool_adapter, chain_api, ) .await; } - Direction::ResubmitAndDone(intent) => { + intent::Outcome::WaitingforMoreTipChanges => { + trace!("active message tracker waiting for more tip changes before status check"); + } + intent::Outcome::Exhausted => { + debug!("active message tracker exhausted: dropping the tracker"); + self.active_message_tracker = None; + } + } + } + + async fn handle_active_message_status( + &self, + message: ActiveMessage, + status: IntentStatus, + wallet_adapter: &WalletAdapter, + mempool_adapter: &MempoolAdapter, + chain_api: &CryptarchiaServiceApi, + ) { + match status { + IntentStatus::NotApplied => { + trace!("active message status: not applied in the tip ledger: resubmitting it"); self.submit_activity( - intent.declaration_id, - intent.metadata, + message.declaration_id, + message.metadata, wallet_adapter, mempool_adapter, chain_api, ) .await; } - Direction::Track(tracker) => { - self.active_message_tracker = Some(tracker); - trace!("no active message to resubmit"); - } - Direction::Done => { - debug!("active message tracker has completed"); + IntentStatus::Applied => { + trace!("active message status: applied in the tip ledger: keep tracking it"); } } }