Skip to main content

meta_srv/procedure/repartition/
gc_requirement.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::sync::Arc;
16
17use common_meta::key::REPARTITION_GC_REQUIRED_KEY;
18use common_meta::key::table_repart::TableRepartManager;
19use common_meta::kv_backend::KvBackendRef;
20use common_meta::rpc::store::PutRequest;
21use common_procedure::ProcedureManagerRef;
22use snafu::{ResultExt, ensure};
23
24use crate::error::{self, Result};
25use crate::procedure::repartition::RepartitionProcedure;
26use crate::procedure::repartition::group::RepartitionGroupProcedure;
27
28const REPARTITION_GC_REQUIREMENT_VALUE: &[u8] = b"repartition";
29
30pub type RepartitionGcRequirementManagerRef = Arc<RepartitionGcRequirementManager>;
31
32/// Persists and enforces the cluster-level GC requirement introduced by repartition.
33pub struct RepartitionGcRequirementManager {
34    kv_backend: KvBackendRef,
35}
36
37impl RepartitionGcRequirementManager {
38    pub fn new(kv_backend: KvBackendRef) -> Self {
39        Self { kv_backend }
40    }
41
42    /// Persists the requirement before a repartition procedure can be submitted.
43    pub async fn require_gc(&self) -> Result<()> {
44        self.kv_backend
45            .put(
46                PutRequest::new()
47                    .with_key(REPARTITION_GC_REQUIRED_KEY.as_bytes().to_vec())
48                    .with_value(REPARTITION_GC_REQUIREMENT_VALUE.to_vec()),
49            )
50            .await
51            .context(error::RepartitionGcRequirementSnafu)?;
52        Ok(())
53    }
54
55    pub async fn is_gc_required(&self) -> Result<bool> {
56        self.kv_backend
57            .get(REPARTITION_GC_REQUIRED_KEY.as_bytes())
58            .await
59            .map(|value| value.is_some())
60            .context(error::RepartitionGcRequirementSnafu)
61    }
62
63    /// Backfills the requirement from durable state written by older versions.
64    ///
65    /// An unfinished procedure covers the crash window before repartition mappings
66    /// are written. A non-empty mapping covers completed repartitions whose
67    /// manifests can still contain cross-region file references.
68    pub async fn reconcile_legacy_state(
69        &self,
70        procedure_manager: &ProcedureManagerRef,
71    ) -> Result<bool> {
72        if self.is_gc_required().await? {
73            return Ok(true);
74        }
75
76        let table_reparts = TableRepartManager::new(self.kv_backend.clone())
77            .table_reparts()
78            .await
79            .context(error::RepartitionGcRequirementSnafu)?;
80        let has_cross_region_references = table_reparts
81            .iter()
82            .any(|(_, value)| !value.src_to_dst.is_empty());
83        let has_unfinished_repartition = procedure_manager
84            .has_unfinished_procedure(&[
85                RepartitionProcedure::TYPE_NAME,
86                RepartitionGroupProcedure::TYPE_NAME,
87            ])
88            .await
89            .context(error::InspectRepartitionProceduresSnafu)?;
90
91        if has_cross_region_references || has_unfinished_repartition {
92            self.require_gc().await?;
93            return Ok(true);
94        }
95
96        Ok(false)
97    }
98
99    pub async fn ensure_gc_enabled(
100        &self,
101        gc_enabled: bool,
102        procedure_manager: &ProcedureManagerRef,
103    ) -> Result<()> {
104        // Reconciliation is also a legacy migration and must run before GC can
105        // remove the repartition mappings used to backfill the durable marker.
106        let required = self.reconcile_legacy_state(procedure_manager).await?;
107        ensure!(!required || gc_enabled, error::RepartitionGcRequiredSnafu);
108        Ok(())
109    }
110}
111
112#[cfg(test)]
113mod tests {
114    use std::collections::HashMap;
115    use std::time::Duration;
116
117    use async_trait::async_trait;
118    use common_meta::kv_backend::memory::MemoryKvBackend;
119    use common_meta::state_store::KvStateStore;
120    use common_procedure::local::{LocalManager, ManagerConfig};
121    use common_procedure::{
122        Context, LockKey, Procedure, ProcedureId, ProcedureManager, ProcedureWithId, Status,
123    };
124    use store_api::storage::RegionId;
125    use tokio::sync::oneshot;
126    use tokio::time::timeout;
127
128    use super::*;
129
130    fn local_procedure_manager(kv_backend: KvBackendRef) -> Arc<LocalManager> {
131        let state_store = Arc::new(KvStateStore::new(kv_backend));
132        Arc::new(LocalManager::new(
133            ManagerConfig::default(),
134            state_store.clone(),
135            state_store,
136            None,
137            None,
138        ))
139    }
140
141    fn procedure_manager(kv_backend: KvBackendRef) -> ProcedureManagerRef {
142        local_procedure_manager(kv_backend)
143    }
144
145    #[derive(Debug)]
146    struct LegacyRepartitionProcedure {
147        persisted: bool,
148        block_after_persist: bool,
149        persisted_tx: Option<oneshot::Sender<()>>,
150    }
151
152    #[async_trait]
153    impl Procedure for LegacyRepartitionProcedure {
154        fn type_name(&self) -> &str {
155            RepartitionProcedure::TYPE_NAME
156        }
157
158        async fn execute(&mut self, _ctx: &Context) -> common_procedure::error::Result<Status> {
159            if !self.persisted {
160                self.persisted = true;
161                return Ok(Status::executing(true));
162            }
163
164            if self.block_after_persist {
165                if let Some(tx) = self.persisted_tx.take() {
166                    let _ = tx.send(());
167                }
168                return std::future::pending().await;
169            }
170
171            Ok(Status::done())
172        }
173
174        fn dump(&self) -> common_procedure::error::Result<String> {
175            Ok("{}".to_string())
176        }
177
178        fn lock_key(&self) -> LockKey {
179            LockKey::default()
180        }
181    }
182
183    #[tokio::test]
184    async fn test_completed_repartition_requirement_rejects_gc_disabled_restart() {
185        let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::new());
186        let manager = RepartitionGcRequirementManager::new(kv_backend.clone());
187        let procedure_manager = procedure_manager(kv_backend.clone());
188
189        manager.require_gc().await.unwrap();
190        let marker = kv_backend
191            .get(REPARTITION_GC_REQUIRED_KEY.as_bytes())
192            .await
193            .unwrap()
194            .unwrap();
195        assert_eq!(REPARTITION_GC_REQUIREMENT_VALUE, marker.value.as_slice());
196
197        let err = manager
198            .ensure_gc_enabled(false, &procedure_manager)
199            .await
200            .unwrap_err();
201        assert!(matches!(err, error::Error::RepartitionGcRequired { .. }));
202        assert!(manager.is_gc_required().await.unwrap());
203        manager
204            .ensure_gc_enabled(true, &procedure_manager)
205            .await
206            .unwrap();
207    }
208
209    #[tokio::test]
210    async fn test_reconcile_legacy_repartition_mapping() {
211        let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::new());
212        let manager = RepartitionGcRequirementManager::new(kv_backend.clone());
213        let procedure_manager = procedure_manager(kv_backend.clone());
214        let table_id = 1024;
215        let source = RegionId::new(table_id, 1);
216        let destination = RegionId::new(table_id, 2);
217        TableRepartManager::new(kv_backend)
218            .update_mappings(table_id, &HashMap::from([(source, vec![destination])]))
219            .await
220            .unwrap();
221
222        assert!(
223            manager
224                .reconcile_legacy_state(&procedure_manager)
225                .await
226                .unwrap()
227        );
228        assert!(manager.is_gc_required().await.unwrap());
229    }
230
231    #[tokio::test]
232    async fn test_legacy_unfinished_repartition_fails_closed_and_recovers() {
233        let kv_backend: KvBackendRef = Arc::new(MemoryKvBackend::new());
234        let first_manager = local_procedure_manager(kv_backend.clone());
235        first_manager.start().await.unwrap();
236
237        let procedure_id = ProcedureId::random();
238        let (persisted_tx, persisted_rx) = oneshot::channel();
239        first_manager
240            .submit(ProcedureWithId {
241                id: procedure_id,
242                procedure: Box::new(LegacyRepartitionProcedure {
243                    persisted: false,
244                    block_after_persist: true,
245                    persisted_tx: Some(persisted_tx),
246                }),
247                context: Default::default(),
248            })
249            .await
250            .unwrap();
251        timeout(Duration::from_secs(10), persisted_rx)
252            .await
253            .unwrap()
254            .unwrap();
255        first_manager.stop().await.unwrap();
256
257        let requirement_manager = RepartitionGcRequirementManager::new(kv_backend.clone());
258        let disabled_manager = local_procedure_manager(kv_backend.clone());
259        let disabled_manager_ref: ProcedureManagerRef = disabled_manager.clone();
260        let err = requirement_manager
261            .ensure_gc_enabled(false, &disabled_manager_ref)
262            .await
263            .unwrap_err();
264        assert!(matches!(err, error::Error::RepartitionGcRequired { .. }));
265        assert!(requirement_manager.is_gc_required().await.unwrap());
266
267        assert!(
268            disabled_manager
269                .has_unfinished_procedure(&[RepartitionProcedure::TYPE_NAME])
270                .await
271                .unwrap()
272        );
273        assert!(
274            disabled_manager
275                .procedure_state(procedure_id)
276                .await
277                .unwrap()
278                .is_none()
279        );
280
281        let enabled_manager = local_procedure_manager(kv_backend);
282        enabled_manager
283            .register_loader(
284                RepartitionProcedure::TYPE_NAME,
285                Box::new(|_| {
286                    Ok(Box::new(LegacyRepartitionProcedure {
287                        persisted: true,
288                        block_after_persist: false,
289                        persisted_tx: None,
290                    }) as _)
291                }),
292            )
293            .unwrap();
294        enabled_manager.start().await.unwrap();
295        let mut watcher = enabled_manager.procedure_watcher(procedure_id).unwrap();
296        timeout(Duration::from_secs(10), async {
297            while !watcher.borrow().is_done() {
298                watcher.changed().await.unwrap();
299            }
300        })
301        .await
302        .unwrap();
303        assert!(
304            !enabled_manager
305                .has_unfinished_procedure(&[RepartitionProcedure::TYPE_NAME])
306                .await
307                .unwrap()
308        );
309        enabled_manager.stop().await.unwrap();
310    }
311}