rings_core/message/handlers/
mod.rs1#![deny(missing_docs)]
2use 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
24pub mod connection;
26pub mod custom;
28pub mod e2e;
30pub mod stabilization;
32pub mod storage;
34#[derive(Clone)]
39pub struct MessageHandler {
40 transport: Arc<SwarmTransport>,
41 dht: Arc<PeerRing>,
42 swarm_callback: SharedSwarmCallback,
43}
44
45#[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 async fn handle(&self, ctx: &MessagePayload, msg: &T) -> Result<()>;
51}
52
53impl MessageHandler {
54 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 pub(crate) async fn connect_dht_peer(&self, peer: Did) -> Result<()> {
83 self.run_effects([CoreEffect::connect_dht_peer(peer)]).await
84 }
85
86 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 let Some(dht_ev) = self.transport.join_routable_peer(peer)? else {
106 return Err(Error::SwarmMissDidInTable(peer));
107 };
108 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 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}