Skip to main content

rings_core/dht/stabilization/
storage_repair.rs

1use super::Stabilizer;
2use super::STORAGE_REPAIR_FRESH_CONNECTION_GRACE_MS;
3use super::STORAGE_REPAIR_MAX_DELIVERIES_PER_STEP;
4use crate::dht::types::ChordStorageRepair;
5use crate::dht::Did;
6use crate::dht::PeerRingAction;
7use crate::dht::StorageSyncDelivery;
8use crate::error::Error;
9use crate::error::Result;
10use crate::message::SyncEntriesWithSuccessor;
11use crate::swarm::transport::TrackedStorageSyncOutcome;
12use crate::swarm::transport::TransportReadiness;
13use crate::utils::get_epoch_ms_i64;
14
15#[derive(Clone, Copy, Debug, Eq, PartialEq)]
16/// Observable completion state for one bounded storage repair pass.
17pub enum StorageRepairOutcome {
18    /// Every selected delivery completed.
19    Complete,
20    /// At least one delivery was deferred; repair remains pending for a later pass.
21    Deferred,
22}
23
24impl StorageRepairOutcome {
25    /// Return whether this pass completed all selected deliveries.
26    pub const fn is_complete(self) -> bool {
27        matches!(self, Self::Complete)
28    }
29}
30
31struct PlannedStorageRepairDelivery {
32    delivery: StorageSyncDelivery,
33}
34
35#[derive(Clone, Copy, Debug)]
36enum StorageRepairDeferReason {
37    MissingNextHop,
38    NextHopNotAdmitted,
39    NextHopTransportMissing,
40    NextHopTransportNotReady(TransportReadiness),
41    NextHopFresh { connected_for_ms: i64 },
42}
43
44#[derive(Clone, Copy, Debug, Eq, PartialEq)]
45enum RepairDeliveryResult {
46    Sent,
47    Deferred,
48}
49
50impl StorageRepairDeferReason {
51    const fn as_str(self) -> &'static str {
52        match self {
53            Self::MissingNextHop => "missing_next_hop",
54            Self::NextHopNotAdmitted => "next_hop_not_admitted",
55            Self::NextHopTransportMissing => "next_hop_transport_missing",
56            Self::NextHopTransportNotReady(_) => "next_hop_transport_not_ready",
57            Self::NextHopFresh { .. } => "next_hop_fresh",
58        }
59    }
60
61    const fn connected_for_ms(self) -> Option<i64> {
62        match self {
63            Self::NextHopFresh { connected_for_ms } => Some(connected_for_ms),
64            _ => None,
65        }
66    }
67
68    const fn transport_readiness(self) -> Option<TransportReadiness> {
69        match self {
70            Self::NextHopTransportNotReady(readiness) => Some(readiness),
71            _ => None,
72        }
73    }
74}
75
76fn is_storage_repair_deferral(error: &Error) -> bool {
77    error.is_deferrable_data_plane_send()
78}
79
80impl Stabilizer {
81    async fn handle_storage_repair_action(
82        &self,
83        act: PeerRingAction,
84    ) -> Result<StorageRepairOutcome> {
85        let deliveries = self.storage_repair_window(act.coalesced_storage_sync_deliveries()?)?;
86        if deliveries.is_empty() {
87            tracing::debug!(
88                target: "rings_core::dht::stabilization",
89                local = %self.dht.did,
90                "STABILIZATION storage repair has no deliveries"
91            );
92            return Ok(StorageRepairOutcome::Complete);
93        }
94
95        tracing::debug!(
96            target: "rings_core::dht::stabilization",
97            local = %self.dht.did,
98            deliveries = deliveries.len(),
99            "STABILIZATION storage repair deliveries prepared"
100        );
101
102        let mut sent = 0usize;
103        let mut deferred = 0usize;
104        for planned in deliveries {
105            match self
106                .send_planned_storage_repair(planned, get_epoch_ms_i64())
107                .await?
108            {
109                RepairDeliveryResult::Sent => sent = sent.saturating_add(1),
110                RepairDeliveryResult::Deferred => deferred = deferred.saturating_add(1),
111            }
112        }
113
114        if deferred > 0 {
115            tracing::debug!(
116                target: "rings_core::dht::stabilization",
117                local = %self.dht.did,
118                sent,
119                deferred,
120                "STABILIZATION storage repair deliveries finished with deferrals"
121            );
122            return Ok(StorageRepairOutcome::Deferred);
123        }
124        Ok(StorageRepairOutcome::Complete)
125    }
126
127    async fn send_planned_storage_repair(
128        &self,
129        planned: PlannedStorageRepairDelivery,
130        now_ms: i64,
131    ) -> Result<RepairDeliveryResult> {
132        let msg = SyncEntriesWithSuccessor::from_delivery(planned.delivery);
133        let purpose = msg.purpose;
134        let destination = msg.destination;
135        let entries = msg.data.len();
136        let next_hop = self.dht.next_hop_for_storage_sync(destination)?;
137        if let Some(reason) = self.storage_repair_defer_reason(next_hop, now_ms)? {
138            self.log_storage_repair_deferred(destination, next_hop, entries, reason, None);
139            return Ok(RepairDeliveryResult::Deferred);
140        }
141        tracing::debug!(
142            target: "rings_core::dht::stabilization",
143            local = %self.dht.did,
144            purpose = ?purpose,
145            destination = ?destination,
146            next_hop = ?next_hop,
147            entries,
148            "STABILIZATION storage repair send start"
149        );
150        match self.transport.send_storage_sync_tracked(msg).await {
151            Ok(TrackedStorageSyncOutcome::Delivered(tx_id)) => {
152                #[cfg(all(test, feature = "dummy", not(target_family = "wasm")))]
153                crate::simulation::record_repair_entries(entries);
154                tracing::debug!(
155                    target: "rings_core::dht::stabilization",
156                    local = %self.dht.did,
157                    tx_id = %tx_id,
158                    destination = ?destination,
159                    next_hop = ?next_hop,
160                    entries,
161                    "STABILIZATION storage repair send complete"
162                );
163                Ok(RepairDeliveryResult::Sent)
164            }
165            Ok(TrackedStorageSyncOutcome::PersistedLocally) => {
166                tracing::debug!(
167                    target: "rings_core::dht::stabilization",
168                    local = %self.dht.did,
169                    destination = ?destination,
170                    entries,
171                    "STABILIZATION storage repair persisted locally"
172                );
173                Ok(RepairDeliveryResult::Sent)
174            }
175            Ok(TrackedStorageSyncOutcome::Deferred) => {
176                tracing::warn!(
177                    target: "rings_core::dht::stabilization",
178                    local = %self.dht.did,
179                    destination = ?destination,
180                    next_hop = ?next_hop,
181                    entries,
182                    "STABILIZATION storage repair delivery cancelled and deferred"
183                );
184                Ok(RepairDeliveryResult::Deferred)
185            }
186            Err(error) if is_storage_repair_deferral(&error) => {
187                tracing::warn!(
188                    target: "rings_core::dht::stabilization",
189                    local = %self.dht.did,
190                    destination = ?destination,
191                    next_hop = ?next_hop,
192                    entries,
193                    error = ?error,
194                    "STABILIZATION storage repair deferred by transport readiness"
195                );
196                Ok(RepairDeliveryResult::Deferred)
197            }
198            Err(error) => Err(error),
199        }
200    }
201
202    fn log_storage_repair_deferred(
203        &self,
204        destination: crate::dht::StorageSyncDestination,
205        next_hop: Option<Did>,
206        entries: usize,
207        reason: StorageRepairDeferReason,
208        error: Option<&Error>,
209    ) {
210        tracing::debug!(
211            target: "rings_core::dht::stabilization",
212            local = %self.dht.did,
213            destination = ?destination,
214            next_hop = ?next_hop,
215            entries,
216            reason = reason.as_str(),
217            transport_readiness = ?reason.transport_readiness(),
218            connected_for_ms = ?reason.connected_for_ms(),
219            grace_ms = STORAGE_REPAIR_FRESH_CONNECTION_GRACE_MS,
220            error = ?error,
221            "STABILIZATION storage repair deferred"
222        );
223    }
224
225    fn storage_repair_window(
226        &self,
227        deliveries: Vec<StorageSyncDelivery>,
228    ) -> Result<Vec<PlannedStorageRepairDelivery>> {
229        let total = deliveries.len();
230        if total == 0 {
231            return Ok(Vec::new());
232        }
233
234        let mut keyed = deliveries
235            .into_iter()
236            .map(|delivery| (delivery.cursor_key(), delivery))
237            .collect::<Vec<_>>();
238        keyed.sort_by(|(left, _), (right, _)| left.cmp(right));
239        let ordered = keyed
240            .iter()
241            .map(|(cursor, _)| cursor.clone())
242            .collect::<Vec<_>>();
243        let start = self
244            .transport
245            .storage_repair_window_start(&ordered, STORAGE_REPAIR_MAX_DELIVERIES_PER_STEP)?;
246        keyed.rotate_left(start);
247        keyed.truncate(STORAGE_REPAIR_MAX_DELIVERIES_PER_STEP);
248        let next_cursor = keyed.last().map(|(cursor, _)| cursor.clone());
249        let selected = keyed
250            .into_iter()
251            .map(|(_, delivery)| PlannedStorageRepairDelivery { delivery })
252            .collect::<Vec<_>>();
253        if let Some(cursor) = next_cursor {
254            self.transport.advance_storage_repair_cursor(cursor)?;
255        }
256        tracing::debug!(
257            target: "rings_core::dht::stabilization",
258            local = %self.dht.did,
259            total_deliveries = total,
260            selected_deliveries = selected.len(),
261            start,
262            "STABILIZATION storage repair delivery window selected"
263        );
264        Ok(selected)
265    }
266
267    fn storage_repair_defer_reason(
268        &self,
269        next_hop: Option<Did>,
270        now_ms: i64,
271    ) -> Result<Option<StorageRepairDeferReason>> {
272        let Some(next_hop) = next_hop else {
273            return Ok(Some(StorageRepairDeferReason::MissingNextHop));
274        };
275        if !self.transport.is_admitted_connection(next_hop) {
276            return Ok(Some(StorageRepairDeferReason::NextHopNotAdmitted));
277        }
278        let Some(next_hop_connection) = self.transport.admitted_connection(next_hop)? else {
279            return Ok(Some(StorageRepairDeferReason::NextHopTransportMissing));
280        };
281        let readiness = next_hop_connection.readiness();
282        if !readiness.can_make_progress() {
283            return Ok(Some(StorageRepairDeferReason::NextHopTransportNotReady(
284                readiness,
285            )));
286        }
287        if let Some(connected_for_ms) = self.peer_connected_for_ms(next_hop, now_ms) {
288            if connected_for_ms < STORAGE_REPAIR_FRESH_CONNECTION_GRACE_MS {
289                return Ok(Some(StorageRepairDeferReason::NextHopFresh {
290                    connected_for_ms,
291                }));
292            }
293        }
294
295        Ok(None)
296    }
297
298    fn peer_connected_for_ms(&self, peer: Did, now_ms: i64) -> Option<i64> {
299        match self.transport.peer_connected_for_ms(peer, now_ms) {
300            Ok(age) => age,
301            Err(error) => {
302                tracing::warn!(
303                    target: "rings_core::dht::stabilization",
304                    local = %self.dht.did,
305                    peer = %peer,
306                    error = %error,
307                    "STABILIZATION storage repair connection age check failed"
308                );
309                None
310            }
311        }
312    }
313
314    /// Republish locally-held entries to their current affine owners.
315    pub async fn repair_storage(&self) -> Result<StorageRepairOutcome> {
316        tracing::debug!(
317            target: "rings_core::dht::stabilization",
318            local = %self.dht.did,
319            redundancy = self.transport.storage_redundancy(),
320            "STABILIZATION repair_storage republish start"
321        );
322        let action = self
323            .dht
324            .republish_local_entries(self.transport.storage_redundancy())
325            .await?;
326        let (action_kind, action_count) = match &action {
327            PeerRingAction::None => ("None", 0),
328            PeerRingAction::Some(_) => ("Some", 1),
329            PeerRingAction::SomeEntry(_) => ("SomeEntry", 1),
330            PeerRingAction::EntryMisses(misses) => ("EntryMisses", misses.len()),
331            PeerRingAction::RemoteAction(_, _) => ("RemoteAction", 1),
332            PeerRingAction::MultiActions(actions) => ("MultiActions", actions.len()),
333        };
334        tracing::debug!(
335            target: "rings_core::dht::stabilization",
336            local = %self.dht.did,
337            action_kind,
338            action_count,
339            "STABILIZATION repair_storage republish action prepared"
340        );
341        let outcome = self.handle_storage_repair_action(action).await?;
342        tracing::debug!(
343            target: "rings_core::dht::stabilization",
344            local = %self.dht.did,
345            "STABILIZATION repair_storage republish complete"
346        );
347        Ok(outcome)
348    }
349}
350
351#[cfg(test)]
352mod tests {
353    use std::sync::Arc;
354
355    use super::*;
356    use crate::dht::entry::Entry;
357    use crate::dht::entry::EntryKind;
358    use crate::dht::entry::PlacedEntry;
359    use crate::dht::StorageSyncDestination;
360    use crate::ecc::SecretKey;
361    use crate::session::SessionSk;
362    use crate::storage::MemStorage;
363    use crate::swarm::SwarmBuilder;
364
365    fn repair_deliveries(values: &[u32]) -> Result<Vec<StorageSyncDelivery>> {
366        let actions = values
367            .iter()
368            .copied()
369            .map(|value| {
370                let destination = StorageSyncDestination::PlacementKey(Did::from(value));
371                let entry = Entry::new(Did::from(value + 100), vec![], EntryKind::Data);
372                PeerRingAction::sync_entries_for_repair(destination, vec![PlacedEntry::new(
373                    destination.did(),
374                    entry,
375                )])
376            })
377            .collect();
378        PeerRingAction::MultiActions(actions).coalesced_storage_sync_deliveries()
379    }
380
381    fn selected_destination(stabilizer: &Stabilizer, values: &[u32]) -> Result<Did> {
382        let mut selected = stabilizer.storage_repair_window(repair_deliveries(values)?)?;
383        let Some(planned) = selected.pop() else {
384            return Err(crate::error::Error::InvalidMessage(
385                "repair window selected no delivery".to_string(),
386            ));
387        };
388        let destination = planned.delivery.into_message_parts().1.did();
389        Ok(destination)
390    }
391
392    #[test]
393    fn test_changing_delivery_sets_preserve_repair_progress_across_stabilizers() -> Result<()> {
394        let session = SessionSk::new_with_seckey(&SecretKey::random())?;
395        let swarm = Arc::new(
396            SwarmBuilder::new(
397                0,
398                "stun://stun.l.google.com:19302",
399                Box::new(MemStorage::new()),
400                session,
401            )
402            .build(),
403        );
404
405        let selected = [
406            selected_destination(&swarm.stabilizer(), &[1, 2, 3])?,
407            selected_destination(&swarm.stabilizer(), &[2, 3])?,
408            selected_destination(&swarm.stabilizer(), &[1, 2, 3])?,
409            selected_destination(&swarm.stabilizer(), &[1, 2, 3])?,
410        ];
411
412        assert_eq!(selected, [
413            Did::from(1u32),
414            Did::from(2u32),
415            Did::from(3u32),
416            Did::from(1u32),
417        ]);
418        Ok(())
419    }
420
421    #[test]
422    fn test_deferred_delivery_rotates_without_losing_its_retry() -> Result<()> {
423        let session = SessionSk::new_with_seckey(&SecretKey::random())?;
424        let swarm = Arc::new(
425            SwarmBuilder::new(
426                0,
427                "stun://stun.l.google.com:19302",
428                Box::new(MemStorage::new()),
429                session,
430            )
431            .build(),
432        );
433        let stabilizer = swarm.stabilizer();
434
435        let selected = [
436            selected_destination(&stabilizer, &[1, 2])?,
437            selected_destination(&stabilizer, &[1, 2])?,
438            selected_destination(&stabilizer, &[1, 2])?,
439        ];
440
441        assert_eq!(selected, [
442            Did::from(1u32),
443            Did::from(2u32),
444            Did::from(1u32)
445        ]);
446        Ok(())
447    }
448}