protocol/fetcher: Remove `from` from `QueuedFetch`
This commit is contained in:
parent
53a1286901
commit
9081c11b98
|
|
@ -368,7 +368,6 @@ impl radicle::node::Handle for Handle {
|
||||||
"queue": queue.iter().map(|fetch| {
|
"queue": queue.iter().map(|fetch| {
|
||||||
json!({
|
json!({
|
||||||
"rid": fetch.rid,
|
"rid": fetch.rid,
|
||||||
"from": fetch.from,
|
|
||||||
"refsAt": fetch.refs_at,
|
"refsAt": fetch.refs_at,
|
||||||
})
|
})
|
||||||
}).collect::<Vec<_>>()
|
}).collect::<Vec<_>>()
|
||||||
|
|
|
||||||
|
|
@ -191,13 +191,11 @@ impl FetcherState {
|
||||||
.or_insert(Queue::new(self.config.maximum_queue_size));
|
.or_insert(Queue::new(self.config.maximum_queue_size));
|
||||||
match queue.enqueue(QueuedFetch {
|
match queue.enqueue(QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from,
|
|
||||||
refs_at,
|
refs_at,
|
||||||
timeout,
|
timeout,
|
||||||
}) {
|
}) {
|
||||||
Enqueue::CapacityReached(QueuedFetch {
|
Enqueue::CapacityReached(QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from,
|
|
||||||
refs_at,
|
refs_at,
|
||||||
timeout,
|
timeout,
|
||||||
}) => event::Fetch::QueueAtCapacity {
|
}) => event::Fetch::QueueAtCapacity {
|
||||||
|
|
@ -301,9 +299,6 @@ impl ActiveFetch {
|
||||||
pub struct QueuedFetch {
|
pub struct QueuedFetch {
|
||||||
/// The repository that will be fetched.
|
/// The repository that will be fetched.
|
||||||
pub rid: RepoId,
|
pub rid: RepoId,
|
||||||
// TODO(finto): this might be redundant, since queues are per node
|
|
||||||
/// The peer from which the repository will be fetched from.
|
|
||||||
pub from: NodeId,
|
|
||||||
/// The references that the fetch is being performed for.
|
/// The references that the fetch is being performed for.
|
||||||
pub refs_at: Vec<RefsAt>,
|
pub refs_at: Vec<RefsAt>,
|
||||||
/// The timeout given for the fetch request.
|
/// The timeout given for the fetch request.
|
||||||
|
|
|
||||||
|
|
@ -8,7 +8,7 @@ use std::time::Duration;
|
||||||
use qcheck::Arbitrary;
|
use qcheck::Arbitrary;
|
||||||
|
|
||||||
use radicle::storage::refs::RefsAt;
|
use radicle::storage::refs::RefsAt;
|
||||||
use radicle_core::{NodeId, RepoId};
|
use radicle_core::RepoId;
|
||||||
|
|
||||||
use crate::fetcher::state::{MaxQueueSize, QueuedFetch};
|
use crate::fetcher::state::{MaxQueueSize, QueuedFetch};
|
||||||
|
|
||||||
|
|
@ -20,7 +20,6 @@ impl Arbitrary for QueuedFetch {
|
||||||
|
|
||||||
QueuedFetch {
|
QueuedFetch {
|
||||||
rid: RepoId::arbitrary(g),
|
rid: RepoId::arbitrary(g),
|
||||||
from: NodeId::arbitrary(g),
|
|
||||||
refs_at,
|
refs_at,
|
||||||
timeout: Duration::from_secs(u64::arbitrary(g) % 3600),
|
timeout: Duration::from_secs(u64::arbitrary(g) % 3600),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,6 @@ pub fn create_queue(capacity: usize) -> Queue {
|
||||||
pub fn create_fetch() -> QueuedFetch {
|
pub fn create_fetch() -> QueuedFetch {
|
||||||
QueuedFetch {
|
QueuedFetch {
|
||||||
rid: arbitrary::gen(1),
|
rid: arbitrary::gen(1),
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -30,7 +30,6 @@ fn same_rid_merges_anywhere_in_queue(max_size: MaxQueueSize, merge_index: usize)
|
||||||
let target_index = merge_index % items.len();
|
let target_index = merge_index % items.len();
|
||||||
let same_rid_item = QueuedFetch {
|
let same_rid_item = QueuedFetch {
|
||||||
rid: items[target_index].rid,
|
rid: items[target_index].rid,
|
||||||
from: arbitrary::gen(1), // Different from
|
|
||||||
refs_at: vec![arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(60),
|
timeout: Duration::from_secs(60),
|
||||||
};
|
};
|
||||||
|
|
@ -51,14 +50,12 @@ fn combines_refs(base_refs_count: u8, merge_refs_count: u8) -> bool {
|
||||||
|
|
||||||
let base_item = QueuedFetch {
|
let base_item = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: base_refs.clone(),
|
refs_at: base_refs.clone(),
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
||||||
let merge_item = QueuedFetch {
|
let merge_item = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: merge_refs.clone(),
|
refs_at: merge_refs.clone(),
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
@ -89,7 +86,6 @@ fn empty_refs_fetches_all() -> bool {
|
||||||
// First enqueue with specific refs
|
// First enqueue with specific refs
|
||||||
let item_with_refs = QueuedFetch {
|
let item_with_refs = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![arbitrary::gen(1), arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1), arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
@ -97,7 +93,6 @@ fn empty_refs_fetches_all() -> bool {
|
||||||
// Second enqueue with empty refs (fetch everything)
|
// Second enqueue with empty refs (fetch everything)
|
||||||
let item_empty_refs = QueuedFetch {
|
let item_empty_refs = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
@ -119,14 +114,12 @@ fn longer_timeout_preserved(short_secs: u16, long_secs: u16) -> bool {
|
||||||
|
|
||||||
let item_short = QueuedFetch {
|
let item_short = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: short,
|
timeout: short,
|
||||||
};
|
};
|
||||||
|
|
||||||
let item_long = QueuedFetch {
|
let item_long = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: long,
|
timeout: long,
|
||||||
};
|
};
|
||||||
|
|
@ -151,14 +144,12 @@ fn does_not_increase_queue_length() -> bool {
|
||||||
|
|
||||||
let item1 = QueuedFetch {
|
let item1 = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
||||||
let item2 = QueuedFetch {
|
let item2 = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(60),
|
timeout: Duration::from_secs(60),
|
||||||
};
|
};
|
||||||
|
|
@ -194,21 +185,18 @@ fn succeed_when_at_capacity() -> bool {
|
||||||
|
|
||||||
let item1 = QueuedFetch {
|
let item1 = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
||||||
let item2 = QueuedFetch {
|
let item2 = QueuedFetch {
|
||||||
rid: arbitrary::gen(1), // Different rid
|
rid: arbitrary::gen(1), // Different rid
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
};
|
};
|
||||||
|
|
||||||
let merge_item = QueuedFetch {
|
let merge_item = QueuedFetch {
|
||||||
rid, // Same as item1
|
rid, // Same as item1
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(60),
|
timeout: Duration::from_secs(60),
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
use radicle::test::arbitrary;
|
use radicle::test::arbitrary;
|
||||||
use radicle_core::{NodeId, RepoId};
|
use radicle_core::RepoId;
|
||||||
|
|
||||||
use crate::fetcher::state::Enqueue;
|
use crate::fetcher::state::Enqueue;
|
||||||
use crate::fetcher::test::queue::helpers::*;
|
use crate::fetcher::test::queue::helpers::*;
|
||||||
|
|
@ -12,7 +12,6 @@ fn zero_timeout_accepted() {
|
||||||
let mut queue = create_queue(10);
|
let mut queue = create_queue(10);
|
||||||
let item = QueuedFetch {
|
let item = QueuedFetch {
|
||||||
rid: arbitrary::gen(1),
|
rid: arbitrary::gen(1),
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::ZERO,
|
timeout: Duration::ZERO,
|
||||||
};
|
};
|
||||||
|
|
@ -24,7 +23,6 @@ fn max_timeout_accepted() {
|
||||||
let mut queue = create_queue(10);
|
let mut queue = create_queue(10);
|
||||||
let item = QueuedFetch {
|
let item = QueuedFetch {
|
||||||
rid: arbitrary::gen(1),
|
rid: arbitrary::gen(1),
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::MAX,
|
timeout: Duration::MAX,
|
||||||
};
|
};
|
||||||
|
|
@ -34,18 +32,15 @@ fn max_timeout_accepted() {
|
||||||
#[test]
|
#[test]
|
||||||
fn empty_refs_at_items_can_be_equal() {
|
fn empty_refs_at_items_can_be_equal() {
|
||||||
let rid: RepoId = arbitrary::gen(1);
|
let rid: RepoId = arbitrary::gen(1);
|
||||||
let from: NodeId = arbitrary::gen(1);
|
|
||||||
let timeout = Duration::from_secs(30);
|
let timeout = Duration::from_secs(30);
|
||||||
|
|
||||||
let item1 = QueuedFetch {
|
let item1 = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from,
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout,
|
timeout,
|
||||||
};
|
};
|
||||||
let item2 = QueuedFetch {
|
let item2 = QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from,
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout,
|
timeout,
|
||||||
};
|
};
|
||||||
|
|
@ -64,19 +59,16 @@ fn merge_preserves_position_in_queue() {
|
||||||
// Enqueue three items
|
// Enqueue three items
|
||||||
let _ = queue.enqueue(QueuedFetch {
|
let _ = queue.enqueue(QueuedFetch {
|
||||||
rid: rid_first,
|
rid: rid_first,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
});
|
});
|
||||||
let _ = queue.enqueue(QueuedFetch {
|
let _ = queue.enqueue(QueuedFetch {
|
||||||
rid: rid_second,
|
rid: rid_second,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
});
|
});
|
||||||
let _ = queue.enqueue(QueuedFetch {
|
let _ = queue.enqueue(QueuedFetch {
|
||||||
rid: rid_third,
|
rid: rid_third,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![],
|
refs_at: vec![],
|
||||||
timeout: Duration::from_secs(30),
|
timeout: Duration::from_secs(30),
|
||||||
});
|
});
|
||||||
|
|
@ -84,7 +76,6 @@ fn merge_preserves_position_in_queue() {
|
||||||
// Merge into the second item
|
// Merge into the second item
|
||||||
let result = queue.enqueue(QueuedFetch {
|
let result = queue.enqueue(QueuedFetch {
|
||||||
rid: rid_second,
|
rid: rid_second,
|
||||||
from: arbitrary::gen(1),
|
|
||||||
refs_at: vec![arbitrary::gen(1)],
|
refs_at: vec![arbitrary::gen(1)],
|
||||||
timeout: Duration::from_secs(60),
|
timeout: Duration::from_secs(60),
|
||||||
});
|
});
|
||||||
|
|
|
||||||
|
|
@ -81,7 +81,6 @@ fn complete_then_dequeue_fifo() {
|
||||||
assert!(queued.is_some());
|
assert!(queued.is_some());
|
||||||
let queued = queued.unwrap();
|
let queued = queued.unwrap();
|
||||||
assert_eq!(queued.rid, repo_2);
|
assert_eq!(queued.rid, repo_2);
|
||||||
assert_eq!(queued.from, node_a);
|
|
||||||
assert_eq!(queued.refs_at, refs_at_2);
|
assert_eq!(queued.refs_at, refs_at_2);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -41,7 +41,6 @@ fn cannot_dequeue_while_node_at_capacity() {
|
||||||
let result = state.dequeue(&node_a);
|
let result = state.dequeue(&node_a);
|
||||||
let queued = result.unwrap();
|
let queued = result.unwrap();
|
||||||
assert_eq!(queued.rid, repo_2);
|
assert_eq!(queued.rid, repo_2);
|
||||||
assert_eq!(queued.from, node_a);
|
|
||||||
assert_eq!(queued.refs_at, refs_at_2);
|
assert_eq!(queued.refs_at, refs_at_2);
|
||||||
assert_eq!(queued.timeout, timeout_2);
|
assert_eq!(queued.timeout, timeout_2);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1184,7 +1184,6 @@ where
|
||||||
|
|
||||||
let Some(fetcher::QueuedFetch {
|
let Some(fetcher::QueuedFetch {
|
||||||
rid,
|
rid,
|
||||||
from,
|
|
||||||
refs_at,
|
refs_at,
|
||||||
timeout,
|
timeout,
|
||||||
}) = self.fetcher.dequeue(&nid)
|
}) = self.fetcher.dequeue(&nid)
|
||||||
|
|
@ -1203,14 +1202,14 @@ where
|
||||||
continue;
|
continue;
|
||||||
};
|
};
|
||||||
|
|
||||||
debug!(target: "service", "Dequeued fetch for {} from {}", rid, from);
|
debug!(target: "service", "Dequeued fetch for {} from {}", rid, nid);
|
||||||
|
|
||||||
if let Some(refs) = NonEmpty::from_vec(refs_at.clone()) {
|
if let Some(refs) = NonEmpty::from_vec(refs_at.clone()) {
|
||||||
self.fetch_refs_at(rid, from, refs, scope, timeout);
|
self.fetch_refs_at(rid, nid, refs, scope, timeout);
|
||||||
} else {
|
} else {
|
||||||
// Channel is `None` since they will already be
|
// Channel is `None` since they will already be
|
||||||
// registered with the fetcher service.
|
// registered with the fetcher service.
|
||||||
self.fetch(rid, from, refs_at, timeout, None);
|
self.fetch(rid, nid, refs_at, timeout, None);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue