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