1pub(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 pub client: KafkaClientRef,
51 pub table_metadata_manager: TableMetadataManagerRef,
53 pub leader_region_registry: LeaderRegionRegistryRef,
55}
56
57#[derive(Serialize, Deserialize)]
59pub struct WalPruneData {
60 pub topic: String,
62 pub prunable_entry_id: EntryId,
64 #[serde(default)]
66 pub logical_delete: bool,
67}
68
69#[derive(Debug, Clone, Copy, PartialEq, Eq)]
70struct WalPruneOutcome;
71
72pub struct WalPruneProcedure {
74 pub data: WalPruneData,
75 pub context: Context,
76 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 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(
148 &partition_client,
149 &self.data.topic,
150 self.data.prunable_entry_id,
151 )
152 .await?;
153 }
154
155 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 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 } 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 use crate::procedure::wal_prune::test_util::TestEnv;
264
265 async fn mock_test_data(context: Context, topic: &str) -> u64 {
270 let n_region = 10;
271 let n_table = 5;
272 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 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_name = format!("test_procedure_execution-{}", topic_name);
382 let env = TestEnv::new();
383 let context = env.build_wal_prune_context(broker_endpoints).await;
384 TestEnv::prepare_topic(&context.client, &topic_name).await;
386
387 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_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_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 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_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 TestEnv::prepare_topic(&context.client, &topic_name).await;
443
444 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 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 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}