common_function/
handlers.rs1use std::sync::Arc;
16
17use api::v1::meta::ReconcileRequest;
18use async_trait::async_trait;
19use catalog::CatalogManagerRef;
20use common_base::AffectedRows;
21use common_meta::rpc::procedure::{
22 GcRegionsRequest as MetaGcRegionsRequest, GcResponse as MetaGcResponse,
23 GcTableRequest as MetaGcTableRequest, ManageRegionFollowerRequest, MigrateRegionRequest,
24 ProcedureStateResponse,
25};
26use common_query::Output;
27use common_query::error::Result;
28use session::context::QueryContextRef;
29use store_api::storage::RegionId;
30use table::requests::{
31 BuildIndexTableRequest, CompactTableRequest, DeleteRequest, FlushTableRequest, InsertRequest,
32};
33use table::table_name::TableName;
34
35#[async_trait]
37pub trait TableMutationHandler: Send + Sync {
38 async fn insert(&self, request: InsertRequest, ctx: QueryContextRef) -> Result<Output>;
40
41 async fn delete(&self, request: DeleteRequest, ctx: QueryContextRef) -> Result<AffectedRows>;
43
44 async fn flush(&self, request: FlushTableRequest, ctx: QueryContextRef)
46 -> Result<AffectedRows>;
47
48 async fn compact(
50 &self,
51 request: CompactTableRequest,
52 ctx: QueryContextRef,
53 ) -> Result<AffectedRows>;
54
55 async fn build_index(
57 &self,
58 request: BuildIndexTableRequest,
59 ctx: QueryContextRef,
60 ) -> Result<AffectedRows>;
61
62 async fn flush_region(&self, region_id: RegionId, ctx: QueryContextRef)
64 -> Result<AffectedRows>;
65
66 async fn compact_region(
68 &self,
69 region_id: RegionId,
70 ctx: QueryContextRef,
71 ) -> Result<AffectedRows>;
72
73 async fn discard_unflushed_data(
75 &self,
76 region_id: RegionId,
77 ctx: QueryContextRef,
78 ) -> Result<AffectedRows>;
79
80 async fn discard_unflushed_data_by_table(
82 &self,
83 table_name: TableName,
84 ctx: QueryContextRef,
85 ) -> Result<AffectedRows>;
86}
87
88#[async_trait]
90pub trait ProcedureServiceHandler: Send + Sync {
91 async fn purge_table(&self, query_ctx: QueryContextRef, table_name: TableName) -> Result<()>;
93
94 async fn migrate_region(
96 &self,
97 query_ctx: QueryContextRef,
98 request: MigrateRegionRequest,
99 ) -> Result<Option<String>>;
100
101 async fn reconcile(&self, request: ReconcileRequest) -> Result<Option<String>>;
103
104 async fn query_procedure_state(&self, pid: &str) -> Result<ProcedureStateResponse>;
106
107 async fn manage_region_follower(&self, request: ManageRegionFollowerRequest) -> Result<()>;
109
110 fn catalog_manager(&self) -> &CatalogManagerRef;
112
113 async fn gc_regions(
115 &self,
116 query_ctx: QueryContextRef,
117 request: MetaGcRegionsRequest,
118 ) -> Result<MetaGcResponse>;
119
120 async fn gc_table(
122 &self,
123 query_ctx: QueryContextRef,
124 request: MetaGcTableRequest,
125 ) -> Result<MetaGcResponse>;
126}
127
128#[async_trait]
130pub trait FlowServiceHandler: Send + Sync {
131 async fn flush(
132 &self,
133 catalog: &str,
134 flow: &str,
135 ctx: QueryContextRef,
136 ) -> Result<api::v1::flow::FlowResponse>;
137}
138
139pub type TableMutationHandlerRef = Arc<dyn TableMutationHandler>;
140
141pub type ProcedureServiceHandlerRef = Arc<dyn ProcedureServiceHandler>;
142
143pub type FlowServiceHandlerRef = Arc<dyn FlowServiceHandler>;