common_query/
promql_annotations.rs1use 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}