meta_srv/procedure/repartition/
gc_requirement.rs1use 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
32pub 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 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 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 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}