Skip to main content

meta_srv/handler/
flow_state_handler.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 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
24/// Extracts the flownode identity from a heartbeat request.
25///
26/// Prefers `header.member_id` — the canonical identity metasrv uses for
27/// flownodes (see `get_node_id` in `src/meta-srv/src/service/heartbeat.rs`).
28/// Falls back to `peer.id`, the operator-configured node id, which is also
29/// unique within the cluster (a flownode requires `node_id` in its config, see
30/// `src/cmd/src/flownode.rs`). Returns `None` when neither is present.
31fn 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            // TODO(#7987-followup): start_time_map is not yet propagated through the heartbeat
74            // wire format (`api::v1::meta::FlowStat`); it will always be empty in distributed
75            // mode until a follow-up PR adds heartbeat propagation.
76            let value: FlowStateValue =
77                FlowStateValue::new(state_size, last_exec_time_map, Default::default());
78            match node_identity(req) {
79                Some(node_id) => {
80                    // Merge by node so that reports from different flownodes
81                    // don't overwrite each other in the global state.
82                    self.flow_state_manager
83                        .merge(node_id, value)
84                        .await
85                        .context(FlowStateHandlerSnafu)?;
86                }
87                // No usable identity in the request: ignore the report instead
88                // of falling back to a whole-map replace, which would clobber
89                // other nodes' reports.
90                // Normal flownodes always carry header.member_id/peer.id; a
91                // report without either indicates an old or malformed client.
92                // Log at debug to avoid an anomalous sender spamming warn.
93                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        // Node 42 (header.member_id) reports flow 1.
160        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        // Node 7 (only peer.id present) reports flow 2.
174        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        // No header and no peer: the report must be ignored and no KV written.
203        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}