1use 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
370fn 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 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}