1use std::collections::HashMap;
16
17use common_telemetry::debug;
18use futures_util::stream::BoxStream;
19use snafu::ensure;
20
21use crate::error::{self, Result};
22use crate::key::TABLE_NAME_KEY_PREFIX;
23use crate::key::txn_helper::TxnOpGetResponseSet;
24use crate::kv_backend::KvBackendRef;
25use crate::kv_backend::txn::{Compare, CompareOp, Txn, TxnOp};
26use crate::range_stream::{DEFAULT_PAGE_SIZE, PaginationStream};
27use crate::rpc::KeyValue;
28use crate::rpc::store::{BatchDeleteRequest, BatchGetRequest, RangeRequest};
29
30pub struct TombstoneManager {
54 kv_backend: KvBackendRef,
55 tombstone_prefix: String,
56 #[cfg(test)]
58 max_txn_ops: Option<usize>,
59}
60
61const TOMBSTONE_PREFIX: &str = "__tombstone/";
62const MOVE_VALUE_TXN_OPS_PER_KEY: usize = 4;
63const RESTORE_VALUE_TXN_OPS_PER_KEY: usize = 6;
64
65pub(crate) fn to_tombstone_key(key: &[u8]) -> Vec<u8> {
66 [TOMBSTONE_PREFIX.as_bytes(), key].concat()
67}
68
69impl TombstoneManager {
70 pub fn new(kv_backend: KvBackendRef) -> Self {
72 Self::new_with_prefix(kv_backend, TOMBSTONE_PREFIX)
73 }
74
75 pub fn new_with_prefix(kv_backend: KvBackendRef, prefix: &str) -> Self {
77 Self {
78 kv_backend,
79 tombstone_prefix: prefix.to_string(),
80 #[cfg(test)]
81 max_txn_ops: None,
82 }
83 }
84
85 pub fn to_tombstone(&self, key: &[u8]) -> Vec<u8> {
86 [self.tombstone_prefix.as_bytes(), key].concat()
87 }
88
89 pub fn strip_tombstone_prefix<'a>(&self, tombstone_key: &'a [u8]) -> Result<&'a [u8]> {
91 ensure!(
92 tombstone_key.starts_with(self.tombstone_prefix.as_bytes()),
93 error::UnexpectedSnafu {
94 err_msg: format!(
95 "The key '{}' does not start with tombstone prefix '{}'.",
96 String::from_utf8_lossy(tombstone_key),
97 self.tombstone_prefix
98 ),
99 }
100 );
101
102 Ok(&tombstone_key[self.tombstone_prefix.len()..])
103 }
104
105 pub async fn get(&self, key: &[u8]) -> Result<Option<KeyValue>> {
107 let tombstone_key = self.to_tombstone(key);
108 let response = self
109 .kv_backend
110 .range(RangeRequest::new().with_key(tombstone_key))
111 .await?;
112 Ok(response.kvs.into_iter().next())
113 }
114
115 pub async fn batch_get(&self, keys: &[Vec<u8>]) -> Result<HashMap<Vec<u8>, KeyValue>> {
117 let tombstone_keys = keys
118 .iter()
119 .map(|key| self.to_tombstone(key))
120 .collect::<Vec<_>>();
121 let resp = self
122 .kv_backend
123 .batch_get(BatchGetRequest::new().with_keys(tombstone_keys))
124 .await?;
125
126 resp.kvs
127 .into_iter()
128 .map(|kv| Ok((self.strip_tombstone_prefix(&kv.key)?.to_vec(), kv)))
129 .collect::<Result<HashMap<_, _>>>()
130 }
131
132 pub fn tombstones(&self) -> BoxStream<'static, Result<KeyValue>> {
134 self.scan_prefix(self.tombstone_prefix.as_bytes().to_vec())
135 }
136
137 pub fn tombstoned_table_names(&self) -> BoxStream<'static, Result<KeyValue>> {
139 self.scan_prefix(
140 format!("{}{}/", self.tombstone_prefix, TABLE_NAME_KEY_PREFIX).into_bytes(),
141 )
142 }
143
144 pub fn tombstoned_table_names_by_catalog(
146 &self,
147 catalog: &str,
148 ) -> BoxStream<'static, Result<KeyValue>> {
149 self.scan_prefix(
150 format!(
151 "{}{}/{catalog}/",
152 self.tombstone_prefix, TABLE_NAME_KEY_PREFIX
153 )
154 .into_bytes(),
155 )
156 }
157
158 fn scan_prefix(&self, prefix: Vec<u8>) -> BoxStream<'static, Result<KeyValue>> {
160 let req = RangeRequest::new().with_prefix(prefix);
161 let stream = PaginationStream::new(self.kv_backend.clone(), req, DEFAULT_PAGE_SIZE, Ok)
162 .into_stream();
163
164 Box::pin(stream)
165 }
166
167 #[cfg(test)]
168 pub fn set_max_txn_ops(&mut self, max_txn_ops: usize) {
169 self.max_txn_ops = Some(max_txn_ops);
170 }
171
172 fn build_move_value_txn(
178 &self,
179 src_key: Vec<u8>,
180 value: Vec<u8>,
181 dest_key: Vec<u8>,
182 require_dest_not_exists: bool,
183 ) -> Txn {
184 let mut compares = vec![Compare::with_value(
185 src_key.clone(),
186 CompareOp::Equal,
187 value.clone(),
188 )];
189 if require_dest_not_exists {
190 compares.push(Compare::with_value_not_exists(
191 dest_key.clone(),
192 CompareOp::Equal,
193 ));
194 }
195
196 let mut failure = vec![TxnOp::Get(src_key.clone())];
197 if require_dest_not_exists {
198 failure.push(TxnOp::Get(dest_key.clone()));
199 }
200
201 Txn::new()
202 .when(compares)
203 .and_then(vec![
204 TxnOp::Put(dest_key.clone(), value.clone()),
205 TxnOp::Delete(src_key.clone()),
206 ])
207 .or_else(failure)
208 }
209
210 async fn move_values_inner(
211 &self,
212 keys: &[Vec<u8>],
213 dest_keys: &[Vec<u8>],
214 require_dest_not_exists: bool,
215 extra_ops: Vec<TxnOp>,
216 ) -> Result<usize> {
217 ensure!(
218 keys.len() == dest_keys.len(),
219 error::UnexpectedSnafu {
220 err_msg: format!(
221 "The length of keys({}) does not match the length of dest_keys({}).",
222 keys.len(),
223 dest_keys.len()
224 ),
225 }
226 );
227 let lookup_table = keys.iter().zip(dest_keys.iter()).collect::<HashMap<_, _>>();
229
230 let resp = self
231 .kv_backend
232 .batch_get(BatchGetRequest::new().with_keys(keys.to_vec()))
233 .await?;
234 let mut results = resp
235 .kvs
236 .into_iter()
237 .map(|kv| (kv.key, kv.value))
238 .collect::<HashMap<_, _>>();
239 if results.is_empty() && extra_ops.is_empty() {
240 return Ok(0);
241 }
242
243 const MAX_RETRIES: usize = 8;
244 for _ in 0..MAX_RETRIES {
245 let (txns, keys): (Vec<_>, Vec<_>) = results
246 .iter()
247 .map(|(key, value)| {
248 let txn = self.build_move_value_txn(
249 key.clone(),
250 value.clone(),
251 lookup_table[&key].clone(),
252 require_dest_not_exists,
253 );
254 (txn, key.clone())
255 })
256 .unzip();
257 let mut txns = txns;
258 if !extra_ops.is_empty() {
259 txns.push(Txn::new().and_then(extra_ops.clone()));
260 }
261 let mut resp = self.kv_backend.txn(Txn::merge_all(txns)).await?;
262 if resp.succeeded {
263 return Ok(keys.len());
264 }
265 let mut set = TxnOpGetResponseSet::from(&mut resp.responses);
266 for key in &keys {
268 let dest_key = lookup_table[&key].clone();
269 if require_dest_not_exists {
270 let mut filter = TxnOpGetResponseSet::filter(dest_key.clone());
271 if filter(&mut set).is_some() {
272 return error::TombstoneTargetAlreadyExistsSnafu {
273 key: String::from_utf8_lossy(&dest_key).to_string(),
274 }
275 .fail();
276 }
277 }
278
279 let mut filter = TxnOpGetResponseSet::filter(key.clone());
280 if let Some(value) = filter(&mut set) {
281 results.insert(key.clone(), value);
282 } else {
283 results.remove(key);
284 }
285 }
286 }
287
288 error::MoveValuesSnafu {
289 err_msg: format!(
290 "keys: {:?}",
291 keys.iter().map(|key| String::from_utf8_lossy(key)),
292 ),
293 }
294 .fail()
295 }
296
297 fn max_txn_ops(&self) -> usize {
298 #[cfg(test)]
299 if let Some(max_txn_ops) = self.max_txn_ops {
300 return max_txn_ops;
301 }
302 self.kv_backend.max_txn_ops()
303 }
304
305 async fn execute_extra_ops(&self, extra_ops: Vec<TxnOp>) -> Result<()> {
306 let max_txn_ops = self.max_txn_ops();
307 ensure!(
308 max_txn_ops > 0,
309 error::UnexpectedSnafu {
310 err_msg: "max_txn_ops must be greater than 0".to_string(),
311 }
312 );
313 for chunk in extra_ops.chunks(max_txn_ops) {
314 self.move_values_inner(&[], &[], false, chunk.to_vec())
315 .await?;
316 }
317 Ok(())
318 }
319
320 async fn move_values_with_extra(
324 &self,
325 keys: Vec<Vec<u8>>,
326 dest_keys: Vec<Vec<u8>>,
327 require_dest_not_exists: bool,
328 mut extra_ops: Vec<TxnOp>,
329 ) -> Result<usize> {
330 ensure!(
331 keys.len() == dest_keys.len(),
332 error::UnexpectedSnafu {
333 err_msg: format!(
334 "The length of keys({}) does not match the length of dest_keys({}).",
335 keys.len(),
336 dest_keys.len()
337 ),
338 }
339 );
340 if keys.is_empty() {
341 if !extra_ops.is_empty() {
342 self.execute_extra_ops(extra_ops).await?;
343 }
344 return Ok(0);
345 }
346 let txn_ops_per_key = if require_dest_not_exists {
347 RESTORE_VALUE_TXN_OPS_PER_KEY
348 } else {
349 MOVE_VALUE_TXN_OPS_PER_KEY
350 };
351 let max_txn_ops = self.max_txn_ops();
352 ensure!(
353 max_txn_ops >= txn_ops_per_key,
354 error::UnexpectedSnafu {
355 err_msg: format!(
356 "max_txn_ops {max_txn_ops} is smaller than required txn ops per key {txn_ops_per_key}"
357 )
358 }
359 );
360 let merge_extra_ops =
361 !extra_ops.is_empty() && max_txn_ops.saturating_sub(txn_ops_per_key) >= extra_ops.len();
362 let reserved_ops = if merge_extra_ops { extra_ops.len() } else { 0 };
363 let chunk_size = (max_txn_ops - reserved_ops) / txn_ops_per_key;
364 let total_ops = keys
365 .len()
366 .saturating_mul(txn_ops_per_key)
367 .saturating_add(extra_ops.len());
368 if keys.len() > chunk_size || total_ops > max_txn_ops {
369 debug!(
370 "Moving values with multiple chunks, keys len: {}, chunk_size: {}",
371 keys.len(),
372 chunk_size
373 );
374 let mut moved_keys = 0;
375 let keys_chunks = keys.chunks(chunk_size).collect::<Vec<_>>();
376 let dest_keys_chunks = dest_keys.chunks(chunk_size).collect::<Vec<_>>();
377 let num_chunks = keys_chunks.len();
378 for (index, (keys, dest_keys)) in
379 keys_chunks.into_iter().zip(dest_keys_chunks).enumerate()
380 {
381 let chunk_extra_ops = if merge_extra_ops && index + 1 == num_chunks {
382 std::mem::take(&mut extra_ops)
383 } else {
384 vec![]
385 };
386 moved_keys += self
387 .move_values_inner(keys, dest_keys, require_dest_not_exists, chunk_extra_ops)
388 .await?;
389 }
390 if !extra_ops.is_empty() {
391 self.execute_extra_ops(extra_ops).await?;
392 }
393 Ok(moved_keys)
394 } else {
395 self.move_values_inner(&keys, &dest_keys, require_dest_not_exists, extra_ops)
396 .await
397 }
398 }
399
400 async fn move_values(
401 &self,
402 keys: Vec<Vec<u8>>,
403 dest_keys: Vec<Vec<u8>>,
404 require_dest_not_exists: bool,
405 ) -> Result<usize> {
406 self.move_values_with_extra(keys, dest_keys, require_dest_not_exists, vec![])
407 .await
408 }
409
410 pub async fn create(&self, keys: Vec<Vec<u8>>) -> Result<usize> {
418 let (keys, dest_keys): (Vec<_>, Vec<_>) = keys
419 .into_iter()
420 .map(|key| {
421 let tombstone_key = self.to_tombstone(&key);
422 (key, tombstone_key)
423 })
424 .unzip();
425
426 self.move_values(keys, dest_keys, false).await
427 }
428
429 pub async fn create_with_marker(
431 &self,
432 keys: Vec<Vec<u8>>,
433 marker_key: Vec<u8>,
434 marker_value: Vec<u8>,
435 ) -> Result<usize> {
436 self.create_with_markers(keys, vec![(marker_key, marker_value)])
437 .await
438 }
439
440 pub async fn create_with_markers(
442 &self,
443 keys: Vec<Vec<u8>>,
444 markers: Vec<(Vec<u8>, Vec<u8>)>,
445 ) -> Result<usize> {
446 let (keys, dest_keys): (Vec<_>, Vec<_>) = keys
447 .into_iter()
448 .map(|key| {
449 let tombstone_key = self.to_tombstone(&key);
450 (key, tombstone_key)
451 })
452 .unzip();
453 let marker_ops = markers
454 .into_iter()
455 .map(|(key, value)| TxnOp::Put(self.to_tombstone(&key), value))
456 .collect();
457
458 self.move_values_with_extra(keys, dest_keys, false, marker_ops)
459 .await
460 }
461
462 pub async fn restore(&self, keys: Vec<Vec<u8>>) -> Result<usize> {
470 let (keys, dest_keys): (Vec<_>, Vec<_>) = keys
471 .into_iter()
472 .map(|key| {
473 let tombstone_key = self.to_tombstone(&key);
474 (tombstone_key, key)
475 })
476 .unzip();
477
478 self.move_values(keys, dest_keys, true).await
479 }
480
481 pub async fn restore_with_marker(
483 &self,
484 keys: Vec<Vec<u8>>,
485 marker_key: Vec<u8>,
486 ) -> Result<usize> {
487 let (keys, dest_keys): (Vec<_>, Vec<_>) = keys
488 .into_iter()
489 .map(|key| {
490 let tombstone_key = self.to_tombstone(&key);
491 (tombstone_key, key)
492 })
493 .unzip();
494 let marker_key = self.to_tombstone(&marker_key);
495
496 self.move_values_with_extra(keys, dest_keys, true, vec![TxnOp::Delete(marker_key)])
497 .await
498 }
499
500 pub async fn restore_with_markers(
502 &self,
503 keys: Vec<Vec<u8>>,
504 marker_keys: Vec<Vec<u8>>,
505 ) -> Result<usize> {
506 let (keys, dest_keys): (Vec<_>, Vec<_>) = keys
507 .into_iter()
508 .map(|key| {
509 let tombstone_key = self.to_tombstone(&key);
510 (tombstone_key, key)
511 })
512 .unzip();
513 let marker_ops = marker_keys
514 .into_iter()
515 .map(|key| TxnOp::Delete(self.to_tombstone(&key)))
516 .collect();
517 self.move_values_with_extra(keys, dest_keys, true, marker_ops)
518 .await
519 }
520
521 pub async fn delete(&self, keys: Vec<Vec<u8>>) -> Result<usize> {
525 let keys = keys
526 .iter()
527 .map(|key| self.to_tombstone(key))
528 .collect::<Vec<_>>();
529
530 let num_keys = keys.len();
531 let _ = self
532 .kv_backend
533 .batch_delete(BatchDeleteRequest::new().with_keys(keys))
534 .await?;
535
536 Ok(num_keys)
537 }
538
539 pub async fn delete_with_marker(
541 &self,
542 mut keys: Vec<Vec<u8>>,
543 marker_key: Vec<u8>,
544 ) -> Result<usize> {
545 keys.push(marker_key);
546 self.delete(keys).await
547 }
548
549 pub async fn delete_with_markers(
551 &self,
552 mut keys: Vec<Vec<u8>>,
553 marker_keys: Vec<Vec<u8>>,
554 ) -> Result<usize> {
555 keys.extend(marker_keys);
556 self.delete(keys).await
557 }
558}
559
560#[cfg(test)]
561mod tests {
562
563 use std::any::Any;
564 use std::collections::HashMap;
565 use std::sync::atomic::{AtomicUsize, Ordering};
566 use std::sync::{Arc, Mutex};
567
568 use crate::error::{Error, Result};
569 use crate::key::tombstone::{
570 MOVE_VALUE_TXN_OPS_PER_KEY, RESTORE_VALUE_TXN_OPS_PER_KEY, TombstoneManager,
571 };
572 use crate::kv_backend::memory::MemoryKvBackend;
573 use crate::kv_backend::txn::{Txn, TxnRequest, TxnResponse};
574 use crate::kv_backend::{KvBackend, TxnService};
575 use crate::rpc::KeyValue;
576 use crate::rpc::store::{
577 BatchDeleteRequest, BatchDeleteResponse, BatchGetRequest, BatchGetResponse,
578 BatchPutRequest, BatchPutResponse, DeleteRangeRequest, DeleteRangeResponse, PutRequest,
579 PutResponse, RangeRequest, RangeResponse,
580 };
581
582 struct TxnOpLimitKvBackend {
583 inner: Arc<MemoryKvBackend<Error>>,
584 max_txn_ops: usize,
585 txn_op_counts: Mutex<Vec<usize>>,
586 }
587
588 #[async_trait::async_trait]
589 impl TxnService for TxnOpLimitKvBackend {
590 type Error = Error;
591
592 async fn txn(&self, txn: Txn) -> Result<TxnResponse> {
593 let TxnRequest {
594 compare,
595 success,
596 failure,
597 } = txn.req();
598 let txn_ops = compare.len() + success.len() + failure.len();
599 assert!(
600 txn_ops <= self.max_txn_ops,
601 "txn ops {txn_ops} exceeds limit {}",
602 self.max_txn_ops
603 );
604 self.txn_op_counts.lock().unwrap().push(txn_ops);
605 self.inner.txn(txn).await
606 }
607
608 fn max_txn_ops(&self) -> usize {
609 self.max_txn_ops
610 }
611 }
612
613 #[async_trait::async_trait]
614 impl KvBackend for TxnOpLimitKvBackend {
615 fn name(&self) -> &str {
616 "txn_op_limit"
617 }
618
619 fn as_any(&self) -> &dyn Any {
620 self
621 }
622
623 async fn range(&self, req: RangeRequest) -> Result<RangeResponse> {
624 self.inner.range(req).await
625 }
626
627 async fn put(&self, req: PutRequest) -> Result<PutResponse> {
628 self.inner.put(req).await
629 }
630
631 async fn batch_put(&self, req: BatchPutRequest) -> Result<BatchPutResponse> {
632 self.inner.batch_put(req).await
633 }
634
635 async fn batch_get(&self, req: BatchGetRequest) -> Result<BatchGetResponse> {
636 self.inner.batch_get(req).await
637 }
638
639 async fn delete_range(&self, req: DeleteRangeRequest) -> Result<DeleteRangeResponse> {
640 self.inner.delete_range(req).await
641 }
642
643 async fn batch_delete(&self, req: BatchDeleteRequest) -> Result<BatchDeleteResponse> {
644 self.inner.batch_delete(req).await
645 }
646 }
647
648 struct StalePointGetKvBackend {
649 inner: Arc<MemoryKvBackend<Error>>,
650 get_calls: AtomicUsize,
651 range_calls: AtomicUsize,
652 }
653
654 #[async_trait::async_trait]
655 impl TxnService for StalePointGetKvBackend {
656 type Error = Error;
657
658 async fn txn(&self, txn: Txn) -> Result<TxnResponse> {
659 self.inner.txn(txn).await
660 }
661 }
662
663 #[async_trait::async_trait]
664 impl KvBackend for StalePointGetKvBackend {
665 fn name(&self) -> &str {
666 "stale_point_get"
667 }
668
669 fn as_any(&self) -> &dyn Any {
670 self
671 }
672
673 async fn range(&self, req: RangeRequest) -> Result<RangeResponse> {
674 assert_eq!(b"__tombstone/foo", req.key.as_slice());
675 assert!(req.range_end.is_empty());
676 self.range_calls.fetch_add(1, Ordering::SeqCst);
677 self.inner.range(req).await
678 }
679
680 async fn get(&self, key: &[u8]) -> Result<Option<KeyValue>> {
681 self.get_calls.fetch_add(1, Ordering::SeqCst);
682 Ok(Some((key.to_vec(), b"stale".to_vec()).into()))
683 }
684
685 async fn put(&self, req: PutRequest) -> Result<PutResponse> {
686 self.inner.put(req).await
687 }
688
689 async fn batch_put(&self, req: BatchPutRequest) -> Result<BatchPutResponse> {
690 self.inner.batch_put(req).await
691 }
692
693 async fn batch_get(&self, req: BatchGetRequest) -> Result<BatchGetResponse> {
694 self.inner.batch_get(req).await
695 }
696
697 async fn delete_range(&self, req: DeleteRangeRequest) -> Result<DeleteRangeResponse> {
698 self.inner.delete_range(req).await
699 }
700
701 async fn batch_delete(&self, req: BatchDeleteRequest) -> Result<BatchDeleteResponse> {
702 self.inner.batch_delete(req).await
703 }
704 }
705
706 #[derive(Debug, Clone)]
707 struct MoveValue {
708 key: Vec<u8>,
709 dest_key: Vec<u8>,
710 value: Vec<u8>,
711 }
712
713 async fn check_moved_values(
714 kv_backend: Arc<MemoryKvBackend<Error>>,
715 move_values: &[MoveValue],
716 ) {
717 for MoveValue {
718 key,
719 dest_key,
720 value,
721 } in move_values
722 {
723 assert!(kv_backend.get(key).await.unwrap().is_none());
724 assert_eq!(
725 &kv_backend.get(dest_key).await.unwrap().unwrap().value,
726 value,
727 );
728 }
729 }
730
731 #[tokio::test]
732 async fn test_create_tombstone() {
733 let kv_backend = Arc::new(MemoryKvBackend::default());
734 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
735 kv_backend
736 .put(PutRequest::new().with_key("bar").with_value("baz"))
737 .await
738 .unwrap();
739 kv_backend
740 .put(PutRequest::new().with_key("foo").with_value("hi"))
741 .await
742 .unwrap();
743 tombstone_manager
744 .create(vec![b"bar".to_vec(), b"foo".to_vec()])
745 .await
746 .unwrap();
747 assert!(!kv_backend.exists(b"bar").await.unwrap());
748 assert!(!kv_backend.exists(b"foo").await.unwrap());
749 assert_eq!(
750 kv_backend
751 .get(&tombstone_manager.to_tombstone(b"bar"))
752 .await
753 .unwrap()
754 .unwrap()
755 .value,
756 b"baz"
757 );
758 assert_eq!(
759 kv_backend
760 .get(&tombstone_manager.to_tombstone(b"foo"))
761 .await
762 .unwrap()
763 .unwrap()
764 .value,
765 b"hi"
766 );
767 assert_eq!(kv_backend.len(), 2);
768 }
769
770 #[tokio::test]
771 async fn test_get_tombstone_bypasses_stale_point_get() {
772 let inner = Arc::new(MemoryKvBackend::default());
773 let backend = Arc::new(StalePointGetKvBackend {
774 inner: inner.clone(),
775 get_calls: AtomicUsize::new(0),
776 range_calls: AtomicUsize::new(0),
777 });
778 let manager = TombstoneManager::new(backend.clone());
779 inner
780 .put(
781 PutRequest::new()
782 .with_key(manager.to_tombstone(b"foo"))
783 .with_value(b"authoritative"),
784 )
785 .await
786 .unwrap();
787
788 let value = manager.get(b"foo").await.unwrap().unwrap();
789
790 assert_eq!(b"authoritative", value.value.as_slice());
791 assert_eq!(0, backend.get_calls.load(Ordering::SeqCst));
792 assert_eq!(1, backend.range_calls.load(Ordering::SeqCst));
793 }
794
795 #[tokio::test]
796 async fn test_create_tombstone_with_non_exist_values() {
797 let kv_backend = Arc::new(MemoryKvBackend::default());
798 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
799
800 kv_backend
801 .put(PutRequest::new().with_key("bar").with_value("baz"))
802 .await
803 .unwrap();
804 kv_backend
805 .put(PutRequest::new().with_key("foo").with_value("hi"))
806 .await
807 .unwrap();
808
809 tombstone_manager
810 .create(vec![b"bar".to_vec(), b"baz".to_vec()])
811 .await
812 .unwrap();
813 check_moved_values(
814 kv_backend.clone(),
815 &[MoveValue {
816 key: b"bar".to_vec(),
817 dest_key: tombstone_manager.to_tombstone(b"bar"),
818 value: b"baz".to_vec(),
819 }],
820 )
821 .await;
822 }
823
824 #[tokio::test]
825 async fn test_restore_tombstone() {
826 let kv_backend = Arc::new(MemoryKvBackend::default());
827 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
828 kv_backend
829 .put(PutRequest::new().with_key("bar").with_value("baz"))
830 .await
831 .unwrap();
832 kv_backend
833 .put(PutRequest::new().with_key("foo").with_value("hi"))
834 .await
835 .unwrap();
836 let expected_kvs = kv_backend.dump();
837 tombstone_manager
838 .create(vec![b"bar".to_vec(), b"foo".to_vec()])
839 .await
840 .unwrap();
841 tombstone_manager
842 .restore(vec![b"bar".to_vec(), b"foo".to_vec()])
843 .await
844 .unwrap();
845 assert_eq!(expected_kvs, kv_backend.dump());
846 }
847
848 #[tokio::test]
849 async fn test_delete_tombstone() {
850 let kv_backend = Arc::new(MemoryKvBackend::default());
851 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
852 kv_backend
853 .put(PutRequest::new().with_key("bar").with_value("baz"))
854 .await
855 .unwrap();
856 kv_backend
857 .put(PutRequest::new().with_key("foo").with_value("hi"))
858 .await
859 .unwrap();
860 tombstone_manager
861 .create(vec![b"bar".to_vec(), b"foo".to_vec()])
862 .await
863 .unwrap();
864 tombstone_manager
865 .delete(vec![b"bar".to_vec(), b"foo".to_vec()])
866 .await
867 .unwrap();
868 assert!(kv_backend.is_empty());
869 }
870
871 #[tokio::test]
872 async fn test_batch_get_tombstones() {
873 let kv_backend = Arc::new(MemoryKvBackend::default());
874 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
875 kv_backend
876 .put(PutRequest::new().with_key("bar").with_value("baz"))
877 .await
878 .unwrap();
879 kv_backend
880 .put(PutRequest::new().with_key("foo").with_value("hi"))
881 .await
882 .unwrap();
883
884 tombstone_manager
885 .create(vec![b"bar".to_vec(), b"foo".to_vec()])
886 .await
887 .unwrap();
888
889 let kvs = tombstone_manager
890 .batch_get(&[b"bar".to_vec(), b"foo".to_vec(), b"missing".to_vec()])
891 .await
892 .unwrap();
893
894 assert_eq!(kvs.len(), 2);
895 assert_eq!(kvs.get(b"bar".as_slice()).unwrap().value, b"baz");
896 assert_eq!(kvs.get(b"foo".as_slice()).unwrap().value, b"hi");
897 assert!(!kvs.contains_key(b"missing".as_slice()));
898 }
899
900 #[tokio::test]
901 async fn test_move_values() {
902 let kv_backend = Arc::new(MemoryKvBackend::default());
903 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
904 let kvs = HashMap::from([
905 (b"bar".to_vec(), b"baz".to_vec()),
906 (b"foo".to_vec(), b"hi".to_vec()),
907 (b"baz".to_vec(), b"hello".to_vec()),
908 ]);
909 for (key, value) in &kvs {
910 kv_backend
911 .put(
912 PutRequest::new()
913 .with_key(key.clone())
914 .with_value(value.clone()),
915 )
916 .await
917 .unwrap();
918 }
919 let move_values = kvs
920 .iter()
921 .map(|(key, value)| MoveValue {
922 key: key.clone(),
923 dest_key: tombstone_manager.to_tombstone(key),
924 value: value.clone(),
925 })
926 .collect::<Vec<_>>();
927 let (keys, dest_keys): (Vec<_>, Vec<_>) = move_values
928 .clone()
929 .into_iter()
930 .map(|kv| (kv.key, kv.dest_key))
931 .unzip();
932 let moved_keys = tombstone_manager
933 .move_values(keys.clone(), dest_keys.clone(), false)
934 .await
935 .unwrap();
936 assert_eq!(kvs.len(), moved_keys);
937 check_moved_values(kv_backend.clone(), &move_values).await;
938 let moved_keys = tombstone_manager
940 .move_values(keys.clone(), dest_keys.clone(), false)
941 .await
942 .unwrap();
943 assert_eq!(0, moved_keys);
944 check_moved_values(kv_backend.clone(), &move_values).await;
945 }
946
947 #[tokio::test]
948 async fn test_move_values_with_max_txn_ops() {
949 common_telemetry::init_default_ut_logging();
950 let kv_backend = Arc::new(MemoryKvBackend::default());
951 let mut tombstone_manager = TombstoneManager::new(kv_backend.clone());
952 tombstone_manager.set_max_txn_ops(4);
953 let kvs = HashMap::from([
954 (b"bar".to_vec(), b"baz".to_vec()),
955 (b"foo".to_vec(), b"hi".to_vec()),
956 (b"baz".to_vec(), b"hello".to_vec()),
957 (b"qux".to_vec(), b"world".to_vec()),
958 (b"quux".to_vec(), b"world".to_vec()),
959 (b"quuux".to_vec(), b"world".to_vec()),
960 (b"quuuux".to_vec(), b"world".to_vec()),
961 (b"quuuuux".to_vec(), b"world".to_vec()),
962 (b"quuuuuux".to_vec(), b"world".to_vec()),
963 ]);
964 for (key, value) in &kvs {
965 kv_backend
966 .put(
967 PutRequest::new()
968 .with_key(key.clone())
969 .with_value(value.clone()),
970 )
971 .await
972 .unwrap();
973 }
974 let move_values = kvs
975 .iter()
976 .map(|(key, value)| MoveValue {
977 key: key.clone(),
978 dest_key: tombstone_manager.to_tombstone(key),
979 value: value.clone(),
980 })
981 .collect::<Vec<_>>();
982 let (keys, dest_keys): (Vec<_>, Vec<_>) = move_values
983 .clone()
984 .into_iter()
985 .map(|kv| (kv.key, kv.dest_key))
986 .unzip();
987 let moved_keys = tombstone_manager
988 .move_values(keys.clone(), dest_keys.clone(), false)
989 .await
990 .unwrap();
991 assert_eq!(kvs.len(), moved_keys);
992 check_moved_values(kv_backend.clone(), &move_values).await;
993 let moved_keys = tombstone_manager
995 .move_values(keys.clone(), dest_keys.clone(), false)
996 .await
997 .unwrap();
998 assert_eq!(0, moved_keys);
999 check_moved_values(kv_backend.clone(), &move_values).await;
1000 }
1001
1002 #[tokio::test]
1003 async fn test_restore_chunks_by_total_txn_ops_limit() {
1004 let inner = Arc::new(MemoryKvBackend::default());
1005 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1006 inner: inner.clone(),
1007 max_txn_ops: 6,
1008 txn_op_counts: Mutex::new(Vec::new()),
1009 });
1010 let tombstone_manager = TombstoneManager::new(kv_backend);
1011 let kvs = HashMap::from([
1012 (b"bar".to_vec(), b"baz".to_vec()),
1013 (b"foo".to_vec(), b"hi".to_vec()),
1014 (b"baz".to_vec(), b"hello".to_vec()),
1015 ]);
1016 for (key, value) in &kvs {
1017 inner
1018 .put(
1019 PutRequest::new()
1020 .with_key(tombstone_manager.to_tombstone(key))
1021 .with_value(value.clone()),
1022 )
1023 .await
1024 .unwrap();
1025 }
1026
1027 let restored = tombstone_manager
1028 .restore(kvs.keys().cloned().collect())
1029 .await
1030 .unwrap();
1031
1032 assert_eq!(kvs.len(), restored);
1033 }
1034
1035 #[tokio::test]
1036 async fn test_restore_fails_fast_when_txn_op_limit_too_small() {
1037 let inner = Arc::new(MemoryKvBackend::default());
1038 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1039 inner: inner.clone(),
1040 max_txn_ops: 5,
1041 txn_op_counts: Mutex::new(Vec::new()),
1042 });
1043 let tombstone_manager = TombstoneManager::new(kv_backend);
1044 let key = b"foo".to_vec();
1045 inner
1046 .put(
1047 PutRequest::new()
1048 .with_key(tombstone_manager.to_tombstone(&key))
1049 .with_value(b"hi".to_vec()),
1050 )
1051 .await
1052 .unwrap();
1053
1054 let err = tombstone_manager.restore(vec![key]).await.unwrap_err();
1055
1056 assert!(matches!(err, Error::Unexpected { .. }));
1057 assert!(
1058 err.to_string()
1059 .contains("max_txn_ops 5 is smaller than required txn ops per key 6")
1060 );
1061 }
1062
1063 #[tokio::test]
1064 async fn test_create_chunks_by_total_txn_ops_limit() {
1065 let inner = Arc::new(MemoryKvBackend::default());
1066 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1067 inner: inner.clone(),
1068 max_txn_ops: 4,
1069 txn_op_counts: Mutex::new(Vec::new()),
1070 });
1071 let tombstone_manager = TombstoneManager::new(kv_backend);
1072 let kvs = HashMap::from([
1073 (b"bar".to_vec(), b"baz".to_vec()),
1074 (b"foo".to_vec(), b"hi".to_vec()),
1075 (b"baz".to_vec(), b"hello".to_vec()),
1076 ]);
1077 for (key, value) in &kvs {
1078 inner
1079 .put(
1080 PutRequest::new()
1081 .with_key(key.clone())
1082 .with_value(value.clone()),
1083 )
1084 .await
1085 .unwrap();
1086 }
1087
1088 let moved = tombstone_manager
1089 .create(kvs.keys().cloned().collect())
1090 .await
1091 .unwrap();
1092
1093 assert_eq!(kvs.len(), moved);
1094 }
1095
1096 #[tokio::test]
1097 async fn test_create_with_marker_uses_final_transaction_at_exact_limit() {
1098 let inner = Arc::new(MemoryKvBackend::default());
1099 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1100 inner: inner.clone(),
1101 max_txn_ops: MOVE_VALUE_TXN_OPS_PER_KEY,
1102 txn_op_counts: Mutex::new(Vec::new()),
1103 });
1104 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1105 for (key, value) in [(b"bar".as_slice(), b"baz".as_slice()), (b"foo", b"hi")] {
1106 inner
1107 .put(PutRequest::new().with_key(key).with_value(value))
1108 .await
1109 .unwrap();
1110 }
1111
1112 tombstone_manager
1113 .create_with_marker(
1114 vec![b"bar".to_vec(), b"foo".to_vec()],
1115 b"__dropped_at/42".to_vec(),
1116 b"1234".to_vec(),
1117 )
1118 .await
1119 .unwrap();
1120 tombstone_manager
1121 .create_with_marker(
1122 vec![b"bar".to_vec(), b"foo".to_vec()],
1123 b"__dropped_at/42".to_vec(),
1124 b"1234".to_vec(),
1125 )
1126 .await
1127 .unwrap();
1128
1129 assert_eq!(
1130 kv_backend.txn_op_counts.lock().unwrap().as_slice(),
1131 &[4, 4, 1, 1]
1132 );
1133 assert_eq!(
1134 inner
1135 .get(&tombstone_manager.to_tombstone(b"__dropped_at/42"))
1136 .await
1137 .unwrap()
1138 .unwrap()
1139 .value,
1140 b"1234"
1141 );
1142 }
1143
1144 #[tokio::test]
1145 async fn test_create_with_markers_chunks_extra_ops_after_key_move() {
1146 let inner = Arc::new(MemoryKvBackend::default());
1147 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1148 inner: inner.clone(),
1149 max_txn_ops: MOVE_VALUE_TXN_OPS_PER_KEY + 1,
1150 txn_op_counts: Mutex::new(Vec::new()),
1151 });
1152 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1153 inner
1154 .put(PutRequest::new().with_key(b"foo").with_value(b"bar"))
1155 .await
1156 .unwrap();
1157
1158 tombstone_manager
1159 .create_with_markers(
1160 vec![b"foo".to_vec()],
1161 (0..6)
1162 .map(|index| (format!("marker-{index}").into_bytes(), vec![index]))
1163 .collect(),
1164 )
1165 .await
1166 .unwrap();
1167
1168 assert_eq!(
1169 kv_backend.txn_op_counts.lock().unwrap().as_slice(),
1170 &[MOVE_VALUE_TXN_OPS_PER_KEY, 5, 1]
1171 );
1172 }
1173
1174 #[tokio::test]
1175 async fn test_create_with_markers_chunks_marker_only_transactions() {
1176 let inner = Arc::new(MemoryKvBackend::default());
1177 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1178 inner,
1179 max_txn_ops: 2,
1180 txn_op_counts: Mutex::new(Vec::new()),
1181 });
1182 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1183 let markers = (0..5)
1184 .map(|index| (format!("marker-{index}").into_bytes(), vec![index]))
1185 .collect();
1186
1187 tombstone_manager
1188 .create_with_markers(vec![], markers)
1189 .await
1190 .unwrap();
1191
1192 assert_eq!(
1193 kv_backend.txn_op_counts.lock().unwrap().as_slice(),
1194 &[2, 2, 1]
1195 );
1196 }
1197
1198 #[tokio::test]
1199 async fn test_restore_with_marker_uses_final_transaction_at_exact_limit() {
1200 let inner = Arc::new(MemoryKvBackend::default());
1201 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1202 inner: inner.clone(),
1203 max_txn_ops: RESTORE_VALUE_TXN_OPS_PER_KEY,
1204 txn_op_counts: Mutex::new(Vec::new()),
1205 });
1206 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1207 for (key, value) in [(b"bar".as_slice(), b"baz".as_slice()), (b"foo", b"hi")] {
1208 inner
1209 .put(
1210 PutRequest::new()
1211 .with_key(tombstone_manager.to_tombstone(key))
1212 .with_value(value),
1213 )
1214 .await
1215 .unwrap();
1216 }
1217 inner
1218 .put(
1219 PutRequest::new()
1220 .with_key(tombstone_manager.to_tombstone(b"__dropped_at/42"))
1221 .with_value(b"1234"),
1222 )
1223 .await
1224 .unwrap();
1225
1226 tombstone_manager
1227 .restore_with_marker(
1228 vec![b"bar".to_vec(), b"foo".to_vec()],
1229 b"__dropped_at/42".to_vec(),
1230 )
1231 .await
1232 .unwrap();
1233 tombstone_manager
1234 .restore_with_marker(
1235 vec![b"bar".to_vec(), b"foo".to_vec()],
1236 b"__dropped_at/42".to_vec(),
1237 )
1238 .await
1239 .unwrap();
1240
1241 assert_eq!(
1242 kv_backend.txn_op_counts.lock().unwrap().as_slice(),
1243 &[6, 6, 1, 1]
1244 );
1245 assert!(
1246 inner
1247 .get(&tombstone_manager.to_tombstone(b"__dropped_at/42"))
1248 .await
1249 .unwrap()
1250 .is_none()
1251 );
1252 }
1253
1254 #[tokio::test]
1255 async fn test_delete_with_marker_removes_marker_with_limited_backend() {
1256 let inner = Arc::new(MemoryKvBackend::default());
1257 let kv_backend = Arc::new(TxnOpLimitKvBackend {
1258 inner: inner.clone(),
1259 max_txn_ops: MOVE_VALUE_TXN_OPS_PER_KEY,
1260 txn_op_counts: Mutex::new(Vec::new()),
1261 });
1262 let tombstone_manager = TombstoneManager::new(kv_backend);
1263 for key in [b"bar".as_slice(), b"foo", b"__dropped_at/42"] {
1264 inner
1265 .put(
1266 PutRequest::new()
1267 .with_key(tombstone_manager.to_tombstone(key))
1268 .with_value(b"value"),
1269 )
1270 .await
1271 .unwrap();
1272 }
1273
1274 tombstone_manager
1275 .delete_with_marker(
1276 vec![b"bar".to_vec(), b"foo".to_vec()],
1277 b"__dropped_at/42".to_vec(),
1278 )
1279 .await
1280 .unwrap();
1281
1282 assert!(inner.is_empty());
1283 }
1284
1285 #[tokio::test]
1286 async fn test_move_values_with_non_exists_values() {
1287 let kv_backend = Arc::new(MemoryKvBackend::default());
1288 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1289 let kvs = HashMap::from([
1290 (b"bar".to_vec(), b"baz".to_vec()),
1291 (b"foo".to_vec(), b"hi".to_vec()),
1292 (b"baz".to_vec(), b"hello".to_vec()),
1293 ]);
1294 for (key, value) in &kvs {
1295 kv_backend
1296 .put(
1297 PutRequest::new()
1298 .with_key(key.clone())
1299 .with_value(value.clone()),
1300 )
1301 .await
1302 .unwrap();
1303 }
1304 let move_values = kvs
1305 .iter()
1306 .map(|(key, value)| MoveValue {
1307 key: key.clone(),
1308 dest_key: tombstone_manager.to_tombstone(key),
1309 value: value.clone(),
1310 })
1311 .collect::<Vec<_>>();
1312 let (mut keys, mut dest_keys): (Vec<_>, Vec<_>) = move_values
1313 .clone()
1314 .into_iter()
1315 .map(|kv| (kv.key, kv.dest_key))
1316 .unzip();
1317 keys.push(b"non-exists".to_vec());
1318 dest_keys.push(b"hi/non-exists".to_vec());
1319 let moved_keys = tombstone_manager
1320 .move_values(keys.clone(), dest_keys.clone(), false)
1321 .await
1322 .unwrap();
1323 check_moved_values(kv_backend.clone(), &move_values).await;
1324 assert_eq!(3, moved_keys);
1325 let moved_keys = tombstone_manager
1327 .move_values(keys.clone(), dest_keys.clone(), false)
1328 .await
1329 .unwrap();
1330 check_moved_values(kv_backend.clone(), &move_values).await;
1331 assert_eq!(0, moved_keys);
1332 }
1333
1334 #[tokio::test]
1335 async fn test_move_values_changed() {
1336 let kv_backend = Arc::new(MemoryKvBackend::default());
1337 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1338 let kvs = HashMap::from([
1339 (b"bar".to_vec(), b"baz".to_vec()),
1340 (b"foo".to_vec(), b"hi".to_vec()),
1341 (b"baz".to_vec(), b"hello".to_vec()),
1342 ]);
1343 for (key, value) in &kvs {
1344 kv_backend
1345 .put(
1346 PutRequest::new()
1347 .with_key(key.clone())
1348 .with_value(value.clone()),
1349 )
1350 .await
1351 .unwrap();
1352 }
1353
1354 kv_backend
1355 .put(PutRequest::new().with_key("baz").with_value("changed"))
1356 .await
1357 .unwrap();
1358
1359 let move_values = kvs
1360 .iter()
1361 .map(|(key, value)| MoveValue {
1362 key: key.clone(),
1363 dest_key: tombstone_manager.to_tombstone(key),
1364 value: value.clone(),
1365 })
1366 .collect::<Vec<_>>();
1367 let (keys, dest_keys): (Vec<_>, Vec<_>) = move_values
1368 .clone()
1369 .into_iter()
1370 .map(|kv| (kv.key, kv.dest_key))
1371 .unzip();
1372 let moved_keys = tombstone_manager
1373 .move_values(keys, dest_keys, false)
1374 .await
1375 .unwrap();
1376 assert_eq!(kvs.len(), moved_keys);
1377 }
1378
1379 #[tokio::test]
1380 async fn test_move_values_overwrite_dest_values() {
1381 let kv_backend = Arc::new(MemoryKvBackend::default());
1382 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1383 let kvs = HashMap::from([
1384 (b"bar".to_vec(), b"baz".to_vec()),
1385 (b"foo".to_vec(), b"hi".to_vec()),
1386 (b"baz".to_vec(), b"hello".to_vec()),
1387 ]);
1388 for (key, value) in &kvs {
1389 kv_backend
1390 .put(
1391 PutRequest::new()
1392 .with_key(key.clone())
1393 .with_value(value.clone()),
1394 )
1395 .await
1396 .unwrap();
1397 }
1398
1399 let move_values = kvs
1401 .iter()
1402 .map(|(key, value)| MoveValue {
1403 key: key.clone(),
1404 dest_key: tombstone_manager.to_tombstone(key),
1405 value: value.clone(),
1406 })
1407 .collect::<Vec<_>>();
1408 let (keys, dest_keys): (Vec<_>, Vec<_>) = move_values
1409 .clone()
1410 .into_iter()
1411 .map(|kv| (kv.key, kv.dest_key))
1412 .unzip();
1413 tombstone_manager
1414 .move_values(keys, dest_keys, false)
1415 .await
1416 .unwrap();
1417 check_moved_values(kv_backend.clone(), &move_values).await;
1418
1419 let kvs = HashMap::from([
1421 (b"bar".to_vec(), b"new baz".to_vec()),
1422 (b"foo".to_vec(), b"new hi".to_vec()),
1423 (b"baz".to_vec(), b"new baz".to_vec()),
1424 ]);
1425 for (key, value) in &kvs {
1426 kv_backend
1427 .put(
1428 PutRequest::new()
1429 .with_key(key.clone())
1430 .with_value(value.clone()),
1431 )
1432 .await
1433 .unwrap();
1434 }
1435 let move_values = kvs
1436 .iter()
1437 .map(|(key, value)| MoveValue {
1438 key: key.clone(),
1439 dest_key: tombstone_manager.to_tombstone(key),
1440 value: value.clone(),
1441 })
1442 .collect::<Vec<_>>();
1443 let (keys, dest_keys): (Vec<_>, Vec<_>) = move_values
1444 .clone()
1445 .into_iter()
1446 .map(|kv| (kv.key, kv.dest_key))
1447 .unzip();
1448 tombstone_manager
1449 .move_values(keys, dest_keys, false)
1450 .await
1451 .unwrap();
1452 check_moved_values(kv_backend.clone(), &move_values).await;
1453 }
1454
1455 #[tokio::test]
1456 async fn test_move_values_with_different_lengths() {
1457 let kv_backend = Arc::new(MemoryKvBackend::default());
1458 let tombstone_manager = TombstoneManager::new(kv_backend.clone());
1459
1460 let keys = vec![b"bar".to_vec(), b"foo".to_vec()];
1461 let dest_keys = vec![b"bar".to_vec(), b"foo".to_vec(), b"baz".to_vec()];
1462
1463 let err = tombstone_manager
1464 .move_values(keys, dest_keys, false)
1465 .await
1466 .unwrap_err();
1467 assert!(
1468 err.to_string()
1469 .contains("The length of keys(2) does not match the length of dest_keys(3)."),
1470 );
1471
1472 let moved_keys = tombstone_manager
1473 .move_values(vec![], vec![], false)
1474 .await
1475 .unwrap();
1476 assert_eq!(0, moved_keys);
1477 }
1478}