rings_core/dht/storage/
sync.rs1use std::str::FromStr;
2
3use async_trait::async_trait;
4use rings_transport::core::transport::MAX_DATA_CHANNEL_MESSAGE_SIZE;
5use serde::Serialize;
6
7use super::StorageSyncDestination;
8use super::StorageSyncPurpose;
9use super::StorageSyncTarget;
10use crate::consts::MAX_CHUNK_ENVELOPE_OVERHEAD;
11use crate::consts::TRANSPORT_CUSTOM_OVERHEAD;
12use crate::dht::chord::PeerRing;
13use crate::dht::chord::PeerRingAction;
14use crate::dht::did::BiasId;
15use crate::dht::entry::Entry;
16use crate::dht::entry::PlacedEntry;
17use crate::dht::entry::SyncedEntryAck;
18use crate::dht::ChordStorageSync;
19use crate::dht::Did;
20use crate::error::Error;
21use crate::error::Result;
22use crate::message::types::Message;
23use crate::message::types::SyncEntriesWithSuccessor;
24
25pub(crate) const SYNC_BATCH_MAX_BYTES: usize = MAX_DATA_CHANNEL_MESSAGE_SIZE / 4;
31
32const SYNC_BATCH_ENVELOPE_HEADROOM_BYTES: usize =
33 MAX_CHUNK_ENVELOPE_OVERHEAD + TRANSPORT_CUSTOM_OVERHEAD;
34
35fn serialized_wire_size<T: Serialize>(value: &T) -> Result<usize> {
36 let bytes = rings_codec::serialized_size(value).map_err(Error::CodecSerialize)?;
37 usize::try_from(bytes).map_err(|_| Error::MessageSizeOverflow)
38}
39
40fn add_wire_cost(total: usize, next: usize) -> Result<usize> {
41 total.checked_add(next).ok_or(Error::MessageSizeOverflow)
42}
43
44fn sync_entries_fixed_wire_cost() -> Result<usize> {
45 let empty_message = Message::SyncEntriesWithSuccessor(SyncEntriesWithSuccessor {
46 purpose: StorageSyncPurpose::OwnershipHandoff,
47 destination: StorageSyncDestination::PhysicalOwner(Did::from(0u32)),
48 data: Vec::new(),
49 });
50 add_wire_cost(
51 serialized_wire_size(&empty_message)?,
52 SYNC_BATCH_ENVELOPE_HEADROOM_BYTES,
53 )
54}
55
56fn placed_entry_wire_cost(placed: &PlacedEntry) -> Result<usize> {
57 serialized_wire_size(placed)
58}
59
60#[cfg(all(test, not(all(feature = "wasm", target_family = "wasm"))))]
61pub(super) fn sync_entries_batch_wire_cost(data: &[PlacedEntry]) -> Result<usize> {
62 let mut cost = sync_entries_fixed_wire_cost()?;
63 for placed in data {
64 cost = add_wire_cost(cost, placed_entry_wire_cost(placed)?)?;
65 }
66 Ok(cost)
67}
68
69pub(super) fn sync_entries_batches(
70 data: Vec<PlacedEntry>,
71 max_batch_bytes: usize,
72) -> Result<Vec<Vec<PlacedEntry>>> {
73 let mut batches = Vec::new();
74 let mut current = Vec::new();
75 let fixed_cost = sync_entries_fixed_wire_cost()?;
76 let mut current_cost = fixed_cost;
77
78 for placed in data {
88 let placed_cost = placed_entry_wire_cost(&placed)?;
89 let candidate_cost = add_wire_cost(current_cost, placed_cost)?;
90 if current.is_empty() {
91 current.push(placed);
92 current_cost = candidate_cost;
93 continue;
94 }
95
96 if candidate_cost <= max_batch_bytes {
97 current.push(placed);
98 current_cost = candidate_cost;
99 } else {
100 batches.push(current);
101 current = vec![placed];
102 current_cost = add_wire_cost(fixed_cost, placed_cost)?;
103 }
104 }
105
106 if !current.is_empty() {
107 batches.push(current);
108 }
109
110 Ok(batches)
111}
112
113#[cfg_attr(all(feature = "wasm", target_family = "wasm"), async_trait(?Send))]
114#[cfg_attr(not(all(feature = "wasm", target_family = "wasm")), async_trait)]
115impl ChordStorageSync<PeerRingAction> for PeerRing {
116 async fn sync_entries_with_successor(&self, new_successor: Did) -> Result<PeerRingAction> {
120 if self.storage_virtual_nodes_enabled()? {
121 return self.copy_entries_to_observed_virtual_storage_owners().await;
122 }
123
124 let mut data = Vec::<PlacedEntry>::new();
125 let all_items: Vec<(String, Entry)> = self.storage.get_all().await?;
126
127 for (entry_key_str, entry) in all_items {
137 let entry_key = Did::from_str(&entry_key_str)?;
138 if BiasId::cmp_from_observer(self.did, entry_key, new_successor)
139 == std::cmp::Ordering::Greater
140 {
141 data.push(PlacedEntry::new(entry_key, entry));
142 }
143 }
144
145 let batches = sync_entries_batches(data, SYNC_BATCH_MAX_BYTES)?;
146 Ok(batches
147 .into_iter()
148 .map(|batch| {
149 PeerRingAction::sync_entries_for_handoff(
150 StorageSyncDestination::PhysicalOwner(new_successor),
151 batch,
152 )
153 })
154 .collect::<Vec<_>>()
155 .into())
156 }
157
158 async fn acknowledge_synced_entries(&self, acks: &[SyncedEntryAck]) -> Result<PeerRingAction> {
159 for ack in acks {
170 let Some(local_entry) = self.storage.get(&ack.key.to_string()).await? else {
171 continue;
172 };
173 if ack.confirms_local_value(&local_entry)? {
174 self.storage.remove(&ack.key.to_string()).await?;
175 }
176 }
177
178 Ok(PeerRingAction::None)
179 }
180}
181
182impl PeerRing {
183 async fn copy_entries_to_observed_virtual_storage_owners(&self) -> Result<PeerRingAction> {
184 let all_items: Vec<(String, Entry)> = self.storage.get_all().await?;
185 let mut by_target =
186 std::collections::BTreeMap::<StorageSyncDestination, Vec<PlacedEntry>>::new();
187
188 for (entry_key_str, entry) in all_items {
199 let entry_key = Did::from_str(&entry_key_str)?;
200 if let StorageSyncTarget::Remote(target) = self.storage_sync_target(entry_key)? {
201 by_target
202 .entry(target)
203 .or_default()
204 .push(PlacedEntry::new(entry_key, entry));
205 }
206 }
207
208 let mut actions = Vec::new();
209 for (target, data) in by_target {
210 for batch in sync_entries_batches(data, SYNC_BATCH_MAX_BYTES)? {
211 actions.push(PeerRingAction::sync_entries_for_repair(target, batch));
212 }
213 }
214 Ok(actions.into())
215 }
216}