Skip to main content

rings_core/message/handlers/stabilization/
mod.rs

1use 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;