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
74#[async_trait]
76pub trait ProcedureServiceHandler: Send + Sync {
77 async fn purge_table(&self, table_name: TableName, query_ctx: QueryContextRef) -> Result<()>;
79
80 async fn migrate_region(&self, request: MigrateRegionRequest) -> Result<Option<String>>;
82
83 async fn reconcile(&self, request: ReconcileRequest) -> Result<Option<String>>;
85
86 async fn query_procedure_state(&self, pid: &str) -> Result<ProcedureStateResponse>;
88
89 async fn manage_region_follower(&self, request: ManageRegionFollowerRequest) -> Result<()>;
91
92 fn catalog_manager(&self) -> &CatalogManagerRef;
94
95 async fn gc_regions(&self, request: MetaGcRegionsRequest) -> Result<MetaGcResponse>;
97
98 async fn gc_table(&self, request: MetaGcTableRequest) -> Result<MetaGcResponse>;
100}
101
102#[async_trait]
104pub trait FlowServiceHandler: Send + Sync {
105 async fn flush(
106 &self,
107 catalog: &str,
108 flow: &str,
109 ctx: QueryContextRef,
110 ) -> Result<api::v1::flow::FlowResponse>;
111}
112
113pub type TableMutationHandlerRef = Arc<dyn TableMutationHandler>;
114
115pub type ProcedureServiceHandlerRef = Arc<dyn ProcedureServiceHandler>;
116
117pub type FlowServiceHandlerRef = Arc<dyn FlowServiceHandler>;