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::kv_backend::memory::MemoryKvBackend;
222 use common_meta::sequence::SequenceBuilder;
223 use common_procedure::{
224 EventContext, EventTrigger, Procedure, ProcedureEvent, ProcedureId, ProcedureState,
225 RetryPhase,
226 };
227 use store_api::storage::FileId;
228
229 use super::*;
230 use crate::gc::BatchGcProcedure;
231 use crate::procedure::test_util::MailboxContext;
232
233 #[test]
234 fn test_batch_gc_lifecycle_event_contract() {
235 let first = RegionId::new(1024, 1);
236 let second = RegionId::new(1024, 2);
237 let event =
238 BatchGcEvent::with_config(&[second, first, second], true, Duration::from_secs(10));
239
240 assert_event_contract(&event, BATCH_GC_EVENT_TYPE, &schema(), &[null_row()]);
241 assert_eq!(
242 event.json_payload().unwrap(),
243 serde_json::json!({
244 "version": 1,
245 "regions": [second.as_u64(), first.as_u64(), second.as_u64()],
246 "full_file_listing": true,
247 "timeout": "10s",
248 })
249 );
250 }
251
252 #[test]
253 fn test_batch_gc_report_event_contract() {
254 let first = RegionId::new(1024, 1);
255 let second = RegionId::new(1024, 2);
256 let third = RegionId::new(1024, 3);
257 let first_file = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
258 let second_file = FileId::parse_str("00000000-0000-0000-0000-000000000002").unwrap();
259 let index_file = FileId::parse_str("00000000-0000-0000-0000-000000000003").unwrap();
260 let report = GcReport {
261 deleted_files: HashMap::from([(first, vec![second_file, first_file])]),
262 deleted_indexes: HashMap::from([(second, vec![(index_file, 42)])]),
263 processed_regions: HashSet::from([first, second]),
264 need_retry_regions: HashSet::from([third]),
265 };
266
267 let event = BatchGcEvent::with_report(&report).unwrap();
268
269 assert_event_contract(
270 &event,
271 BATCH_GC_EVENT_TYPE,
272 &schema(),
273 &[
274 region_row(
275 first,
276 Some(serde_json::json!({
277 "deleted_files": [
278 "00000000-0000-0000-0000-000000000002",
279 "00000000-0000-0000-0000-000000000001",
280 ],
281 "need_retry": false,
282 })),
283 ),
284 region_row(
285 second,
286 Some(serde_json::json!({
287 "deleted_indexes": [{
288 "file_id": "00000000-0000-0000-0000-000000000003",
289 "index_version": 42,
290 }],
291 "need_retry": false,
292 })),
293 ),
294 region_row(
295 third,
296 Some(serde_json::json!({
297 "need_retry": true,
298 })),
299 ),
300 ],
301 );
302 assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null);
303 }
304
305 #[test]
306 fn test_batch_gc_report_event_skips_noop_regions() {
307 let noop = RegionId::new(1024, 1);
308 let retry = RegionId::new(1024, 2);
309 let event = BatchGcEvent::with_report(&GcReport {
310 processed_regions: HashSet::from([noop]),
311 need_retry_regions: HashSet::from([retry]),
312 ..Default::default()
313 })
314 .unwrap();
315
316 assert_event_contract(
317 &event,
318 BATCH_GC_EVENT_TYPE,
319 &schema(),
320 &[region_row(
321 retry,
322 Some(serde_json::json!({
323 "need_retry": true,
324 })),
325 )],
326 );
327 }
328
329 #[test]
330 fn test_batch_gc_event_filter() {
331 let procedure = batch_gc_procedure();
332 let running = ProcedureState::Running;
333 let manual_context = PersistentEventContext::new(TriggerReason::Manual);
334 let event_context = |trigger, lifecycle_state, event_type_filter| EventContext {
335 procedure_id: ProcedureId::random(),
336 lifecycle_state,
337 trigger,
338 event_type_filter: Arc::new(event_type_filter),
339 event_context: None,
340 };
341
342 assert!(
343 procedure
344 .event(&event_context(
345 EventTrigger::Submitted,
346 &running,
347 EventTypeFilter::All,
348 ))
349 .is_none()
350 );
351 assert!(
352 procedure
353 .event(&EventContext {
354 procedure_id: ProcedureId::random(),
355 lifecycle_state: &running,
356 trigger: EventTrigger::Submitted,
357 event_type_filter: Arc::new(EventTypeFilter::All),
358 event_context: Some(&manual_context),
359 })
360 .is_some()
361 );
362
363 let report = GcReport {
364 processed_regions: HashSet::from([RegionId::new(1024, 1)]),
365 ..Default::default()
366 };
367 let done = ProcedureState::Done {
368 output: Some(Arc::new(report)),
369 };
370 assert!(
371 procedure
372 .event(&event_context(
373 EventTrigger::Succeeded,
374 &done,
375 EventTypeFilter::All,
376 ))
377 .is_none()
378 );
379
380 assert!(
381 procedure
382 .event(&event_context(
383 EventTrigger::Recovered,
384 &running,
385 EventTypeFilter::All,
386 ))
387 .is_none()
388 );
389
390 let retrying = procedure
391 .event(&event_context(
392 EventTrigger::Retrying {
393 phase: RetryPhase::Execute,
394 attempt: 1,
395 },
396 &running,
397 EventTypeFilter::All,
398 ))
399 .unwrap();
400 assert_eq!(
401 retrying.json_payload().unwrap(),
402 serde_json::json!({
403 "version": 1,
404 "regions": [RegionId::new(1024, 1).as_u64()],
405 "full_file_listing": true,
406 "timeout": "10s",
407 })
408 );
409 assert_eq!(retrying.extra_rows().unwrap(), vec![null_row()]);
410
411 assert!(
412 procedure
413 .event(&event_context(
414 EventTrigger::Recovered,
415 &running,
416 EventTypeFilter::Only(HashSet::from(["another_event".to_string()])),
417 ))
418 .is_none()
419 );
420 assert!(
421 procedure
422 .event(&event_context(
423 EventTrigger::Recovered,
424 &running,
425 EventTypeFilter::Only(HashSet::new()),
426 ))
427 .is_none()
428 );
429
430 let missing = ProcedureState::Done { output: None };
431 assert!(
432 procedure
433 .event(&event_context(
434 EventTrigger::Succeeded,
435 &missing,
436 EventTypeFilter::All,
437 ))
438 .is_none()
439 );
440 let wrong = ProcedureState::Done {
441 output: Some(Arc::new("not a GC report".to_string())),
442 };
443 assert!(
444 procedure
445 .event(&event_context(
446 EventTrigger::Succeeded,
447 &wrong,
448 EventTypeFilter::All,
449 ))
450 .is_none()
451 );
452 }
453
454 #[test]
455 fn test_batch_gc_terminal_event_with_report() {
456 let mut procedure = batch_gc_procedure();
457 let running = ProcedureState::Running;
458 let event_context = |trigger| EventContext {
459 procedure_id: ProcedureId::random(),
460 lifecycle_state: &running,
461 trigger,
462 event_type_filter: Arc::new(EventTypeFilter::All),
463 event_context: None,
464 };
465 let report = GcReport {
466 need_retry_regions: HashSet::from([RegionId::new(1024, 2)]),
467 ..Default::default()
468 };
469 procedure.set_gc_report_for_test(report);
470
471 let retrying = procedure
472 .event(&event_context(EventTrigger::Retrying {
473 phase: RetryPhase::Execute,
474 attempt: 2,
475 }))
476 .unwrap();
477 assert_eq!(
478 retrying.json_payload().unwrap(),
479 serde_json::json!({
480 "version": 1,
481 "regions": [RegionId::new(1024, 1).as_u64()],
482 "full_file_listing": true,
483 "timeout": "10s",
484 })
485 );
486 assert_eq!(retrying.extra_rows().unwrap(), vec![null_row()]);
487
488 for trigger in [EventTrigger::Failed, EventTrigger::Poisoned] {
489 let event = procedure.event(&event_context(trigger)).unwrap();
490 assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null);
491 assert_eq!(
492 event.extra_rows().unwrap(),
493 vec![region_row(
494 RegionId::new(1024, 2),
495 Some(serde_json::json!({
496 "need_retry": true,
497 })),
498 )],
499 );
500 }
501
502 let procedure = batch_gc_procedure();
503 let failed_without_report = procedure
504 .event(&event_context(EventTrigger::Failed))
505 .unwrap();
506 assert_eq!(
507 failed_without_report.json_payload().unwrap(),
508 serde_json::json!({
509 "version": 1,
510 "regions": [RegionId::new(1024, 1).as_u64()],
511 "full_file_listing": true,
512 "timeout": "10s",
513 })
514 );
515 assert_eq!(
516 failed_without_report.extra_rows().unwrap(),
517 vec![null_row()]
518 );
519 }
520
521 #[test]
522 fn test_batch_gc_procedure_event_contract() {
523 let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
524 let region_id = RegionId::new(1024, 1);
525 let file_id = FileId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
526 let internal = BatchGcEvent::with_report(&GcReport {
527 deleted_files: HashMap::from([(region_id, vec![file_id])]),
528 ..Default::default()
529 })
530 .unwrap();
531 let event = ProcedureEvent::new(
532 procedure_id,
533 Box::new(internal),
534 ProcedureState::Done { output: None },
535 EventTrigger::Succeeded,
536 );
537 let mut event_schema = procedure_schema();
538 event_schema.extend(schema());
539 event_schema.push(ACTOR_COLUMN.column_schema());
540 event_schema.push(EVENT_CONTEXT_COLUMN.column_schema());
541 let mut values = vec![
542 ValueData::StringValue(procedure_id.to_string()).into(),
543 ValueData::StringValue("Done".to_string()).into(),
544 ValueData::StringValue(String::new()).into(),
545 jsonb_value(&serde_json::json!({"type": "Succeeded"})),
546 ];
547 values.extend(
548 region_row(
549 region_id,
550 Some(serde_json::json!({
551 "deleted_files": ["00000000-0000-0000-0000-000000000001"],
552 "need_retry": false,
553 })),
554 )
555 .values,
556 );
557 values.push(Value { value_data: None });
558 values.push(Value { value_data: None });
559
560 assert_event_contract(
561 &event,
562 BATCH_GC_EVENT_TYPE,
563 &event_schema,
564 &[Row { values }],
565 );
566 }
567
568 fn procedure_schema() -> Vec<ColumnSchema> {
569 column_schemas([
570 &PROCEDURE_ID_COLUMN,
571 &PROCEDURE_STATE_COLUMN,
572 &PROCEDURE_ERROR_COLUMN,
573 &PROCEDURE_TRIGGER_COLUMN,
574 ])
575 }
576
577 fn region_row(region_id: RegionId, report: Option<serde_json::Value>) -> Row {
578 Row {
579 values: vec![
580 ValueData::U64Value(region_id.as_u64()).into(),
581 ValueData::U32Value(region_id.table_id()).into(),
582 ValueData::U32Value(region_id.region_number()).into(),
583 nullable_json(report.as_ref()),
584 ],
585 }
586 }
587
588 fn batch_gc_procedure() -> BatchGcProcedure {
589 let kv_backend = Arc::new(MemoryKvBackend::new());
590 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
591 let mailbox_sequence = SequenceBuilder::new("test_batch_gc_event", kv_backend).build();
592 let mailbox = MailboxContext::new(mailbox_sequence);
593 BatchGcProcedure::new(
594 mailbox.mailbox().clone(),
595 table_metadata_manager,
596 "localhost".to_string(),
597 vec![RegionId::new(1024, 1)],
598 true,
599 Duration::from_secs(10),
600 HashMap::new(),
601 )
602 }
603}