Skip to main content

meta_srv/procedure/wal_prune/
manager.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::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/// Tracks running [WalPruneProcedure]s and the resources they hold.
35/// A [WalPruneProcedure] is holding a semaphore permit to limit the number of concurrent procedures.
36///
37/// TODO(CookiePie): Similar to [RegionMigrationProcedureTracker], maybe can refactor to a unified framework.
38#[derive(Clone)]
39pub struct WalPruneProcedureTracker {
40    running_procedures: Arc<RwLock<HashSet<String>>>,
41}
42
43impl WalPruneProcedureTracker {
44    /// Insert a running [WalPruneProcedure] for the given topic name and
45    /// consume acquire a semaphore permit for the given topic name.
46    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    /// Number of running [WalPruneProcedure]s.
59    pub fn len(&self) -> usize {
60        self.running_procedures.read().unwrap().len()
61    }
62}
63
64/// [WalPruneProcedureGuard] is a guard for [WalPruneProcedure].
65/// It is used to track the running [WalPruneProcedure]s.
66/// When the guard is dropped, it will remove the topic name from the running procedures and release the semaphore.
67pub 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
79/// Event is used to notify the [WalPruneManager] to do some work.
80///
81/// - `Tick`: Trigger a submission of [WalPruneProcedure] to prune remote WAL.
82pub 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] is a ticker that periodically sends [Event]s to the [WalPruneManager].
96    /// It is used to trigger the [WalPruneManager] to submit [WalPruneProcedure]s.
97    WalPruneTicker,
98    event_type = Event,
99    event_value = Event::Tick
100);
101
102/// [WalPruneManager] manages all remote WAL related tasks in metasrv.
103///
104/// [WalPruneManager] is responsible for:
105/// 1. Registering [WalPruneProcedure] loader in the procedure manager.
106/// 2. Periodically receive [Event::Tick] to submit [WalPruneProcedure] to prune remote WAL.
107/// 3. Use a semaphore to limit the number of concurrent [WalPruneProcedure]s.
108pub(crate) struct WalPruneManager {
109    /// Receives [Event]s.
110    receiver: Receiver<Event>,
111    /// Procedure manager.
112    procedure_manager: ProcedureManagerRef,
113    /// Tracker for running [WalPruneProcedure]s.
114    tracker: WalPruneProcedureTracker,
115    /// Semaphore to limit the number of concurrent [WalPruneProcedure]s.
116    semaphore: Arc<Semaphore>,
117    /// Whether pruning only updates metadata and skips Kafka DeleteRecords.
118    logical_delete: bool,
119
120    /// Context for [WalPruneProcedure].
121    wal_prune_context: WalPruneContext,
122}
123
124impl WalPruneManager {
125    /// Returns a new empty [`WalPruneManager`].
126    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    /// Start the [WalPruneManager]. It will register [WalPruneProcedure] loader in the procedure manager.
146    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    /// Returns a mpsc channel with a buffer capacity of 1024 for sending and receiving `Event` messages.
168    pub(crate) fn channel() -> (Sender<Event>, Receiver<Event>) {
169        tokio::sync::mpsc::channel(1024)
170    }
171
172    /// Runs the main loop. Performs actions on received events.
173    ///
174    /// - `Tick`: Submit `limit` [WalPruneProcedure]s to prune remote WAL.
175    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    /// Submits a [WalPruneProcedure] for the given topic name.
186    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    /// Retrieve topics from the table metadata manager.
282    /// Since [WalPruneManager] submits procedures depending on the order of the topics, we should sort the topics.
283    /// TODO(CookiePie): Can register topics in memory instead of retrieving from the table metadata manager every time.
284    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            // wait a bit longer to make sure not all ticks are skipped
320            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        // Wait for at least one procedure to be submitted.
399        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}