use std::{
fmt,
sync::atomic::{AtomicI64, Ordering},
time::Instant,
};
#[derive(Debug)]
pub struct SequenceNumberGenerator {
base_timestamp: Instant,
starting_offset: AtomicI64,
last_sequence_number: AtomicI64,
}
impl SequenceNumberGenerator {
pub fn new(starting_offset: i64) -> Self {
Self {
base_timestamp: Instant::now(),
starting_offset: AtomicI64::new(starting_offset),
last_sequence_number: AtomicI64::new(starting_offset.saturating_sub(1)),
}
}
#[inline]
pub fn get_sequence_number(&self) -> i64 {
let elapsed_nanos = self.base_timestamp.elapsed().as_nanos();
let elapsed = i64::try_from(elapsed_nanos).unwrap_or(i64::MAX);
let offset = self.starting_offset.load(Ordering::Acquire);
let candidate = elapsed.saturating_add(offset);
let mut prev = self.last_sequence_number.load(Ordering::Acquire);
loop {
let next = if candidate > prev {
candidate
} else {
prev.saturating_add(1)
};
match self.last_sequence_number.compare_exchange_weak(
prev,
next,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => return next,
Err(actual) => prev = actual,
}
}
}
pub fn set_starting_offset(&self, starting_offset: i64) {
self
.starting_offset
.store(starting_offset, Ordering::Release);
self
.last_sequence_number
.fetch_max(starting_offset.saturating_sub(1), Ordering::Release);
}
}
impl fmt::Display for SequenceNumberGenerator {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"{},{},{}",
self.starting_offset.load(Ordering::Relaxed),
self.last_sequence_number.load(Ordering::Relaxed),
self.base_timestamp.elapsed().as_nanos()
)
}
}
#[cfg(test)]
mod tests {
use std::{sync::Arc, thread};
use gxhash::HashSet;
use super::SequenceNumberGenerator;
#[test]
fn monotonic_and_offset_lift() {
let generator = SequenceNumberGenerator::new(0);
let first = generator.get_sequence_number();
let second = generator.get_sequence_number();
let third = generator.get_sequence_number();
assert!(
second > first,
"单调时钟结合原子递增保证连续取号严格单调递增"
);
assert!(third > second);
generator.set_starting_offset(1_000_000);
assert!(generator.get_sequence_number() >= 1_000_000);
}
#[test]
fn concurrent_strict_monotonicity() {
let generator = Arc::new(SequenceNumberGenerator::new(100));
let thread_count = 8;
let iterations_per_thread = 5000;
let handles: Vec<_> = (0..thread_count)
.map(|_| {
let g = Arc::clone(&generator);
thread::spawn(move || {
let mut numbers = Vec::with_capacity(iterations_per_thread);
for _ in 0..iterations_per_thread {
numbers.push(g.get_sequence_number());
}
numbers
})
})
.collect();
let mut all_numbers = Vec::with_capacity(thread_count * iterations_per_thread);
for h in handles {
let nums = h.join().unwrap();
for window in nums.windows(2) {
assert!(
window[1] > window[0],
"单线程视角严格递增: {} <= {}",
window[1],
window[0]
);
}
all_numbers.extend(nums);
}
let total_count = all_numbers.len();
let unique_set: HashSet<i64> = all_numbers.into_iter().collect();
assert_eq!(
unique_set.len(),
total_count,
"并发全序严格唯一,无任何重复序列号"
);
}
#[test]
fn display_and_boundary_offsets() {
let generator = SequenceNumberGenerator::new(-100);
let s = format!("{generator}");
assert!(s.starts_with("-100,"));
let num = generator.get_sequence_number();
assert!(num >= -100);
generator.set_starting_offset(i64::MAX - 1000);
assert!(generator.get_sequence_number() >= i64::MAX - 1000);
}
}