1use std::any::Any;
16use std::collections::BTreeSet;
17use std::time::Duration;
18
19use api::v1::value::ValueData;
20use api::v1::{ColumnSchema, Row, Value};
21use common_event_recorder::Event;
22use common_event_recorder::error::{Result, SerializeEventSnafu};
23use common_event_recorder::event_table::{
24 GC_REPORT_COLUMN, REGION_ID_COLUMN, REGION_NUMBER_COLUMN, TABLE_ID_COLUMN, column_schemas,
25 nullable_json,
26};
27use serde::Serialize;
28use snafu::ResultExt;
29use store_api::storage::{GcReport, IndexVersion, RegionId};
30
31pub(crate) const BATCH_GC_EVENT_TYPE: &str = "batch_gc";
33
34const PAYLOAD_VERSION: u8 = 1;
35
36#[derive(Debug, Serialize)]
37struct BatchGcPayload {
38 version: u8,
39 regions: Vec<RegionId>,
40 full_file_listing: bool,
41 #[serde(with = "humantime_serde")]
42 timeout: Duration,
43}
44
45#[derive(Debug, Serialize)]
46struct DeletedIndexPayload {
47 file_id: String,
48 index_version: IndexVersion,
49}
50
51#[derive(Debug, Serialize)]
52struct BatchGcRegionReport {
53 #[serde(skip_serializing_if = "Vec::is_empty")]
54 deleted_files: Vec<String>,
55 #[serde(skip_serializing_if = "Vec::is_empty")]
56 deleted_indexes: Vec<DeletedIndexPayload>,
57 need_retry: bool,
58}
59
60#[derive(Debug)]
61struct BatchGcRegionRow {
62 region_id: RegionId,
63 report: BatchGcRegionReport,
64}
65
66#[derive(Debug)]
67pub(crate) struct BatchGcEvent {
68 payload: Option<BatchGcPayload>,
69 regions: Option<Vec<BatchGcRegionRow>>,
71}
72
73impl BatchGcEvent {
74 pub(crate) fn with_config(
75 regions: &[RegionId],
76 full_file_listing: bool,
77 timeout: Duration,
78 ) -> Self {
79 Self {
80 payload: Some(BatchGcPayload {
81 version: PAYLOAD_VERSION,
82 regions: regions.to_vec(),
83 full_file_listing,
84 timeout,
85 }),
86 regions: None,
87 }
88 }
89
90 pub(crate) fn with_report(report: &GcReport) -> Option<Self> {
92 let mut region_ids = report
93 .deleted_files
94 .keys()
95 .copied()
96 .collect::<BTreeSet<_>>();
97 region_ids.extend(report.deleted_indexes.keys().copied());
98 region_ids.extend(report.need_retry_regions.iter().copied());
99
100 let regions: Vec<_> = region_ids
101 .into_iter()
102 .filter_map(|region_id| {
103 region_report(region_id, report)
104 .map(|report| BatchGcRegionRow { region_id, report })
105 })
106 .collect();
107
108 (!regions.is_empty()).then_some(Self {
109 payload: None,
110 regions: Some(regions),
111 })
112 }
113}
114
115impl Event for BatchGcEvent {
116 fn event_type(&self) -> &str {
117 BATCH_GC_EVENT_TYPE
118 }
119
120 fn json_payload(&self) -> Result<serde_json::Value> {
121 self.payload
122 .as_ref()
123 .map(serde_json::to_value)
124 .transpose()
125 .context(SerializeEventSnafu)
126 .map(|payload| payload.unwrap_or(serde_json::Value::Null))
127 }
128
129 fn extra_schema(&self) -> Vec<ColumnSchema> {
130 schema()
131 }
132
133 fn extra_rows(&self) -> Result<Vec<Row>> {
134 let Some(regions) = &self.regions else {
135 return Ok(vec![null_row()]);
136 };
137
138 regions
139 .iter()
140 .map(|region| {
141 let report = serde_json::to_value(®ion.report).context(SerializeEventSnafu)?;
142 Ok(Row {
143 values: vec![
144 ValueData::U64Value(region.region_id.as_u64()).into(),
145 ValueData::U32Value(region.region_id.table_id()).into(),
146 ValueData::U32Value(region.region_id.region_number()).into(),
147 nullable_json(Some(&report)),
148 ],
149 })
150 })
151 .collect()
152 }
153
154 fn as_any(&self) -> &dyn Any {
155 self
156 }
157}
158
159fn schema() -> Vec<ColumnSchema> {
160 column_schemas([
161 ®ION_ID_COLUMN,
162 &TABLE_ID_COLUMN,
163 ®ION_NUMBER_COLUMN,
164 &GC_REPORT_COLUMN,
165 ])
166}
167
168fn null_row() -> Row {
169 Row {
170 values: (0..schema().len())
171 .map(|_| Value { value_data: None })
172 .collect(),
173 }
174}
175
176fn region_report(region_id: RegionId, report: &GcReport) -> Option<BatchGcRegionReport> {
177 let deleted_files = report
178 .deleted_files
179 .get(®ion_id)
180 .into_iter()
181 .flatten()
182 .map(ToString::to_string)
183 .collect::<Vec<_>>();
184
185 let deleted_indexes = report
186 .deleted_indexes
187 .get(®ion_id)
188 .into_iter()
189 .flatten()
190 .map(|(file_id, index_version)| DeletedIndexPayload {
191 file_id: file_id.to_string(),
192 index_version: *index_version,
193 })
194 .collect::<Vec<_>>();
195
196 let need_retry = report.need_retry_regions.contains(®ion_id);
197 if deleted_files.is_empty() && deleted_indexes.is_empty() && !need_retry {
198 return None;
199 }
200
201 Some(BatchGcRegionReport {
202 deleted_files,
203 deleted_indexes,
204 need_retry,
205 })
206}
207
208#[cfg(test)]
209mod tests {
210 use std::collections::{HashMap, HashSet};
211 use std::sync::Arc;
212
213 use api::v1::ColumnSchema;
214 use common_event_recorder::event_table::{
215 ACTOR_COLUMN, EVENT_CONTEXT_COLUMN, PROCEDURE_ERROR_COLUMN, PROCEDURE_ID_COLUMN,
216 PROCEDURE_STATE_COLUMN, PROCEDURE_TRIGGER_COLUMN, jsonb_value,
217 };
218 use common_event_recorder::testing::assert_event_contract;
219 use common_event_recorder::{EventTypeFilter, PersistentEventContext, TriggerReason};
220 use common_meta::key::TableMetadataManager;
221 use common_meta::key::runtime_switch::RuntimeSwitchManager;
222 use common_meta::kv_backend::memory::MemoryKvBackend;
223 use common_meta::sequence::SequenceBuilder;
224 use common_procedure::{
225 EventContext, EventTrigger, Procedure, ProcedureEvent, ProcedureId, ProcedureState,
226 RetryPhase,
227 };
228 use store_api::storage::FileId;
229
230 use super::*;
231 use crate::gc::BatchGcProcedure;
232 use crate::procedure::test_util::MailboxContext;
233
234 #[test]
235 fn test_batch_gc_lifecycle_event_contract() {
236 let first = RegionId::new(1024, 1);
237 let second = RegionId::new(1024, 2);
238 let event =
239 BatchGcEvent::with_config(&[second, first, second], true, Duration::from_secs(10));
240
241 assert_event_contract(&event, BATCH_GC_EVENT_TYPE, &schema(), &[null_row()]);
242 assert_eq!(
243 event.json_payload().unwrap(),
244 serde_json::json!({
245 "version": 1,
246 "regions": [second.as_u64(), first.as_u64(), second.as_u64()],
247 "full_file_listing": true,
248 "timeout": "10s",
249 })
250 );
251 }
252
253 #[test]
254 fn test_batch_gc_report_event_contract() {
255 let first = RegionId::new(1024, 1);
256 let second = RegionId::new(1024, 2);
257 let third = RegionId::new(1024, 3);
258 let first_file = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
259 let second_file = FileId::parse_str("00000000-0000-0000-0000-000000000002").unwrap();
260 let index_file = FileId::parse_str("00000000-0000-0000-0000-000000000003").unwrap();
261 let report = GcReport {
262 deleted_files: HashMap::from([(first, vec![second_file, first_file])]),
263 deleted_indexes: HashMap::from([(second, vec![(index_file, 42)])]),
264 processed_regions: HashSet::from([first, second]),
265 need_retry_regions: HashSet::from([third]),
266 };
267
268 let event = BatchGcEvent::with_report(&report).unwrap();
269
270 assert_event_contract(
271 &event,
272 BATCH_GC_EVENT_TYPE,
273 &schema(),
274 &[
275 region_row(
276 first,
277 Some(serde_json::json!({
278 "deleted_files": [
279 "00000000-0000-0000-0000-000000000002",
280 "00000000-0000-0000-0000-000000000001",
281 ],
282 "need_retry": false,
283 })),
284 ),
285 region_row(
286 second,
287 Some(serde_json::json!({
288 "deleted_indexes": [{
289 "file_id": "00000000-0000-0000-0000-000000000003",
290 "index_version": 42,
291 }],
292 "need_retry": false,
293 })),
294 ),
295 region_row(
296 third,
297 Some(serde_json::json!({
298 "need_retry": true,
299 })),
300 ),
301 ],
302 );
303 assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null);
304 }
305
306 #[test]
307 fn test_batch_gc_report_event_skips_noop_regions() {
308 let noop = RegionId::new(1024, 1);
309 let retry = RegionId::new(1024, 2);
310 let event = BatchGcEvent::with_report(&GcReport {
311 processed_regions: HashSet::from([noop]),
312 need_retry_regions: HashSet::from([retry]),
313 ..Default::default()
314 })
315 .unwrap();
316
317 assert_event_contract(
318 &event,
319 BATCH_GC_EVENT_TYPE,
320 &schema(),
321 &[region_row(
322 retry,
323 Some(serde_json::json!({
324 "need_retry": true,
325 })),
326 )],
327 );
328 }
329
330 #[test]
331 fn test_batch_gc_event_filter() {
332 let procedure = batch_gc_procedure();
333 let running = ProcedureState::Running;
334 let manual_context = PersistentEventContext::new(TriggerReason::Manual);
335 let event_context = |trigger, lifecycle_state, event_type_filter| EventContext {
336 procedure_id: ProcedureId::random(),
337 lifecycle_state,
338 trigger,
339 event_type_filter: Arc::new(event_type_filter),
340 event_context: None,
341 };
342
343 assert!(
344 procedure
345 .event(&event_context(
346 EventTrigger::Submitted,
347 &running,
348 EventTypeFilter::All,
349 ))
350 .is_none()
351 );
352 assert!(
353 procedure
354 .event(&EventContext {
355 procedure_id: ProcedureId::random(),
356 lifecycle_state: &running,
357 trigger: EventTrigger::Submitted,
358 event_type_filter: Arc::new(EventTypeFilter::All),
359 event_context: Some(&manual_context),
360 })
361 .is_some()
362 );
363
364 let report = GcReport {
365 processed_regions: HashSet::from([RegionId::new(1024, 1)]),
366 ..Default::default()
367 };
368 let done = ProcedureState::Done {
369 output: Some(Arc::new(report)),
370 };
371 assert!(
372 procedure
373 .event(&event_context(
374 EventTrigger::Succeeded,
375 &done,
376 EventTypeFilter::All,
377 ))
378 .is_none()
379 );
380
381 assert!(
382 procedure
383 .event(&event_context(
384 EventTrigger::Recovered,
385 &running,
386 EventTypeFilter::All,
387 ))
388 .is_none()
389 );
390
391 let retrying = procedure
392 .event(&event_context(
393 EventTrigger::Retrying {
394 phase: RetryPhase::Execute,
395 attempt: 1,
396 },
397 &running,
398 EventTypeFilter::All,
399 ))
400 .unwrap();
401 assert_eq!(
402 retrying.json_payload().unwrap(),
403 serde_json::json!({
404 "version": 1,
405 "regions": [RegionId::new(1024, 1).as_u64()],
406 "full_file_listing": true,
407 "timeout": "10s",
408 })
409 );
410 assert_eq!(retrying.extra_rows().unwrap(), vec![null_row()]);
411
412 assert!(
413 procedure
414 .event(&event_context(
415 EventTrigger::Recovered,
416 &running,
417 EventTypeFilter::Only(HashSet::from(["another_event".to_string()])),
418 ))
419 .is_none()
420 );
421 assert!(
422 procedure
423 .event(&event_context(
424 EventTrigger::Recovered,
425 &running,
426 EventTypeFilter::Only(HashSet::new()),
427 ))
428 .is_none()
429 );
430
431 let missing = ProcedureState::Done { output: None };
432 assert!(
433 procedure
434 .event(&event_context(
435 EventTrigger::Succeeded,
436 &missing,
437 EventTypeFilter::All,
438 ))
439 .is_none()
440 );
441 let wrong = ProcedureState::Done {
442 output: Some(Arc::new("not a GC report".to_string())),
443 };
444 assert!(
445 procedure
446 .event(&event_context(
447 EventTrigger::Succeeded,
448 &wrong,
449 EventTypeFilter::All,
450 ))
451 .is_none()
452 );
453 }
454
455 #[test]
456 fn test_batch_gc_terminal_event_with_report() {
457 let mut procedure = batch_gc_procedure();
458 let running = ProcedureState::Running;
459 let event_context = |trigger| EventContext {
460 procedure_id: ProcedureId::random(),
461 lifecycle_state: &running,
462 trigger,
463 event_type_filter: Arc::new(EventTypeFilter::All),
464 event_context: None,
465 };
466 let report = GcReport {
467 need_retry_regions: HashSet::from([RegionId::new(1024, 2)]),
468 ..Default::default()
469 };
470 procedure.set_gc_report_for_test(report);
471
472 let retrying = procedure
473 .event(&event_context(EventTrigger::Retrying {
474 phase: RetryPhase::Execute,
475 attempt: 2,
476 }))
477 .unwrap();
478 assert_eq!(
479 retrying.json_payload().unwrap(),
480 serde_json::json!({
481 "version": 1,
482 "regions": [RegionId::new(1024, 1).as_u64()],
483 "full_file_listing": true,
484 "timeout": "10s",
485 })
486 );
487 assert_eq!(retrying.extra_rows().unwrap(), vec![null_row()]);
488
489 for trigger in [EventTrigger::Failed, EventTrigger::Poisoned] {
490 let event = procedure.event(&event_context(trigger)).unwrap();
491 assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null);
492 assert_eq!(
493 event.extra_rows().unwrap(),
494 vec![region_row(
495 RegionId::new(1024, 2),
496 Some(serde_json::json!({
497 "need_retry": true,
498 })),
499 )],
500 );
501 }
502
503 let procedure = batch_gc_procedure();
504 let failed_without_report = procedure
505 .event(&event_context(EventTrigger::Failed))
506 .unwrap();
507 assert_eq!(
508 failed_without_report.json_payload().unwrap(),
509 serde_json::json!({
510 "version": 1,
511 "regions": [RegionId::new(1024, 1).as_u64()],
512 "full_file_listing": true,
513 "timeout": "10s",
514 })
515 );
516 assert_eq!(
517 failed_without_report.extra_rows().unwrap(),
518 vec![null_row()]
519 );
520 }
521
522 #[test]
523 fn test_batch_gc_procedure_event_contract() {
524 let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
525 let region_id = RegionId::new(1024, 1);
526 let file_id = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
527 let internal = BatchGcEvent::with_report(&GcReport {
528 deleted_files: HashMap::from([(region_id, vec![file_id])]),
529 ..Default::default()
530 })
531 .unwrap();
532 let event = ProcedureEvent::new(
533 procedure_id,
534 Box::new(internal),
535 ProcedureState::Done { output: None },
536 EventTrigger::Succeeded,
537 );
538 let mut event_schema = procedure_schema();
539 event_schema.extend(schema());
540 event_schema.push(ACTOR_COLUMN.column_schema());
541 event_schema.push(EVENT_CONTEXT_COLUMN.column_schema());
542 let mut values = vec![
543 ValueData::StringValue(procedure_id.to_string()).into(),
544 ValueData::StringValue("Done".to_string()).into(),
545 ValueData::StringValue(String::new()).into(),
546 jsonb_value(&serde_json::json!({"type": "Succeeded"})),
547 ];
548 values.extend(
549 region_row(
550 region_id,
551 Some(serde_json::json!({
552 "deleted_files": ["00000000-0000-0000-0000-000000000001"],
553 "need_retry": false,
554 })),
555 )
556 .values,
557 );
558 values.push(Value { value_data: None });
559 values.push(Value { value_data: None });
560
561 assert_event_contract(
562 &event,
563 BATCH_GC_EVENT_TYPE,
564 &event_schema,
565 &[Row { values }],
566 );
567 }
568
569 fn procedure_schema() -> Vec<ColumnSchema> {
570 column_schemas([
571 &PROCEDURE_ID_COLUMN,
572 &PROCEDURE_STATE_COLUMN,
573 &PROCEDURE_ERROR_COLUMN,
574 &PROCEDURE_TRIGGER_COLUMN,
575 ])
576 }
577
578 fn region_row(region_id: RegionId, report: Option<serde_json::Value>) -> Row {
579 Row {
580 values: vec![
581 ValueData::U64Value(region_id.as_u64()).into(),
582 ValueData::U32Value(region_id.table_id()).into(),
583 ValueData::U32Value(region_id.region_number()).into(),
584 nullable_json(report.as_ref()),
585 ],
586 }
587 }
588
589 fn batch_gc_procedure() -> BatchGcProcedure {
590 let kv_backend = Arc::new(MemoryKvBackend::new());
591 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
592 let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend.clone()));
593 let mailbox_sequence = SequenceBuilder::new("test_batch_gc_event", kv_backend).build();
594 let mailbox = MailboxContext::new(mailbox_sequence);
595 BatchGcProcedure::new(
596 mailbox.mailbox().clone(),
597 table_metadata_manager,
598 runtime_switch_manager,
599 "localhost".to_string(),
600 vec![RegionId::new(1024, 1)],
601 true,
602 Duration::from_secs(10),
603 HashMap::new(),
604 )
605 }
606}