Skip to main content

servers/batcher/
flow_notifier.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::num::NonZeroUsize;
16
17use common_batcher::notifier::{Notifier, run_notifier};
18use common_meta::cache::TableFlownodeSetCacheRef;
19use common_meta::node_manager::NodeManagerRef;
20use common_runtime::spawn_global;
21use common_telemetry::{error, warn};
22use prometheus::IntCounterVec;
23use tokio::sync::mpsc::Receiver;
24use tokio::sync::mpsc::error::TrySendError;
25
26use crate::batcher::flow_sender::{FlowNotification, FlowSender};
27
28/// Best-effort queue admission and diagnostics, independent of payload extraction.
29#[derive(Clone)]
30pub(in crate::batcher) struct FlowNotifier {
31    notifier: Notifier<FlowNotification>,
32    dropped: IntCounterVec,
33}
34
35impl FlowNotifier {
36    /// Each construction creates an independent queue with caller-owned metrics.
37    pub fn try_new(
38        capacity: usize,
39        dropped: IntCounterVec,
40    ) -> Option<(Self, Receiver<FlowNotification>)> {
41        let (notifier, receiver) = Notifier::try_new(capacity)?;
42        Some((Self { notifier, dropped }, receiver))
43    }
44
45    /// Never waits for queue capacity or delivery and never changes the write result.
46    pub fn try_notify(&self, notification: FlowNotification) -> bool {
47        match self.notifier.try_notify(notification) {
48            Ok(()) => true,
49            Err(TrySendError::Full(notification)) => {
50                self.dropped.with_label_values(&["full"]).inc();
51                warn!(
52                    "Dropping flow notification because queue is full, table_id: {}, queue_capacity: {}",
53                    notification.table_id,
54                    self.notifier.max_capacity()
55                );
56                false
57            }
58            Err(TrySendError::Closed(notification)) => {
59                self.dropped.with_label_values(&["closed"]).inc();
60                error!(
61                    "Dropping flow notification because queue is closed, table_id: {}, queue_capacity: {}",
62                    notification.table_id,
63                    self.notifier.max_capacity()
64                );
65                false
66            }
67        }
68    }
69}
70
71/// Composes queue consumption and delivery without giving the notifier task ownership.
72pub(in crate::batcher) fn start_flow_notification_worker(
73    receiver: Receiver<FlowNotification>,
74    cache: TableFlownodeSetCacheRef,
75    node_manager: NodeManagerRef,
76) {
77    let sender = FlowSender::new(cache, node_manager);
78    let concurrency = NonZeroUsize::new(8).unwrap();
79    spawn_global(run_notifier(receiver, concurrency, move |notification| {
80        let sender = sender.clone();
81        async move { sender.send(notification).await }
82    }));
83}
84
85#[cfg(test)]
86mod tests {
87
88    use prometheus::{IntCounterVec, Opts};
89
90    use crate::batcher::flow_notifier::FlowNotifier;
91    use crate::batcher::flow_sender::FlowNotification;
92
93    fn dropped_counter() -> IntCounterVec {
94        // Unregistered collectors remain independent even with identical names.
95        IntCounterVec::new(
96            Opts::new("test_flow_dropped", "Dropped notifications"),
97            &["reason"],
98        )
99        .unwrap()
100    }
101
102    fn notification(table_id: u32) -> FlowNotification {
103        FlowNotification {
104            table_id,
105            timestamps: vec![-1, 42, 1_700_000_000_000_000_001, 42],
106        }
107    }
108
109    fn dropped_counts(counter: &IntCounterVec) -> (u64, u64) {
110        (
111            counter.with_label_values(&["full"]).get(),
112            counter.with_label_values(&["closed"]).get(),
113        )
114    }
115
116    #[test]
117    fn test_queues_and_metrics_are_independent() {
118        let first_dropped = dropped_counter();
119        let second_dropped = dropped_counter();
120        let (first, mut first_rx) = FlowNotifier::try_new(1, first_dropped.clone()).unwrap();
121        let (second, mut second_rx) = FlowNotifier::try_new(2, second_dropped.clone()).unwrap();
122
123        assert!(first.try_notify(notification(1)));
124        assert!(!first.try_notify(notification(2)));
125        assert_eq!(dropped_counts(&first_dropped), (1, 0));
126        assert_eq!(dropped_counts(&second_dropped), (0, 0));
127        assert!(second.try_notify(notification(3)));
128        assert!(second.try_notify(notification(4)));
129        assert!(!second.try_notify(notification(5)));
130        assert_eq!(dropped_counts(&second_dropped), (1, 0));
131
132        let payload = first_rx.try_recv().unwrap();
133        assert_eq!(payload.table_id, 1);
134        assert_eq!(payload.timestamps, notification(1).timestamps);
135        assert!(first_rx.try_recv().is_err());
136        for table_id in [3, 4] {
137            let payload = second_rx.try_recv().unwrap();
138            assert_eq!(payload.table_id, table_id);
139            assert_eq!(payload.timestamps, notification(table_id).timestamps);
140        }
141        assert!(second_rx.try_recv().is_err());
142
143        drop(first_rx);
144        assert!(!first.try_notify(notification(6)));
145        assert_eq!(dropped_counts(&first_dropped), (1, 1));
146        assert_eq!(dropped_counts(&second_dropped), (1, 0));
147        assert!(second.try_notify(notification(7)));
148        assert_eq!(second_rx.try_recv().unwrap().table_id, 7);
149        drop(second_rx);
150        assert!(!second.try_notify(notification(8)));
151        assert_eq!(dropped_counts(&second_dropped), (1, 1));
152    }
153
154    #[test]
155    fn test_clones_share_queue_and_drop_metrics() {
156        let dropped = dropped_counter();
157        let (notifier, mut receiver) = FlowNotifier::try_new(1, dropped.clone()).unwrap();
158        let clone = notifier.clone();
159        assert!(notifier.try_notify(notification(1)));
160        assert!(!clone.try_notify(notification(2)));
161        assert_eq!(dropped_counts(&dropped), (1, 0));
162        assert_eq!(receiver.try_recv().unwrap().table_id, 1);
163        drop(notifier);
164        assert!(clone.try_notify(notification(3)));
165        let payload = receiver.try_recv().unwrap();
166        assert_eq!(payload.table_id, 3);
167        assert_eq!(payload.timestamps, notification(3).timestamps);
168        drop(receiver);
169        assert!(!clone.try_notify(notification(4)));
170        assert_eq!(dropped_counts(&dropped), (1, 1));
171    }
172}