Skip to main content

common_query/
promql_annotations.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::{BTreeSet, HashMap};
16use std::sync::{Arc, Mutex, Weak};
17
18use datafusion::config::{ConfigEntry, ConfigExtension, ExtensionOptions};
19use datafusion_common::{DataFusionError, Result as DfResult};
20use once_cell::sync::Lazy;
21
22type PromqlAnnotationStateRef = Mutex<PromqlAnnotationState>;
23
24static PROMQL_ANNOTATION_REGISTRY: Lazy<Mutex<HashMap<String, Weak<PromqlAnnotationStateRef>>>> =
25    Lazy::new(|| Mutex::new(HashMap::new()));
26
27#[derive(Debug, Default)]
28struct PromqlAnnotationState {
29    warnings: BTreeSet<String>,
30    infos: BTreeSet<String>,
31}
32
33#[derive(Clone, Debug)]
34pub struct PromqlAnnotationCollector {
35    inner: Arc<PromqlAnnotationStateRef>,
36}
37
38impl Default for PromqlAnnotationCollector {
39    fn default() -> Self {
40        Self {
41            inner: Arc::new(Mutex::new(PromqlAnnotationState::default())),
42        }
43    }
44}
45
46impl PromqlAnnotationCollector {
47    pub fn record_warning<S: Into<String>>(&self, warning: S) {
48        self.inner
49            .lock()
50            .expect("promql annotation collector poisoned")
51            .warnings
52            .insert(warning.into());
53    }
54
55    pub fn record_info<S: Into<String>>(&self, info: S) {
56        self.inner
57            .lock()
58            .expect("promql annotation collector poisoned")
59            .infos
60            .insert(info.into());
61    }
62
63    pub fn append_to(&self, warnings: &mut Vec<String>, infos: &mut Vec<String>) {
64        let inner = self
65            .inner
66            .lock()
67            .expect("promql annotation collector poisoned");
68        append_unique(warnings, inner.warnings.iter().cloned());
69        append_unique(infos, inner.infos.iter().cloned());
70    }
71}
72
73impl ConfigExtension for PromqlAnnotationCollector {
74    const PREFIX: &'static str = "greptime_promql_annotations";
75}
76
77impl ExtensionOptions for PromqlAnnotationCollector {
78    fn as_any(&self) -> &dyn std::any::Any {
79        self
80    }
81
82    fn as_any_mut(&mut self) -> &mut dyn std::any::Any {
83        self
84    }
85
86    fn cloned(&self) -> Box<dyn ExtensionOptions> {
87        Box::new(self.clone())
88    }
89
90    fn set(&mut self, key: &str, value: &str) -> DfResult<()> {
91        Err(DataFusionError::NotImplemented(format!(
92            "PromqlAnnotationCollector does not support set key: {key} with value: {value}"
93        )))
94    }
95
96    fn entries(&self) -> Vec<ConfigEntry> {
97        vec![]
98    }
99}
100
101pub fn promql_annotation_collector(query_id: &str) -> PromqlAnnotationCollector {
102    let mut registry = PROMQL_ANNOTATION_REGISTRY
103        .lock()
104        .expect("promql annotation registry poisoned");
105    registry.retain(|_, collector| collector.strong_count() > 0);
106
107    if let Some(inner) = registry.get(query_id).and_then(Weak::upgrade) {
108        return PromqlAnnotationCollector { inner };
109    }
110
111    let collector = PromqlAnnotationCollector::default();
112    registry.insert(query_id.to_string(), Arc::downgrade(&collector.inner));
113    collector
114}
115
116pub fn get_promql_annotation_collector(query_id: &str) -> Option<PromqlAnnotationCollector> {
117    let mut registry = PROMQL_ANNOTATION_REGISTRY
118        .lock()
119        .expect("promql annotation registry poisoned");
120    registry.retain(|_, collector| collector.strong_count() > 0);
121    registry
122        .get(query_id)
123        .and_then(Weak::upgrade)
124        .map(|inner| PromqlAnnotationCollector { inner })
125}
126
127fn append_unique(target: &mut Vec<String>, values: impl Iterator<Item = String>) {
128    let mut seen = target.iter().cloned().collect::<BTreeSet<_>>();
129    for value in values {
130        if seen.insert(value.clone()) {
131            target.push(value);
132        }
133    }
134    target.sort();
135}
136
137#[cfg(test)]
138mod tests {
139    use super::*;
140
141    #[test]
142    fn collector_deduplicates_and_orders_annotations() {
143        let collector = PromqlAnnotationCollector::default();
144        collector.record_info("z-info");
145        collector.record_info("a-info");
146        collector.record_info("z-info");
147        collector.record_warning("b-warning");
148        collector.record_warning("a-warning");
149
150        let mut warnings = vec!["c-warning".to_string()];
151        let mut infos = vec![];
152        collector.append_to(&mut warnings, &mut infos);
153
154        assert_eq!(warnings, vec!["a-warning", "b-warning", "c-warning"]);
155        assert_eq!(infos, vec!["a-info", "z-info"]);
156    }
157
158    #[test]
159    fn registry_prunes_dead_collectors() {
160        let key = "registry_prunes_dead_collectors";
161        {
162            let collector = promql_annotation_collector(key);
163            assert!(get_promql_annotation_collector(key).is_some());
164            drop(collector);
165        }
166
167        assert!(get_promql_annotation_collector(key).is_none());
168        assert!(!PROMQL_ANNOTATION_REGISTRY.lock().unwrap().contains_key(key));
169    }
170}