Skip to content
45 changes: 42 additions & 3 deletions src/pointer/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -845,9 +845,21 @@ impl Inner {
// have committed while it ran.
let replacing = index.get(&address).is_some_and(|entry| entry.on_disk);
let outcome = match index.get(&address) {
// Nothing held, or a record this node lost: either way the
// write must happen, whatever state it carries.
None | Some(IndexEntry { on_disk: false, .. }) => PutOutcome::Changed,
// Nothing held: any state lands.
None => PutOutcome::Changed,
// A record this node lost: the lost state itself lands, to
// restore it, and so does anything that replaces it, but
// nothing older, so an arrival verified before the loss cannot
// roll the node back past it (the rule `admits` states).
Some(entry @ IndexEntry { on_disk: false, .. }) => {
if entry.state.state_id == record.state_id()
|| record.state().replaces(&entry.state)
{
PutOutcome::Changed
} else {
PutOutcome::Stale
}
}
Some(entry) if entry.state.state_id == record.state_id() => PutOutcome::Unchanged,
Some(entry) if record.state().replaces(&entry.state) => PutOutcome::Changed,
Some(_) => PutOutcome::Stale,
Expand Down Expand Up @@ -1754,6 +1766,33 @@ mod tests {
assert!(store.get(&record.address()).await.expect("get").is_some());
}

/// The commit itself holds the line `admits` states for a lost record, so
/// an older state that reaches it without that gate, or was admitted
/// before the loss, cannot roll the node back.
#[tokio::test]
async fn a_commit_after_a_loss_takes_nothing_older_than_what_was_lost() {
let (store, _dir) = store().await;
let held = signed(1, 3, 1);
store.put_bytes(&held.to_bytes()).await.expect("put");
std::fs::remove_file(store.file_for(&held.address())).expect("remove");
assert!(store.get(&held.address()).await.expect("get").is_none());

assert_eq!(
store
.put_bytes(&signed(1, 2, 2).to_bytes())
.await
.expect("put"),
PutOutcome::Stale,
"an older state must not replace the one that was lost"
);
assert!(!store.contains(&held.address()), "nothing was written");
assert_eq!(
store.put_bytes(&held.to_bytes()).await.expect("put"),
PutOutcome::Changed,
"the lost state itself is restored"
);
}

#[tokio::test]
async fn a_lost_record_is_restored_and_cannot_be_rolled_back() {
// The node remembers what it lost. Without that the address would look
Expand Down
44 changes: 25 additions & 19 deletions src/replication/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5549,23 +5549,27 @@ async fn handle_replication_message(
}
ReplicationMessageBody::PointerFetchRequest(request) => {
if let Some(pointers) = &ctx.pointers {
pointers.serve_fetch_detached(
*source,
request,
msg.request_id,
rr_message_id.map(ToOwned::to_owned),
);
pointers
.serve_fetch_detached(
*source,
request,
msg.request_id,
rr_message_id.map(ToOwned::to_owned),
)
.await;
}
Ok(())
}
ReplicationMessageBody::PointerStateRequest(request) => {
if let Some(pointers) = &ctx.pointers {
pointers.serve_state_detached(
*source,
request,
msg.request_id,
rr_message_id.map(ToOwned::to_owned),
);
pointers
.serve_state_detached(
*source,
request,
msg.request_id,
rr_message_id.map(ToOwned::to_owned),
)
.await;
}
Ok(())
}
Expand Down Expand Up @@ -6536,12 +6540,6 @@ async fn dispatch_neighbor_sync_request(
received_at: Instant,
rr_message_id: Option<&str>,
) -> Result<()> {
// A peer syncing with us gets our pointer hints as well — including a
// node that is bootstrapping, which is how it learns the pointers it
// should hold.
if let Some(pointers) = &ctx.pointers {
pointers.push_hints_detached(vec![source]);
}
let guard = match admit_bounded_responder(
&ctx.neighbor_sync_responder_admission_semaphore,
&ctx.neighbor_sync_responder_inflight,
Expand All @@ -6565,7 +6563,7 @@ async fn dispatch_neighbor_sync_request(
return Ok(());
}
};

let pointers = ctx.pointers.clone();
let worker_semaphore = Arc::clone(&ctx.neighbor_sync_responder_worker_semaphore);
let p2p_node = Arc::clone(&ctx.p2p_node);
let storage = Arc::clone(&ctx.storage);
Expand Down Expand Up @@ -6598,6 +6596,14 @@ async fn dispatch_neighbor_sync_request(
);
return;
}
// A peer whose sync request is admitted and fresh gets our pointer
// hints as well, including a node that is bootstrapping, which is how
// it learns the pointers it should hold. Only such a request: each
// answer scans every pointer held, and a refused or stale one must
// cost nothing.
if let Some(pointers) = &pointers {
pointers.answer_sync_with_hints(source);
}
if let Err(e) = handle_neighbor_sync_request(
&source,
&request,
Expand Down
Loading
Loading