meta_srv/handler/
flow_state_handler.rs1use api::v1::meta::{FlowStat, HeartbeatRequest, Role};
16use common_meta::key::flow::flow_state::{FlowStateManager, FlowStateValue};
17use common_telemetry::debug;
18use snafu::ResultExt;
19
20use crate::error::{FlowStateHandlerSnafu, Result};
21use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
22use crate::metasrv::Context;
23
24fn node_identity(req: &HeartbeatRequest) -> Option<u64> {
32 req.header
33 .as_ref()
34 .map(|header| header.member_id)
35 .or_else(|| req.peer.as_ref().map(|peer| peer.id))
36}
37
38pub struct FlowStateHandler {
39 flow_state_manager: FlowStateManager,
40}
41
42impl FlowStateHandler {
43 pub fn new(flow_state_manager: FlowStateManager) -> Self {
44 Self { flow_state_manager }
45 }
46}
47
48#[async_trait::async_trait]
49impl HeartbeatHandler for FlowStateHandler {
50 fn is_acceptable(&self, role: Role) -> bool {
51 role == Role::Flownode
52 }
53
54 async fn handle(
55 &self,
56 req: &HeartbeatRequest,
57 _ctx: &mut Context,
58 _acc: &mut HeartbeatAccumulator,
59 ) -> Result<HandleControl> {
60 if let Some(FlowStat {
61 flow_stat_size,
62 flow_last_exec_time_map,
63 }) = &req.flow_stat
64 {
65 let state_size = flow_stat_size
66 .iter()
67 .map(|(k, v)| (*k, *v as usize))
68 .collect();
69 let last_exec_time_map = flow_last_exec_time_map
70 .iter()
71 .map(|(k, v)| (*k, *v))
72 .collect();
73 let value: FlowStateValue =
77 FlowStateValue::new(state_size, last_exec_time_map, Default::default());
78 match node_identity(req) {
79 Some(node_id) => {
80 self.flow_state_manager
83 .merge(node_id, value)
84 .await
85 .context(FlowStateHandlerSnafu)?;
86 }
87 None => {
94 debug!(
95 "Ignore flow state report without node identity (no header.member_id and no peer.id): {value:?}"
96 );
97 }
98 }
99 }
100 Ok(HandleControl::Continue)
101 }
102}
103
104#[cfg(test)]
105mod tests {
106 use std::collections::HashMap;
107
108 use api::v1::meta::{HeartbeatRequest, Peer, RequestHeader, Role};
109
110 use super::*;
111 use crate::handler::test_utils::TestEnv;
112
113 #[test]
114 fn test_node_identity_prefers_header_member_id() {
115 let req = HeartbeatRequest {
116 header: Some(RequestHeader::new(42, Role::Flownode, HashMap::new())),
117 peer: Some(Peer {
118 id: 99,
119 addr: "127.0.0.1:4001".to_string(),
120 }),
121 ..Default::default()
122 };
123 assert_eq!(node_identity(&req), Some(42));
124 }
125
126 #[test]
127 fn test_node_identity_falls_back_to_peer_id() {
128 let req = HeartbeatRequest {
129 header: None,
130 peer: Some(Peer {
131 id: 99,
132 addr: "127.0.0.1:4001".to_string(),
133 }),
134 ..Default::default()
135 };
136 assert_eq!(node_identity(&req), Some(99));
137 }
138
139 #[test]
140 fn test_node_identity_none_without_header_and_peer() {
141 let req = HeartbeatRequest::default();
142 assert_eq!(node_identity(&req), None);
143 }
144
145 fn flow_stat() -> api::v1::meta::FlowStat {
146 api::v1::meta::FlowStat {
147 flow_stat_size: HashMap::from([(1, 1024)]),
148 flow_last_exec_time_map: HashMap::from([(1, 100)]),
149 }
150 }
151
152 #[tokio::test]
153 async fn test_handle_merges_reports_from_different_nodes() {
154 let env = TestEnv::new();
155 let ctx = env.ctx();
156 let flow_state_manager = FlowStateManager::new(ctx.in_memory.clone().as_kv_backend_ref());
157 let handler = FlowStateHandler::new(flow_state_manager);
158
159 let req_a = HeartbeatRequest {
161 header: Some(RequestHeader::new(42, Role::Flownode, HashMap::new())),
162 peer: Some(Peer {
163 id: 42,
164 addr: "127.0.0.1:4001".to_string(),
165 }),
166 flow_stat: Some(flow_stat()),
167 ..Default::default()
168 };
169 let mut ctx = env.ctx();
170 let mut acc = HeartbeatAccumulator::default();
171 handler.handle(&req_a, &mut ctx, &mut acc).await.unwrap();
172
173 let req_b = HeartbeatRequest {
175 header: None,
176 peer: Some(Peer {
177 id: 7,
178 addr: "127.0.0.1:4007".to_string(),
179 }),
180 flow_stat: Some(api::v1::meta::FlowStat {
181 flow_stat_size: HashMap::from([(2, 2048)]),
182 flow_last_exec_time_map: HashMap::from([(2, 200)]),
183 }),
184 ..Default::default()
185 };
186 let mut ctx = env.ctx();
187 let mut acc = HeartbeatAccumulator::default();
188 handler.handle(&req_b, &mut ctx, &mut acc).await.unwrap();
189
190 let value = handler.flow_state_manager.get().await.unwrap().unwrap();
191 assert_eq!(value.last_exec_time_map.get(&1), Some(&100));
192 assert_eq!(value.last_exec_time_map.get(&2), Some(&200));
193 }
194
195 #[tokio::test]
196 async fn test_handle_ignores_report_without_identity() {
197 let env = TestEnv::new();
198 let ctx = env.ctx();
199 let flow_state_manager = FlowStateManager::new(ctx.in_memory.clone().as_kv_backend_ref());
200 let handler = FlowStateHandler::new(flow_state_manager);
201
202 let req = HeartbeatRequest {
204 header: None,
205 peer: None,
206 flow_stat: Some(flow_stat()),
207 ..Default::default()
208 };
209 let mut ctx = env.ctx();
210 let mut acc = HeartbeatAccumulator::default();
211 handler.handle(&req, &mut ctx, &mut acc).await.unwrap();
212
213 assert!(handler.flow_state_manager.get().await.unwrap().is_none());
214 }
215}