servers/batcher/
flow_sender.rs1use 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
23pub(in crate::batcher) struct FlowNotification {
25 pub table_id: TableId,
26 pub timestamps: Vec<i64>,
27}
28
29#[derive(Clone)]
32pub(in crate::batcher) struct FlowSender {
33 cache: TableFlownodeSetCacheRef,
34 node_manager: NodeManagerRef,
35}
36
37impl FlowSender {
38 pub fn new(cache: TableFlownodeSetCacheRef, node_manager: NodeManagerRef) -> Self {
40 Self {
41 cache,
42 node_manager,
43 }
44 }
45
46 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 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}