diff --git a/Makefile b/Makefile index 022481738..c3fd05f3c 100644 --- a/Makefile +++ b/Makefile @@ -48,6 +48,7 @@ SCYLLA_TEST_FILTER := $(subst ${SPACE},${EMPTY},ClusterTests.*\ :ServerSideFailureTests.*\ :ServerSideFailureThreeNodeTests.*\ :TimestampTests.*\ +:UuidTests.*\ :HostFilterTest.*\ :ExecutionProfileTest.*\ :DCExecutionProfileTest.*\ @@ -105,6 +106,7 @@ CASSANDRA_TEST_FILTER := $(subst ${SPACE},${EMPTY},ClusterTests.*\ :ServerSideFailureTests.*\ :ServerSideFailureThreeNodeTests.*\ :TimestampTests.*\ +:UuidTests.*\ :HostFilterTest.*\ :ExecutionProfileTest.*\ :DCExecutionProfileTest.*\ diff --git a/docs/source/topics/using/data-types/uuids.md b/docs/source/topics/using/data-types/uuids.md index f6f4747ea..98489ad83 100644 --- a/docs/source/topics/using/data-types/uuids.md +++ b/docs/source/topics/using/data-types/uuids.md @@ -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(); @@ -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 diff --git a/scylla-rust-wrapper/src/cql_types/uuid.rs b/scylla-rust-wrapper/src/cql_types/uuid.rs index 4b4dcb06e..e8341917b 100644 --- a/scylla-rust-wrapper/src/cql_types/uuid.rs +++ b/scylla-rust-wrapper/src/cql_types/uuid.rs @@ -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 { + 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), + ); + } + + // 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 +} + +// 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; } } } @@ -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, + uuid_gen: CassBorrowedSharedPtr, 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, }; @@ -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, + uuid_gen: CassBorrowedSharedPtr, 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; }; @@ -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) { 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); + } +} diff --git a/tests/src/integration/tests/test_uuids.cpp b/tests/src/integration/tests/test_uuids.cpp new file mode 100644 index 000000000..ae82d0492 --- /dev/null +++ b/tests/src/integration/tests/test_uuids.cpp @@ -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. +*/ + +#include "integration.hpp" + +#include +#include +#include +#include +#include + +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 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 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 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 lock(start_mutex); + start_condition.wait(lock, [&ready_threads]() { return ready_threads == THREAD_COUNT; }); + start = true; + } + start_condition.notify_all(); + + for (std::vector::iterator it = threads.begin(); it != threads.end(); ++it) { + it->join(); + } + + std::vector 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())); +}