rings_core/message/handlers/stabilization/
mod.rs1use async_trait::async_trait;
2
3use crate::dht::ChordStorageSync;
4use crate::error::Error;
5use crate::error::Result;
6use crate::message::effects::CoreEffect;
7use crate::message::types::Message;
8use crate::message::types::NotifyPredecessorReport;
9use crate::message::types::NotifyPredecessorSend;
10use crate::message::types::SyncEntriesWithSuccessor;
11use crate::message::HandleMsg;
12use crate::message::MessageHandler;
13use crate::message::MessagePayload;
14
15#[cfg_attr(all(feature = "wasm", target_family = "wasm"), async_trait(?Send))]
16#[cfg_attr(not(all(feature = "wasm", target_family = "wasm")), async_trait)]
17impl HandleMsg<NotifyPredecessorSend> for MessageHandler {
18 async fn handle(&self, ctx: &MessagePayload, msg: &NotifyPredecessorSend) -> Result<()> {
19 if ctx.should_forward_from(self.dht.did) {
20 return self
21 .run_effects([CoreEffect::forward_payload(ctx, None)])
22 .await;
23 }
24
25 let origin = self.verified_notify_predecessor_origin(ctx, msg)?;
26 let Some(predecessor) = self.transport.notify_admitted_predecessor(origin)? else {
27 return Err(Error::NotifyPredecessorOriginNotAdmitted { origin });
28 };
29
30 if predecessor != origin {
31 return self
32 .run_effects([CoreEffect::send_report_message(
33 ctx,
34 Message::NotifyPredecessorReport(NotifyPredecessorReport { did: predecessor }),
35 )])
36 .await;
37 }
38
39 Ok(())
40 }
41}
42
43impl MessageHandler {
44 fn verified_notify_predecessor_origin(
45 &self,
46 ctx: &MessagePayload,
47 msg: &NotifyPredecessorSend,
48 ) -> Result<crate::dht::Did> {
49 let origin = ctx.relay.try_origin_sender()?;
50 if msg.did != origin {
51 return Err(Error::NotifyPredecessorOriginMismatch {
52 claimed: msg.did,
53 origin,
54 });
55 }
56 Ok(origin)
57 }
58}
59
60#[cfg_attr(all(feature = "wasm", target_family = "wasm"), async_trait(?Send))]
61#[cfg_attr(not(all(feature = "wasm", target_family = "wasm")), async_trait)]
62impl HandleMsg<NotifyPredecessorReport> for MessageHandler {
63 async fn handle(&self, _ctx: &MessagePayload, msg: &NotifyPredecessorReport) -> Result<()> {
64 self.run_effects([CoreEffect::connect_dht_peer(msg.did)])
65 .await?;
66
67 let deliveries = self
68 .dht
69 .sync_entries_with_successor(msg.did)
70 .await?
71 .coalesced_storage_sync_deliveries()?;
72 let effects = deliveries.into_iter().map(|delivery| {
73 let msg = SyncEntriesWithSuccessor::from_delivery(delivery);
74 CoreEffect::send_storage_sync(msg)
75 });
76 self.run_effects(effects).await?;
77
78 Ok(())
79 }
80}
81
82#[cfg(not(all(feature = "wasm", target_family = "wasm")))]
83#[cfg(test)]
84mod tests;