1use std::collections::BTreeMap;
24
25use api::v1::RowInsertRequests;
26use common_catalog::consts::SEMANTIC_GRAPH_WINDOW_NANOS;
27use common_grpc::precision::Precision;
28use common_query::prelude::{greptime_timestamp, greptime_value};
29use otel_arrow_rust::proto::opentelemetry::common::v1::KeyValue;
30use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
31 AggregationTemporality, ResourceMetrics, metric,
32};
33
34use crate::error::Result;
35use crate::otlp::metrics::{
36 INSTANCE_KEY, JOB_KEY, ServiceIdentity, exponential_histogram_gate,
37 exponential_histogram_value, histogram_data_point_rejection, scalar_value_string,
38 service_identity,
39};
40use crate::otlp::trace::{
41 KEY_CONTAINER_ID, KEY_CONTAINER_NAME, KEY_HOST_ID, KEY_HOST_NAME, KEY_K8S_CONTAINER_NAME,
42 KEY_K8S_NAMESPACE_NAME, KEY_K8S_NODE_NAME, KEY_K8S_POD_NAME, KEY_K8S_POD_UID, KEY_SERVICE_NAME,
43 KEY_SERVICE_NAMESPACE,
44};
45use crate::row_writer::{self, MultiTableData};
46
47pub const OTEL_RESOURCE_INFO_TABLE_NAME: &str = "greptime_otel_resource_info";
50
51fn is_projected_attr(key: &str) -> bool {
57 matches!(
58 key,
59 KEY_SERVICE_NAME
60 | KEY_SERVICE_NAMESPACE
61 | KEY_HOST_ID
62 | KEY_HOST_NAME
63 | KEY_CONTAINER_ID
64 | KEY_CONTAINER_NAME
65 | KEY_K8S_POD_UID
66 | KEY_K8S_POD_NAME
67 | KEY_K8S_CONTAINER_NAME
68 | KEY_K8S_NAMESPACE_NAME
69 | KEY_K8S_NODE_NAME
70 )
71}
72
73const MAX_PROJECTED_TAGS: usize = 13;
76
77#[derive(Debug, Default)]
86pub struct ResourceInfoData {
87 rows: BTreeMap<Vec<(String, String)>, BTreeMap<i64, i64>>,
88}
89
90impl ResourceInfoData {
91 pub fn observe(&mut self, raw_attrs: &[KeyValue], resource: &ResourceMetrics) {
93 let mut tags = Vec::with_capacity(MAX_PROJECTED_TAGS);
94 let ServiceIdentity { job, instance } = service_identity(raw_attrs);
95 if let Some(job) = job {
96 tags.push((JOB_KEY.to_string(), job));
97 }
98 if let Some(instance) = instance {
99 tags.push((INSTANCE_KEY.to_string(), instance));
100 }
101 for kv in raw_attrs {
102 if is_projected_attr(&kv.key)
103 && let Some(value) = scalar_value_string(kv.value.as_ref())
104 {
105 tags.push((kv.key.clone(), value));
106 }
107 }
108 if tags.is_empty() {
109 return;
110 }
111 tags.sort_unstable();
114
115 let mut observed: BTreeMap<i64, i64> = BTreeMap::new();
116 for_each_encoded_time(resource, |ts| {
117 let window = ts - ts.rem_euclid(SEMANTIC_GRAPH_WINDOW_NANOS);
118 observed
119 .entry(window)
120 .and_modify(|newest| *newest = (*newest).max(ts))
121 .or_insert(ts);
122 });
123 if observed.is_empty() {
124 return;
125 }
126
127 let windows = self.rows.entry(tags).or_default();
128 for (window, newest) in observed {
129 windows
130 .entry(window)
131 .and_modify(|seen| *seen = (*seen).max(newest))
132 .or_insert(newest);
133 }
134 }
135
136 pub fn into_row_insert_requests(self) -> Result<Option<RowInsertRequests>> {
139 if self.rows.is_empty() {
140 return Ok(None);
141 }
142
143 let mut writer = MultiTableData::default();
144 let table = writer.get_or_default_table_data(
145 OTEL_RESOURCE_INFO_TABLE_NAME,
146 MAX_PROJECTED_TAGS + 2,
147 self.rows.values().map(BTreeMap::len).sum(),
148 );
149 for (tags, windows) in &self.rows {
150 for ts_nanos in windows.values().copied() {
151 let mut row = table.alloc_one_row();
152 row_writer::write_tags(table, tags.iter().cloned(), &mut row)?;
153 row_writer::write_f64(table, greptime_value(), 1.0, &mut row)?;
154 row_writer::write_ts_to_millis(
155 table,
156 greptime_timestamp(),
157 Some(ts_nanos),
158 Precision::Nanosecond,
159 &mut row,
160 )?;
161 table.add_row(row);
162 }
163 }
164
165 let (requests, _) = writer.into_row_insert_requests();
166 Ok(Some(requests))
167 }
168}
169
170fn for_each_encoded_time(resource: &ResourceMetrics, mut visit: impl FnMut(i64)) {
174 fn visit_all(points: impl Iterator<Item = u64>, visit: &mut impl FnMut(i64)) {
175 for ts in points {
176 visit(ts as i64);
177 }
178 }
179 for scope in &resource.scope_metrics {
180 for m in &scope.metrics {
181 match &m.data {
182 Some(metric::Data::Gauge(g)) => {
183 visit_all(g.data_points.iter().map(|p| p.time_unix_nano), &mut visit)
184 }
185 Some(metric::Data::Sum(s)) => {
186 visit_all(s.data_points.iter().map(|p| p.time_unix_nano), &mut visit)
187 }
188 Some(metric::Data::Histogram(h)) => {
189 let is_delta = matches!(
190 AggregationTemporality::try_from(h.aggregation_temporality),
191 Ok(AggregationTemporality::Delta)
192 );
193 visit_all(
194 h.data_points
195 .iter()
196 .filter(|point| {
197 histogram_data_point_rejection(point, is_delta).is_none()
198 })
199 .map(|point| point.time_unix_nano),
200 &mut visit,
201 )
202 }
203 Some(metric::Data::Summary(s)) => {
204 visit_all(s.data_points.iter().map(|p| p.time_unix_nano), &mut visit)
205 }
206 Some(metric::Data::ExponentialHistogram(h))
207 if exponential_histogram_gate(h).is_ok() =>
208 {
209 for point in &h.data_points {
210 if let Ok((_, ts)) = exponential_histogram_value(point) {
211 visit(ts);
212 }
213 }
214 }
215 Some(metric::Data::ExponentialHistogram(_)) => {}
216 None => {}
217 }
218 }
219 }
220}
221
222#[cfg(test)]
223mod tests {
224 use api::v1::SemanticType;
225 use api::v1::value::ValueData;
226 use common_query::prelude::set_default_prefix;
227 use otel_arrow_rust::proto::opentelemetry::common::v1::{AnyValue, any_value};
228 use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
229 AggregationTemporality, ExponentialHistogram, ExponentialHistogramDataPoint, Gauge, Metric,
230 NumberDataPoint, ScopeMetrics,
231 };
232
233 use super::*;
234
235 mod delta;
236
237 fn kv(key: &str, value: &str) -> KeyValue {
238 KeyValue {
239 key: key.into(),
240 value: Some(AnyValue {
241 value: Some(any_value::Value::StringValue(value.into())),
242 }),
243 }
244 }
245
246 fn gauge_at(times: &[i64]) -> ResourceMetrics {
247 ResourceMetrics {
248 scope_metrics: vec![ScopeMetrics {
249 metrics: vec![Metric {
250 data: Some(metric::Data::Gauge(Gauge {
251 data_points: times
252 .iter()
253 .map(|ts| NumberDataPoint {
254 time_unix_nano: *ts as u64,
255 ..Default::default()
256 })
257 .collect(),
258 })),
259 ..Default::default()
260 }],
261 ..Default::default()
262 }],
263 ..Default::default()
264 }
265 }
266
267 #[test]
268 fn observe_projects_allowlist_and_dedups_per_request() {
269 let mut data = ResourceInfoData::default();
270 let attrs = vec![
271 kv("service.name", "api"),
272 kv("service.namespace", "shop"),
273 kv("service.instance.id", "inst-1"),
274 kv("host.id", "h-1"),
275 kv("k8s.node.name", "node-a"),
276 kv("os.type", "linux"),
277 ];
278 data.observe(&attrs, &gauge_at(&[100, 50]));
279 assert_eq!(data.rows.len(), 1);
280 let (tags, windows) = data.rows.iter().next().unwrap();
281 assert_eq!(windows.values().copied().collect::<Vec<_>>(), vec![100]);
282 assert!(tags.contains(&("job".to_string(), "shop/api".to_string())));
283 assert!(tags.contains(&("instance".to_string(), "inst-1".to_string())));
284 assert!(tags.contains(&("service.name".to_string(), "api".to_string())));
285 assert!(tags.contains(&("k8s.node.name".to_string(), "node-a".to_string())));
286 assert!(
287 tags.iter()
288 .all(|(k, _)| k != "os.type" && k != "service.instance.id")
289 );
290
291 data.observe(&[kv("host.id", "h-2")], &gauge_at(&[10]));
292 assert_eq!(data.rows.len(), 2);
293
294 let mut empty = ResourceInfoData::default();
295 empty.observe(&[kv("os.type", "linux")], &gauge_at(&[100]));
296 assert!(empty.into_row_insert_requests().unwrap().is_none());
297 }
298
299 #[test]
301 fn observe_keeps_one_row_per_graph_window() {
302 let window = SEMANTIC_GRAPH_WINDOW_NANOS;
303 let mut data = ResourceInfoData::default();
304 data.observe(
305 &[kv("service.name", "api")],
306 &gauge_at(&[window + 1, window + 2, 3 * window + 7]),
307 );
308
309 let windows = data.rows.values().next().unwrap();
310 assert_eq!(
311 windows.iter().collect::<Vec<_>>(),
312 vec![(&window, &(window + 2)), (&(3 * window), &(3 * window + 7))]
313 );
314 }
315
316 #[test]
319 fn observe_ignores_data_the_encoder_drops() {
320 let exponential = |temporality: AggregationTemporality| ResourceMetrics {
321 scope_metrics: vec![ScopeMetrics {
322 metrics: vec![Metric {
323 data: Some(metric::Data::ExponentialHistogram(ExponentialHistogram {
324 data_points: vec![ExponentialHistogramDataPoint {
325 time_unix_nano: 100,
326 ..Default::default()
327 }],
328 aggregation_temporality: temporality as i32,
329 })),
330 ..Default::default()
331 }],
332 ..Default::default()
333 }],
334 ..Default::default()
335 };
336 for temporality in [
337 AggregationTemporality::Delta,
338 AggregationTemporality::Unspecified,
339 ] {
340 let resource = exponential(temporality);
341 let mut data = ResourceInfoData::default();
342 data.observe(&[kv("service.name", "api")], &resource);
343 assert!(data.into_row_insert_requests().unwrap().is_none());
344 }
345 }
346
347 #[test]
348 fn rows_carry_raw_key_tags_value_and_millis_timestamp() {
349 set_default_prefix(None).unwrap();
350 let mut data = ResourceInfoData::default();
351 data.observe(
352 &[kv("service.name", "api"), kv("host.id", "h-1")],
353 &gauge_at(&[1_700_000_000_123_456_789]),
354 );
355 let requests = data.into_row_insert_requests().unwrap().unwrap();
356 assert_eq!(requests.inserts.len(), 1);
357 let insert = &requests.inserts[0];
358 assert_eq!(insert.table_name, OTEL_RESOURCE_INFO_TABLE_NAME);
359
360 let rows = insert.rows.as_ref().unwrap();
361 let names = rows
362 .schema
363 .iter()
364 .map(|c| c.column_name.as_str())
365 .collect::<Vec<_>>();
366 assert_eq!(
367 names,
368 vec![
369 "host.id",
370 "job",
371 "service.name",
372 greptime_value(),
373 greptime_timestamp()
374 ]
375 );
376 for column in &rows.schema[..3] {
377 assert_eq!(column.semantic_type, SemanticType::Tag as i32);
378 }
379
380 assert_eq!(rows.rows.len(), 1);
381 let values = &rows.rows[0].values;
382 assert_eq!(values[3].value_data, Some(ValueData::F64Value(1.0)));
383 assert_eq!(
384 values[4].value_data,
385 Some(ValueData::TimestampMillisecondValue(1_700_000_000_123))
386 );
387 }
388}