meta_srv/procedure/wal_prune/
manager.rs1use std::collections::HashSet;
16use std::fmt::{Debug, Formatter};
17use std::sync::{Arc, RwLock};
18
19use common_procedure::{ProcedureId, ProcedureManagerRef, ProcedureWithId, watcher};
20use common_telemetry::{debug, error, info, warn};
21use futures::future::join_all;
22use snafu::{OptionExt, ResultExt};
23use tokio::sync::Semaphore;
24use tokio::sync::mpsc::{Receiver, Sender};
25
26use crate::define_ticker;
27use crate::error::{self, Result};
28use crate::metrics::METRIC_META_REMOTE_WAL_PRUNE_EXECUTE;
29use crate::procedure::wal_prune::utils::{find_pruneable_entry_id_for_topic, should_trigger_prune};
30use crate::procedure::wal_prune::{Context as WalPruneContext, WalPruneProcedure};
31
32pub type WalPruneTickerRef = Arc<WalPruneTicker>;
33
34#[derive(Clone)]
39pub struct WalPruneProcedureTracker {
40 running_procedures: Arc<RwLock<HashSet<String>>>,
41}
42
43impl WalPruneProcedureTracker {
44 pub fn insert_running_procedure(&self, topic_name: String) -> Option<WalPruneProcedureGuard> {
47 let mut running_procedures = self.running_procedures.write().unwrap();
48 if running_procedures.insert(topic_name.clone()) {
49 Some(WalPruneProcedureGuard {
50 topic_name,
51 running_procedures: self.running_procedures.clone(),
52 })
53 } else {
54 None
55 }
56 }
57
58 pub fn len(&self) -> usize {
60 self.running_procedures.read().unwrap().len()
61 }
62}
63
64pub struct WalPruneProcedureGuard {
68 topic_name: String,
69 running_procedures: Arc<RwLock<HashSet<String>>>,
70}
71
72impl Drop for WalPruneProcedureGuard {
73 fn drop(&mut self) {
74 let mut running_procedures = self.running_procedures.write().unwrap();
75 running_procedures.remove(&self.topic_name);
76 }
77}
78
79pub enum Event {
83 Tick,
84}
85
86impl Debug for Event {
87 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
88 match self {
89 Event::Tick => write!(f, "Tick"),
90 }
91 }
92}
93
94define_ticker!(
95 WalPruneTicker,
98 event_type = Event,
99 event_value = Event::Tick
100);
101
102pub(crate) struct WalPruneManager {
109 receiver: Receiver<Event>,
111 procedure_manager: ProcedureManagerRef,
113 tracker: WalPruneProcedureTracker,
115 semaphore: Arc<Semaphore>,
117 logical_delete: bool,
119
120 wal_prune_context: WalPruneContext,
122}
123
124impl WalPruneManager {
125 pub fn new(
127 parallelism: usize,
128 logical_delete: bool,
129 receiver: Receiver<Event>,
130 procedure_manager: ProcedureManagerRef,
131 wal_prune_context: WalPruneContext,
132 ) -> Self {
133 Self {
134 receiver,
135 procedure_manager,
136 wal_prune_context,
137 tracker: WalPruneProcedureTracker {
138 running_procedures: Arc::new(RwLock::new(HashSet::new())),
139 },
140 semaphore: Arc::new(Semaphore::new(parallelism)),
141 logical_delete,
142 }
143 }
144
145 pub async fn try_start(mut self) -> Result<()> {
147 let context = self.wal_prune_context.clone();
148 let tracker = self.tracker.clone();
149 self.procedure_manager
150 .register_loader(
151 WalPruneProcedure::TYPE_NAME,
152 Box::new(move |json| {
153 let tracker = tracker.clone();
154 WalPruneProcedure::from_json(json, &context, tracker).map(|p| Box::new(p) as _)
155 }),
156 )
157 .context(error::RegisterProcedureLoaderSnafu {
158 type_name: WalPruneProcedure::TYPE_NAME,
159 })?;
160 common_runtime::spawn_global(async move {
161 self.run().await;
162 });
163 info!("WalPruneProcedureManager Started.");
164 Ok(())
165 }
166
167 pub(crate) fn channel() -> (Sender<Event>, Receiver<Event>) {
169 tokio::sync::mpsc::channel(1024)
170 }
171
172 pub(crate) async fn run(&mut self) {
176 while let Some(event) = self.receiver.recv().await {
177 match event {
178 Event::Tick => self.handle_tick_request().await.unwrap_or_else(|e| {
179 error!(e; "Failed to handle tick request");
180 }),
181 }
182 }
183 }
184
185 pub async fn wait_procedure(
187 &self,
188 topic_name: &str,
189 prunable_entry_id: u64,
190 ) -> Result<ProcedureId> {
191 let guard = self
192 .tracker
193 .insert_running_procedure(topic_name.to_string())
194 .with_context(|| error::PruneTaskAlreadyRunningSnafu { topic: topic_name })?;
195
196 let procedure = WalPruneProcedure::new(
197 self.wal_prune_context.clone(),
198 Some(guard),
199 topic_name.to_string(),
200 prunable_entry_id,
201 self.logical_delete,
202 );
203 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure));
204 let procedure_id = procedure_with_id.id;
205 METRIC_META_REMOTE_WAL_PRUNE_EXECUTE
206 .with_label_values(&[topic_name])
207 .inc();
208 let procedure_manager = self.procedure_manager.clone();
209 let mut watcher = procedure_manager
210 .submit(procedure_with_id)
211 .await
212 .context(error::SubmitProcedureSnafu)?;
213 watcher::wait(&mut watcher)
214 .await
215 .context(error::WaitProcedureSnafu)?;
216
217 Ok(procedure_id)
218 }
219
220 async fn try_prune(&self, topic_name: &str) -> Result<()> {
221 let table_metadata_manager = self.wal_prune_context.table_metadata_manager.clone();
222 let leader_region_registry = self.wal_prune_context.leader_region_registry.clone();
223 let prunable_entry_id = find_pruneable_entry_id_for_topic(
224 &table_metadata_manager,
225 &leader_region_registry,
226 topic_name,
227 )
228 .await?;
229 let Some(prunable_entry_id) = prunable_entry_id else {
230 debug!(
231 "No prunable entry id found for topic {}, skipping prune",
232 topic_name
233 );
234 return Ok(());
235 };
236 let current = table_metadata_manager
237 .topic_name_manager()
238 .get(topic_name)
239 .await
240 .context(error::TableMetadataManagerSnafu)?
241 .map(|v| v.into_inner().pruned_entry_id);
242 debug!(
243 "Found prunable entry id {} for topic {}, current pruned entry id: {:?}",
244 prunable_entry_id, topic_name, current
245 );
246 if !should_trigger_prune(current, prunable_entry_id) {
247 debug!(
248 "No need to prune topic {}, current pruned entry id: {:?}, prunable entry id: {}",
249 topic_name, current, prunable_entry_id
250 );
251 return Ok(());
252 }
253
254 self.wait_procedure(topic_name, prunable_entry_id)
255 .await
256 .map(|_| ())
257 }
258
259 async fn handle_tick_request(&self) -> Result<()> {
260 let topics = self.retrieve_sorted_topics().await?;
261 let mut tasks = Vec::with_capacity(topics.len());
262 for topic_name in topics.iter() {
263 tasks.push(async {
264 let _permit = self.semaphore.acquire().await.unwrap();
265 match self.try_prune(topic_name).await {
266 Ok(_) => {}
267 Err(error::Error::PruneTaskAlreadyRunning { topic, .. }) => {
268 warn!("Prune task for topic {} is already running", topic);
269 }
270 Err(err) => {
271 error!(err; "Failed to prune remote WAL for topic {}", topic_name.as_str());
272 }
273 }
274 });
275 }
276
277 join_all(tasks).await;
278 Ok(())
279 }
280
281 async fn retrieve_sorted_topics(&self) -> Result<Vec<String>> {
285 self.wal_prune_context
286 .table_metadata_manager
287 .topic_name_manager()
288 .range()
289 .await
290 .context(error::TableMetadataManagerSnafu)
291 }
292}
293
294#[cfg(test)]
295mod test {
296 use std::assert_matches;
297 use std::time::Duration;
298
299 use common_meta::key::topic_name::TopicNameKey;
300 use common_meta::leadership_notifier::LeadershipChangeListener;
301 use common_wal::maybe_skip_kafka_integration_test;
302 use common_wal::test_util::get_kafka_endpoints;
303 use tokio::time::{sleep, timeout};
304
305 use super::*;
306 use crate::procedure::test_util::new_wal_prune_metadata;
307 use crate::procedure::wal_prune::test_util::TestEnv;
308
309 #[tokio::test]
310 async fn test_wal_prune_ticker() {
311 common_telemetry::init_default_ut_logging();
312 let (tx, mut rx) = WalPruneManager::channel();
313 let interval = Duration::from_millis(50);
314 let ticker = WalPruneTicker::new(interval, tx);
315 assert_eq!(ticker.name(), "WalPruneTicker");
316
317 for _ in 0..2 {
318 ticker.start();
319 sleep(4 * interval).await;
321 assert!(!rx.is_empty());
322 while let Ok(event) = rx.try_recv() {
323 assert_matches!(event, Event::Tick);
324 }
325 }
326 ticker.stop();
327 }
328
329 #[tokio::test]
330 async fn test_wal_prune_tracker_and_guard() {
331 let tracker = WalPruneProcedureTracker {
332 running_procedures: Arc::new(RwLock::new(HashSet::new())),
333 };
334 let topic_name = uuid::Uuid::new_v4().to_string();
335 {
336 let guard = tracker
337 .insert_running_procedure(topic_name.clone())
338 .unwrap();
339 assert_eq!(guard.topic_name, topic_name);
340 assert_eq!(guard.running_procedures.read().unwrap().len(), 1);
341
342 let result = tracker.insert_running_procedure(topic_name.clone());
343 assert!(result.is_none());
344 }
345 assert_eq!(tracker.running_procedures.read().unwrap().len(), 0);
346 }
347
348 async fn mock_wal_prune_manager(
349 broker_endpoints: Vec<String>,
350 limit: usize,
351 ) -> (Sender<Event>, WalPruneManager) {
352 let test_env = TestEnv::new();
353 let (tx, rx) = WalPruneManager::channel();
354 let wal_prune_context = test_env.build_wal_prune_context(broker_endpoints).await;
355 (
356 tx,
357 WalPruneManager::new(
358 limit,
359 false,
360 rx,
361 test_env.procedure_manager.clone(),
362 wal_prune_context,
363 ),
364 )
365 }
366
367 async fn mock_topics(manager: &WalPruneManager, topics: &[String]) {
368 let topic_name_keys = topics
369 .iter()
370 .map(|topic| TopicNameKey::new(topic))
371 .collect::<Vec<_>>();
372 manager
373 .wal_prune_context
374 .table_metadata_manager
375 .topic_name_manager()
376 .batch_put(topic_name_keys)
377 .await
378 .unwrap();
379 }
380
381 #[tokio::test]
382 async fn test_wal_prune_manager() {
383 maybe_skip_kafka_integration_test!();
384 let broker_endpoints = get_kafka_endpoints();
385 let limit = 6;
386 let (tx, manager) = mock_wal_prune_manager(broker_endpoints, limit).await;
387 let topics = (0..limit * 2)
388 .map(|_| uuid::Uuid::new_v4().to_string())
389 .collect::<Vec<_>>();
390 mock_topics(&manager, &topics).await;
391
392 let tracker = manager.tracker.clone();
393 let handler =
394 common_runtime::spawn_global(async move { manager.try_start().await.unwrap() });
395 handler.await.unwrap();
396
397 tx.send(Event::Tick).await.unwrap();
398 timeout(Duration::from_millis(100), async move { tracker.len() > 0 })
400 .await
401 .unwrap();
402 }
403
404 #[tokio::test]
405 async fn test_find_pruneable_entry_id_for_topic_none() {
406 let test_env = TestEnv::new();
407 let prunable_entry_id = find_pruneable_entry_id_for_topic(
408 &test_env.table_metadata_manager,
409 &test_env.leader_region_registry,
410 "test_topic",
411 )
412 .await
413 .unwrap();
414 assert!(prunable_entry_id.is_none());
415 }
416
417 #[tokio::test]
418 async fn test_find_pruneable_entry_id_for_topic_some() {
419 let test_env = TestEnv::new();
420 let topic = "test_topic";
421 let expected_prunable_entry_id = new_wal_prune_metadata(
422 test_env.table_metadata_manager.clone(),
423 test_env.leader_region_registry.clone(),
424 2,
425 5,
426 &[3, 10, 23, 50, 52, 82, 130],
427 topic.to_string(),
428 )
429 .await;
430 let prunable_entry_id = find_pruneable_entry_id_for_topic(
431 &test_env.table_metadata_manager,
432 &test_env.leader_region_registry,
433 topic,
434 )
435 .await
436 .unwrap()
437 .unwrap();
438 assert_eq!(prunable_entry_id, expected_prunable_entry_id);
439 }
440}