Skip to main content

servers/batcher/
flow_sender.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;
16
17use api::v1::flow::{DirtyWindowRequest, DirtyWindowRequests};
18use common_meta::cache::TableFlownodeSetCacheRef;
19use common_meta::node_manager::NodeManagerRef;
20use common_telemetry::error;
21use store_api::storage::TableId;
22
23/// One source table's timestamps, expressed in its native time unit.
24pub(in crate::batcher) struct FlowNotification {
25    pub table_id: TableId,
26    pub timestamps: Vec<i64>,
27}
28
29/// Resolves Flow targets and delivers dirty-window notifications best-effort.
30/// Queue admission, delivery concurrency and write completion belong to callers.
31#[derive(Clone)]
32pub(in crate::batcher) struct FlowSender {
33    cache: TableFlownodeSetCacheRef,
34    node_manager: NodeManagerRef,
35}
36
37impl FlowSender {
38    /// Binds delivery dependencies without creating a queue or starting a task.
39    pub fn new(cache: TableFlownodeSetCacheRef, node_manager: NodeManagerRef) -> Self {
40        Self {
41            cache,
42            node_manager,
43        }
44    }
45
46    /// Sends once per distinct peer, without retrying or failing the completed write.
47    /// Peer RPCs stay within the caller's notification concurrency budget.
48    pub async fn send(&self, notification: FlowNotification) {
49        let table_id = notification.table_id;
50        let flownodes = match self.cache.get(table_id).await {
51            Ok(Some(flownodes)) => flownodes,
52            Ok(None) => return,
53            Err(e) => {
54                error!(e; "Failed to get flownodes for table id: {}", table_id);
55                return;
56            }
57        };
58        let peers = flownodes.values().cloned().collect::<HashSet<_>>();
59
60        for peer in peers {
61            if let Err(e) = self
62                .node_manager
63                .flownode(&peer)
64                .await
65                .handle_mark_window_dirty(DirtyWindowRequests {
66                    requests: vec![DirtyWindowRequest {
67                        table_id,
68                        timestamps: notification.timestamps.clone(),
69                        time_ranges: Vec::new(),
70                    }],
71                })
72                .await
73            {
74                error!(
75                    e;
76                    "Failed to mark timestamps as dirty, table_id: {}, peer_id: {}, peer_addr: {}",
77                    table_id,
78                    peer.id,
79                    peer.addr
80                );
81            }
82        }
83    }
84}
85
86#[cfg(test)]
87mod tests {
88    use std::sync::{Arc, Mutex};
89
90    use api::v1::flow::{DirtyWindowRequest, DirtyWindowRequests, FlowRequest, FlowResponse};
91    use api::v1::region::InsertRequests;
92    use async_trait::async_trait;
93    use common_meta::error::{Result as MetaResult, UnexpectedSnafu};
94    use common_meta::node_manager::{
95        DatanodeManager, DatanodeRef, Flownode, FlownodeManager, FlownodeRef,
96    };
97    use common_meta::peer::Peer;
98
99    use crate::batcher::flow_sender::{FlowNotification, FlowSender};
100    use crate::batcher::test_util::mock_table_flownode_cache;
101
102    #[derive(Default)]
103    struct RecordingNodeManager {
104        requests: Arc<Mutex<Vec<(u64, DirtyWindowRequests)>>>,
105        fail_first: bool,
106    }
107
108    struct RecordingFlownode {
109        peer_id: u64,
110        requests: Arc<Mutex<Vec<(u64, DirtyWindowRequests)>>>,
111        fail_first: bool,
112    }
113
114    #[async_trait]
115    impl Flownode for RecordingFlownode {
116        async fn handle(&self, _: FlowRequest) -> MetaResult<FlowResponse> {
117            unreachable!("notifications must use the dirty-window RPC")
118        }
119
120        async fn handle_inserts(&self, _: InsertRequests) -> MetaResult<FlowResponse> {
121            unreachable!("notifications must not mirror row inserts")
122        }
123
124        async fn handle_mark_window_dirty(
125            &self,
126            request: DirtyWindowRequests,
127        ) -> MetaResult<FlowResponse> {
128            let mut requests = self.requests.lock().unwrap();
129            requests.push((self.peer_id, request));
130            if self.fail_first && requests.len() == 1 {
131                return UnexpectedSnafu {
132                    err_msg: "injected first notification failure",
133                }
134                .fail();
135            }
136            Ok(FlowResponse::default())
137        }
138    }
139
140    #[async_trait]
141    impl DatanodeManager for RecordingNodeManager {
142        async fn datanode(&self, _: &Peer) -> DatanodeRef {
143            unreachable!("flow notifications must not contact datanodes")
144        }
145    }
146
147    #[async_trait]
148    impl FlownodeManager for RecordingNodeManager {
149        async fn flownode(&self, peer: &Peer) -> FlownodeRef {
150            Arc::new(RecordingFlownode {
151                peer_id: peer.id,
152                requests: self.requests.clone(),
153                fail_first: self.fail_first,
154            })
155        }
156    }
157
158    #[tokio::test]
159    async fn test_deduplicated_delivery_preserves_payload_and_continues_after_error() {
160        // HashSet peer order is unspecified, so fail the first attempted RPC rather
161        // than a chosen peer. The second peer must still receive the notification.
162        for fail_first in [false, true] {
163            let first = Peer {
164                id: 1,
165                addr: "flow-1".to_string(),
166            };
167            let second = Peer {
168                id: 2,
169                addr: "flow-2".to_string(),
170            };
171            let cache =
172                mock_table_flownode_cache(42, vec![(0, first.clone()), (1, first), (2, second)])
173                    .await;
174            let manager = Arc::new(RecordingNodeManager {
175                fail_first,
176                ..Default::default()
177            });
178            let sender = FlowSender::new(cache, manager.clone());
179            let timestamps = vec![-1, 42, 1_700_000_000_000_000_001, 42];
180            sender
181                .send(FlowNotification {
182                    table_id: 42,
183                    timestamps: timestamps.clone(),
184                })
185                .await;
186            let expected = DirtyWindowRequests {
187                requests: vec![DirtyWindowRequest {
188                    table_id: 42,
189                    timestamps,
190                    time_ranges: vec![],
191                }],
192            };
193            let mut actual = manager.requests.lock().unwrap().clone();
194            actual.sort_by_key(|(peer, _)| *peer);
195            assert_eq!(actual, vec![(1, expected.clone()), (2, expected)]);
196        }
197    }
198
199    #[tokio::test]
200    async fn test_missing_targets_do_not_send() {
201        let cache = mock_table_flownode_cache(42, vec![]).await;
202        let manager = Arc::new(RecordingNodeManager::default());
203        let sender = FlowSender::new(cache, manager.clone());
204        sender
205            .send(FlowNotification {
206                table_id: 42,
207                timestamps: vec![42],
208            })
209            .await;
210        assert!(manager.requests.lock().unwrap().is_empty());
211    }
212}