Skip to main content

operator/
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 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/// The operator for procedures which implements [`ProcedureServiceHandler`].
38#[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}