Skip to content
Open
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
38 changes: 36 additions & 2 deletions src/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,11 @@ pub(crate) fn next_nonce() -> u64 {
}
// more than 300 seconds behind
if nonce + 300000 < now_ms {
CUR_NONCE.fetch_max(now_ms + 1, Ordering::Relaxed);
return now_ms;
// Catch the counter up to the wall clock, then claim a slot from it. Returning
// `now_ms` directly hands the same nonce to every caller that races into this
// branch, and the exchange rejects the duplicates.
CUR_NONCE.fetch_max(now_ms, Ordering::Relaxed);
return CUR_NONCE.fetch_add(1, Ordering::Relaxed);
}
nonce
}
Expand Down Expand Up @@ -95,6 +98,37 @@ lazy_static! {
mod tests {
use super::*;

#[test]
fn next_nonce_is_unique_after_catch_up() {
for _ in 0..50 {
// Park the counter far enough behind the wall clock that every caller takes the
// catch-up branch, which is what an idle client hits before its next burst.
let now = now_timestamp_ms();
CUR_NONCE.store(now - 400_000, Ordering::Relaxed);

let barrier = std::sync::Arc::new(std::sync::Barrier::new(8));
let handles: Vec<_> = (0..8)
.map(|_| {
let barrier = barrier.clone();
std::thread::spawn(move || {
barrier.wait();
next_nonce()
})
})
.collect();
let nonces: Vec<u64> = handles.into_iter().map(|h| h.join().unwrap()).collect();

let mut sorted = nonces.clone();
sorted.sort_unstable();
sorted.dedup();
assert_eq!(
sorted.len(),
nonces.len(),
"duplicate nonces handed out: {nonces:?}"
);
}
}

#[test]
fn float_to_string_for_hashing_test() {
assert_eq!(float_to_string_for_hashing(0.), "0".to_string());
Expand Down