Skip to main content

common_meta/key/
tombstone.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
30/// [TombstoneManager] provides the ability to:
31/// - logically delete values
32/// - restore the deleted values
33///
34/// The tombstone mechanism is primarily used for table metadata deletion.
35/// When a table is logically deleted, its associated metadata keys are moved
36/// to a tombstone area by prepending a prefix (default is `__tombstone/`).
37///
38/// All involved keys in the tombstone mechanism:
39/// - `__table_name/{catalog}/{schema}/{table_name}`: Maps table name to table ID.
40/// - `__table_info/{table_id}`: Stores table information/metadata.
41/// - `__table_route/{table_id}`: Stores table routing information.
42/// - `__table_repart/{table_id}`: Stores table repartition information.
43/// - `__dn_table/{datanode_id}/{table_id}`: Maps datanodes to table regions.
44/// - `__topic_region/{topic_name}/{region_id}`: Maps regions to WAL topics.
45///
46/// These keys are moved to:
47/// - `__tombstone/__table_name/...`
48/// - `__tombstone/__table_info/...`
49/// - `__tombstone/__table_route/...`
50/// - `__tombstone/__table_repart/...`
51/// - `__tombstone/__dn_table/...`
52/// - `__tombstone/__topic_region/...`
53pub struct TombstoneManager {
54    kv_backend: KvBackendRef,
55    tombstone_prefix: String,
56    // Only used for testing.
57    #[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    /// Returns [TombstoneManager].
71    pub fn new(kv_backend: KvBackendRef) -> Self {
72        Self::new_with_prefix(kv_backend, TOMBSTONE_PREFIX)
73    }
74
75    /// Returns [TombstoneManager] with a custom tombstone prefix.
76    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    /// Removes the tombstone prefix from a tombstoned key.
90    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    /// Gets a single tombstoned value by its original key.
106    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    /// Gets tombstoned values by their original keys.
116    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    /// Streams all tombstoned key-value pairs.
133    pub fn tombstones(&self) -> BoxStream<'static, Result<KeyValue>> {
134        self.scan_prefix(self.tombstone_prefix.as_bytes().to_vec())
135    }
136
137    /// Streams tombstoned table-name entries only.
138    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    /// Streams tombstoned table-name entries in the provided catalog.
145    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    /// Streams tombstoned entries under the provided prefix.
159    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    /// Moves value to `dest_key`.
173    ///
174    /// Puts `value` to `dest_key` if the value of `src_key` equals `value`.
175    ///
176    /// Otherwise retrieves the value of `src_key`.
177    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        // The key -> dest key mapping.
228        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            // Updates results.
267            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    /// Moves values to `dest_key`.
321    ///
322    /// Returns the number of keys that were moved.
323    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    /// Creates tombstones for keys.
411    ///
412    /// Preforms to:
413    /// - deletes origin values.
414    /// - stores tombstone values.
415    ///
416    /// Returns the number of keys that were moved.
417    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    /// Creates tombstones and stores an additional marker in the tombstone namespace.
430    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    /// Creates tombstones and stores additional markers in the tombstone namespace.
441    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    /// Restores tombstones for keys.
463    ///
464    /// Preforms to:
465    /// - restore origin value.
466    /// - deletes tombstone values.
467    ///
468    /// Returns the number of keys that were restored.
469    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    /// Restores tombstones and deletes an associated marker.
482    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    /// Restores tombstones and deletes associated markers.
501    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    /// Deletes tombstones values for the specified `keys`.
522    ///
523    /// Returns the number of keys that were deleted.
524    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    /// Deletes tombstones and an associated marker.
540    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    /// Deletes tombstones and associated markers.
550    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        // Moves again
939        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        // Moves again
994        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        // Moves again
1326        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        // Prepares
1400        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        // Overwrites existing dest keys.
1420        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}