Skip to main content

servers/batcher/logical_table/
batch.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashSet;
16use std::sync::Arc;
17use std::time::Instant;
18
19use catalog::CatalogManagerRef;
20use common_batcher::flush_limiter::FlushLimiter;
21use common_meta::node_manager::NodeManagerRef;
22use common_query::prelude::GREPTIME_PHYSICAL_TABLE;
23use common_telemetry::{debug, warn};
24use partition::manager::PartitionRuleManagerRef;
25use session::context::QueryContextRef;
26use snafu::OptionExt;
27
28use crate::batcher::flow_notifier::FlowNotifier;
29use crate::batcher::logical_table::PHYSICAL_TABLE_KEY;
30use crate::batcher::logical_table::batch_convert::{
31    TableBatch, concat_modified_batches, transform_logical_batches_to_physical,
32};
33use crate::batcher::logical_table::flow_notifier::enqueue_flow_notifications;
34use crate::batcher::logical_table::pending_worker::FlushWaiter;
35use crate::batcher::logical_table::region_write::{
36    CatalogManagerPhysicalFlushAdapter, NodeManagerPhysicalFlushAdapter,
37    PartitionManagerPhysicalFlushAdapter, PhysicalFlushCatalogProvider, PhysicalFlushNodeRequester,
38    PhysicalFlushPartitionProvider, encode_region_write_requests, flush_region_writes_concurrently,
39    plan_region_batches, resolve_region_targets,
40};
41use crate::error;
42use crate::error::Result;
43use crate::metrics::{
44    FLUSH_DROPPED_ROWS, FLUSH_ELAPSED, FLUSH_FAILURES, FLUSH_ROWS, FLUSH_TOTAL,
45    PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED,
46};
47
48pub(in crate::batcher::logical_table) struct Batch {
49    pub(in crate::batcher::logical_table) table_batches: Vec<TableBatch>,
50    pub(in crate::batcher::logical_table) total_row_count: usize,
51    pub(in crate::batcher::logical_table) db_string: String,
52    pub(in crate::batcher::logical_table) ctx: QueryContextRef,
53    pub(in crate::batcher::logical_table) waiters: Vec<FlushWaiter>,
54}
55
56pub(in crate::batcher::logical_table) async fn spawn_flush(
57    flush: Batch,
58    partition_manager: PartitionRuleManagerRef,
59    node_manager: NodeManagerRef,
60    catalog_manager: CatalogManagerRef,
61    flow_notification_tx: FlowNotifier,
62    flush_limiter: FlushLimiter,
63) {
64    match flush_limiter.acquire().await {
65        Ok(permit) => {
66            tokio::spawn(async move {
67                let _permit = permit;
68                flush_batch_with_managers(
69                    flush,
70                    partition_manager,
71                    node_manager,
72                    catalog_manager,
73                    flow_notification_tx,
74                )
75                .await;
76            });
77        }
78        Err(err) => {
79            warn!(err; "Flush semaphore closed, flushing inline");
80            flush_batch_with_managers(
81                flush,
82                partition_manager,
83                node_manager,
84                catalog_manager,
85                flow_notification_tx,
86            )
87            .await;
88        }
89    }
90}
91
92pub(in crate::batcher::logical_table) async fn flush_batch_with_managers(
93    flush: Batch,
94    partition_manager: PartitionRuleManagerRef,
95    node_manager: NodeManagerRef,
96    catalog_manager: CatalogManagerRef,
97    flow_notification_tx: FlowNotifier,
98) {
99    let partition_provider = PartitionManagerPhysicalFlushAdapter { partition_manager };
100    let node_requester = NodeManagerPhysicalFlushAdapter {
101        node_manager: node_manager.clone(),
102    };
103    let catalog_provider = CatalogManagerPhysicalFlushAdapter { catalog_manager };
104    flush_batch(
105        flush,
106        &partition_provider,
107        &node_requester,
108        &catalog_provider,
109        flow_notification_tx,
110    )
111    .await;
112}
113
114pub(in crate::batcher::logical_table) async fn flush_batch(
115    flush: Batch,
116    partition_manager: &(impl PhysicalFlushPartitionProvider + ?Sized),
117    node_manager: &(impl PhysicalFlushNodeRequester + ?Sized),
118    catalog_manager: &(impl PhysicalFlushCatalogProvider + ?Sized),
119    flow_notification_tx: FlowNotifier,
120) {
121    let Batch {
122        table_batches,
123        total_row_count,
124        db_string,
125        ctx,
126        waiters,
127    } = flush;
128    let start = Instant::now();
129
130    // Physical-table-level flush: transform all logical table batches
131    // into physical format and write them together.
132    let physical_table_name = ctx
133        .extension(PHYSICAL_TABLE_KEY)
134        .unwrap_or(GREPTIME_PHYSICAL_TABLE)
135        .to_string();
136    let result = flush_batch_physical(
137        &table_batches,
138        &physical_table_name,
139        &ctx,
140        partition_manager,
141        node_manager,
142        catalog_manager,
143    )
144    .await;
145
146    let elapsed = start.elapsed().as_secs_f64();
147    FLUSH_ELAPSED.observe(elapsed);
148
149    debug!(
150        "Pending rows batch flushed, total rows: {}, elapsed time: {}s",
151        total_row_count, elapsed
152    );
153
154    match result {
155        Ok(affected_rows) => {
156            FLUSH_TOTAL.inc();
157            FLUSH_ROWS.observe(total_row_count as f64);
158            operator::metrics::DIST_INGEST_ROW_COUNT
159                .with_label_values(&[db_string.as_str()])
160                .inc_by(affected_rows as u64);
161
162            notify_waiters(waiters, Ok(()));
163            enqueue_flow_notifications(table_batches, &flow_notification_tx);
164        }
165        Err(err) => {
166            FLUSH_FAILURES.inc();
167            FLUSH_DROPPED_ROWS.inc_by(total_row_count as u64);
168            notify_waiters(waiters, Err(err));
169        }
170    }
171}
172
173/// Flushes a batch of logical table rows by transforming them into the physical table format
174/// and writing them to the appropriate datanode regions.
175///
176/// This function performs the end-to-end physical flush pipeline:
177/// 1. Resolves the physical table metadata and column ID mapping.
178/// 2. Fetches the physical table's partition rule.
179/// 3. Transforms each logical table batch into the physical (sparse primary key) format.
180/// 4. Concatenates all transformed batches into a single combined batch.
181/// 5. Splits the combined batch by partition rule and sends region write requests
182///    concurrently to the target datanodes.
183pub async fn flush_batch_physical(
184    table_batches: &[TableBatch],
185    physical_table_name: &str,
186    ctx: &QueryContextRef,
187    partition_manager: &(impl PhysicalFlushPartitionProvider + ?Sized),
188    node_manager: &(impl PhysicalFlushNodeRequester + ?Sized),
189    catalog_manager: &(impl PhysicalFlushCatalogProvider + ?Sized),
190) -> Result<usize> {
191    // 1. Resolve the physical table and get column ID mapping
192    let physical_table = {
193        let _timer = PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED
194            .with_label_values(&["flush_physical_resolve_table"])
195            .start_timer();
196        catalog_manager
197            .physical_table(
198                ctx.current_catalog(),
199                &ctx.current_schema(),
200                physical_table_name,
201                ctx.as_ref(),
202            )
203            .await?
204            .with_context(|| error::InternalSnafu {
205                err_msg: format!(
206                    "Physical table '{}' not found during pending flush",
207                    physical_table_name
208                ),
209            })?
210    };
211
212    let physical_table_info = physical_table.table_info;
213    let name_to_ids = physical_table
214        .col_name_to_ids
215        .with_context(|| error::InternalSnafu {
216            err_msg: format!(
217                "Physical table '{}' has no column IDs for pending flush",
218                physical_table_name
219            ),
220        })?;
221
222    // 2. Get the physical table's partition rule (one lookup instead of N)
223    let partition_rule = {
224        let _timer = PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED
225            .with_label_values(&["flush_physical_fetch_partition_rule"])
226            .start_timer();
227        partition_manager
228            .find_table_partition_rule(physical_table_info.as_ref())
229            .await?
230    };
231    let partition_columns = partition_rule.partition_columns();
232    let partition_columns_set: HashSet<&str> =
233        partition_columns.iter().map(String::as_str).collect();
234
235    // 3. Transform each logical table batch into physical format
236    let modified_batches =
237        transform_logical_batches_to_physical(table_batches, &name_to_ids, &partition_columns_set)?;
238
239    // 4. Concatenate all modified batches (all share the same physical schema)
240    let combined_batch = concat_modified_batches(&modified_batches)?;
241
242    // 5. Split by physical partition rule and send to regions
243    let physical_table_id = physical_table_info.table_id();
244    let planned_batches = plan_region_batches(
245        combined_batch,
246        physical_table_id,
247        partition_rule.as_ref(),
248        partition_columns,
249    )?;
250
251    let resolved_batches = resolve_region_targets(planned_batches, partition_manager).await?;
252    let region_writes = encode_region_write_requests(resolved_batches, ctx.skip_wal())?;
253    flush_region_writes_concurrently(node_manager, region_writes).await
254}
255
256pub(in crate::batcher::logical_table) fn notify_waiters(
257    waiters: Vec<FlushWaiter>,
258    result: Result<()>,
259) {
260    let shared_result = result.map_err(Arc::new);
261    for waiter in waiters {
262        let _ = waiter.response_tx.send(match &shared_result {
263            Ok(()) => Ok(()),
264            Err(error) => Err(Arc::clone(error)),
265        });
266        // waiter._permit is dropped here, releasing the inflight semaphore slot
267    }
268}
269
270#[cfg(test)]
271mod tests {
272    use std::collections::HashMap;
273    use std::future::poll_fn;
274    use std::sync::Arc;
275    use std::sync::atomic::{AtomicUsize, Ordering};
276    use std::task::Poll;
277    use std::time::Duration;
278
279    use api::region::RegionResponse;
280    use api::v1::meta::Peer;
281    use api::v1::region::RegionRequest;
282    use arrow::record_batch::RecordBatch;
283    use async_trait::async_trait;
284    use catalog::error::Result as CatalogResult;
285    use common_batcher::request_limiter::RequestLimiter;
286    use common_meta::cache::TableFlownodeSetCacheRef;
287    use common_meta::node_manager::NodeManagerRef;
288    use datatypes::schema::{ColumnSchema as DtColumnSchema, Schema as DtSchema};
289    use partition::error::Result as PartitionResult;
290    use partition::partition::{PartitionRule, PartitionRuleRef, RegionMask};
291    use store_api::storage::RegionId;
292    use table::metadata::TableId;
293    use table::test_util::table_info::test_table_info;
294    use tokio::sync::{mpsc, oneshot};
295
296    use crate::batcher::flow_notifier::{FlowNotifier, start_flow_notification_worker};
297    use crate::batcher::logical_table::batch::{
298        Batch, flush_batch, flush_batch_physical, notify_waiters,
299    };
300    use crate::batcher::logical_table::batch_convert::TableBatch;
301    use crate::batcher::logical_table::pending_worker::FlushWaiter;
302    use crate::batcher::logical_table::region_write::{
303        PhysicalFlushCatalogProvider, PhysicalFlushNodeRequester, PhysicalFlushPartitionProvider,
304        PhysicalTableMetadata,
305    };
306    use crate::batcher::logical_table::test_util::{
307        FlowNotificationMockNodeManager, RecordingFlownode, mock_aligned_tag_batch,
308    };
309    use crate::batcher::test_util::mock_table_flownode_cache;
310    use crate::error;
311    use crate::error::Error;
312    use crate::metrics::FLOW_NOTIFICATION_DROPPED;
313
314    #[tokio::test]
315    async fn test_flush_batch_notifies_flownode_after_successful_physical_write() {
316        let table_id = 42;
317        let peer = Peer {
318            id: 7,
319            addr: "flow-7".to_string(),
320        };
321        let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
322        let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
323        let _requests_tx = requests_tx.clone();
324        let flow_node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
325            flownode: Arc::new(RecordingFlownode { requests_tx }),
326        });
327        let flow_notification_tx = mock_flow_notification_sender(cache, flow_node_manager.clone());
328        let ctx = session::context::QueryContext::arc();
329        let table_batches = vec![TableBatch {
330            table_name: "cpu".to_string(),
331            table_id,
332            batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
333            row_count: 1,
334        }];
335        let writes = Arc::new(AtomicUsize::new(0));
336
337        flush_batch(
338            Batch {
339                table_batches,
340                total_row_count: 1,
341                db_string: ctx.get_db_string(),
342                ctx,
343                waiters: Vec::new(),
344            },
345            &MockFlushPartitionProvider {
346                partition_rule_calls: Arc::new(AtomicUsize::new(0)),
347                region_leader_calls: Arc::new(AtomicUsize::new(0)),
348            },
349            &MockFlushNodeRequester {
350                writes: writes.clone(),
351                fail: false,
352            },
353            &MockFlushCatalogProvider {
354                table: Some(mock_physical_table_metadata(1024)),
355            },
356            flow_notification_tx,
357        )
358        .await;
359
360        assert_eq!(1, writes.load(Ordering::SeqCst));
361        let requests = tokio::time::timeout(Duration::from_secs(1), requests_rx.recv())
362            .await
363            .unwrap()
364            .unwrap();
365        assert_eq!(
366            vec![api::v1::flow::DirtyWindowRequest {
367                table_id,
368                timestamps: vec![1000],
369                time_ranges: Vec::new(),
370            }],
371            requests.requests
372        );
373    }
374
375    #[tokio::test]
376    async fn test_flush_batch_does_not_notify_flownode_after_physical_write_error() {
377        let table_id = 42;
378        let peer = Peer {
379            id: 7,
380            addr: "flow-7".to_string(),
381        };
382        let cache = mock_table_flownode_cache(table_id, vec![(0, peer.clone()), (1, peer)]).await;
383        let (requests_tx, mut requests_rx) = mpsc::unbounded_channel();
384        let _requests_tx = requests_tx.clone();
385        let flow_node_manager: NodeManagerRef = Arc::new(FlowNotificationMockNodeManager {
386            flownode: Arc::new(RecordingFlownode { requests_tx }),
387        });
388        let flow_notification_tx = mock_flow_notification_sender(cache, flow_node_manager.clone());
389        let ctx = session::context::QueryContext::arc();
390        let table_batches = vec![TableBatch {
391            table_name: "cpu".to_string(),
392            table_id,
393            batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
394            row_count: 1,
395        }];
396        let writes = Arc::new(AtomicUsize::new(0));
397
398        flush_batch(
399            Batch {
400                table_batches,
401                total_row_count: 1,
402                db_string: ctx.get_db_string(),
403                ctx,
404                waiters: Vec::new(),
405            },
406            &MockFlushPartitionProvider {
407                partition_rule_calls: Arc::new(AtomicUsize::new(0)),
408                region_leader_calls: Arc::new(AtomicUsize::new(0)),
409            },
410            &MockFlushNodeRequester {
411                writes: writes.clone(),
412                fail: true,
413            },
414            &MockFlushCatalogProvider {
415                table: Some(mock_physical_table_metadata(1024)),
416            },
417            flow_notification_tx,
418        )
419        .await;
420
421        assert_eq!(1, writes.load(Ordering::SeqCst));
422        assert!(
423            tokio::time::timeout(Duration::from_millis(50), requests_rx.recv())
424                .await
425                .is_err()
426        );
427    }
428
429    #[tokio::test]
430    async fn test_cancelled_waiter_retains_request_slot_until_notification() {
431        let limiter = RequestLimiter::try_new(1).unwrap();
432        let (response_tx, response_rx) = oneshot::channel();
433        let waiter = FlushWaiter {
434            response_tx,
435            _permit: limiter.acquire().await.unwrap(),
436        };
437        drop(response_rx);
438        let next = limiter.acquire();
439        tokio::pin!(next);
440        assert!(poll_fn(|cx| Poll::Ready(next.as_mut().poll(cx).is_pending())).await);
441        notify_waiters(vec![waiter], Ok(()));
442        let _permit = next.await.unwrap();
443    }
444
445    #[tokio::test]
446    async fn test_flush_batch_physical_uses_mockable_trait_dependencies() {
447        let table_batches = vec![TableBatch {
448            table_name: "t1".to_string(),
449            table_id: 11,
450            batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
451            row_count: 1,
452        }];
453        let partition_calls = Arc::new(AtomicUsize::new(0));
454        let leader_calls = Arc::new(AtomicUsize::new(0));
455        let node = MockFlushNodeRequester::default();
456        let ctx = session::context::QueryContext::arc();
457
458        flush_batch_physical(
459            &table_batches,
460            "phy",
461            &ctx,
462            &MockFlushPartitionProvider {
463                partition_rule_calls: partition_calls.clone(),
464                region_leader_calls: leader_calls.clone(),
465            },
466            &node,
467            &MockFlushCatalogProvider {
468                table: Some(mock_physical_table_metadata(1024)),
469            },
470        )
471        .await
472        .unwrap();
473
474        assert_eq!(1, partition_calls.load(Ordering::SeqCst));
475        assert_eq!(1, leader_calls.load(Ordering::SeqCst));
476        assert_eq!(1, node.writes.load(Ordering::SeqCst));
477    }
478
479    #[tokio::test]
480    async fn test_flush_batch_physical_returns_actual_affected_rows() {
481        let table_batches = vec![TableBatch {
482            table_name: "t1".to_string(),
483            table_id: 11,
484            batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
485            row_count: 1,
486        }];
487        let ctx = session::context::QueryContext::arc();
488
489        let affected_rows = flush_batch_physical(
490            &table_batches,
491            "phy",
492            &ctx,
493            &MockFlushPartitionProvider {
494                partition_rule_calls: Arc::new(AtomicUsize::new(0)),
495                region_leader_calls: Arc::new(AtomicUsize::new(0)),
496            },
497            &AffectedRowsFlushNodeRequester { affected_rows: 7 },
498            &MockFlushCatalogProvider {
499                table: Some(mock_physical_table_metadata(1024)),
500            },
501        )
502        .await
503        .unwrap();
504
505        assert_eq!(7, affected_rows);
506    }
507
508    #[tokio::test]
509    async fn test_flush_batch_physical_stops_before_partition_and_node_when_table_missing() {
510        let table_batches = vec![TableBatch {
511            table_name: "t1".to_string(),
512            table_id: 11,
513            batches: vec![mock_aligned_tag_batch("tag1", "host-1", 1000, 1.0)],
514            row_count: 1,
515        }];
516        let partition_calls = Arc::new(AtomicUsize::new(0));
517        let leader_calls = Arc::new(AtomicUsize::new(0));
518        let node = MockFlushNodeRequester::default();
519        let ctx = session::context::QueryContext::arc();
520
521        let err = flush_batch_physical(
522            &table_batches,
523            "missing_phy",
524            &ctx,
525            &MockFlushPartitionProvider {
526                partition_rule_calls: partition_calls.clone(),
527                region_leader_calls: leader_calls.clone(),
528            },
529            &node,
530            &MockFlushCatalogProvider { table: None },
531        )
532        .await
533        .unwrap_err();
534
535        assert!(
536            err.to_string()
537                .contains("Physical table 'missing_phy' not found")
538        );
539        assert_eq!(0, partition_calls.load(Ordering::SeqCst));
540        assert_eq!(0, leader_calls.load(Ordering::SeqCst));
541        assert_eq!(0, node.writes.load(Ordering::SeqCst));
542    }
543
544    #[tokio::test]
545    async fn test_flush_batch_physical_aborts_immediately_on_transform_error() {
546        let table_batches = vec![
547            TableBatch {
548                table_name: "broken".to_string(),
549                table_id: 11,
550                batches: vec![mock_aligned_tag_batch("unknown_tag", "host-1", 1000, 1.0)],
551                row_count: 1,
552            },
553            TableBatch {
554                table_name: "healthy".to_string(),
555                table_id: 12,
556                batches: vec![mock_aligned_tag_batch("tag1", "host-2", 2000, 2.0)],
557                row_count: 1,
558            },
559        ];
560        let partition_calls = Arc::new(AtomicUsize::new(0));
561        let leader_calls = Arc::new(AtomicUsize::new(0));
562        let node = MockFlushNodeRequester::default();
563        let ctx = session::context::QueryContext::arc();
564
565        let err = flush_batch_physical(
566            &table_batches,
567            "phy",
568            &ctx,
569            &MockFlushPartitionProvider {
570                partition_rule_calls: partition_calls.clone(),
571                region_leader_calls: leader_calls.clone(),
572            },
573            &node,
574            &MockFlushCatalogProvider {
575                table: Some(mock_physical_table_metadata(1024)),
576            },
577        )
578        .await
579        .unwrap_err();
580
581        assert!(err.to_string().contains("unknown_tag"));
582        assert_eq!(1, partition_calls.load(Ordering::SeqCst));
583        assert_eq!(0, leader_calls.load(Ordering::SeqCst));
584        assert_eq!(0, node.writes.load(Ordering::SeqCst));
585    }
586
587    fn mock_physical_table_metadata(table_id: TableId) -> PhysicalTableMetadata {
588        let schema = Arc::new(
589            DtSchema::try_new(vec![
590                DtColumnSchema::new(
591                    "__primary_key",
592                    datatypes::prelude::ConcreteDataType::binary_datatype(),
593                    false,
594                ),
595                DtColumnSchema::new(
596                    "greptime_timestamp",
597                    datatypes::prelude::ConcreteDataType::timestamp_millisecond_datatype(),
598                    false,
599                ),
600                DtColumnSchema::new(
601                    "greptime_value",
602                    datatypes::prelude::ConcreteDataType::float64_datatype(),
603                    true,
604                ),
605                DtColumnSchema::new(
606                    "tag1",
607                    datatypes::prelude::ConcreteDataType::string_datatype(),
608                    true,
609                ),
610            ])
611            .unwrap(),
612        );
613        let mut table_info = test_table_info(table_id, "phy", "public", "greptime", schema);
614        table_info.meta.column_ids = vec![0, 1, 2, 3];
615
616        PhysicalTableMetadata {
617            table_info: Arc::new(table_info),
618            col_name_to_ids: Some(HashMap::from([("tag1".to_string(), 3)])),
619        }
620    }
621
622    struct MockFlushCatalogProvider {
623        table: Option<PhysicalTableMetadata>,
624    }
625
626    #[async_trait]
627    impl PhysicalFlushCatalogProvider for MockFlushCatalogProvider {
628        async fn physical_table(
629            &self,
630            _catalog: &str,
631            _schema: &str,
632            _table_name: &str,
633            _query_ctx: &session::context::QueryContext,
634        ) -> CatalogResult<Option<PhysicalTableMetadata>> {
635            Ok(self.table.clone())
636        }
637    }
638
639    struct SingleRegionPartitionRule;
640
641    impl PartitionRule for SingleRegionPartitionRule {
642        fn as_any(&self) -> &dyn std::any::Any {
643            self
644        }
645
646        fn partition_columns(&self) -> &[String] {
647            &[]
648        }
649
650        fn find_region(
651            &self,
652            _values: &[datatypes::prelude::Value],
653        ) -> partition::error::Result<store_api::storage::RegionNumber> {
654            unimplemented!()
655        }
656
657        fn split_record_batch(
658            &self,
659            record_batch: &RecordBatch,
660        ) -> partition::error::Result<HashMap<store_api::storage::RegionNumber, RegionMask>>
661        {
662            Ok(HashMap::from([(
663                1,
664                RegionMask::new(
665                    arrow::array::BooleanArray::from(vec![true; record_batch.num_rows()]),
666                    record_batch.num_rows(),
667                ),
668            )]))
669        }
670    }
671
672    struct MockFlushPartitionProvider {
673        partition_rule_calls: Arc<AtomicUsize>,
674        region_leader_calls: Arc<AtomicUsize>,
675    }
676
677    #[async_trait]
678    impl PhysicalFlushPartitionProvider for MockFlushPartitionProvider {
679        async fn find_table_partition_rule(
680            &self,
681            _table_info: &table::metadata::TableInfo,
682        ) -> PartitionResult<PartitionRuleRef> {
683            self.partition_rule_calls.fetch_add(1, Ordering::SeqCst);
684            Ok(Arc::new(SingleRegionPartitionRule))
685        }
686
687        async fn find_region_leader(&self, _region_id: RegionId) -> error::Result<Peer> {
688            self.region_leader_calls.fetch_add(1, Ordering::SeqCst);
689            Ok(Peer {
690                id: 1,
691                addr: "node-1".to_string(),
692            })
693        }
694    }
695
696    #[derive(Default)]
697    struct MockFlushNodeRequester {
698        writes: Arc<AtomicUsize>,
699        fail: bool,
700    }
701
702    #[async_trait]
703    impl PhysicalFlushNodeRequester for MockFlushNodeRequester {
704        async fn handle(
705            &self,
706            _peer: &Peer,
707            _request: RegionRequest,
708        ) -> error::Result<RegionResponse> {
709            self.writes.fetch_add(1, Ordering::SeqCst);
710            if self.fail {
711                return Err(Error::Internal {
712                    err_msg: "physical write failed".to_string(),
713                });
714            }
715            Ok(RegionResponse::new(0))
716        }
717    }
718
719    fn mock_flow_notification_sender(
720        cache: TableFlownodeSetCacheRef,
721        node_manager: NodeManagerRef,
722    ) -> FlowNotifier {
723        let (tx, rx) = FlowNotifier::try_new(16, FLOW_NOTIFICATION_DROPPED.clone()).unwrap();
724        start_flow_notification_worker(rx, cache, node_manager);
725        tx
726    }
727
728    #[derive(Default)]
729    struct AffectedRowsFlushNodeRequester {
730        affected_rows: usize,
731    }
732
733    #[async_trait]
734    impl PhysicalFlushNodeRequester for AffectedRowsFlushNodeRequester {
735        async fn handle(
736            &self,
737            _peer: &Peer,
738            _request: RegionRequest,
739        ) -> error::Result<RegionResponse> {
740            Ok(RegionResponse::new(self.affected_rows))
741        }
742    }
743}