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
2 changes: 2 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ SCYLLA_TEST_FILTER := $(subst ${SPACE},${EMPTY},ClusterTests.*\
:ServerSideFailureTests.*\
:ServerSideFailureThreeNodeTests.*\
:TimestampTests.*\
:UuidTests.*\
:HostFilterTest.*\
:ExecutionProfileTest.*\
:DCExecutionProfileTest.*\
Expand Down Expand Up @@ -105,6 +106,7 @@ CASSANDRA_TEST_FILTER := $(subst ${SPACE},${EMPTY},ClusterTests.*\
:ServerSideFailureTests.*\
:ServerSideFailureThreeNodeTests.*\
:TimestampTests.*\
:UuidTests.*\
:HostFilterTest.*\
:ExecutionProfileTest.*\
:DCExecutionProfileTest.*\
Expand Down
12 changes: 12 additions & 0 deletions docs/source/topics/using/data-types/uuids.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,17 @@ Version 4 can be used with ScyllaDB/Cassandra's `uuid` type for unique identific
A UUID generator object is used to create new UUIDs. The [`CassUuidGen`] object
is thread-safe. It should only be created once per application and reused.

Within one wall-clock millisecond, one generator can allocate 10,000 distinct
version 1 UUID timestamps. When the generated timestamp is in the current
millisecond and this capacity is exhausted, calls to `cass_uuid_gen_time()`
busy-wait until the system clock advances, consuming CPU and adding latency.

If the system clock moves backward, or a thread uses a stale clock sample after
another thread advances the timestamp, the generator preserves monotonicity by
incrementing the last timestamp. Generated timestamps can therefore move ahead
of wall time. At sustained rates above 10 million UUIDs per second, this drift
can grow indefinitely. See [issue #507] for possible improvements.

```c
CassUuidGen* uuid_gen = cass_uuid_gen_new();

Expand Down Expand Up @@ -65,3 +76,4 @@ cass_uuid_string(uuid, uuid_str);
```
[`cass_uuid_timestamp()`]: https://cpp-rs-driver.docs.scylladb.com/stable/api/struct.CassUuid#1a3980467a0bb6642054ecf37d49aebf1a
[`CassUuidGen`]: https://cpp-rs-driver.docs.scylladb.com/stable/api/struct.CassUuidGen
[issue #507]: https://github.com/scylladb/cpp-rs-driver/issues/507
149 changes: 134 additions & 15 deletions scylla-rust-wrapper/src/cql_types/uuid.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,21 +49,49 @@ fn rand_clock_seq_and_node(node: u64) -> u64 {
result
}

// Ported from UuidGen::monotonic_timestamp, but simplified at
// a cost of performance.
fn monotonic_timestamp(last_timestamp: &mut AtomicU64) -> u64 {
fn try_monotonic_timestamp(last_timestamp: &AtomicU64, now: u64) -> Option<u64> {
let last = last_timestamp.load(Ordering::SeqCst);

// The wall clock advanced, so use its current millisecond as the new baseline.
if now > last {
return last_timestamp
.compare_exchange(last, now, Ordering::SeqCst, Ordering::SeqCst)
.map(|_| now)
.ok();
}

let last_ms = to_milliseconds(last);
// Preserve monotonicity after clock rollback or when another thread advanced the timestamp
// after this thread sampled the clock.
if to_milliseconds(now) < last_ms {
return Some(
last_timestamp
.fetch_add(1, Ordering::SeqCst)
.wrapping_add(1),
);
}
Comment thread
dkropachev marked this conversation as resolved.

// Allocate the next 100-nanosecond tick within the current wall-clock millisecond.
let candidate = last.wrapping_add(1);
if to_milliseconds(candidate) == last_ms {
return last_timestamp
.compare_exchange(last, candidate, Ordering::SeqCst, Ordering::SeqCst)
.map(|_| candidate)
.ok();
}

None
Comment thread
dkropachev marked this conversation as resolved.
}

// Ported from UuidGen::monotonic_timestamp.
fn monotonic_timestamp(last_timestamp: &AtomicU64) -> u64 {
loop {
let now = SystemTime::now();
let now = now.duration_since(UNIX_EPOCH).unwrap();
let now = from_unix_timestamp(now.as_millis() as u64);

let last = last_timestamp.load(Ordering::SeqCst);
if last < now
&& last_timestamp
.compare_exchange(last, now, Ordering::SeqCst, Ordering::SeqCst)
.is_ok()
{
return now;
if let Some(timestamp) = try_monotonic_timestamp(last_timestamp, now) {
return timestamp;
}
}
}
Expand Down Expand Up @@ -135,16 +163,16 @@ pub unsafe extern "C" fn cass_uuid_gen_new_with_node(

#[unsafe(no_mangle)]
pub unsafe extern "C" fn cass_uuid_gen_time(
uuid_gen: CassBorrowedExclusivePtr<CassUuidGen, CMut>,
uuid_gen: CassBorrowedSharedPtr<CassUuidGen, CMut>,
output: *mut CassUuid,
) {
let Some(uuid_gen) = BoxFFI::as_mut_ref(uuid_gen) else {
let Some(uuid_gen) = BoxFFI::as_ref(uuid_gen) else {
tracing::error!("Provided null uuid generator pointer to cass_uuid_gen_time!");
return;
};

let uuid = CassUuid {
time_and_version: set_version(monotonic_timestamp(&mut uuid_gen.last_timestamp), 1),
time_and_version: set_version(monotonic_timestamp(&uuid_gen.last_timestamp), 1),
clock_seq_and_node: uuid_gen.clock_seq_and_node,
};

Expand All @@ -167,11 +195,11 @@ pub unsafe extern "C" fn cass_uuid_gen_random(_uuid_gen: *mut CassUuidGen, outpu

#[unsafe(no_mangle)]
pub unsafe extern "C" fn cass_uuid_gen_from_time(
uuid_gen: CassBorrowedExclusivePtr<CassUuidGen, CMut>,
uuid_gen: CassBorrowedSharedPtr<CassUuidGen, CMut>,
timestamp: cass_uint64_t,
output: *mut CassUuid,
) {
let Some(uuid_gen) = BoxFFI::as_mut_ref(uuid_gen) else {
let Some(uuid_gen) = BoxFFI::as_ref(uuid_gen) else {
tracing::error!("Provided null uuid generator pointer to cass_uuid_gen_from_time!");
return;
};
Expand Down Expand Up @@ -273,3 +301,94 @@ pub unsafe extern "C" fn cass_uuid_from_string_n(
pub unsafe extern "C" fn cass_uuid_gen_free(uuid_gen: CassOwnedExclusivePtr<CassUuidGen, CMut>) {
BoxFFI::free(uuid_gen);
}

#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Barrier};
use std::thread;

#[test]
fn monotonic_timestamp_uses_submillisecond_ticks() {
let now = from_unix_timestamp(1_700_000_000_000);
let last_timestamp = AtomicU64::new(now);

for expected_offset in 1..10_000 {
assert_eq!(
try_monotonic_timestamp(&last_timestamp, now),
Some(now + expected_offset)
);
}

assert_eq!(try_monotonic_timestamp(&last_timestamp, now), None);
}

#[test]
fn monotonic_timestamp_remains_monotonic_during_clock_rollback() {
let now = from_unix_timestamp(1_700_000_000_000);
let future = from_unix_timestamp(1_700_000_000_001);
let last_timestamp = AtomicU64::new(future);

assert_eq!(
try_monotonic_timestamp(&last_timestamp, now),
Some(future + 1)
);
assert_eq!(last_timestamp.load(Ordering::SeqCst), future + 1);
}

#[test]
fn monotonic_timestamp_handles_stale_sample_across_millisecond_boundary() {
let stale_now = from_unix_timestamp(1_700_000_000_000);
let next_millisecond = from_unix_timestamp(1_700_000_000_001);
let last_timestamp = AtomicU64::new(stale_now);

assert_eq!(
try_monotonic_timestamp(&last_timestamp, next_millisecond),
Some(next_millisecond)
);
assert_eq!(
try_monotonic_timestamp(&last_timestamp, stale_now),
Some(next_millisecond + 1)
);
}

#[test]
fn monotonic_timestamp_is_unique_under_concurrency() {
const THREADS: usize = 8;
const UUIDS_PER_THREAD: usize = 1_000;

let now = from_unix_timestamp(1_700_000_000_000);
let last_timestamp = Arc::new(AtomicU64::new(0));
let start = Arc::new(Barrier::new(THREADS));
let mut handles = Vec::with_capacity(THREADS);

for _ in 0..THREADS {
let last_timestamp = Arc::clone(&last_timestamp);
let start = Arc::clone(&start);
handles.push(thread::spawn(move || {
let mut timestamps = Vec::with_capacity(UUIDS_PER_THREAD);
start.wait();
for _ in 0..UUIDS_PER_THREAD {
loop {
if let Some(timestamp) = try_monotonic_timestamp(&last_timestamp, now) {
timestamps.push(timestamp);
break;
}
}
}
timestamps
}));
}

let mut timestamps: Vec<_> = handles
.into_iter()
.flat_map(|handle| handle.join().unwrap())
.collect();
timestamps.sort_unstable();
timestamps.dedup();

assert_eq!(timestamps.len(), THREADS * UUIDS_PER_THREAD);
assert_eq!(timestamps[0], now);
assert_eq!(timestamps[timestamps.len() - 1], now + 7_999);
}
}
100 changes: 100 additions & 0 deletions tests/src/integration/tests/test_uuids.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
/*
Copyright (c) 2026 ScyllaDB Ltd.

Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at

http://www.apache.org/licenses/LICENSE-2.0

Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
Comment thread
dkropachev marked this conversation as resolved.

#include "integration.hpp"

#include <algorithm>
#include <condition_variable>
#include <mutex>
#include <thread>
#include <vector>

namespace {

const size_t THREAD_COUNT = 8;
const size_t UUIDS_PER_THREAD = 1000;
const cass_uint64_t UUID_TIMESTAMP_MASK = 0x0FFFFFFFFFFFFFFFULL;

} // namespace

class UuidTests : public Integration {
public:
UuidTests() { is_ccm_requested_ = false; }
};

/**
* Generates time UUIDs concurrently through the public C API.
*
* @test_category uuid
* @expected_result Every generated UUID has a unique, monotonically allocated timestamp.
*/
CASSANDRA_INTEGRATION_TEST_F(UuidTests, ConcurrentTimeGeneration) {
UuidGen generator;
CassUuidGen* const uuid_gen = generator.get();
std::vector<CassUuid> uuids(THREAD_COUNT * UUIDS_PER_THREAD);

std::mutex start_mutex;
std::condition_variable start_condition;
size_t ready_threads = 0;
bool start = false;

std::vector<std::thread> threads;
threads.reserve(THREAD_COUNT);
for (size_t thread_index = 0; thread_index < THREAD_COUNT; ++thread_index) {
threads.push_back(std::thread([&, thread_index]() {
{
std::unique_lock<std::mutex> lock(start_mutex);
++ready_threads;
start_condition.notify_all();
start_condition.wait(lock, [&start]() { return start; });
}

const size_t offset = thread_index * UUIDS_PER_THREAD;
for (size_t i = 0; i < UUIDS_PER_THREAD; ++i) {
cass_uuid_gen_time(uuid_gen, &uuids[offset + i]);
}
}));
}

{
std::unique_lock<std::mutex> lock(start_mutex);
start_condition.wait(lock, [&ready_threads]() { return ready_threads == THREAD_COUNT; });
start = true;
}
start_condition.notify_all();

for (std::vector<std::thread>::iterator it = threads.begin(); it != threads.end(); ++it) {
it->join();
}

std::vector<cass_uint64_t> timestamps;
timestamps.reserve(uuids.size());
for (size_t i = 0; i < uuids.size(); ++i) {
EXPECT_EQ(1u, cass_uuid_version(uuids[i]));
timestamps.push_back(uuids[i].time_and_version & UUID_TIMESTAMP_MASK);
}

for (size_t thread_index = 0; thread_index < THREAD_COUNT; ++thread_index) {
const size_t offset = thread_index * UUIDS_PER_THREAD;
for (size_t i = 1; i < UUIDS_PER_THREAD; ++i) {
EXPECT_LT(uuids[offset + i - 1].time_and_version & UUID_TIMESTAMP_MASK,
uuids[offset + i].time_and_version & UUID_TIMESTAMP_MASK);
}
}

std::sort(timestamps.begin(), timestamps.end());
EXPECT_EQ(timestamps.end(), std::unique(timestamps.begin(), timestamps.end()));
}
Loading