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)]
16pub enum StorageRepairOutcome {
18 Complete,
20 Deferred,
22}
23
24impl StorageRepairOutcome {
25 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 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}