1use crate::handlers::{FlowServiceHandlerRef, ProcedureServiceHandlerRef, TableMutationHandlerRef};
16
17#[derive(Clone, Default)]
20pub struct FunctionState {
21 pub table_mutation_handler: Option<TableMutationHandlerRef>,
23 pub procedure_service_handler: Option<ProcedureServiceHandlerRef>,
25 pub flow_service_handler: Option<FlowServiceHandlerRef>,
27}
28
29impl FunctionState {
30 #[cfg(any(test, feature = "testing"))]
32 pub fn mock() -> Self {
33 use std::sync::Arc;
34
35 use api::v1::meta::{ProcedureStatus, ReconcileRequest};
36 use async_trait::async_trait;
37 use catalog::CatalogManagerRef;
38 use common_base::AffectedRows;
39 use common_meta::rpc::procedure::{
40 GcRegionsRequest, GcResponse, GcTableRequest, ManageRegionFollowerRequest,
41 MigrateRegionRequest, ProcedureStateResponse,
42 };
43 use common_query::Output;
44 use common_query::error::Result;
45 use session::context::QueryContextRef;
46 use store_api::storage::RegionId;
47 use table::requests::{
48 BuildIndexTableRequest, CompactTableRequest, DeleteRequest, FlushTableRequest,
49 InsertRequest,
50 };
51
52 use crate::handlers::{FlowServiceHandler, ProcedureServiceHandler, TableMutationHandler};
53 struct MockProcedureServiceHandler;
54 struct MockTableMutationHandler;
55 struct MockFlowServiceHandler;
56 const ROWS: usize = 42;
57
58 #[async_trait]
59 impl ProcedureServiceHandler for MockProcedureServiceHandler {
60 async fn purge_table(
61 &self,
62 _table_name: table::table_name::TableName,
63 _query_ctx: QueryContextRef,
64 ) -> Result<()> {
65 Ok(())
66 }
67
68 async fn migrate_region(
69 &self,
70 _request: MigrateRegionRequest,
71 ) -> Result<Option<String>> {
72 Ok(Some("test_pid".to_string()))
73 }
74
75 async fn reconcile(&self, _request: ReconcileRequest) -> Result<Option<String>> {
76 Ok(Some("test_pid".to_string()))
77 }
78
79 async fn query_procedure_state(&self, _pid: &str) -> Result<ProcedureStateResponse> {
80 Ok(ProcedureStateResponse {
81 status: ProcedureStatus::Done.into(),
82 error: "OK".to_string(),
83 ..Default::default()
84 })
85 }
86
87 async fn manage_region_follower(
88 &self,
89 _request: ManageRegionFollowerRequest,
90 ) -> Result<()> {
91 Ok(())
92 }
93
94 async fn gc_regions(&self, _request: GcRegionsRequest) -> Result<GcResponse> {
95 Ok(GcResponse {
96 processed_regions: 1,
97 need_retry_regions: vec![],
98 deleted_files: 0,
99 deleted_indexes: 0,
100 })
101 }
102
103 async fn gc_table(&self, _request: GcTableRequest) -> Result<GcResponse> {
104 Ok(GcResponse {
105 processed_regions: 1,
106 need_retry_regions: vec![],
107 deleted_files: 0,
108 deleted_indexes: 0,
109 })
110 }
111
112 fn catalog_manager(&self) -> &CatalogManagerRef {
113 unimplemented!()
114 }
115 }
116
117 #[async_trait]
118 impl TableMutationHandler for MockTableMutationHandler {
119 async fn insert(
120 &self,
121 _request: InsertRequest,
122 _ctx: QueryContextRef,
123 ) -> Result<Output> {
124 Ok(Output::new_with_affected_rows(ROWS))
125 }
126
127 async fn delete(
128 &self,
129 _request: DeleteRequest,
130 _ctx: QueryContextRef,
131 ) -> Result<AffectedRows> {
132 Ok(ROWS)
133 }
134
135 async fn flush(
136 &self,
137 _request: FlushTableRequest,
138 _ctx: QueryContextRef,
139 ) -> Result<AffectedRows> {
140 Ok(ROWS)
141 }
142
143 async fn compact(
144 &self,
145 _request: CompactTableRequest,
146 _ctx: QueryContextRef,
147 ) -> Result<AffectedRows> {
148 Ok(ROWS)
149 }
150
151 async fn build_index(
152 &self,
153 _request: BuildIndexTableRequest,
154 _ctx: QueryContextRef,
155 ) -> Result<AffectedRows> {
156 Ok(ROWS)
157 }
158
159 async fn flush_region(
160 &self,
161 _region_id: RegionId,
162 _ctx: QueryContextRef,
163 ) -> Result<AffectedRows> {
164 Ok(ROWS)
165 }
166
167 async fn compact_region(
168 &self,
169 _region_id: RegionId,
170 _ctx: QueryContextRef,
171 ) -> Result<AffectedRows> {
172 Ok(ROWS)
173 }
174 }
175
176 #[async_trait]
177 impl FlowServiceHandler for MockFlowServiceHandler {
178 async fn flush(
179 &self,
180 _catalog: &str,
181 _flow: &str,
182 _ctx: QueryContextRef,
183 ) -> Result<api::v1::flow::FlowResponse> {
184 todo!()
185 }
186 }
187
188 Self {
189 table_mutation_handler: Some(Arc::new(MockTableMutationHandler)),
190 procedure_service_handler: Some(Arc::new(MockProcedureServiceHandler)),
191 flow_service_handler: Some(Arc::new(MockFlowServiceHandler)),
192 }
193 }
194}