1use api::v1::meta::ReconcileRequest;
16use async_trait::async_trait;
17use catalog::CatalogManagerRef;
18use common_error::ext::BoxedError;
19use common_function::handlers::ProcedureServiceHandler;
20use common_meta::key::TableMetadataManagerRef;
21use common_meta::procedure_executor::{ExecutorContext, ProcedureExecutorRef};
22use common_meta::rpc::ddl::{DdlTask, SubmitDdlTaskRequest, TriggerReason};
23use common_meta::rpc::procedure::{
24 GcRegionsRequest as MetaGcRegionsRequest, GcResponse as MetaGcResponse,
25 GcTableRequest as MetaGcTableRequest, ManageRegionFollowerRequest, MigrateRegionRequest,
26 ProcedureStateResponse,
27};
28use common_query::error as query_error;
29use common_query::error::Result as QueryResult;
30use session::context::QueryContextRef;
31use snafu::ResultExt;
32use table::table_name::TableName;
33
34use crate::error;
35use crate::utils::to_executor_context;
36
37#[derive(Clone)]
39pub struct ProcedureServiceOperator {
40 procedure_executor: ProcedureExecutorRef,
41 catalog_manager: CatalogManagerRef,
42 table_metadata_manager: TableMetadataManagerRef,
43}
44
45impl ProcedureServiceOperator {
46 pub fn new(
47 procedure_executor: ProcedureExecutorRef,
48 catalog_manager: CatalogManagerRef,
49 table_metadata_manager: TableMetadataManagerRef,
50 ) -> Self {
51 Self {
52 procedure_executor,
53 catalog_manager,
54 table_metadata_manager,
55 }
56 }
57}
58
59#[async_trait]
60impl ProcedureServiceHandler for ProcedureServiceOperator {
61 async fn purge_table(
62 &self,
63 query_ctx: QueryContextRef,
64 table_name: TableName,
65 ) -> QueryResult<()> {
66 let dropped = self
67 .table_metadata_manager
68 .get_dropped_table(&table_name)
69 .await
70 .map_err(BoxedError::new)
71 .context(query_error::ProcedureServiceSnafu)?
72 .ok_or_else(|| {
73 error::TableNotFoundSnafu {
74 table_name: table_name.to_string(),
75 }
76 .build()
77 })
78 .map_err(BoxedError::new)
79 .context(query_error::ProcedureServiceSnafu)?;
80 let executor_context = to_executor_context(query_ctx, TriggerReason::Manual);
81 let request = SubmitDdlTaskRequest::new(DdlTask::new_purge_dropped_table(dropped.table_id));
82 self.procedure_executor
83 .submit_ddl_task(executor_context, request)
84 .await
85 .map_err(BoxedError::new)
86 .context(query_error::ProcedureServiceSnafu)?;
87 Ok(())
88 }
89
90 async fn migrate_region(
91 &self,
92 query_ctx: QueryContextRef,
93 request: MigrateRegionRequest,
94 ) -> QueryResult<Option<String>> {
95 let executor_context = to_executor_context(query_ctx, TriggerReason::Manual);
96 Ok(self
97 .procedure_executor
98 .migrate_region(&executor_context, request)
99 .await
100 .map_err(BoxedError::new)
101 .context(query_error::ProcedureServiceSnafu)?
102 .pid
103 .map(|pid| String::from_utf8_lossy(&pid.key).to_string()))
104 }
105
106 async fn reconcile(&self, request: ReconcileRequest) -> QueryResult<Option<String>> {
107 Ok(self
108 .procedure_executor
109 .reconcile(&ExecutorContext::default(), request)
110 .await
111 .map_err(BoxedError::new)
112 .context(query_error::ProcedureServiceSnafu)?
113 .pid
114 .map(|pid| String::from_utf8_lossy(&pid.key).to_string()))
115 }
116
117 async fn query_procedure_state(&self, pid: &str) -> QueryResult<ProcedureStateResponse> {
118 self.procedure_executor
119 .query_procedure_state(&ExecutorContext::default(), pid)
120 .await
121 .map_err(BoxedError::new)
122 .context(query_error::ProcedureServiceSnafu)
123 }
124
125 async fn manage_region_follower(
126 &self,
127 request: ManageRegionFollowerRequest,
128 ) -> QueryResult<()> {
129 self.procedure_executor
130 .manage_region_follower(&ExecutorContext::default(), request)
131 .await
132 .map_err(BoxedError::new)
133 .context(query_error::ProcedureServiceSnafu)
134 }
135
136 fn catalog_manager(&self) -> &CatalogManagerRef {
137 &self.catalog_manager
138 }
139
140 async fn gc_regions(
141 &self,
142 query_ctx: QueryContextRef,
143 request: MetaGcRegionsRequest,
144 ) -> QueryResult<MetaGcResponse> {
145 let executor_context = to_executor_context(query_ctx, TriggerReason::Manual);
146 self.procedure_executor
147 .gc_regions(&executor_context, request)
148 .await
149 .map_err(BoxedError::new)
150 .context(query_error::ProcedureServiceSnafu)
151 }
152
153 async fn gc_table(
154 &self,
155 query_ctx: QueryContextRef,
156 request: MetaGcTableRequest,
157 ) -> QueryResult<MetaGcResponse> {
158 let executor_context = to_executor_context(query_ctx, TriggerReason::Manual);
159 self.procedure_executor
160 .gc_table(&executor_context, request)
161 .await
162 .map_err(BoxedError::new)
163 .context(query_error::ProcedureServiceSnafu)
164 }
165}
166
167#[cfg(test)]
168mod tests {
169 use std::collections::HashMap;
170 use std::sync::{Arc, Mutex};
171
172 use api::v1::meta::{ProcedureDetailResponse, ReconcileResponse};
173 use common_meta::key::TableMetadataManager;
174 use common_meta::key::table_route::TableRouteValue;
175 use common_meta::key::test_utils::new_test_table_info_with_name;
176 use common_meta::kv_backend::memory::MemoryKvBackend;
177 use common_meta::procedure_executor::{
178 ExecutorContext, ProcedureExecutor, ProcedureExecutorRef,
179 };
180 use common_meta::rpc::ddl::{DdlTask, SubmitDdlTaskRequest, SubmitDdlTaskResponse};
181 use common_meta::rpc::procedure::{MigrateRegionResponse, ProcedureStateResponse};
182 use session::context::QueryContextBuilder;
183 use table::table_name::TableName;
184
185 use super::*;
186
187 #[derive(Default)]
188 struct RecordingProcedureExecutor {
189 requests: Mutex<Vec<SubmitDdlTaskRequest>>,
190 contexts: Mutex<Vec<ExecutorContext>>,
191 fail: bool,
192 }
193
194 #[async_trait]
195 impl ProcedureExecutor for RecordingProcedureExecutor {
196 async fn submit_ddl_task(
197 &self,
198 context: ExecutorContext,
199 request: SubmitDdlTaskRequest,
200 ) -> common_meta::error::Result<SubmitDdlTaskResponse> {
201 self.contexts.lock().unwrap().push(context);
202 self.requests.lock().unwrap().push(request);
203 if self.fail {
204 return common_meta::error::UnsupportedSnafu {
205 operation: "test purge failure",
206 }
207 .fail();
208 }
209 Ok(SubmitDdlTaskResponse::default())
210 }
211 async fn migrate_region(
212 &self,
213 _: &ExecutorContext,
214 _: MigrateRegionRequest,
215 ) -> common_meta::error::Result<MigrateRegionResponse> {
216 unimplemented!()
217 }
218 async fn reconcile(
219 &self,
220 _: &ExecutorContext,
221 _: ReconcileRequest,
222 ) -> common_meta::error::Result<ReconcileResponse> {
223 unimplemented!()
224 }
225 async fn query_procedure_state(
226 &self,
227 _: &ExecutorContext,
228 _: &str,
229 ) -> common_meta::error::Result<ProcedureStateResponse> {
230 unimplemented!()
231 }
232 async fn list_procedures(
233 &self,
234 _: &ExecutorContext,
235 ) -> common_meta::error::Result<ProcedureDetailResponse> {
236 unimplemented!()
237 }
238 }
239
240 async fn operator_with_tombstone(
241 table_id: u32,
242 name: &TableName,
243 fail: bool,
244 ) -> (ProcedureServiceOperator, Arc<RecordingProcedureExecutor>) {
245 let manager = Arc::new(TableMetadataManager::new(Arc::new(
246 MemoryKvBackend::default(),
247 )));
248 let mut info = new_test_table_info_with_name(table_id, &name.table_name);
249 info.catalog_name = name.catalog_name.clone();
250 info.schema_name = name.schema_name.clone();
251 let route = TableRouteValue::physical(vec![]);
252 manager
253 .create_table_metadata(info, route.clone(), HashMap::new())
254 .await
255 .unwrap();
256 manager
257 .delete_table_metadata(table_id, name, &route, &HashMap::new(), None)
258 .await
259 .unwrap();
260 let executor = Arc::new(RecordingProcedureExecutor {
261 fail,
262 ..Default::default()
263 });
264 let operator = ProcedureServiceOperator::new(
265 executor.clone() as ProcedureExecutorRef,
266 catalog::memory::MemoryCatalogManager::new(),
267 manager,
268 );
269 (operator, executor)
270 }
271
272 #[tokio::test]
273 async fn test_purge_table_submits_tombstone_id_and_query_context() {
274 let name = TableName::new("catalog", "schema", "metrics");
275 let (operator, executor) = operator_with_tombstone(42, &name, false).await;
276 let manager = operator.table_metadata_manager.clone();
277 let mut live = new_test_table_info_with_name(99, "metrics");
278 live.catalog_name = "catalog".to_string();
279 live.schema_name = "schema".to_string();
280 manager
281 .create_table_metadata(live, TableRouteValue::physical(vec![]), HashMap::new())
282 .await
283 .unwrap();
284 let query_ctx = QueryContextBuilder::default()
285 .current_catalog("catalog".to_string())
286 .current_schema("schema".to_string())
287 .build()
288 .into();
289
290 operator.purge_table(query_ctx, name.clone()).await.unwrap();
291
292 let requests = executor.requests.lock().unwrap();
293 assert!(
294 matches!(&requests[0].task, DdlTask::PurgeDroppedTable(task) if task.table_id == 42)
295 );
296 let contexts = executor.contexts.lock().unwrap();
297 let query_context = contexts[0].query_context.as_ref().unwrap();
298 assert_eq!(query_context.current_catalog, "catalog");
299 assert_eq!(query_context.current_schema, "schema");
300 assert_eq!(
301 contexts[0].event_input.as_ref().map(|input| input.reason),
302 Some(TriggerReason::Manual)
303 );
304 }
305
306 #[tokio::test]
307 async fn test_purge_table_missing_tombstone_is_table_not_found() {
308 let manager = Arc::new(TableMetadataManager::new(Arc::new(
309 MemoryKvBackend::default(),
310 )));
311 let executor = Arc::new(RecordingProcedureExecutor::default());
312 let operator = ProcedureServiceOperator::new(
313 executor.clone(),
314 catalog::memory::MemoryCatalogManager::new(),
315 manager,
316 );
317 let error = operator
318 .purge_table(
319 QueryContextBuilder::default()
320 .current_catalog("catalog".to_string())
321 .current_schema("schema".to_string())
322 .build()
323 .into(),
324 TableName::new("catalog", "schema", "missing"),
325 )
326 .await
327 .unwrap_err();
328 assert_eq!(
329 common_error::status_code::StatusCode::TableNotFound,
330 common_error::ext::ErrorExt::status_code(&error)
331 );
332 assert!(executor.requests.lock().unwrap().is_empty());
333 }
334
335 #[tokio::test]
336 async fn test_purge_table_propagates_submission_failure() {
337 let name = TableName::new("catalog", "schema", "metrics");
338 let (operator, executor) = operator_with_tombstone(42, &name, true).await;
339 assert!(
340 operator
341 .purge_table(
342 QueryContextBuilder::default()
343 .current_catalog("catalog".to_string())
344 .current_schema("schema".to_string())
345 .build()
346 .into(),
347 name,
348 )
349 .await
350 .is_err()
351 );
352 assert_eq!(1, executor.requests.lock().unwrap().len());
353 }
354}