use ::core::borrow::Borrow;
use ::core::fmt::{Debug, Display, Formatter};
use ::core::iter::FusedIterator;
use ::core::ops::RangeBounds;
use alloc::vec::Vec;
use super::set::BTreeSet;
use crate::core::node::NodeLike;
use crate::{cdc::change::ChangeEvent, core::pair::Pair};
#[cfg(feature = "cdc")]
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Topology<T> {
pub node_capacity: usize,
pub nodes: Vec<Vec<T>>,
}
#[cfg(feature = "cdc")]
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum TopologyError {
ZeroNodeCapacity,
EmptyNode { index: usize },
OversizedNode { index: usize, len: usize, capacity: usize },
UnsortedNode { index: usize },
OverlappingNodes { left: usize, right: usize },
}
#[cfg(feature = "cdc")]
impl Display for TopologyError {
fn fmt(&self, formatter: &mut Formatter<'_>) -> ::core::fmt::Result {
match self {
Self::ZeroNodeCapacity => formatter.write_str("topology node capacity must be non-zero"),
Self::EmptyNode { index } => write!(formatter, "topology node {index} is empty"),
Self::OversizedNode { index, len, capacity } => write!(
formatter,
"topology node {index} has {len} entries but capacity is {capacity}"
),
Self::UnsortedNode { index } => write!(formatter, "topology node {index} is not strictly ordered"),
Self::OverlappingNodes { left, right } => {
write!(
formatter,
"topology nodes {left} and {right} overlap or are out of order"
)
}
}
}
}
#[cfg(feature = "cdc")]
impl ::core::error::Error for TopologyError {}
#[derive(Debug)]
pub struct BTreeMap<K, V, Node = Vec<Pair<K, V>>>
where
K: Send + Ord + Clone + 'static,
V: Send + Clone + 'static,
Node: NodeLike<Pair<K, V>>,
{
pub(crate) set: BTreeSet<Pair<K, V>, Node>,
}
impl<K, V, Node> Default for BTreeMap<K, V, Node>
where
K: Send + Ord + Clone,
V: Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
fn default() -> Self {
Self {
set: Default::default(),
}
}
}
pub struct Iter<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
inner: super::set::Iter<'a, Pair<K, V>, Node>,
}
impl<'a, K, V, Node> Iterator for Iter<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
type Item = (K, V);
fn next(&mut self) -> Option<Self::Item> {
if let Some(entry) = self.inner.next() {
return Some((entry.key, entry.value));
}
None
}
}
impl<'a, K, V, Node> DoubleEndedIterator for Iter<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
fn next_back(&mut self) -> Option<Self::Item> {
if let Some(entry) = self.inner.next_back() {
return Some((entry.key, entry.value));
}
None
}
}
impl<'a, K, V, Node> FusedIterator for Iter<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
}
pub struct Range<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
inner: super::set::Range<'a, Pair<K, V>, Node>,
}
impl<'a, K, V, Node> Iterator for Range<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
type Item = (K, V);
fn next(&mut self) -> Option<Self::Item> {
if let Some(entry) = self.inner.next() {
return Some((entry.key, entry.value));
}
None
}
}
impl<'a, K, V, Node> DoubleEndedIterator for Range<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
fn next_back(&mut self) -> Option<Self::Item> {
if let Some(entry) = self.inner.next_back() {
return Some((entry.key, entry.value));
}
None
}
}
impl<'a, K, V, Node> FusedIterator for Range<'a, K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
}
impl<K, V, Node> BTreeMap<K, V, Node>
where
K: Debug + Send + Ord + Clone + 'static,
V: Debug + Send + Clone + 'static,
Node: NodeLike<Pair<K, V>> + Send + 'static,
{
pub fn new() -> Self {
Self {
set: Default::default(),
}
}
pub fn with_maximum_node_size(node_capacity: usize) -> Self {
Self {
set: BTreeSet::with_maximum_node_size(node_capacity),
}
}
#[cfg(feature = "cdc")]
pub fn attach_node(&self, node: Node) {
self.set.attach_node(node)
}
#[cfg(feature = "cdc")]
pub fn attach_nodes(&self, nodes: impl IntoIterator<Item = Node>) {
self.set.attach_nodes(nodes)
}
#[cfg(feature = "cdc")]
pub fn snapshot_nodes(&self) -> Vec<Node>
where
Node: Clone,
{
self.set
.index
.read()
.values()
.map(|node| (*node.read()).clone())
.collect()
}
#[cfg(feature = "cdc")]
pub fn export_topology(&self) -> Topology<Pair<K, V>> {
let (node_capacity, nodes) = self.set.export_topology();
Topology { node_capacity, nodes }
}
#[cfg(feature = "cdc")]
pub fn from_topology(topology: Topology<Pair<K, V>>) -> Result<Self, TopologyError> {
if topology.node_capacity == 0 {
return Err(TopologyError::ZeroNodeCapacity);
}
for (index, node) in topology.nodes.iter().enumerate() {
if node.is_empty() {
return Err(TopologyError::EmptyNode { index });
}
if node.len() > topology.node_capacity {
return Err(TopologyError::OversizedNode {
index,
len: node.len(),
capacity: topology.node_capacity,
});
}
if node.windows(2).any(|pair| pair[0] >= pair[1]) {
return Err(TopologyError::UnsortedNode { index });
}
if index > 0
&& topology.nodes[index - 1]
.last()
.expect("non-empty node validated above")
>= node.first().expect("non-empty node validated above")
{
return Err(TopologyError::OverlappingNodes {
left: index - 1,
right: index,
});
}
}
let map = Self::with_maximum_node_size(topology.node_capacity);
map.attach_nodes(topology.nodes.into_iter().map(|values| {
let mut node = Node::with_capacity(topology.node_capacity);
for value in values {
let (inserted, _) = NodeLike::insert(&mut node, value);
debug_assert!(inserted, "validated topology contains unique values");
}
node
}));
Ok(map)
}
pub fn contains_key<Q>(&self, key: &Q) -> bool
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
{
self.set.contains(key)
}
pub fn get<Q>(&self, key: &Q) -> Option<super::r#ref::Ref<Pair<K, V>, Node>>
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
{
self.set.get(key)
}
#[inline(always)]
pub fn lookup_for_select<Q>(&self, key: &Q) -> Option<V>
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
V: Clone,
{
self.set.get_with(key, |pair| pair.value.clone())
}
#[inline(always)]
pub fn lookup_for_select_optimistic<Q>(&self, key: &Q) -> Option<V>
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
V: Clone,
{
self.set.get_with_optimistic(key, |pair| pair.value.clone())
}
pub fn insert(&self, key: K, value: V) -> Option<V> {
let new_entry = Pair { key, value };
self.set.put(new_entry).map(|pair| pair.value)
}
pub fn checked_insert(&self, key: K, value: V) -> Option<()> {
let new_entry = Pair { key, value };
self.set.put_checked(new_entry).ok().map(|_| ())
}
pub fn insert_cdc(&self, key: K, value: V) -> (Option<V>, Vec<ChangeEvent<Pair<K, V>>>) {
let new_entry = Pair { key, value };
let (old_value, cdc) = self.set.put_cdc(new_entry);
(old_value.map(|pair| pair.value), cdc)
}
pub fn checked_insert_cdc(&self, key: K, value: V) -> Option<Vec<ChangeEvent<Pair<K, V>>>> {
let new_entry = Pair { key, value };
self.set.put_cdc_checked(new_entry).ok().map(|(_, evs)| evs)
}
pub fn remove<Q>(&self, key: &Q) -> Option<(K, V)>
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
{
self.set.remove(key).map(|pair| (pair.key, pair.value))
}
#[allow(clippy::type_complexity)]
pub fn remove_cdc<Q>(&self, key: &Q) -> (Option<(K, V)>, Vec<ChangeEvent<Pair<K, V>>>)
where
Pair<K, V>: Borrow<Q> + Ord,
Q: Ord + ?Sized,
{
let (old_value, cdc) = self.set.remove_cdc(key);
(old_value.map(|pair| (pair.key, pair.value)), cdc)
}
pub fn len(&self) -> usize {
self.set.len()
}
pub fn is_empty(&self) -> bool {
self.set.is_empty()
}
pub fn capacity(&self) -> usize {
self.set.capacity()
}
pub fn node_count(&self) -> usize {
self.set.node_count()
}
pub fn iter(&self) -> Iter<'_, K, V, Node> {
Iter { inner: self.set.iter() }
}
pub fn range<Q, R>(&self, range: R) -> Range<'_, K, V, Node>
where
Pair<K, V>: Borrow<Q>,
Q: Ord + ?Sized,
R: RangeBounds<Q>,
{
Range {
inner: BTreeSet::range(&self.set, range),
}
}
}
#[cfg(test)]
mod tests {
use super::BTreeMap;
use super::ChangeEvent;
use super::Pair;
#[cfg(feature = "cdc")]
use super::{Topology, TopologyError};
use crate::core::constants::DEFAULT_INNER_SIZE;
use crate::BTreeSet;
use rand::Rng;
use scc::HashMap;
use std::fmt::Debug;
use std::sync::{Arc, Mutex};
use std::thread;
#[cfg(feature = "cdc")]
#[test]
fn pointer_free_topology_round_trip_preserves_nodes() {
let map = BTreeMap::<u64, u64>::with_maximum_node_size(4);
for key in 0..37 {
map.insert(key, key * 10);
}
let topology = map.export_topology();
assert!(topology.nodes.len() > 1);
let restored = BTreeMap::<u64, u64>::from_topology(topology.clone()).unwrap();
assert_eq!(restored.export_topology(), topology);
assert_eq!(
restored.iter().collect::<Vec<_>>(),
(0..37).map(|k| (k, k * 10)).collect::<Vec<_>>()
);
}
#[cfg(feature = "cdc")]
#[test]
fn pointer_free_topology_rejects_invalid_boundaries() {
let error = BTreeMap::<u64, u64>::from_topology(Topology {
node_capacity: 4,
nodes: vec![
vec![Pair { key: 1, value: 10 }, Pair { key: 3, value: 30 }],
vec![Pair { key: 2, value: 20 }],
],
})
.unwrap_err();
assert_eq!(error, TopologyError::OverlappingNodes { left: 0, right: 1 });
}
#[test]
fn test_range_edge_cast() {
let maximum_node_size = 3;
let map = BTreeMap::<usize, &str>::with_maximum_node_size(maximum_node_size);
map.insert(1usize, "a");
map.insert(2usize, "b");
map.insert(3usize, "c");
map.insert(4usize, "d");
map.insert(5usize, "e");
map.insert(6usize, "f");
map.insert(7usize, "g");
let mid_range = map.range::<usize, _>(3..5).collect::<BTreeSet<_>>();
assert_eq!(
mid_range,
vec![(3usize, "c"), (4usize, "d"),].into_iter().collect::<BTreeSet<_>>()
);
}
#[test]
fn test_split_insert_replaces_left_split_max() {
let maximum_node_size = 4;
let map = BTreeMap::<usize, &str>::with_maximum_node_size(maximum_node_size);
for key in 0..maximum_node_size {
assert_eq!(map.insert(key, "old"), None);
}
let split_left_max = maximum_node_size / 2 - 1;
assert_eq!(map.insert(split_left_max, "new"), Some("old"));
assert_eq!(map.len(), maximum_node_size);
assert_eq!(map.get(&split_left_max).map(|entry| entry.get().value), Some("new"));
assert_eq!(map.iter().filter(|(key, _)| *key == split_left_max).count(), 1);
}
#[derive(Debug, Default)]
struct PersistedBTreeMap<K, V>
where
K: Debug + Ord + Clone,
V: Debug + Clone + PartialEq,
{
nodes: std::collections::BTreeMap<K, Vec<Pair<K, V>>>,
}
impl<K: Debug + Ord + Clone, V: Debug + Clone + PartialEq> PersistedBTreeMap<K, V> {
fn persist(&mut self, event: &ChangeEvent<Pair<K, V>>) {
match event {
ChangeEvent::CreateNode { max_value, event_id: _ } => {
let node = vec![max_value.clone()];
self.nodes.insert(max_value.key.clone(), node);
}
ChangeEvent::RemoveNode { max_value, event_id: _ } => {
self.nodes.remove(&max_value.key);
}
ChangeEvent::InsertAt {
max_value,
index,
value,
event_id: _,
} => {
if let Some(node) = self.nodes.get_mut(&max_value.key) {
node.insert(*index, value.clone());
}
if max_value.key < value.key {
let node = self.nodes.remove(&max_value.key).unwrap();
self.nodes.insert(value.key.clone(), node);
}
}
ChangeEvent::RemoveAt {
max_value,
index,
value,
event_id: _,
} => {
let mut max_removed = false;
if let Some(node) = self.nodes.get_mut(&max_value.key) {
node.remove(*index);
max_removed = max_value.key == value.key;
}
if max_removed {
let node = self.nodes.remove(&max_value.key).unwrap();
if let Some(new_max) = node.last() {
self.nodes.insert(new_max.key.clone(), node);
}
}
}
ChangeEvent::SplitNode {
max_value,
split_index,
event_id: _,
} => {
if let Some(mut old_node) = self.nodes.remove(&max_value.key) {
let new_node = old_node.split_off(*split_index);
let new_max_value = new_node.last().unwrap();
self.nodes.insert(new_max_value.key.clone(), new_node);
let old_max_value = old_node.last().unwrap();
self.nodes.insert(old_max_value.key.clone(), old_node);
}
}
}
}
fn contains_pair(&self, key: &K, value: &V) -> bool {
for node in self.nodes.values() {
if let Ok(pos) = node.binary_search(&Pair {
key: key.clone(),
value: value.clone(),
}) {
if node[pos].value == *value {
return true;
}
}
}
false
}
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_single_insert() {
let map = BTreeMap::<usize, &str>::new();
let mut mock_state = PersistedBTreeMap::default();
let (_, events) = map.insert_cdc(1, "a");
for event in events {
mock_state.persist(&event);
}
assert!(mock_state.contains_pair(&1, &"a"));
assert!(map.contains_key(&1));
assert_eq!(map.get(&1).unwrap().get().value, "a");
let expected_state = map
.set
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_multiple_inserts() {
let map = BTreeMap::<usize, String>::new();
let mut mock_state = PersistedBTreeMap::default();
for i in 0..1024 {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
for event in events {
mock_state.persist(&event);
}
}
for i in 0..1024 {
assert!(mock_state.contains_pair(&i, &format!("val{}", i)));
assert!(map.contains_key(&i));
assert_eq!(map.get(&i).unwrap().get().value, format!("val{}", i));
}
let expected_state = map
.set
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_updates() {
let map = BTreeMap::<usize, &str>::new();
let mut mock_state = PersistedBTreeMap::default();
let (_, events) = map.insert_cdc(1, "a");
for event in events {
mock_state.persist(&event);
}
let (_, events) = map.insert_cdc(1, "b");
for event in events {
mock_state.persist(&event);
}
assert!(mock_state.contains_pair(&1, &"b"));
assert!(!mock_state.contains_pair(&1, &"a"));
assert!(map.contains_key(&1));
assert_eq!(map.get(&1).unwrap().get().value, "b");
let expected_state = map
.set
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_node_splits() {
let map = BTreeMap::<usize, String>::new();
let mut mock_state = PersistedBTreeMap::default();
let n = crate::core::constants::DEFAULT_INNER_SIZE + 10;
for i in 0..n {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
for event in events {
mock_state.persist(&event);
}
}
for i in 0..n {
assert!(mock_state.contains_pair(&i, &format!("val{}", i)));
assert!(map.contains_key(&i));
assert_eq!(map.get(&i).unwrap().get().value, format!("val{}", i));
}
assert!(mock_state.nodes.len() > 1);
let expected_state = map
.set
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
#[cfg(feature = "cdc")]
#[test]
fn test_concurrent_insert_cdc() {
let map = Arc::new(BTreeMap::<usize, String>::new());
let num_threads = 8;
let operations_per_thread = 1000;
let mut handles = vec![];
let test_data: Vec<Vec<(i32, (usize, String))>> = (0..num_threads)
.map(|_| {
let mut rng = rand::rng();
(0..operations_per_thread)
.map(|_| {
let value = rng.random_range(0..100000);
let operation = rng.random_range(0..2);
(operation, (value, format!("val{value}")))
})
.collect()
})
.collect();
let expected_values = Arc::new(Mutex::new(HashMap::new()));
for thread_idx in 0..num_threads {
let map_clone = Arc::clone(&map);
let expected_values = Arc::clone(&expected_values);
let thread_data = test_data[thread_idx].clone();
let handle = thread::spawn(move || {
let mut events = Vec::new();
for (operation, (k, v)) in thread_data {
if operation == 0 {
let (_, evs) = map_clone.insert_cdc(k, v.clone());
events.extend(evs);
let _ = expected_values.lock().unwrap().insert(k, v);
}
}
events
});
handles.push(handle);
}
let mut final_events = Vec::new();
for handle in handles {
let thread_events = handle.join().unwrap();
final_events.extend(thread_events)
}
final_events.sort_by(|ev1, ev2| ev1.id().cmp(&ev2.id()));
let mut mock_state = PersistedBTreeMap::default();
for ev in final_events {
mock_state.persist(&ev);
}
let expected_state = map
.set
.index
.read()
.iter()
.map(|(key, node)| (key.clone().key, node.read_arc().clone()))
.collect::<_>();
assert_eq!(mock_state.nodes, expected_state);
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_event_ids_sequential_no_gaps() {
let map = BTreeMap::<usize, String>::new();
let mut all_events = Vec::new();
for i in 0..100 {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
all_events.extend(events);
}
all_events.sort_by_key(|e| e.id());
assert!(!all_events.is_empty(), "Should have at least one event");
for i in 1..all_events.len() {
let prev_id = all_events[i - 1].id().inner();
let curr_id = all_events[i].id().inner();
assert_eq!(
curr_id,
prev_id + 1,
"Event IDs should be consecutive: {} followed by {}, but got gap",
prev_id,
curr_id
);
}
}
#[cfg(feature = "cdc")]
#[test]
fn normal_writes_do_not_consume_cdc_event_ids() {
let map = BTreeMap::<usize, String>::new();
map.insert(1, "first".to_owned());
map.insert(1, "second".to_owned());
map.remove(&1);
let (_, events) = map.insert_cdc(2, "recorded".to_owned());
assert!(!events.is_empty());
assert_eq!(events[0].id().inner(), 0);
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_remove_monotonicity() {
let map = BTreeMap::<usize, String>::new();
let mut all_events = Vec::new();
for i in 0..50 {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
all_events.extend(events);
}
for i in 0..25 {
let (_, events) = map.remove_cdc(&i);
all_events.extend(events);
}
all_events.sort_by_key(|e| e.id());
assert!(!all_events.is_empty(), "Should have at least one event");
for i in 1..all_events.len() {
let prev_id = all_events[i - 1].id().inner();
let curr_id = all_events[i].id().inner();
assert_eq!(
curr_id,
prev_id + 1,
"Event IDs should be consecutive across inserts and removes"
);
}
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_split_no_gaps() {
let map = BTreeMap::<usize, String>::new();
let mut all_events = Vec::new();
let n = DEFAULT_INNER_SIZE + 200;
for i in 0..n {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
all_events.extend(events);
}
all_events.sort_by_key(|e| e.id());
assert!(!all_events.is_empty(), "Should have at least one event");
for i in 1..all_events.len() {
let prev_id = all_events[i - 1].id().inner();
let curr_id = all_events[i].id().inner();
assert_eq!(
curr_id,
prev_id + 1,
"Event IDs should be consecutive even during splits"
);
}
let split_events: Vec<_> = all_events
.iter()
.filter(|e| matches!(e, ChangeEvent::SplitNode { .. }))
.collect();
assert!(!split_events.is_empty(), "Should have at least one split event");
}
#[cfg(feature = "cdc")]
#[test]
fn test_concurrent_cdc_no_gaps() {
let map = Arc::new(BTreeMap::<usize, String>::new());
let num_threads = 16;
let operations_per_thread = 500;
let mut handles = vec![];
for thread_idx in 0..num_threads {
let map_clone = Arc::clone(&map);
let handle = thread::spawn(move || {
let mut events = Vec::new();
let base = thread_idx * 10000;
for i in 0..operations_per_thread {
let value = base + i;
let (_, evs) = map_clone.insert_cdc(value, format!("val{}", value));
events.extend(evs);
}
events
});
handles.push(handle);
}
let mut final_events = Vec::new();
for handle in handles {
let thread_events = handle.join().unwrap();
final_events.extend(thread_events);
}
final_events.sort_by_key(|e| e.id());
assert!(!final_events.is_empty(), "Should have at least one event");
for i in 1..final_events.len() {
let prev_id = final_events[i - 1].id().inner();
let curr_id = final_events[i].id().inner();
assert_eq!(
curr_id,
prev_id + 1,
"Concurrent event IDs should be consecutive with no gaps: {} -> {}",
prev_id,
curr_id
);
}
}
#[cfg(feature = "cdc")]
#[test]
fn test_cdc_mixed_operations() {
let map = BTreeMap::<usize, String>::new();
let mut all_events = Vec::new();
for i in 0..100 {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
all_events.extend(events);
}
for i in 0..50 {
let (_, events) = map.remove_cdc(&i);
all_events.extend(events);
}
for i in 100..125 {
let (_, events) = map.insert_cdc(i, format!("val{}", i));
all_events.extend(events);
}
all_events.sort_by_key(|e| e.id());
assert!(!all_events.is_empty(), "Should have at least one event");
for i in 1..all_events.len() {
let prev_id = all_events[i - 1].id().inner();
let curr_id = all_events[i].id().inner();
assert_eq!(curr_id, prev_id + 1, "Mixed operation event IDs should be consecutive");
}
}
}