Skip to main content

operator/
table.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 async_trait::async_trait;
16use client::Output;
17use common_base::AffectedRows;
18use common_error::ext::BoxedError;
19use common_function::handlers::TableMutationHandler;
20use common_query::error as query_error;
21use common_query::error::Result as QueryResult;
22use session::context::QueryContextRef;
23use snafu::ResultExt;
24use store_api::storage::RegionId;
25use table::requests::{
26    BuildIndexTableRequest, CompactTableRequest, DeleteRequest as TableDeleteRequest,
27    FlushTableRequest, InsertRequest as TableInsertRequest,
28};
29use table::table_name::TableName;
30
31use crate::delete::DeleterRef;
32use crate::insert::InserterRef;
33use crate::request::RequesterRef;
34
35pub struct TableMutationOperator {
36    inserter: InserterRef,
37    deleter: DeleterRef,
38    requester: RequesterRef,
39}
40
41impl TableMutationOperator {
42    pub fn new(inserter: InserterRef, deleter: DeleterRef, requester: RequesterRef) -> Self {
43        Self {
44            inserter,
45            deleter,
46            requester,
47        }
48    }
49}
50
51#[async_trait]
52impl TableMutationHandler for TableMutationOperator {
53    async fn insert(
54        &self,
55        request: TableInsertRequest,
56        ctx: QueryContextRef,
57    ) -> QueryResult<Output> {
58        self.inserter
59            .handle_table_insert(request, ctx)
60            .await
61            .map_err(BoxedError::new)
62            .context(query_error::TableMutationSnafu)
63    }
64
65    async fn delete(
66        &self,
67        request: TableDeleteRequest,
68        ctx: QueryContextRef,
69    ) -> QueryResult<AffectedRows> {
70        self.deleter
71            .handle_table_delete(request, ctx)
72            .await
73            .map_err(BoxedError::new)
74            .context(query_error::TableMutationSnafu)
75    }
76
77    async fn flush(
78        &self,
79        request: FlushTableRequest,
80        ctx: QueryContextRef,
81    ) -> QueryResult<AffectedRows> {
82        self.requester
83            .handle_table_flush(request, ctx)
84            .await
85            .map_err(BoxedError::new)
86            .context(query_error::TableMutationSnafu)
87    }
88
89    async fn compact(
90        &self,
91        request: CompactTableRequest,
92        ctx: QueryContextRef,
93    ) -> QueryResult<AffectedRows> {
94        self.requester
95            .handle_table_compaction(request, ctx)
96            .await
97            .map_err(BoxedError::new)
98            .context(query_error::TableMutationSnafu)
99    }
100
101    async fn build_index(
102        &self,
103        request: BuildIndexTableRequest,
104        ctx: QueryContextRef,
105    ) -> QueryResult<AffectedRows> {
106        self.requester
107            .handle_table_build_index(request, ctx)
108            .await
109            .map_err(BoxedError::new)
110            .context(query_error::TableMutationSnafu)
111    }
112
113    async fn flush_region(
114        &self,
115        region_id: RegionId,
116        ctx: QueryContextRef,
117    ) -> QueryResult<AffectedRows> {
118        self.requester
119            .handle_region_flush(region_id, ctx)
120            .await
121            .map_err(BoxedError::new)
122            .context(query_error::TableMutationSnafu)
123    }
124
125    async fn compact_region(
126        &self,
127        region_id: RegionId,
128        ctx: QueryContextRef,
129    ) -> QueryResult<AffectedRows> {
130        self.requester
131            .handle_region_compaction(region_id, ctx)
132            .await
133            .map_err(BoxedError::new)
134            .context(query_error::TableMutationSnafu)
135    }
136
137    async fn discard_unflushed_data(
138        &self,
139        region_id: RegionId,
140        ctx: QueryContextRef,
141    ) -> QueryResult<AffectedRows> {
142        self.requester
143            .handle_discard_unflushed_data(region_id, ctx)
144            .await
145            .map_err(BoxedError::new)
146            .context(query_error::TableMutationSnafu)
147    }
148
149    async fn discard_unflushed_data_by_table(
150        &self,
151        table_name: TableName,
152        ctx: QueryContextRef,
153    ) -> QueryResult<AffectedRows> {
154        self.requester
155            .handle_discard_unflushed_data_by_table(table_name, ctx)
156            .await
157            .map_err(BoxedError::new)
158            .context(query_error::TableMutationSnafu)
159    }
160}