Skip to main content

rings_core/message/handlers/
mod.rs

1#![deny(missing_docs)]
2//! This module implemented message handler of rings network.
3
4use std::sync::Arc;
5
6use async_trait::async_trait;
7
8use super::effects::core_actor_steps;
9use super::effects::lower_dht_action;
10use super::effects::yield_core_actor_step;
11use super::effects::CoreEffect;
12use super::effects::CoreEffectInterpreter;
13use super::MessagePayload;
14use crate::dht::Did;
15use crate::dht::PeerRing;
16use crate::dht::PeerRingAction;
17use crate::error::Error;
18use crate::error::Result;
19use crate::swarm::callback::InnerSwarmCallback;
20use crate::swarm::callback::SharedSwarmCallback;
21use crate::swarm::transport::PendingConnectionAttempt;
22use crate::swarm::transport::SwarmTransport;
23
24/// Operator and Handler for Connection
25pub mod connection;
26/// Operator and Handler for CustomMessage
27pub mod custom;
28/// Operator and Handler for E2E encrypted messages
29pub mod e2e;
30/// Operator and handler for DHT stabilization
31pub mod stabilization;
32/// Operator and Handler for Storage
33pub mod storage;
34/// Shared message-handler handle.
35///
36/// Clone law: cloning duplicates `Arc` handles to the same transport, DHT
37/// state, and callback. It never forks protocol state or transfers ownership.
38#[derive(Clone)]
39pub struct MessageHandler {
40    transport: Arc<SwarmTransport>,
41    dht: Arc<PeerRing>,
42    swarm_callback: SharedSwarmCallback,
43}
44
45/// Generic trait for handle message ,inspired by Actor-Model.
46#[cfg_attr(all(feature = "wasm", target_family = "wasm"), async_trait(?Send))]
47#[cfg_attr(not(all(feature = "wasm", target_family = "wasm")), async_trait)]
48pub trait HandleMsg<T> {
49    /// Message handler.
50    async fn handle(&self, ctx: &MessagePayload, msg: &T) -> Result<()>;
51}
52
53impl MessageHandler {
54    /// Create a new MessageHandler instance.
55    pub fn new(transport: Arc<SwarmTransport>, swarm_callback: SharedSwarmCallback) -> Self {
56        let dht = transport.dht.clone();
57        Self {
58            transport,
59            dht,
60            swarm_callback,
61        }
62    }
63
64    fn inner_callback(&self) -> InnerSwarmCallback {
65        InnerSwarmCallback::new(self.transport.clone(), self.swarm_callback.clone())
66    }
67
68    pub(crate) async fn run_effects<'payload>(
69        &self,
70        effects: impl IntoIterator<Item = CoreEffect<'payload>>,
71    ) -> Result<()> {
72        CoreEffectInterpreter::new(&self.transport, &self.swarm_callback)
73            .run_all(effects)
74            .await
75    }
76
77    /// Idempotently establish a DHT-driven transport connection.
78    ///
79    /// Self and already-connected peers are no-ops. `AlreadyConnected` is treated
80    /// as success so concurrent DHT actions racing through `MultiActions` do not
81    /// fail the whole handler.
82    pub(crate) async fn connect_dht_peer(&self, peer: Did) -> Result<()> {
83        self.run_effects([CoreEffect::connect_dht_peer(peer)]).await
84    }
85
86    /// Idempotently establish DHT-driven transport connections in local quality order.
87    pub(crate) async fn connect_dht_peers(
88        &self,
89        peers: impl IntoIterator<Item = Did>,
90    ) -> Result<()> {
91        for (peer, has_next) in
92            core_actor_steps(self.transport.order_dht_candidates_by_quality(peers).await)
93        {
94            self.connect_dht_peer(peer).await?;
95            if has_next {
96                yield_core_actor_step().await;
97            }
98        }
99        Ok(())
100    }
101
102    pub(crate) async fn join_dht(&self, peer: Did) -> Result<()> {
103        // Default HMCC/Zave join path: maps to the JoinThenSync operation in
104        // the CorrectChord spec (see tests/default/test_dht_convergence.rs).
105        let Some(dht_ev) = self.transport.join_routable_peer(peer)? else {
106            return Err(Error::SwarmMissDidInTable(peer));
107        };
108        // The local join has completed. Follow-up convergence messages are
109        // best-effort: a peer can churn before these sends complete, and that
110        // must not suppress the application-level Connected event.
111        if let Err(e) = self.handle_dht_events(&dht_ev).await {
112            tracing::warn!("Failed to handle dht events while joining {peer}: {e:?}");
113        }
114        Ok(())
115    }
116
117    pub(crate) async fn admit_dht_attempt(
118        &self,
119        attempt: PendingConnectionAttempt,
120    ) -> Result<bool> {
121        let Some(dht_ev) = self.transport.commit_connection_admission(attempt)? else {
122            return Ok(false);
123        };
124        // Local topology and lifecycle state are committed together. Remote
125        // convergence remains best-effort because the peer may churn immediately.
126        if let Err(error) = self.handle_dht_events(&dht_ev).await {
127            tracing::warn!(
128                peer = %attempt.peer(),
129                generation = attempt.generation(),
130                error = ?error,
131                "failed to handle DHT events after connection admission"
132            );
133        }
134        Ok(true)
135    }
136
137    pub(crate) async fn leave_dht_attempt(&self, attempt: PendingConnectionAttempt) -> Result<()> {
138        let should_repair = self
139            .dht
140            .peer_may_share_storage_responsibility(
141                attempt.peer(),
142                self.transport.storage_redundancy(),
143            )
144            .await?;
145        let removed = if self.transport.disconnect_attempt(attempt).await? {
146            true
147        } else {
148            self.transport.remove_retired_attempt_topology(attempt)?
149        };
150        if removed && should_repair {
151            self.transport.request_storage_repair();
152        }
153        Ok(())
154    }
155
156    #[cfg(all(test, not(all(feature = "wasm", target_family = "wasm"))))]
157    pub(crate) async fn leave_dht(&self, peer: Did) -> Result<()> {
158        let should_repair = self
159            .dht
160            .peer_may_share_storage_responsibility(peer, self.transport.storage_redundancy())
161            .await?;
162        self.dht.remove(peer)?;
163        if should_repair {
164            self.transport.request_storage_repair();
165        }
166        Ok(())
167    }
168
169    fn collect_dht_effects(
170        &self,
171        act: &PeerRingAction,
172        effects: &mut Vec<CoreEffect<'static>>,
173    ) -> Result<()> {
174        match act {
175            PeerRingAction::MultiActions(acts) => {
176                for act in acts {
177                    self.collect_dht_effects(act, effects)?;
178                }
179                Ok(())
180            }
181            act => {
182                if let Some(effect) =
183                    lower_dht_action(act, |did| self.transport.get_connection(did).is_some())?
184                {
185                    effects.push(effect);
186                }
187                Ok(())
188            }
189        }
190    }
191
192    async fn run_prioritized_dht_effects(&self, effects: Vec<CoreEffect<'static>>) -> Result<()> {
193        let mut connection_peers = Vec::new();
194        let mut other_effects = Vec::new();
195        for effect in effects {
196            match effect {
197                CoreEffect::ConnectDhtPeer { peer } => {
198                    connection_peers.push(peer);
199                }
200                effect => other_effects.push(effect),
201            }
202        }
203
204        let ordered_peers = self
205            .transport
206            .order_dht_candidates_by_quality(connection_peers)
207            .await;
208        for (peer, has_next) in core_actor_steps(ordered_peers) {
209            if let Err(e) = self.connect_dht_peer(peer).await {
210                tracing::error!("Failed on handle multi connection action: {e:?}");
211            }
212            if has_next || !other_effects.is_empty() {
213                yield_core_actor_step().await;
214            }
215        }
216
217        for (effect, has_next) in core_actor_steps(other_effects) {
218            if let Err(e) = self.run_effects([effect]).await {
219                tracing::error!("Failed on handle multi action: {e:?}");
220            }
221            if has_next {
222                yield_core_actor_step().await;
223            }
224        }
225
226        Ok(())
227    }
228
229    pub(crate) async fn handle_dht_events(&self, act: &PeerRingAction) -> Result<()> {
230        if matches!(act, PeerRingAction::MultiActions(_)) {
231            let mut effects = Vec::new();
232            self.collect_dht_effects(act, &mut effects)?;
233            self.run_prioritized_dht_effects(effects).await
234        } else {
235            let effects =
236                lower_dht_action(act, |did| self.transport.get_connection(did).is_some())?;
237            self.run_effects(effects).await
238        }
239    }
240}