use std::sync::Arc;
use tokio::time::{Duration, Instant};
use kitsune_p2p_types::{tx2::tx2_utils::ShareOpen, KAgent, KSpace };
use linked_hash_map::{Entry, LinkedHashMap};
use crate::{FetchContext, FetchKey, FetchPoolPush, RoughInt};
mod pool_reader;
pub use pool_reader::*;
const NUM_ITEMS_PER_POLL: usize = 100;
#[derive(Clone)]
pub struct FetchPool {
config: FetchConfig,
state: ShareOpen<State>,
}
impl std::fmt::Debug for FetchPool {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.state
.share_ref(|state| f.debug_struct("FetchPool").field("state", state).finish())
}
}
pub type FetchConfig = Arc<dyn FetchPoolConfig>;
pub trait FetchPoolConfig: 'static + Send + Sync {
fn item_retry_delay(&self) -> std::time::Duration {
std::time::Duration::from_secs(90)
}
fn source_retry_delay(&self) -> std::time::Duration {
std::time::Duration::from_secs(5 * 60)
}
fn merge_fetch_contexts(&self, a: u32, b: u32) -> u32;
}
#[derive(Debug)]
pub struct State {
queue: LinkedHashMap<FetchKey, FetchPoolItem>,
}
#[allow(clippy::derivable_impls)]
impl Default for State {
fn default() -> Self {
Self {
queue: Default::default(),
}
}
}
pub struct StateIter<'a> {
state: &'a mut State,
config: &'a dyn FetchPoolConfig,
}
#[derive(Debug, PartialEq, Eq)]
struct Sources(Vec<SourceRecord>);
#[derive(Debug, PartialEq, Eq)]
pub struct FetchPoolItem {
sources: Sources,
space: KSpace,
size: Option<RoughInt>,
pub context: Option<FetchContext>,
last_fetch: Option<Instant>,
}
#[derive(Debug, PartialEq, Eq)]
struct SourceRecord {
source: FetchSource,
last_request: Option<Instant>,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub enum FetchSource {
Agent(KAgent),
}
struct FetchPoolConfigBitwiseOr;
impl FetchPoolConfig for FetchPoolConfigBitwiseOr {
fn merge_fetch_contexts(&self, a: u32, b: u32) -> u32 {
a | b
}
}
impl FetchPool {
pub fn new(config: FetchConfig) -> Self {
Self {
config,
state: ShareOpen::new(State::default()),
}
}
pub fn new_bitwise_or() -> Self {
Self {
config: Arc::new(FetchPoolConfigBitwiseOr),
state: ShareOpen::new(State::default()),
}
}
pub fn push(&self, args: FetchPoolPush) {
self.state.share_mut(|s| {
tracing::debug!(
"FetchPool (size = {}) item added: {:?}",
s.queue.len() + 1,
args
);
s.push(&*self.config, args);
});
}
pub fn remove(&self, key: &FetchKey) -> Option<FetchPoolItem> {
self.state.share_mut(|s| {
let removed = s.remove(key);
tracing::debug!(
"FetchPool (size = {}) item removed: key={:?} val={:?}",
s.queue.len(),
key,
removed
);
removed
})
}
pub fn get_items_to_fetch(&self) -> Vec<(FetchKey, KSpace, FetchSource, Option<FetchContext>)> {
self.state.share_mut(|s| {
let mut out = Vec::new();
for (key, space, source, context) in s.iter_mut(&*self.config) {
out.push((key, space, source, context));
}
out
})
}
}
impl State {
pub fn push(&mut self, config: &dyn FetchPoolConfig, args: FetchPoolPush) {
let FetchPoolPush {
key,
author,
context,
space,
source,
size,
} = args;
match self.queue.entry(key) {
Entry::Vacant(e) => {
let sources = if let Some(author) = author {
Sources(vec![SourceRecord::new(source), SourceRecord::agent(author)])
} else {
Sources(vec![SourceRecord::new(source)])
};
let item = FetchPoolItem {
sources,
space,
size,
context,
last_fetch: None,
};
e.insert(item);
}
Entry::Occupied(mut e) => {
let v = e.get_mut();
v.sources.0.insert(0, SourceRecord::new(source));
v.context = match (v.context.take(), context) {
(Some(a), Some(b)) => Some(config.merge_fetch_contexts(*a, *b).into()),
(a, b) => a.and(b),
}
}
}
}
pub fn iter_mut<'a>(&'a mut self, config: &'a dyn FetchPoolConfig) -> StateIter {
StateIter {
state: self,
config,
}
}
pub fn remove(&mut self, key: &FetchKey) -> Option<FetchPoolItem> {
self.queue.remove(key)
}
}
impl<'a> Iterator for StateIter<'a> {
type Item = (FetchKey, KSpace, FetchSource, Option<FetchContext>);
fn next(&mut self) -> Option<Self::Item> {
let keys: Vec<_> = self
.state
.queue
.keys()
.take(NUM_ITEMS_PER_POLL)
.cloned()
.collect();
for key in keys {
let item = self.state.queue.get_refresh(&key)?;
let item_not_recently_fetched = item
.last_fetch
.map(|t| t.elapsed() >= self.config.item_retry_delay())
.unwrap_or(true);
if item_not_recently_fetched {
if let Some(source) = item.sources.next(self.config.source_retry_delay()) {
let space = item.space.clone();
item.last_fetch = Some(Instant::now());
return Some((key, space, source, item.context));
}
}
}
None
}
}
impl SourceRecord {
fn new(source: FetchSource) -> Self {
Self {
source,
last_request: None,
}
}
fn agent(agent: KAgent) -> Self {
Self {
source: FetchSource::Agent(agent),
last_request: None,
}
}
}
impl Sources {
fn next(&mut self, interval: Duration) -> Option<FetchSource> {
if let Some((i, agent)) = self
.0
.iter()
.enumerate()
.find(|(_, s)| {
s.last_request
.map(|t| t.elapsed() >= interval)
.unwrap_or(true)
})
.map(|(i, s)| (i, s.source.clone()))
{
self.0[i].last_request = Some(Instant::now());
self.0.rotate_left(i + 1);
Some(agent)
} else {
None
}
}
}
#[cfg(test)]
mod tests {
use pretty_assertions::assert_eq;
use std::{sync::Arc, time::Duration};
use kitsune_p2p_types::bin_types::{KitsuneAgent, KitsuneBinType, KitsuneOpHash, KitsuneSpace};
use super::*;
pub(super) struct Config(pub u32, pub u32);
impl FetchPoolConfig for Config {
fn merge_fetch_contexts(&self, a: u32, b: u32) -> u32 {
(a + b).min(1)
}
fn item_retry_delay(&self) -> Duration {
Duration::from_secs(self.0 as u64)
}
fn source_retry_delay(&self) -> Duration {
Duration::from_secs(self.1 as u64)
}
}
pub(super) fn key_op(n: u8) -> FetchKey {
FetchKey::Op(Arc::new(KitsuneOpHash::new(vec![n; 36])))
}
pub(super) fn req(n: u8, context: Option<FetchContext>, source: FetchSource) -> FetchPoolPush {
FetchPoolPush {
key: key_op(n),
author: None,
context,
space: space(0),
source,
size: None,
}
}
pub(super) fn item(
_cfg: &dyn FetchPoolConfig,
sources: Vec<FetchSource>,
context: Option<FetchContext>,
) -> FetchPoolItem {
FetchPoolItem {
sources: Sources(sources.into_iter().map(|s| SourceRecord::new(s)).collect()),
space: Arc::new(KitsuneSpace::new(vec![0; 36])),
context,
size: None,
last_fetch: None,
}
}
pub(super) fn space(i: u8) -> KSpace {
Arc::new(KitsuneSpace::new(vec![i; 36]))
}
pub(super) fn source(i: u8) -> FetchSource {
FetchSource::Agent(Arc::new(KitsuneAgent::new(vec![i; 36])))
}
pub(super) fn sources(ix: impl IntoIterator<Item = u8>) -> Vec<FetchSource> {
ix.into_iter().map(source).collect()
}
pub(super) fn ctx(c: u32) -> Option<FetchContext> {
Some(c.into())
}
#[tokio::test(start_paused = true)]
async fn source_rotation() {
let sec1 = Duration::from_secs(10);
let mut ss = Sources(vec![
SourceRecord {
source: source(1),
last_request: Some(Instant::now()),
}
.into(),
SourceRecord {
source: source(2),
last_request: None,
}
.into(),
]);
tokio::time::advance(Duration::from_secs(1)).await;
assert_eq!(ss.next(sec1), Some(source(2)));
assert_eq!(ss.next(sec1), None);
tokio::time::advance(Duration::from_secs(9)).await;
assert_eq!(ss.next(sec1), Some(source(1)));
tokio::time::advance(Duration::from_secs(1)).await;
assert_eq!(ss.next(sec1), Some(source(2)));
assert_eq!(ss.next(sec1), None);
tokio::time::advance(Duration::from_secs(20)).await;
assert_eq!(ss.next(sec1), Some(source(1)));
assert_eq!(ss.next(sec1), Some(source(2)));
assert_eq!(ss.next(sec1), None);
}
#[test]
fn queue_push() {
let mut q = State::default();
let c = Config(1, 1);
q.push(&c, req(1, ctx(1), source(1)));
q.push(&c, req(1, ctx(0), source(0)));
q.push(&c, req(2, ctx(0), source(0)));
let expected_ready = [
(key_op(1), item(&c, sources(0..=1), ctx(1))),
(key_op(2), item(&c, sources([0]), ctx(0))),
]
.into_iter()
.collect();
assert_eq!(q.queue, expected_ready);
}
#[tokio::test(start_paused = true)]
async fn queue_next() {
let cfg = Config(1, 10);
let mut q = {
let mut queue = [
(key_op(1), item(&cfg, sources(0..=2), ctx(1))),
(key_op(2), item(&cfg, sources(1..=3), ctx(1))),
(key_op(3), item(&cfg, sources(2..=4), ctx(1))),
];
queue[1].1.sources.0[1].last_request = Some(Instant::now() - Duration::from_secs(3));
let queue = queue.into_iter().collect();
State { queue }
};
assert_eq!(q.iter_mut(&cfg).count(), 3);
tokio::time::advance(Duration::from_secs(1)).await;
assert_eq!(q.iter_mut(&cfg).count(), 3);
tokio::time::advance(Duration::from_secs(1)).await;
assert_eq!(q.iter_mut(&cfg).count(), 2);
tokio::time::advance(Duration::from_secs(5)).await;
assert_eq!(
q.iter_mut(&cfg).collect::<Vec<_>>(),
vec![(key_op(2), space(0), source(2), ctx(1))]
);
assert_eq!(q.iter_mut(&cfg).count(), 0);
tokio::time::advance(Duration::from_secs(4)).await;
assert_eq!(q.iter_mut(&cfg).count(), 3);
}
}