use std::collections::BTreeMap;
use std::collections::BTreeSet;
use super::SwarmTransport;
use crate::dht::entry::PlacementMiss;
use crate::dht::Did;
use crate::error::Error;
use crate::error::Result;
use crate::utils::get_epoch_ms_i64;
const STORAGE_LOOKUP_OBSERVATION_TTL_MS: i64 = 30_000;
pub(crate) const STORAGE_LOOKUP_OBSERVATION_CAPACITY: usize = 1024;
pub(super) type StorageLookupObservationMap =
BTreeMap<StorageLookupObservationKey, StorageLookupObservation>;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd)]
pub(super) struct StorageLookupObservationKey {
resource: Did,
redundancy: u16,
}
pub(super) struct StorageLookupObservation {
observed_at_ms: i64,
misses: BTreeSet<PlacementMiss>,
}
fn storage_lookup_observation_now_ms() -> i64 {
get_epoch_ms_i64()
}
fn oldest_storage_lookup_observation_key(
observations: &StorageLookupObservationMap,
) -> Option<StorageLookupObservationKey> {
observations
.iter()
.min_by_key(|(key, observation)| (observation.observed_at_ms, **key))
.map(|(key, _)| *key)
}
fn evict_storage_lookup_observations(observations: &mut StorageLookupObservationMap, now_ms: i64) {
observations.retain(|_, observation| {
now_ms.saturating_sub(observation.observed_at_ms) <= STORAGE_LOOKUP_OBSERVATION_TTL_MS
});
while observations.len() > STORAGE_LOOKUP_OBSERVATION_CAPACITY {
let Some(stale_key) = oldest_storage_lookup_observation_key(observations) else {
break;
};
observations.remove(&stale_key);
}
}
fn reserve_storage_lookup_observation_slot(observations: &mut StorageLookupObservationMap) {
while observations.len() >= STORAGE_LOOKUP_OBSERVATION_CAPACITY {
let Some(stale_key) = oldest_storage_lookup_observation_key(observations) else {
break;
};
observations.remove(&stale_key);
}
}
impl SwarmTransport {
fn storage_lookup_observation_key(
&self,
resource: Did,
redundancy: u16,
) -> Result<StorageLookupObservationKey> {
self.ensure_storage_redundancy_value(redundancy)?;
Ok(StorageLookupObservationKey {
resource,
redundancy,
})
}
pub(crate) fn start_storage_lookup(&self, resource: Did, redundancy: u16) -> Result<()> {
let key = self.storage_lookup_observation_key(resource, redundancy)?;
let mut observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
let now = storage_lookup_observation_now_ms();
evict_storage_lookup_observations(&mut observations, now);
reserve_storage_lookup_observation_slot(&mut observations);
observations.insert(key, StorageLookupObservation {
observed_at_ms: now,
misses: BTreeSet::new(),
});
Ok(())
}
pub(crate) fn ensure_storage_lookup_active(
&self,
resource: Did,
redundancy: u16,
) -> Result<()> {
let key = self.storage_lookup_observation_key(resource, redundancy)?;
let mut observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
let now = storage_lookup_observation_now_ms();
evict_storage_lookup_observations(&mut observations, now);
if observations.contains_key(&key) {
Ok(())
} else {
Err(Error::InvalidMessage(
"storage lookup response has no active local lookup".to_string(),
))
}
}
pub(crate) fn observe_storage_misses(
&self,
resource: Did,
redundancy: u16,
misses: impl IntoIterator<Item = PlacementMiss>,
) -> Result<()> {
let key = self.storage_lookup_observation_key(resource, redundancy)?;
let mut misses = misses.into_iter().peekable();
if misses.peek().is_none() {
return Ok(());
}
let mut observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
let now = storage_lookup_observation_now_ms();
evict_storage_lookup_observations(&mut observations, now);
let Some(observation) = observations.get_mut(&key) else {
return Err(Error::InvalidMessage(
"storage miss observation has no active local lookup".to_string(),
));
};
observation.observed_at_ms = now;
observation.misses.extend(misses);
evict_storage_lookup_observations(&mut observations, now);
Ok(())
}
pub(crate) fn take_storage_misses(
&self,
resource: Did,
redundancy: u16,
) -> Result<Vec<PlacementMiss>> {
let key = self.storage_lookup_observation_key(resource, redundancy)?;
let mut observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
let now = storage_lookup_observation_now_ms();
evict_storage_lookup_observations(&mut observations, now);
let Some(observation) = observations.get_mut(&key) else {
return Err(Error::InvalidMessage(
"storage repair has no active local lookup".to_string(),
));
};
Ok(std::mem::take(&mut observation.misses)
.into_iter()
.collect())
}
#[cfg(all(test, not(all(feature = "wasm", target_family = "wasm"))))]
pub(crate) fn expire_storage_lookup_observation(
&self,
resource: Did,
redundancy: u16,
) -> Result<()> {
let key = self.storage_lookup_observation_key(resource, redundancy)?;
let mut observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
if let Some(observation) = observations.get_mut(&key) {
observation.observed_at_ms = storage_lookup_observation_now_ms()
.saturating_sub(STORAGE_LOOKUP_OBSERVATION_TTL_MS + 1);
}
Ok(())
}
#[cfg(all(test, not(all(feature = "wasm", target_family = "wasm"))))]
pub(crate) fn storage_lookup_observation_count(&self) -> Result<usize> {
let observations = self
.storage_lookup_observations
.lock()
.map_err(|_| Error::DHTSyncLockError)?;
Ok(observations.len())
}
}