Skip to main content

meta_srv/service/
procedure.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::time::Duration;
16
17use api::v1::meta::reconcile_request::Target;
18use api::v1::meta::{
19    DdlTaskRequest as PbDdlTaskRequest, DdlTaskResponse as PbDdlTaskResponse, GcRegionsRequest,
20    GcRegionsResponse, GcStats, GcTableRequest, GcTableResponse, MigrateRegionRequest,
21    MigrateRegionResponse, ProcedureActor, ProcedureDetailRequest, ProcedureDetailResponse,
22    ProcedureEventContext as PbProcedureEventContext, ProcedureStateResponse,
23    QueryProcedureRequest, ReconcileCatalog, ReconcileDatabase, ReconcileRequest,
24    ReconcileResponse, ReconcileTable, ResolveStrategy, Role, procedure_service_server,
25};
26use common_event_recorder::{PersistentEventContext, ProcedureEventInput};
27use common_meta::key::TableMetadataManagerRef;
28use common_meta::key::table_name::TableNameKey;
29use common_meta::peer::Peer;
30use common_meta::procedure_executor::ExecutorContext;
31use common_meta::rpc::ddl::{
32    CREATE_DATABASE_CREATOR_EXTENSION_KEY, CREATE_DATABASE_CREATOR_METADATA_KEY,
33    CreatorGrantIntent, DdlTask, QueryContext, SubmitDdlTaskRequest,
34};
35use common_meta::rpc::procedure::{
36    self, GcRegionsRequest as MetaGcRegionsRequest, GcResponse,
37    GcTableRequest as MetaGcTableRequest,
38};
39use common_procedure::ProcedureContext;
40use snafu::{OptionExt, ResultExt};
41use store_api::storage::RegionId;
42use table::table_reference::TableReference;
43use tonic::metadata::MetadataMap;
44use tonic::{Request, Status};
45
46use crate::error::{TableMetadataManagerSnafu, TableNotFoundSnafu};
47use crate::metasrv::Metasrv;
48use crate::procedure::region_migration::manager::{
49    RegionMigrationProcedureTask, RegionMigrationTriggerReason,
50};
51use crate::service::GrpcResult;
52use crate::{check_leader, error, gc};
53
54struct ProcedureSubmission {
55    actor: Option<String>,
56    event_context: Option<PbProcedureEventContext>,
57}
58
59impl From<(Option<ProcedureActor>, Option<PbProcedureEventContext>)> for ProcedureSubmission {
60    fn from(
61        (actor, event_context): (Option<ProcedureActor>, Option<PbProcedureEventContext>),
62    ) -> Self {
63        Self {
64            actor: actor.map(|actor| actor.username),
65            event_context,
66        }
67    }
68}
69
70impl From<ProcedureSubmission> for ProcedureContext {
71    fn from(context: ProcedureSubmission) -> Self {
72        Self {
73            actor: context.actor,
74            event_context: context.event_context.map(PersistentEventContext::from),
75        }
76    }
77}
78
79#[async_trait::async_trait]
80impl procedure_service_server::ProcedureService for Metasrv {
81    async fn query(
82        &self,
83        request: Request<QueryProcedureRequest>,
84    ) -> GrpcResult<ProcedureStateResponse> {
85        check_leader!(
86            self,
87            request,
88            ProcedureStateResponse,
89            "`query procedure state`"
90        );
91
92        let QueryProcedureRequest { header, pid, .. } = request.into_inner();
93        let _header = header.context(error::MissingRequestHeaderSnafu)?;
94        let pid = pid.context(error::MissingRequiredParameterSnafu { param: "pid" })?;
95        let pid = procedure::pb_pid_to_pid(&pid).context(error::ConvertProtoDataSnafu)?;
96
97        let state = self
98            .procedure_manager()
99            .procedure_state(pid)
100            .await
101            .context(error::QueryProcedureSnafu)?
102            .context(error::ProcedureNotFoundSnafu {
103                pid: pid.to_string(),
104            })?;
105
106        Ok(Response::new(procedure::procedure_state_to_pb_response(
107            &state,
108        )))
109    }
110
111    async fn ddl(&self, request: Request<PbDdlTaskRequest>) -> GrpcResult<PbDdlTaskResponse> {
112        check_leader!(self, request, PbDdlTaskResponse, "`ddl`");
113
114        let (metadata, _, request) = request.into_parts();
115        let PbDdlTaskRequest {
116            header,
117            query_context,
118            task,
119            wait,
120            timeout_secs,
121            event_context,
122            actor,
123        } = request;
124
125        let header = header.context(error::MissingRequestHeaderSnafu)?;
126        let ProcedureSubmission {
127            actor,
128            event_context,
129        } = ProcedureSubmission::from((actor, event_context));
130        let mut query_context = query_context
131            .context(error::MissingRequiredParameterSnafu {
132                param: "query_context",
133            })?
134            .into();
135        let mut task: DdlTask = task
136            .context(error::MissingRequiredParameterSnafu { param: "task" })?
137            .try_into()
138            .context(error::ConvertProtoDataSnafu)?;
139        restore_create_database_creator(&metadata, header.role, &mut task, &mut query_context)?;
140        let executor_context = ExecutorContext {
141            tracing_context: Some(header.tracing_context),
142            query_context: Some(query_context),
143            actor,
144            event_input: event_context.map(ProcedureEventInput::from),
145        };
146        let resp = self
147            .ddl_manager()
148            .submit_ddl_task(
149                executor_context,
150                SubmitDdlTaskRequest {
151                    wait,
152                    timeout: Duration::from_secs(timeout_secs.into()),
153                    task,
154                },
155            )
156            .await
157            .context(error::SubmitDdlTaskSnafu)?
158            .into();
159
160        Ok(Response::new(resp))
161    }
162
163    async fn migrate(
164        &self,
165        request: Request<MigrateRegionRequest>,
166    ) -> GrpcResult<MigrateRegionResponse> {
167        check_leader!(self, request, MigrateRegionResponse, "`migrate`");
168
169        let MigrateRegionRequest {
170            header,
171            region_id,
172            from_peer,
173            to_peer,
174            timeout_secs,
175            event_context,
176            actor,
177        } = request.into_inner();
178
179        let _header = header.context(error::MissingRequestHeaderSnafu)?;
180        let procedure_context =
181            ProcedureContext::from(ProcedureSubmission::from((actor, event_context)));
182        let from_peer = self
183            .lookup_datanode_peer(from_peer)
184            .await?
185            .unwrap_or_else(|| Peer::empty(from_peer));
186        let to_peer = self
187            .lookup_datanode_peer(to_peer)
188            .await?
189            .context(error::PeerUnavailableSnafu { peer_id: to_peer })?;
190
191        let pid = self
192            .region_migration_manager()
193            .submit_procedure(
194                procedure_context,
195                RegionMigrationProcedureTask {
196                    region_id: region_id.into(),
197                    from_peer,
198                    to_peer,
199                    timeout: Duration::from_secs(timeout_secs.into()),
200                    trigger_reason: RegionMigrationTriggerReason::Manual,
201                },
202            )
203            .await?
204            .map(procedure::pid_to_pb_pid);
205
206        let resp = MigrateRegionResponse {
207            pid,
208            ..Default::default()
209        };
210
211        Ok(Response::new(resp))
212    }
213
214    async fn reconcile(&self, request: Request<ReconcileRequest>) -> GrpcResult<ReconcileResponse> {
215        check_leader!(self, request, ReconcileResponse, "`reconcile`");
216
217        let ReconcileRequest { header, target } = request.into_inner();
218        let _header = header.context(error::MissingRequestHeaderSnafu)?;
219        let target = target.context(error::MissingRequiredParameterSnafu { param: "target" })?;
220        let parse_resolve_strategy = |resolve_strategy: i32| {
221            ResolveStrategy::try_from(resolve_strategy)
222                .ok()
223                .context(error::UnexpectedSnafu {
224                    violated: format!("Invalid resolve strategy: {}", resolve_strategy),
225                })
226        };
227        let procedure_id = match target {
228            Target::ReconcileTable(table) => {
229                let ReconcileTable {
230                    catalog_name,
231                    schema_name,
232                    table_name,
233                    resolve_strategy,
234                } = table;
235                let resolve_strategy = parse_resolve_strategy(resolve_strategy)?;
236                let table_ref = TableReference::full(&catalog_name, &schema_name, &table_name);
237                self.reconciliation_manager()
238                    .reconcile_table(table_ref, resolve_strategy.into())
239                    .await
240                    .context(error::SubmitReconcileProcedureSnafu)?
241            }
242            Target::ReconcileDatabase(database) => {
243                let ReconcileDatabase {
244                    catalog_name,
245                    database_name,
246                    resolve_strategy,
247                    parallelism,
248                } = database;
249                let resolve_strategy = parse_resolve_strategy(resolve_strategy)?;
250                self.reconciliation_manager()
251                    .reconcile_database(
252                        catalog_name,
253                        database_name,
254                        resolve_strategy.into(),
255                        parallelism as usize,
256                    )
257                    .await
258                    .context(error::SubmitReconcileProcedureSnafu)?
259            }
260            Target::ReconcileCatalog(catalog) => {
261                let ReconcileCatalog {
262                    catalog_name,
263                    resolve_strategy,
264                    parallelism,
265                } = catalog;
266                let resolve_strategy = parse_resolve_strategy(resolve_strategy)?;
267                self.reconciliation_manager()
268                    .reconcile_catalog(catalog_name, resolve_strategy.into(), parallelism as usize)
269                    .await
270                    .context(error::SubmitReconcileProcedureSnafu)?
271            }
272        };
273        Ok(Response::new(ReconcileResponse {
274            pid: Some(procedure::pid_to_pb_pid(procedure_id)),
275            ..Default::default()
276        }))
277    }
278
279    async fn details(
280        &self,
281        request: Request<ProcedureDetailRequest>,
282    ) -> GrpcResult<ProcedureDetailResponse> {
283        check_leader!(
284            self,
285            request,
286            ProcedureDetailResponse,
287            "`procedure details`"
288        );
289
290        let ProcedureDetailRequest { header } = request.into_inner();
291        let _header = header.context(error::MissingRequestHeaderSnafu)?;
292        let metas = self
293            .procedure_manager()
294            .list_procedures()
295            .await
296            .context(error::QueryProcedureSnafu)?;
297        Ok(Response::new(procedure::procedure_details_to_pb_response(
298            metas,
299        )))
300    }
301
302    async fn gc_regions(
303        &self,
304        request: Request<GcRegionsRequest>,
305    ) -> GrpcResult<GcRegionsResponse> {
306        check_leader!(self, request, GcRegionsResponse, "`gc_regions`");
307
308        let GcRegionsRequest {
309            header,
310            region_ids,
311            full_file_listing,
312            timeout_secs,
313            event_context,
314            actor,
315        } = request.into_inner();
316
317        let _header = header.context(error::MissingRequestHeaderSnafu)?;
318        let procedure_context =
319            ProcedureContext::from(ProcedureSubmission::from((actor, event_context)));
320
321        let response = self
322            .handle_gc_regions(
323                procedure_context,
324                MetaGcRegionsRequest {
325                    region_ids,
326                    full_file_listing,
327                    timeout: Self::normalize_gc_timeout(Duration::from_secs(timeout_secs as u64)),
328                },
329            )
330            .await?;
331
332        Ok(Response::new(gc_response_to_regions_pb(response)))
333    }
334
335    async fn gc_table(&self, request: Request<GcTableRequest>) -> GrpcResult<GcTableResponse> {
336        check_leader!(self, request, GcTableResponse, "`gc_table`");
337
338        let GcTableRequest {
339            header,
340            catalog_name,
341            schema_name,
342            table_name,
343            full_file_listing,
344            timeout_secs,
345            event_context,
346            actor,
347        } = request.into_inner();
348
349        let _header = header.context(error::MissingRequestHeaderSnafu)?;
350        let procedure_context =
351            ProcedureContext::from(ProcedureSubmission::from((actor, event_context)));
352
353        let response = self
354            .handle_gc_table(
355                procedure_context,
356                MetaGcTableRequest {
357                    catalog_name,
358                    schema_name,
359                    table_name,
360                    full_file_listing,
361                    timeout: Self::normalize_gc_timeout(Duration::from_secs(timeout_secs as u64)),
362                },
363            )
364            .await?;
365
366        Ok(Response::new(gc_response_to_table_pb(response)))
367    }
368}
369
370/// Restores the authenticated creator omitted from the protobuf create-database task.
371///
372/// The frontend sends the serialized [`CreatorGrantIntent`] in both a reserved query-context
373/// extension and internal gRPC metadata. This function removes the extension, requires both
374/// copies to match on a frontend create-database request, and writes the decoded intent into the
375/// task. Missing both preserves legacy and admin behavior; missing or mismatched copies fail
376/// closed.
377fn restore_create_database_creator(
378    metadata: &MetadataMap,
379    role: i32,
380    task: &mut DdlTask,
381    query_context: &mut QueryContext,
382) -> Result<(), Status> {
383    let forwarded = query_context
384        .extensions
385        .remove(CREATE_DATABASE_CREATOR_EXTENSION_KEY);
386    let internal = metadata
387        .get_bin(CREATE_DATABASE_CREATOR_METADATA_KEY)
388        .map(|value| {
389            value
390                .to_bytes()
391                .map_err(|_| Status::invalid_argument("Invalid internal creator metadata"))
392        })
393        .transpose()?;
394
395    let (forwarded, internal) = match (forwarded, internal) {
396        (None, None) => return Ok(()),
397        (Some(forwarded), Some(internal)) => (forwarded, internal),
398        _ => {
399            return Err(Status::invalid_argument(
400                "Untrusted create-database creator context",
401            ));
402        }
403    };
404
405    let DdlTask::CreateDatabase(task) = task else {
406        return Err(Status::invalid_argument(
407            "Untrusted create-database creator context",
408        ));
409    };
410    if Role::try_from(role) != Ok(Role::Frontend) || forwarded.as_bytes() != internal.as_ref() {
411        return Err(Status::invalid_argument(
412            "Untrusted create-database creator context",
413        ));
414    }
415
416    task.creator = Some(
417        serde_json::from_str::<CreatorGrantIntent>(&forwarded)
418            .map_err(|_| Status::invalid_argument("Invalid internal creator metadata"))?,
419    );
420    Ok(())
421}
422
423impl Metasrv {
424    fn normalize_gc_timeout(timeout: Duration) -> Option<Duration> {
425        if timeout.is_zero() {
426            None
427        } else {
428            Some(timeout)
429        }
430    }
431
432    async fn handle_gc_regions(
433        &self,
434        procedure_context: ProcedureContext,
435        request: MetaGcRegionsRequest,
436    ) -> error::Result<GcResponse> {
437        let region_ids: Vec<RegionId> = request
438            .region_ids
439            .into_iter()
440            .map(RegionId::from_u64)
441            .collect();
442        self.trigger_gc_for_regions(
443            procedure_context,
444            region_ids,
445            request.full_file_listing,
446            request.timeout,
447        )
448        .await
449    }
450
451    async fn handle_gc_table(
452        &self,
453        procedure_context: ProcedureContext,
454        request: MetaGcTableRequest,
455    ) -> error::Result<GcResponse> {
456        let table_name_key = TableNameKey::new(
457            &request.catalog_name,
458            &request.schema_name,
459            &request.table_name,
460        );
461
462        let table_metadata_manager: &TableMetadataManagerRef = self.table_metadata_manager();
463        let table_id = table_metadata_manager
464            .table_name_manager()
465            .get(table_name_key)
466            .await
467            .context(TableMetadataManagerSnafu)?
468            .context(TableNotFoundSnafu {
469                name: request.table_name.clone(),
470            })?
471            .table_id();
472
473        let (_phy_table_id, route) = table_metadata_manager
474            .table_route_manager()
475            .get_physical_table_route(table_id)
476            .await
477            .context(TableMetadataManagerSnafu)?;
478
479        let region_ids: Vec<RegionId> = route.region_routes.iter().map(|r| r.region.id).collect();
480        self.trigger_gc_for_regions(
481            procedure_context,
482            region_ids,
483            request.full_file_listing,
484            request.timeout,
485        )
486        .await
487    }
488
489    /// Triggers manual GC for specified regions and returns the GC response.
490    async fn trigger_gc_for_regions(
491        &self,
492        procedure_context: ProcedureContext,
493        region_ids: Vec<RegionId>,
494        full_file_listing: bool,
495        timeout: Option<Duration>,
496    ) -> error::Result<GcResponse> {
497        let gc_ticker = self.gc_ticker().context(error::UnexpectedSnafu {
498            violated: "GC ticker not available".to_string(),
499        })?;
500
501        let (tx, rx) = tokio::sync::oneshot::channel();
502        gc_ticker
503            .sender
504            .send(gc::Event::Manually {
505                sender: tx,
506                region_ids: Some(region_ids),
507                full_file_listing: Some(full_file_listing),
508                timeout,
509                procedure_context,
510            })
511            .await
512            .map_err(|_| {
513                error::UnexpectedSnafu {
514                    violated: "Failed to send GC event".to_string(),
515                }
516                .build()
517            })?;
518
519        let job_report = rx.await.map_err(|_| {
520            error::UnexpectedSnafu {
521                violated: "GC job channel closed unexpectedly".to_string(),
522            }
523            .build()
524        })??;
525
526        let report = gc_job_report_to_gc_report(job_report);
527
528        Ok(gc_report_to_response(&report))
529    }
530}
531
532fn gc_job_report_to_gc_report(job_report: crate::gc::GcJobReport) -> store_api::storage::GcReport {
533    job_report.merge_to_report()
534}
535
536fn gc_report_to_response(report: &store_api::storage::GcReport) -> GcResponse {
537    let deleted_files = report.deleted_files.values().map(|v| v.len() as u64).sum();
538    let deleted_indexes = report
539        .deleted_indexes
540        .values()
541        .map(|v| v.len() as u64)
542        .sum();
543    GcResponse {
544        processed_regions: report.processed_regions.len() as u64,
545        need_retry_regions: report
546            .need_retry_regions
547            .iter()
548            .map(|id| id.as_u64())
549            .collect(),
550        deleted_files,
551        deleted_indexes,
552    }
553}
554
555fn gc_response_to_regions_pb(resp: GcResponse) -> GcRegionsResponse {
556    GcRegionsResponse {
557        stats: Some(GcStats {
558            processed_regions: resp.processed_regions,
559            need_retry_regions: resp.need_retry_regions,
560            deleted_files: resp.deleted_files,
561            deleted_indexes: resp.deleted_indexes,
562        }),
563        ..Default::default()
564    }
565}
566
567fn gc_response_to_table_pb(resp: GcResponse) -> GcTableResponse {
568    GcTableResponse {
569        stats: Some(GcStats {
570            processed_regions: resp.processed_regions,
571            need_retry_regions: resp.need_retry_regions,
572            deleted_files: resp.deleted_files,
573            deleted_indexes: resp.deleted_indexes,
574        }),
575        ..Default::default()
576    }
577}
578
579#[cfg(test)]
580mod tests {
581    use std::time::Duration;
582
583    use api::v1::meta::Role;
584    use common_meta::rpc::ddl::{
585        CREATE_DATABASE_CREATOR_EXTENSION_KEY, CREATE_DATABASE_CREATOR_METADATA_KEY,
586        CreatorGrantIntent, DdlTask, QueryContext,
587    };
588    use tonic::metadata::{MetadataMap, MetadataValue};
589
590    use super::{Metasrv, restore_create_database_creator};
591
592    fn create_database_task() -> DdlTask {
593        DdlTask::new_create_database(
594            "greptime".to_string(),
595            "metrics".to_string(),
596            false,
597            Default::default(),
598            None,
599        )
600    }
601
602    fn creator_context(creator: Option<&str>) -> QueryContext {
603        let mut context = QueryContext::default();
604        if let Some(creator) = creator {
605            context.extensions.insert(
606                CREATE_DATABASE_CREATOR_EXTENSION_KEY.to_string(),
607                creator.to_string(),
608            );
609        }
610        context
611    }
612
613    fn creator_metadata(creator: Option<&str>) -> MetadataMap {
614        let mut metadata = MetadataMap::new();
615        if let Some(creator) = creator {
616            metadata.insert_bin(
617                CREATE_DATABASE_CREATOR_METADATA_KEY,
618                MetadataValue::from_bytes(creator.as_bytes()),
619            );
620        }
621        metadata
622    }
623
624    #[test]
625    fn test_normalize_gc_timeout() {
626        assert_eq!(Metasrv::normalize_gc_timeout(Duration::ZERO), None);
627        assert_eq!(
628            Metasrv::normalize_gc_timeout(Duration::from_secs(10)),
629            Some(Duration::from_secs(10))
630        );
631    }
632
633    #[test]
634    fn test_restore_create_database_creator_requires_matching_internal_metadata() {
635        let mut task = create_database_task();
636        let mut legacy = QueryContext::default();
637        restore_create_database_creator(
638            &MetadataMap::new(),
639            Role::Frontend as i32,
640            &mut task,
641            &mut legacy,
642        )
643        .unwrap();
644
645        let creator = CreatorGrantIntent {
646            username: "alice".to_string(),
647            created_at_ns: 42,
648        };
649        let encoded = serde_json::to_string(&creator).unwrap();
650        let other = serde_json::to_string(&CreatorGrantIntent {
651            username: "mallory".to_string(),
652            created_at_ns: 99,
653        })
654        .unwrap();
655
656        for (name, body, metadata, role, create_database) in [
657            (
658                "body only",
659                Some(encoded.as_str()),
660                None,
661                Role::Frontend,
662                true,
663            ),
664            (
665                "metadata only",
666                None,
667                Some(encoded.as_str()),
668                Role::Frontend,
669                true,
670            ),
671            (
672                "mismatch",
673                Some(other.as_str()),
674                Some(encoded.as_str()),
675                Role::Frontend,
676                true,
677            ),
678            ("malformed", Some("{"), Some("{"), Role::Frontend, true),
679            (
680                "wrong role",
681                Some(encoded.as_str()),
682                Some(encoded.as_str()),
683                Role::Datanode,
684                true,
685            ),
686            (
687                "wrong task",
688                Some(encoded.as_str()),
689                Some(encoded.as_str()),
690                Role::Frontend,
691                false,
692            ),
693        ] {
694            let mut task = if create_database {
695                create_database_task()
696            } else {
697                DdlTask::new_drop_database("greptime".to_string(), "metrics".to_string(), false)
698            };
699            let mut context = creator_context(body);
700            assert!(
701                restore_create_database_creator(
702                    &creator_metadata(metadata),
703                    role as i32,
704                    &mut task,
705                    &mut context,
706                )
707                .is_err(),
708                "{name} must fail closed"
709            );
710        }
711
712        let mut task = create_database_task();
713        let mut trusted = creator_context(Some(&encoded));
714        restore_create_database_creator(
715            &creator_metadata(Some(&encoded)),
716            Role::Frontend as i32,
717            &mut task,
718            &mut trusted,
719        )
720        .unwrap();
721        let DdlTask::CreateDatabase(task) = task else {
722            unreachable!();
723        };
724        assert_eq!(task.creator, Some(creator));
725        assert!(trusted.extensions.is_empty());
726    }
727}