Skip to main content

meta_srv/procedure/
wal_prune.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
15pub(crate) mod manager;
16#[cfg(test)]
17mod test_util;
18pub(crate) mod utils;
19
20use std::sync::Arc;
21
22use common_error::ext::BoxedError;
23use common_meta::key::TableMetadataManagerRef;
24use common_meta::lock_key::RemoteWalLock;
25use common_meta::region_registry::LeaderRegionRegistryRef;
26use common_procedure::error::ToJsonSnafu;
27use common_procedure::{
28    Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey,
29    Procedure, ProcedureState, Result as ProcedureResult, Status, StringKey,
30};
31use common_telemetry::{info, warn};
32use manager::{WalPruneProcedureGuard, WalPruneProcedureTracker};
33use rskafka::client::Client;
34use serde::{Deserialize, Serialize};
35use snafu::ResultExt;
36use store_api::logstore::EntryId;
37
38use crate::Result;
39use crate::error::{self};
40use crate::event::wal_prune::{WAL_PRUNE_EVENT_TYPE, WalPruneEvent};
41use crate::procedure::wal_prune::utils::{
42    delete_records, get_offsets_for_topic, get_partition_client, update_pruned_entry_id,
43};
44
45pub type KafkaClientRef = Arc<Client>;
46
47#[derive(Clone)]
48pub struct Context {
49    /// The Kafka client.
50    pub client: KafkaClientRef,
51    /// The table metadata manager.
52    pub table_metadata_manager: TableMetadataManagerRef,
53    /// The leader region registry.
54    pub leader_region_registry: LeaderRegionRegistryRef,
55}
56
57/// The data of WAL pruning.
58#[derive(Serialize, Deserialize)]
59pub struct WalPruneData {
60    /// The topic name to prune.
61    pub topic: String,
62    /// The minimum flush entry id for topic, which is used to prune the WAL.
63    pub prunable_entry_id: EntryId,
64    /// Whether pruning only updates metadata and skips Kafka DeleteRecords.
65    #[serde(default)]
66    pub logical_delete: bool,
67}
68
69#[derive(Debug, Clone, Copy, PartialEq, Eq)]
70struct WalPruneOutcome;
71
72/// The procedure to prune WAL.
73pub struct WalPruneProcedure {
74    pub data: WalPruneData,
75    pub context: Context,
76    /// The latest offset observed during the current execution attempt.
77    observed_latest_offset: Option<u64>,
78    pub _guard: Option<WalPruneProcedureGuard>,
79}
80
81impl WalPruneProcedure {
82    const TYPE_NAME: &'static str = "metasrv-procedure::WalPrune";
83
84    pub fn new(
85        context: Context,
86        guard: Option<WalPruneProcedureGuard>,
87        topic: String,
88        prunable_entry_id: u64,
89        logical_delete: bool,
90    ) -> Self {
91        Self {
92            data: WalPruneData {
93                topic,
94                prunable_entry_id,
95                logical_delete,
96            },
97            context,
98            observed_latest_offset: None,
99            _guard: guard,
100        }
101    }
102
103    pub fn from_json(
104        json: &str,
105        context: &Context,
106        tracker: WalPruneProcedureTracker,
107    ) -> ProcedureResult<Self> {
108        let data: WalPruneData = serde_json::from_str(json).context(ToJsonSnafu)?;
109        let guard = tracker.insert_running_procedure(data.topic.clone());
110        Ok(Self {
111            data,
112            context: context.clone(),
113            observed_latest_offset: None,
114            _guard: guard,
115        })
116    }
117
118    /// Prune the WAL and persist the minimum prunable entry id.
119    ///
120    /// Retry:
121    /// - Kafka client errors that have exhausted rskafka's internal retry.
122    /// - Failed to update the pruned entry id in the table metadata manager.
123    ///
124    /// WAL prune event delivery is best effort. A physical prune that completes before a process
125    /// restart can recover as a no-op and emit no `Succeeded` event.
126    pub async fn on_prune(&mut self) -> Result<Status> {
127        self.observed_latest_offset = None;
128        let partition_client = get_partition_client(&self.context.client, &self.data.topic).await?;
129        let (earliest_offset, latest_offset) =
130            get_offsets_for_topic(&partition_client, &self.data.topic).await?;
131        self.observed_latest_offset = Some(latest_offset);
132        if self.data.prunable_entry_id <= earliest_offset {
133            warn!(
134                "The prunable entry id is less or equal to the earliest offset, topic: {}, prunable entry id: {}, earliest offset: {}, latest offset: {}",
135                self.data.topic, self.data.prunable_entry_id, earliest_offset, latest_offset
136            );
137            return Ok(Status::done());
138        }
139
140        if self.data.logical_delete {
141            info!(
142                "Skipping physical deletion of records for logical WAL pruning, topic: {}, prunable entry id: {}",
143                self.data.topic, self.data.prunable_entry_id
144            );
145        } else {
146            // Delete records.
147            delete_records(
148                &partition_client,
149                &self.data.topic,
150                self.data.prunable_entry_id,
151            )
152            .await?;
153        }
154
155        // Update the pruned entry id for the topic.
156        update_pruned_entry_id(
157            &self.context.table_metadata_manager,
158            &self.data.topic,
159            self.data.prunable_entry_id,
160        )
161        .await
162        .map_err(BoxedError::new)
163        .with_context(|_| error::RetryLaterWithSourceSnafu {
164            reason: format!(
165                "Failed to update pruned entry id for topic: {}",
166                self.data.topic
167            ),
168        })?;
169
170        info!(
171            "Successfully pruned WAL for topic: {}, prunable entry id: {}, latest offset: {}",
172            self.data.topic, self.data.prunable_entry_id, latest_offset
173        );
174        Ok(Status::done_with_output(WalPruneOutcome))
175    }
176}
177
178#[async_trait::async_trait]
179impl Procedure for WalPruneProcedure {
180    fn type_name(&self) -> &str {
181        Self::TYPE_NAME
182    }
183
184    fn rollback_supported(&self) -> bool {
185        false
186    }
187
188    async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
189        let _guard = ctx
190            .provider
191            .acquire_lock(&(RemoteWalLock::Write(self.data.topic.clone()).into()))
192            .await;
193
194        self.on_prune().await.map_err(|e| {
195            if e.is_retryable() {
196                ProcedureError::retry_later(e)
197            } else {
198                ProcedureError::external(e)
199            }
200        })
201    }
202
203    fn dump(&self) -> ProcedureResult<String> {
204        serde_json::to_string(&self.data).context(ToJsonSnafu)
205    }
206
207    /// WAL prune procedure will read the topic-region map from the table metadata manager,
208    /// which are modified by `DROP [TABLE|DATABASE]` and `CREATE [TABLE]` operations.
209    /// But the modifications are atomic, so it does not conflict with the procedure.
210    /// It only abort the procedure sometimes since the `check_heartbeat_collected_region_ids` fails.
211    fn lock_key(&self) -> LockKey {
212        let lock_key: StringKey = RemoteWalLock::Write(self.data.topic.clone()).into();
213        LockKey::new(vec![lock_key])
214    }
215
216    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
217        if !ctx.event_type_filter.allows(WAL_PRUNE_EVENT_TYPE) {
218            return None;
219        }
220
221        if matches!(&ctx.trigger, EventTrigger::Succeeded) {
222            let ProcedureState::Done {
223                output: Some(output),
224            } = ctx.lifecycle_state
225            else {
226                return None;
227            };
228            output.downcast_ref::<WalPruneOutcome>()?;
229        // `Submitted` and `Recovered` are intentionally omitted. `RollingBack` and
230        // `ChildSubmitted` cannot occur because WAL pruning neither supports rollback nor submits
231        // child procedures.
232        } else if !matches!(
233            &ctx.trigger,
234            EventTrigger::Retrying { .. } | EventTrigger::Failed | EventTrigger::Poisoned
235        ) {
236            return None;
237        }
238
239        Some(Box::new(WalPruneEvent::new(
240            &self.data.topic,
241            self.data.prunable_entry_id,
242            self.observed_latest_offset,
243            self.data.logical_delete,
244        )))
245    }
246}
247
248#[cfg(test)]
249mod tests {
250    use std::assert_matches;
251    use std::collections::HashSet;
252
253    use common_event_recorder::EventTypeFilter;
254    use common_procedure::{Output, ProcedureId};
255    use common_wal::maybe_skip_kafka_integration_test;
256    use common_wal::test_util::get_kafka_endpoints;
257    use rskafka::client::partition::{FetchResult, UnknownTopicHandling};
258    use rskafka::record::Record;
259
260    use super::*;
261    use crate::procedure::test_util::new_wal_prune_metadata;
262    // Fix this import to correctly point to the test_util module
263    use crate::procedure::wal_prune::test_util::TestEnv;
264
265    /// Mock a test env for testing.
266    /// Including:
267    /// 1. Prepare some data in the table metadata manager and in-memory kv backend.
268    /// 2. Return the procedure, the minimum last entry id to prune and the regions to flush.
269    async fn mock_test_data(context: Context, topic: &str) -> u64 {
270        let n_region = 10;
271        let n_table = 5;
272        // 5 entries per region.
273        let offsets = mock_wal_entries(
274            context.client.clone(),
275            topic,
276            (n_region * n_table * 5) as usize,
277        )
278        .await;
279
280        new_wal_prune_metadata(
281            context.table_metadata_manager.clone(),
282            context.leader_region_registry.clone(),
283            n_region,
284            n_table,
285            &offsets[1..],
286            topic.to_string(),
287        )
288        .await
289    }
290
291    fn record(i: usize) -> Record {
292        let key = format!("key_{i}");
293        let value = format!("value_{i}");
294        Record {
295            key: Some(key.into()),
296            value: Some(value.into()),
297            timestamp: chrono::Utc::now(),
298            headers: Default::default(),
299        }
300    }
301
302    async fn mock_wal_entries(
303        client: KafkaClientRef,
304        topic_name: &str,
305        n_entries: usize,
306    ) -> Vec<i64> {
307        let controller_client = client.controller_client().unwrap();
308        let _ = controller_client
309            .create_topic(topic_name, 1, 1, 5_000)
310            .await;
311        let partition_client = client
312            .partition_client(topic_name, 0, UnknownTopicHandling::Retry)
313            .await
314            .unwrap();
315        let mut offsets = Vec::with_capacity(n_entries);
316        for i in 0..n_entries {
317            let record = vec![record(i)];
318            let offset = partition_client
319                .produce(
320                    record,
321                    rskafka::client::partition::Compression::NoCompression,
322                )
323                .await
324                .unwrap()
325                .offsets;
326            offsets.extend(offset);
327        }
328        offsets
329    }
330
331    async fn check_entry_id_existence(
332        client: KafkaClientRef,
333        topic_name: &str,
334        entry_id: i64,
335        expect_success: bool,
336    ) {
337        let partition_client = client
338            .partition_client(topic_name, 0, UnknownTopicHandling::Retry)
339            .await
340            .unwrap();
341        let res = partition_client
342            .fetch_records(entry_id, 0..10001, 5_000)
343            .await;
344        if expect_success {
345            assert!(res.is_ok());
346            let FetchResult { records, .. } = res.unwrap();
347            assert!(!records.is_empty());
348        } else {
349            let err = res.unwrap_err();
350            // The error is in a private module so we check it through `to_string()`.
351            assert!(err.to_string().contains("OffsetOutOfRange"));
352        }
353    }
354
355    async fn delete_topic(client: KafkaClientRef, topic_name: &str) {
356        let controller_client = client.controller_client().unwrap();
357        controller_client
358            .delete_topic(topic_name, 5_000)
359            .await
360            .unwrap();
361    }
362
363    #[test]
364    fn test_wal_prune_data_backward_compatibility() {
365        let data: WalPruneData =
366            serde_json::from_str(r#"{"topic":"test_topic","prunable_entry_id":42}"#).unwrap();
367
368        assert_eq!(data.topic, "test_topic");
369        assert_eq!(data.prunable_entry_id, 42);
370        assert!(!data.logical_delete);
371    }
372
373    #[tokio::test]
374    async fn test_procedure_execution() {
375        maybe_skip_kafka_integration_test!();
376        let broker_endpoints = get_kafka_endpoints();
377
378        common_telemetry::init_default_ut_logging();
379        let mut topic_name = uuid::Uuid::new_v4().to_string();
380        // Topic should start with a letter.
381        topic_name = format!("test_procedure_execution-{}", topic_name);
382        let env = TestEnv::new();
383        let context = env.build_wal_prune_context(broker_endpoints).await;
384        // Prepare the topic.
385        TestEnv::prepare_topic(&context.client, &topic_name).await;
386
387        // Mock the test data.
388        let prunable_entry_id = mock_test_data(context.clone(), &topic_name).await;
389        let mut procedure = WalPruneProcedure::new(
390            context.clone(),
391            None,
392            topic_name.clone(),
393            prunable_entry_id,
394            false,
395        );
396        let status = procedure.on_prune().await.unwrap();
397        assert_eq!(
398            status.downcast_output_ref::<WalPruneOutcome>(),
399            Some(&WalPruneOutcome)
400        );
401        // Check if the entry ids after(include) `prunable_entry_id` still exist.
402        check_entry_id_existence(
403            procedure.context.client.clone(),
404            &topic_name,
405            procedure.data.prunable_entry_id as i64,
406            true,
407        )
408        .await;
409        // Check if the entry ids before `prunable_entry_id` are deleted.
410        check_entry_id_existence(
411            procedure.context.client.clone(),
412            &topic_name,
413            procedure.data.prunable_entry_id as i64 - 1,
414            false,
415        )
416        .await;
417
418        let value = env
419            .table_metadata_manager
420            .topic_name_manager()
421            .get(&topic_name)
422            .await
423            .unwrap()
424            .unwrap();
425        assert_eq!(value.pruned_entry_id, procedure.data.prunable_entry_id);
426        // Clean up the topic.
427        delete_topic(procedure.context.client, &topic_name).await;
428    }
429
430    #[tokio::test]
431    async fn test_procedure_execution_with_logical_delete() {
432        maybe_skip_kafka_integration_test!();
433        let broker_endpoints = get_kafka_endpoints();
434
435        common_telemetry::init_default_ut_logging();
436        let mut topic_name = uuid::Uuid::new_v4().to_string();
437        // Topic should start with a letter.
438        topic_name = format!("test_procedure_execution_with_logical_delete-{topic_name}");
439        let env = TestEnv::new();
440        let context = env.build_wal_prune_context(broker_endpoints).await;
441        // Prepare the topic.
442        TestEnv::prepare_topic(&context.client, &topic_name).await;
443
444        // Mock the test data.
445        let prunable_entry_id = mock_test_data(context.clone(), &topic_name).await;
446        let mut procedure = WalPruneProcedure::new(
447            context.clone(),
448            None,
449            topic_name.clone(),
450            prunable_entry_id,
451            true,
452        );
453        let status = procedure.on_prune().await.unwrap();
454        assert_eq!(
455            status.downcast_output_ref::<WalPruneOutcome>(),
456            Some(&WalPruneOutcome)
457        );
458        // Logical delete should keep the entry ids before `prunable_entry_id`.
459        check_entry_id_existence(
460            procedure.context.client.clone(),
461            &topic_name,
462            procedure.data.prunable_entry_id as i64 - 1,
463            true,
464        )
465        .await;
466
467        let value = env
468            .table_metadata_manager
469            .topic_name_manager()
470            .get(&topic_name)
471            .await
472            .unwrap()
473            .unwrap();
474        assert_eq!(value.pruned_entry_id, procedure.data.prunable_entry_id);
475        // Clean up the topic.
476        delete_topic(procedure.context.client, &topic_name).await;
477    }
478
479    #[tokio::test]
480    async fn test_procedure_noop_has_no_output() {
481        maybe_skip_kafka_integration_test!();
482        let broker_endpoints = get_kafka_endpoints();
483        let topic_name = format!("test_procedure_noop_has_no_output-{}", uuid::Uuid::new_v4());
484        let env = TestEnv::new();
485        let context = env.build_wal_prune_context(broker_endpoints).await;
486        TestEnv::prepare_topic(&context.client, &topic_name).await;
487        let mut procedure =
488            WalPruneProcedure::new(context.clone(), None, topic_name.clone(), 0, false);
489
490        let status = procedure.on_prune().await.unwrap();
491
492        assert_matches!(status, Status::Done { output: None });
493        assert_eq!(procedure.observed_latest_offset, Some(0));
494        delete_topic(context.client, &topic_name).await;
495    }
496
497    #[tokio::test]
498    async fn test_wal_prune_event_trigger_selection() {
499        maybe_skip_kafka_integration_test!();
500        let context = TestEnv::new()
501            .build_wal_prune_context(get_kafka_endpoints())
502            .await;
503        let mut procedure =
504            WalPruneProcedure::new(context, None, "test_topic".to_string(), 42, false);
505        let running = ProcedureState::Running;
506        let runtime_context = |trigger, lifecycle_state, event_type_filter| EventContext {
507            procedure_id: ProcedureId::random(),
508            lifecycle_state,
509            trigger,
510            event_type_filter: Arc::new(event_type_filter),
511            event_context: None,
512        };
513
514        for trigger in [EventTrigger::Submitted, EventTrigger::Recovered] {
515            assert!(
516                procedure
517                    .event(&runtime_context(trigger, &running, EventTypeFilter::All))
518                    .is_none()
519            );
520        }
521
522        let retrying_event = procedure
523            .event(&runtime_context(
524                EventTrigger::Retrying {
525                    phase: common_procedure::RetryPhase::Execute,
526                    attempt: 1,
527                },
528                &running,
529                EventTypeFilter::All,
530            ))
531            .unwrap();
532        assert!(
533            retrying_event.extra_rows().unwrap()[0].values[2]
534                .value_data
535                .is_none()
536        );
537
538        procedure.observed_latest_offset = Some(100);
539        for trigger in [EventTrigger::Failed, EventTrigger::Poisoned] {
540            let event = procedure
541                .event(&runtime_context(trigger, &running, EventTypeFilter::All))
542                .unwrap();
543            assert_eq!(
544                event.extra_rows().unwrap()[0].values[2].value_data,
545                Some(api::v1::value::ValueData::U64Value(100))
546            );
547        }
548
549        let done_without_output = ProcedureState::Done { output: None };
550        assert!(
551            procedure
552                .event(&runtime_context(
553                    EventTrigger::Succeeded,
554                    &done_without_output,
555                    EventTypeFilter::All,
556                ))
557                .is_none()
558        );
559        let done_with_wrong_output = ProcedureState::Done {
560            output: Some(Arc::new(42_u64) as Output),
561        };
562        assert!(
563            procedure
564                .event(&runtime_context(
565                    EventTrigger::Succeeded,
566                    &done_with_wrong_output,
567                    EventTypeFilter::All,
568                ))
569                .is_none()
570        );
571
572        let done = ProcedureState::Done {
573            output: Some(Arc::new(WalPruneOutcome) as Output),
574        };
575        let event = procedure
576            .event(&runtime_context(
577                EventTrigger::Succeeded,
578                &done,
579                EventTypeFilter::All,
580            ))
581            .unwrap();
582        assert_eq!(event.event_type(), WAL_PRUNE_EVENT_TYPE);
583        assert_eq!(
584            event.json_payload().unwrap(),
585            serde_json::json!({
586                "version": 1,
587                "logical_delete": false,
588            })
589        );
590
591        assert!(
592            procedure
593                .event(&runtime_context(
594                    EventTrigger::Succeeded,
595                    &done,
596                    EventTypeFilter::Only(HashSet::new()),
597                ))
598                .is_none()
599        );
600    }
601}