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::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}