Skip to main content

meta_srv/event/
gc.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::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
31/// Procedure event type for batch garbage collection.
32pub(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    // None emits a procedure-level lifecycle row with null Region dimensions.
70    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    /// Returns None when the GC report contains no deleted files, indexes, or retries.
91    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(&region.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        &REGION_ID_COLUMN,
162        &TABLE_ID_COLUMN,
163        &REGION_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(&region_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(&region_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(&region_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}