Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 28 additions & 4 deletions src/replication/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2171,13 +2171,22 @@ impl ReplicationEngine {
/// The only guard left standing is the per-attempt recheck inside
/// `execute_single_fetch` — exactly the gate e2e tests use this seam to
/// exercise.
///
/// A replica hint for `key` that reached pending verification first is
/// dropped, so the key enters the fetch queue as it would have without
/// it. Callers host the chunk on a holder before calling this, which makes
/// it advertisable, and a change to the holder's closest peers starts a
/// neighbour-sync round at once. Left in place, the hint's entry made
/// `enqueue_fetch` refuse the key, and with one holder and no paid-list
/// entry it does not reach a quorum, so the key was not fetched that way
/// either. A key already in the fetch queue or in flight, or a full fetch
/// queue, still returns `false`.
#[cfg(any(test, feature = "test-utils"))]
pub async fn enqueue_fetch_for_test(&self, key: XorName, sources: Vec<PeerId>) -> bool {
let distance = crate::client::xor_distance(&key, self.p2p_node.peer_id().as_bytes());
self.queues
.write()
.await
.enqueue_fetch(key, distance, sources)
let mut queues = self.queues.write().await;
queues.remove_pending(&key);
queues.enqueue_fetch(key, distance, sources)
}

/// Test-only: whether `key` is still tracked in any fetch-pipeline stage
Expand All @@ -2187,6 +2196,21 @@ impl ReplicationEngine {
self.queues.read().await.contains_key(key)
}

/// Test-only: whether `key` is in the fetch queue or in flight, which is
/// where a candidate placed by [`Self::enqueue_fetch_for_test`] stays
/// until it resolves.
///
/// Such a candidate carries no verification retry metadata, so it is never
/// requeued for verification: once it leaves these two stages it is done.
/// Pending verification is left out on purpose. The chunk the test hosted
/// stays advertisable, so a replica hint can put the key there after the
/// candidate resolved, and counting that entry would hold a test's wait
/// loop until its deadline.
#[cfg(any(test, feature = "test-utils"))]
pub async fn fetch_queued_or_in_flight_for_test(&self, key: &XorName) -> bool {
self.queues.read().await.fetch_queued_or_in_flight(key)
}

/// Test-only: place `key` into pending verification as though `hinter` had
/// just advertised it with a replica hint. Returns whether it was admitted.
///
Expand Down
8 changes: 8 additions & 0 deletions src/replication/scheduling.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1414,6 +1414,14 @@ impl ReplicationQueues {
|| self.in_flight_fetch.contains_key(key)
}

/// Test-only: whether a key is in the fetch queue or in flight, leaving
/// pending verification out.
#[cfg(any(test, feature = "test-utils"))]
#[must_use]
pub fn fetch_queued_or_in_flight(&self, key: &XorName) -> bool {
self.fetch_payloads.contains_key(key) || self.in_flight_fetch.contains_key(key)
}

/// Check if all bootstrap-related work is done.
///
/// Returns `true` when none of the given bootstrap keys remain in any queue.
Expand Down
4 changes: 2 additions & 2 deletions tests/e2e/fetch_local_write_guard.rs
Original file line number Diff line number Diff line change
Expand Up @@ -198,7 +198,7 @@ async fn write_blocked_node_neither_probes_nor_dials() {
);

let deadline = tokio::time::Instant::now() + OBSERVATION_WINDOW;
while target_engine.fetch_pipeline_contains_for_test(&key).await {
while target_engine.fetch_queued_or_in_flight_for_test(&key).await {
assert!(
tokio::time::Instant::now() < deadline,
"the candidate never resolved; a key stuck in the fetch pipeline \
Expand Down Expand Up @@ -476,7 +476,7 @@ async fn already_held_key_is_not_fetched_again() {
);
let deadline = tokio::time::Instant::now() + OBSERVATION_WINDOW;
while target_engine
.fetch_pipeline_contains_for_test(&held_key)
.fetch_queued_or_in_flight_for_test(&held_key)
.await
{
assert!(
Expand Down
4 changes: 2 additions & 2 deletions tests/e2e/fetch_responsibility_recheck.rs
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ async fn stale_fetch_candidate_is_declined_at_download_time() {

let deadline = tokio::time::Instant::now() + OBSERVATION_WINDOW;
while target_engine
.fetch_pipeline_contains_for_test(&key_out)
.fetch_queued_or_in_flight_for_test(&key_out)
.await
{
assert!(
Expand Down Expand Up @@ -172,7 +172,7 @@ async fn stale_fetch_candidate_is_declined_at_download_time() {
}
let deadline = tokio::time::Instant::now() + OBSERVATION_WINDOW;
while target_engine
.fetch_pipeline_contains_for_test(&key_in)
.fetch_queued_or_in_flight_for_test(&key_in)
.await
{
assert!(
Expand Down
Loading