1use std::net::SocketAddr;
18use std::sync::Arc;
19
20use api::v1::flow::DirtyWindowRequests;
21use api::v1::{RowDeleteRequests, RowInsertRequests};
22use cache::{PARTITION_INFO_CACHE_NAME, TABLE_FLOWNODE_SET_CACHE_NAME, TABLE_ROUTE_CACHE_NAME};
23use catalog::CatalogManagerRef;
24use common_base::Plugins;
25use common_error::ext::BoxedError;
26use common_meta::cache::{LayeredCacheRegistryRef, TableFlownodeSetCacheRef, TableRouteCacheRef};
27use common_meta::key::TableMetadataManagerRef;
28use common_meta::key::flow::FlowMetadataManagerRef;
29use common_meta::kv_backend::KvBackendRef;
30use common_meta::node_manager::{Flownode, NodeManagerRef};
31use common_meta::procedure_executor::ProcedureExecutorRef;
32use common_query::Output;
33use common_runtime::JoinHandle;
34use common_telemetry::tracing::info;
35use futures::TryStreamExt;
36use greptime_proto::v1::flow::{FlowRequest, FlowResponse, InsertRequests, flow_server};
37use itertools::Itertools;
38use operator::delete::Deleter;
39use operator::insert::Inserter;
40use operator::statement::StatementExecutor;
41use partition::cache::PartitionInfoCacheRef;
42use partition::manager::PartitionRuleManager;
43use query::{QueryEngine, QueryEngineFactory};
44use servers::add_service;
45use servers::grpc::builder::GrpcServerBuilder;
46use servers::grpc::{GrpcServer, GrpcServerConfig};
47use servers::http::HttpServerBuilder;
48use servers::metrics_handler::MetricsHandler;
49use servers::server::{ServerHandler, ServerHandlers};
50use session::context::QueryContextRef;
51use snafu::{OptionExt, ResultExt};
52use tokio::sync::{Mutex, broadcast, oneshot};
53use tonic::codec::CompressionEncoding;
54use tonic::{Request, Response, Status};
55
56use crate::adapter::flownode_impl::{FlowDualEngine, FlowDualEngineRef};
57use crate::adapter::{FlowStreamingEngineRef, create_worker};
58use crate::batching_mode::engine::BatchingEngine;
59use crate::error::{
60 CacheRequiredSnafu, DatafusionSnafu, ExternalSnafu, ListFlowsSnafu, ParseAddrSnafu,
61 ShutdownServerSnafu, StartServerSnafu, UnexpectedSnafu, to_status_with_last_err,
62};
63use crate::heartbeat::HeartbeatTask;
64use crate::metrics::{METRIC_FLOW_PROCESSING_TIME, METRIC_FLOW_ROWS};
65use crate::transform::register_function_to_query_engine;
66use crate::utils::{SizeReportSender, StateReportHandler};
67use crate::{Error, FlownodeOptions, FrontendClient, StreamingEngine};
68
69pub const FLOW_NODE_SERVER_NAME: &str = "FLOW_NODE_SERVER";
70#[derive(Clone)]
72pub struct FlowService {
73 pub dual_engine: FlowDualEngineRef,
74}
75
76impl FlowService {
77 pub fn new(manager: FlowDualEngineRef) -> Self {
78 Self {
79 dual_engine: manager,
80 }
81 }
82}
83
84#[async_trait::async_trait]
85impl flow_server::Flow for FlowService {
86 async fn handle_create_remove(
87 &self,
88 request: Request<FlowRequest>,
89 ) -> Result<Response<FlowResponse>, Status> {
90 let _timer = METRIC_FLOW_PROCESSING_TIME
91 .with_label_values(&["ddl"])
92 .start_timer();
93
94 let request = request.into_inner();
95 self.dual_engine
96 .handle(request)
97 .await
98 .map_err(|err| {
99 common_telemetry::error!(err; "Failed to handle flow request");
100 err
101 })
102 .map(Response::new)
103 .map_err(to_status_with_last_err)
104 }
105
106 async fn handle_mirror_request(
107 &self,
108 request: Request<InsertRequests>,
109 ) -> Result<Response<FlowResponse>, Status> {
110 let _timer = METRIC_FLOW_PROCESSING_TIME
111 .with_label_values(&["insert"])
112 .start_timer();
113
114 let request = request.into_inner();
115 let mut row_count = 0;
117 let request = api::v1::region::InsertRequests {
118 requests: request
119 .requests
120 .into_iter()
121 .map(|insert| {
122 insert.rows.as_ref().inspect(|x| row_count += x.rows.len());
123 api::v1::region::InsertRequest {
124 region_id: insert.region_id,
125 rows: insert.rows,
126 partition_expr_version: insert.partition_expr_version,
127 }
128 })
129 .collect_vec(),
130 };
131
132 METRIC_FLOW_ROWS
133 .with_label_values(&["in"])
134 .inc_by(row_count as u64);
135
136 self.dual_engine
137 .handle_inserts(request)
138 .await
139 .map(Response::new)
140 .map_err(to_status_with_last_err)
141 }
142
143 async fn handle_mark_dirty_time_window(
144 &self,
145 reqs: Request<DirtyWindowRequests>,
146 ) -> Result<Response<FlowResponse>, Status> {
147 self.dual_engine
148 .handle_mark_window_dirty(reqs.into_inner())
149 .await
150 .map(Response::new)
151 .map_err(to_status_with_last_err)
152 }
153}
154
155#[derive(Clone)]
156pub struct FlownodeServer {
157 inner: Arc<FlownodeServerInner>,
158}
159
160struct FlownodeServerInner {
164 worker_shutdown_tx: Mutex<broadcast::Sender<()>>,
166 server_shutdown_tx: Mutex<broadcast::Sender<()>>,
168 streaming_task_handler: Mutex<Option<JoinHandle<()>>>,
170 state_report_task_handler: Mutex<Option<JoinHandle<()>>>,
172 flow_service: FlowService,
173}
174
175impl FlownodeServer {
176 pub fn new(flow_service: FlowService) -> Self {
177 let (tx, _rx) = broadcast::channel::<()>(1);
178 let (server_tx, _server_rx) = broadcast::channel::<()>(1);
179 Self {
180 inner: Arc::new(FlownodeServerInner {
181 flow_service,
182 worker_shutdown_tx: Mutex::new(tx),
183 server_shutdown_tx: Mutex::new(server_tx),
184 streaming_task_handler: Mutex::new(None),
185 state_report_task_handler: Mutex::new(None),
186 }),
187 }
188 }
189
190 async fn start_workers(&self) -> Result<(), Error> {
194 let manager_ref = self.inner.flow_service.dual_engine.clone();
195 let mut state_report_task_handler = self.inner.state_report_task_handler.lock().await;
196 let started_state_report_task = state_report_task_handler.is_none();
197 if state_report_task_handler.is_none() {
198 *state_report_task_handler = manager_ref.clone().start_state_report_task().await;
199 }
200 drop(state_report_task_handler);
201 let handle = manager_ref
202 .streaming_engine()
203 .run_background(Some(self.inner.worker_shutdown_tx.lock().await.subscribe()));
204 self.inner
205 .streaming_task_handler
206 .lock()
207 .await
208 .replace(handle);
209
210 if let Err(err) = self
211 .inner
212 .flow_service
213 .dual_engine
214 .start_flow_consistent_check_task()
215 .await
216 {
217 self.rollback_started_workers(started_state_report_task)
218 .await;
219 return Err(err);
220 }
221
222 Ok(())
223 }
224
225 async fn rollback_started_workers(&self, abort_state_report_task: bool) {
226 let tx = self.inner.worker_shutdown_tx.lock().await;
227 if tx.send(()).is_err() {
228 info!("Receiver dropped, the flow node server has already shutdown");
229 }
230 drop(tx);
231
232 if let Some(handle) = self.inner.streaming_task_handler.lock().await.take() {
233 handle.abort();
234 }
235
236 if abort_state_report_task
237 && let Some(handle) = self.inner.state_report_task_handler.lock().await.take()
238 {
239 handle.abort();
240 }
241 }
242
243 async fn stop_workers(&self) -> Result<(), Error> {
245 let tx = self.inner.worker_shutdown_tx.lock().await;
246 if tx.send(()).is_err() {
247 info!("Receiver dropped, the flow node server has already shutdown");
248 }
249 self.inner
252 .flow_service
253 .dual_engine
254 .stop_flow_consistent_check_task()
255 .await?;
256 Ok(())
257 }
258}
259
260impl FlownodeServer {
261 pub fn create_flow_service(&self) -> flow_server::FlowServer<impl flow_server::Flow> {
262 flow_server::FlowServer::new(self.inner.flow_service.clone())
263 .accept_compressed(CompressionEncoding::Gzip)
264 .send_compressed(CompressionEncoding::Gzip)
265 .accept_compressed(CompressionEncoding::Zstd)
266 .send_compressed(CompressionEncoding::Zstd)
267 }
268}
269
270pub struct FlownodeInstance {
272 flownode_server: FlownodeServer,
273 services: ServerHandlers,
274 heartbeat_task: Option<HeartbeatTask>,
275}
276
277impl FlownodeInstance {
278 pub async fn start(&mut self) -> Result<(), crate::Error> {
279 if let Some(task) = &self.heartbeat_task {
280 task.start().await?;
281 }
282
283 self.flownode_server.start_workers().await?;
284
285 self.services.start_all().await.context(StartServerSnafu)?;
286
287 Ok(())
288 }
289 pub async fn shutdown(&mut self) -> Result<(), Error> {
290 self.services
291 .shutdown_all()
292 .await
293 .context(ShutdownServerSnafu)?;
294
295 self.flownode_server.stop_workers().await?;
296
297 if let Some(task) = &self.heartbeat_task {
298 task.shutdown();
299 }
300
301 Ok(())
302 }
303
304 pub fn flownode_server(&self) -> &FlownodeServer {
305 &self.flownode_server
306 }
307
308 pub fn flow_engine(&self) -> FlowDualEngineRef {
309 self.flownode_server.inner.flow_service.dual_engine.clone()
310 }
311
312 pub fn setup_services(&mut self, services: ServerHandlers) {
313 self.services = services;
314 }
315}
316
317pub struct FlownodeBuilder {
319 opts: FlownodeOptions,
320 plugins: Plugins,
321 table_meta: TableMetadataManagerRef,
322 catalog_manager: CatalogManagerRef,
323 flow_metadata_manager: FlowMetadataManagerRef,
324 heartbeat_task: Option<HeartbeatTask>,
325 state_report_handler: Option<StateReportHandler>,
327 frontend_client: Arc<FrontendClient>,
328}
329
330impl FlownodeBuilder {
331 pub fn new(
333 opts: FlownodeOptions,
334 plugins: Plugins,
335 table_meta: TableMetadataManagerRef,
336 catalog_manager: CatalogManagerRef,
337 flow_metadata_manager: FlowMetadataManagerRef,
338 frontend_client: Arc<FrontendClient>,
339 ) -> Self {
340 Self {
341 opts,
342 plugins,
343 table_meta,
344 catalog_manager,
345 flow_metadata_manager,
346 heartbeat_task: None,
347 state_report_handler: None,
348 frontend_client,
349 }
350 }
351
352 pub fn with_heartbeat_task(self, heartbeat_task: HeartbeatTask) -> Self {
353 let (sender, receiver) = SizeReportSender::new();
354 Self {
355 heartbeat_task: Some(heartbeat_task.with_query_stat_size(sender)),
356 state_report_handler: Some(receiver),
357 ..self
358 }
359 }
360
361 pub fn opts(&self) -> &FlownodeOptions {
362 &self.opts
363 }
364
365 pub fn table_meta(&self) -> &TableMetadataManagerRef {
366 &self.table_meta
367 }
368
369 pub fn catalog_manager(&self) -> &CatalogManagerRef {
370 &self.catalog_manager
371 }
372
373 pub fn flow_metadata_manager(&self) -> &FlowMetadataManagerRef {
374 &self.flow_metadata_manager
375 }
376
377 pub fn frontend_client(&self) -> &Arc<FrontendClient> {
378 &self.frontend_client
379 }
380
381 pub fn set_plugins(&mut self, plugins: Plugins) {
382 self.plugins = plugins;
383 }
384
385 pub async fn build(mut self) -> Result<FlownodeInstance, Error> {
386 let query_engine_factory = QueryEngineFactory::try_new_with_plugins(
388 self.catalog_manager.clone(),
390 None,
391 None,
392 None,
393 None,
394 None,
395 false,
396 Default::default(),
397 self.opts.query.clone(),
398 )
399 .context(DatafusionSnafu {
400 context: "Failed to build query engine",
401 })?;
402 let manager = Arc::new(
403 self.build_manager(query_engine_factory.query_engine())
404 .await?,
405 );
406 let batching = Arc::new(BatchingEngine::new(
407 self.frontend_client.clone(),
408 query_engine_factory.query_engine(),
409 self.flow_metadata_manager.clone(),
410 self.table_meta.clone(),
411 self.catalog_manager.clone(),
412 self.opts.flow.batching_mode.clone(),
413 ));
414 let dual = Arc::new(FlowDualEngine::new(
415 manager.clone(),
416 batching,
417 self.flow_metadata_manager.clone(),
418 self.catalog_manager.clone(),
419 self.plugins.clone(),
420 ));
421 if let Some(handler) = self.state_report_handler.take() {
422 dual.set_state_report_handler(handler).await;
423 }
424
425 let server = FlownodeServer::new(FlowService::new(dual));
426
427 let heartbeat_task = self.heartbeat_task;
428
429 let instance = FlownodeInstance {
430 flownode_server: server,
431 services: ServerHandlers::default(),
432 heartbeat_task,
433 };
434 Ok(instance)
435 }
436
437 async fn build_manager(
440 &mut self,
441 query_engine: Arc<dyn QueryEngine>,
442 ) -> Result<StreamingEngine, Error> {
443 let table_meta = self.table_meta.clone();
444
445 register_function_to_query_engine(&query_engine);
446
447 let num_workers = self.opts.flow.num_workers;
448
449 let node_id = self.opts.node_id.map(|id| id as u32);
450
451 let mut man = StreamingEngine::new(node_id, query_engine, table_meta);
452 for worker_id in 0..num_workers {
453 let (tx, rx) = oneshot::channel();
454
455 let _handle = std::thread::Builder::new()
456 .name(format!("flow-worker-{}", worker_id))
457 .spawn(move || {
458 let (handle, mut worker) = create_worker();
459 let _ = tx.send(handle);
460 info!("Flow Worker started in new thread");
461 worker.run();
462 });
463 let worker_handle = rx.await.map_err(|e| {
464 UnexpectedSnafu {
465 reason: format!("Failed to receive worker handle: {}", e),
466 }
467 .build()
468 })?;
469 man.add_worker_handle(worker_handle);
470 }
471 info!("Flow Node Manager started");
472 Ok(man)
473 }
474}
475
476pub struct FlownodeServiceBuilder<'a> {
478 opts: &'a FlownodeOptions,
479 grpc_server: Option<GrpcServer>,
480 enable_http_service: bool,
481}
482
483impl<'a> FlownodeServiceBuilder<'a> {
484 pub fn new(opts: &'a FlownodeOptions) -> Self {
485 Self {
486 opts,
487 grpc_server: None,
488 enable_http_service: false,
489 }
490 }
491
492 pub fn enable_http_service(self) -> Self {
493 Self {
494 enable_http_service: true,
495 ..self
496 }
497 }
498
499 pub fn with_grpc_server(self, grpc_server: GrpcServer) -> Self {
500 Self {
501 grpc_server: Some(grpc_server),
502 ..self
503 }
504 }
505
506 pub fn with_default_grpc_server(mut self, flownode_server: &FlownodeServer) -> Self {
507 let grpc_server = Self::grpc_server_builder(self.opts, flownode_server).build();
508 self.grpc_server = Some(grpc_server);
509 self
510 }
511
512 pub fn build(mut self) -> Result<ServerHandlers, Error> {
513 let handlers = ServerHandlers::default();
514 if let Some(grpc_server) = self.grpc_server.take() {
515 let addr: SocketAddr = self.opts.grpc.bind_addr.parse().context(ParseAddrSnafu {
516 addr: &self.opts.grpc.bind_addr,
517 })?;
518 let handler: ServerHandler = (Box::new(grpc_server), addr);
519 handlers.insert(handler);
520 }
521
522 if self.enable_http_service {
523 let http_server = HttpServerBuilder::new(self.opts.http.clone())
524 .with_metrics_handler(MetricsHandler)
525 .build();
526 let addr: SocketAddr = self.opts.http.addr.parse().context(ParseAddrSnafu {
527 addr: &self.opts.http.addr,
528 })?;
529 let handler: ServerHandler = (Box::new(http_server), addr);
530 handlers.insert(handler);
531 }
532 Ok(handlers)
533 }
534
535 pub fn grpc_server_builder(
536 opts: &FlownodeOptions,
537 flownode_server: &FlownodeServer,
538 ) -> GrpcServerBuilder {
539 let config = GrpcServerConfig {
540 max_recv_message_size: opts.grpc.max_recv_message_size.as_bytes() as usize,
541 max_send_message_size: opts.grpc.max_send_message_size.as_bytes() as usize,
542 tls: opts.grpc.tls.clone(),
543 max_connection_age: opts.grpc.max_connection_age,
544 };
545 let service = flownode_server.create_flow_service();
546 let runtime = common_runtime::global_runtime();
547 let mut builder = GrpcServerBuilder::new(config, runtime);
548 add_service!(builder, service);
549 builder
550 }
551}
552
553#[derive(Clone)]
558pub struct FrontendInvoker {
559 inserter: Arc<Inserter>,
560 deleter: Arc<Deleter>,
561 statement_executor: Arc<StatementExecutor>,
562}
563
564impl FrontendInvoker {
565 pub fn new(
566 inserter: Arc<Inserter>,
567 deleter: Arc<Deleter>,
568 statement_executor: Arc<StatementExecutor>,
569 ) -> Self {
570 Self {
571 inserter,
572 deleter,
573 statement_executor,
574 }
575 }
576
577 pub async fn build_from(
578 flow_streaming_engine: FlowStreamingEngineRef,
579 catalog_manager: CatalogManagerRef,
580 kv_backend: KvBackendRef,
581 layered_cache_registry: LayeredCacheRegistryRef,
582 procedure_executor: ProcedureExecutorRef,
583 node_manager: NodeManagerRef,
584 origin_frontend_addr: String,
585 ) -> Result<FrontendInvoker, Error> {
586 let table_route_cache: TableRouteCacheRef =
587 layered_cache_registry.get().context(CacheRequiredSnafu {
588 name: TABLE_ROUTE_CACHE_NAME,
589 })?;
590 let partition_info_cache: PartitionInfoCacheRef =
591 layered_cache_registry.get().context(CacheRequiredSnafu {
592 name: PARTITION_INFO_CACHE_NAME,
593 })?;
594
595 let partition_manager = Arc::new(PartitionRuleManager::new(
596 kv_backend.clone(),
597 table_route_cache.clone(),
598 partition_info_cache.clone(),
599 ));
600
601 let table_flownode_cache: TableFlownodeSetCacheRef =
602 layered_cache_registry.get().context(CacheRequiredSnafu {
603 name: TABLE_FLOWNODE_SET_CACHE_NAME,
604 })?;
605
606 let inserter = Arc::new(Inserter::new(
610 catalog_manager.clone(),
611 partition_manager.clone(),
612 node_manager.clone(),
613 table_flownode_cache,
614 true,
615 ));
616
617 let deleter = Arc::new(Deleter::new(
618 catalog_manager.clone(),
619 partition_manager.clone(),
620 node_manager.clone(),
621 ));
622
623 let query_engine = flow_streaming_engine.query_engine.clone();
624
625 let statement_executor = Arc::new(StatementExecutor::new(
626 catalog_manager.clone(),
627 query_engine.clone(),
628 procedure_executor.clone(),
629 kv_backend.clone(),
630 layered_cache_registry.clone(),
631 inserter.clone(),
632 partition_manager,
633 None,
634 origin_frontend_addr,
635 ));
636
637 let invoker = FrontendInvoker::new(inserter, deleter, statement_executor);
638 Ok(invoker)
639 }
640}
641
642impl FrontendInvoker {
643 pub async fn row_inserts(
644 &self,
645 requests: RowInsertRequests,
646 ctx: QueryContextRef,
647 ) -> common_frontend::error::Result<Output> {
648 let _timer = METRIC_FLOW_PROCESSING_TIME
649 .with_label_values(&["output_insert"])
650 .start_timer();
651
652 self.inserter
653 .handle_row_inserts(requests, ctx, &self.statement_executor, false, false)
654 .await
655 .map_err(BoxedError::new)
656 .context(common_frontend::error::ExternalSnafu)
657 }
658
659 pub async fn row_deletes(
660 &self,
661 requests: RowDeleteRequests,
662 ctx: QueryContextRef,
663 ) -> common_frontend::error::Result<Output> {
664 let _timer = METRIC_FLOW_PROCESSING_TIME
665 .with_label_values(&["output_delete"])
666 .start_timer();
667
668 self.deleter
669 .handle_row_deletes(requests, ctx)
670 .await
671 .map_err(BoxedError::new)
672 .context(common_frontend::error::ExternalSnafu)
673 }
674
675 pub fn statement_executor(&self) -> Arc<StatementExecutor> {
676 self.statement_executor.clone()
677 }
678}
679
680pub(crate) async fn get_all_flow_ids(
682 flow_metadata_manager: &FlowMetadataManagerRef,
683 catalog_manager: &CatalogManagerRef,
684 nodeid: Option<u64>,
685) -> Result<Vec<u32>, Error> {
686 let ret = if let Some(nodeid) = nodeid {
687 let flow_ids_one_node = flow_metadata_manager
688 .flownode_flow_manager()
689 .flows(nodeid)
690 .try_collect::<Vec<_>>()
691 .await
692 .context(ListFlowsSnafu { id: Some(nodeid) })?;
693 flow_ids_one_node.into_iter().map(|(id, _)| id).collect()
694 } else {
695 let all_catalogs = catalog_manager
696 .catalog_names()
697 .await
698 .map_err(BoxedError::new)
699 .context(ExternalSnafu)?;
700 let mut all_flow_ids = vec![];
701 for catalog in all_catalogs {
702 let flows = flow_metadata_manager
703 .flow_name_manager()
704 .flow_names(&catalog)
705 .await
706 .try_collect::<Vec<_>>()
707 .await
708 .map_err(BoxedError::new)
709 .context(ExternalSnafu)?;
710
711 all_flow_ids.extend(flows.into_iter().map(|(_, id)| id.flow_id()));
712 }
713 all_flow_ids
714 };
715
716 Ok(ret)
717}
718
719#[cfg(test)]
720mod tests {
721 use std::sync::Arc;
722 use std::time::Duration;
723
724 use api::v1::meta::Role;
725 use catalog::memory::new_memory_catalog_manager;
726 use common_base::Plugins;
727 use common_meta::key::TableMetadataManager;
728 use common_meta::key::flow::FlowMetadataManager;
729 use common_meta::kv_backend::memory::MemoryKvBackend;
730 use meta_client::client::MetaClient;
731 use query::options::QueryOptions;
732
733 use super::*;
734 use crate::adapter::flownode_impl::FlowDualEngine;
735 use crate::batching_mode::BatchingModeOptions;
736 use crate::batching_mode::engine::BatchingEngine;
737 use crate::utils::SizeReportSender;
738
739 async fn new_test_flownode_server() -> (FlownodeServer, SizeReportSender) {
740 let (frontend_client, _handler) =
741 FrontendClient::from_empty_grpc_handler(QueryOptions::default());
742
743 new_test_flownode_server_with_frontend_client(
744 frontend_client,
745 BatchingModeOptions::default(),
746 None,
747 )
748 .await
749 }
750
751 async fn new_test_flownode_server_with_frontend_client(
752 frontend_client: FrontendClient,
753 batching_opts: BatchingModeOptions,
754 node_id: Option<u32>,
755 ) -> (FlownodeServer, SizeReportSender) {
756 let kv_backend = Arc::new(MemoryKvBackend::new());
757 let table_meta = Arc::new(TableMetadataManager::new(kv_backend.clone()));
758 table_meta.init().await.unwrap();
759 let flow_meta = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
760 let catalog_manager = new_memory_catalog_manager().unwrap();
761 let query_engine = crate::test_utils::create_test_query_engine();
762
763 let streaming_engine = Arc::new(StreamingEngine::new(
764 node_id,
765 query_engine.clone(),
766 table_meta.clone(),
767 ));
768 let batching_engine = Arc::new(BatchingEngine::new(
769 Arc::new(frontend_client),
770 query_engine,
771 flow_meta.clone(),
772 table_meta,
773 catalog_manager.clone(),
774 batching_opts,
775 ));
776 let dual_engine = Arc::new(FlowDualEngine::new(
777 streaming_engine,
778 batching_engine,
779 flow_meta,
780 catalog_manager,
781 Plugins::new(),
782 ));
783
784 let (report_sender, report_handler) = SizeReportSender::new();
785 dual_engine.set_state_report_handler(report_handler).await;
786
787 let server = FlownodeServer::new(FlowService::new(dual_engine));
788 (server, report_sender)
789 }
790
791 #[tokio::test]
792 async fn test_state_report_handler_survives_worker_restart() {
793 let (server, report_sender) = new_test_flownode_server().await;
794
795 server.start_workers().await.unwrap();
796 report_sender.query(Duration::from_secs(3)).await.unwrap();
797
798 server.stop_workers().await.unwrap();
799 report_sender.query(Duration::from_secs(3)).await.unwrap();
800
801 server.start_workers().await.unwrap();
802 report_sender.query(Duration::from_secs(3)).await.unwrap();
803
804 server.stop_workers().await.unwrap();
805 }
806
807 #[tokio::test]
808 async fn test_start_workers_rolls_back_on_check_task_start_failure() {
809 let batching_opts = BatchingModeOptions {
810 experimental_frontend_scan_timeout: Duration::from_millis(1),
811 ..Default::default()
812 };
813 let frontend_client = FrontendClient::from_meta_client(
814 Arc::new(MetaClient::new(0, Role::Frontend)),
815 QueryOptions::default(),
816 batching_opts.clone(),
817 )
818 .unwrap();
819 let (server, _report_sender) =
820 new_test_flownode_server_with_frontend_client(frontend_client, batching_opts, Some(1))
821 .await;
822
823 server.start_workers().await.unwrap_err();
824
825 assert!(server.inner.streaming_task_handler.lock().await.is_none());
826 assert!(
827 server
828 .inner
829 .state_report_task_handler
830 .lock()
831 .await
832 .is_none()
833 );
834 }
835}