Skip to main content

servers/otlp/metrics/
resource_info.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
15//! Resource descriptor synthesized from OTLP metrics requests, so the entity
16//! graph reads one table instead of scanning every logical metric table for
17//! attributes the promote filter may have dropped.
18//!
19//! Columns are a fixed allowlist under the raw OTel attribute names: the
20//! conventions whitelist matches fixed names, so they must not follow the
21//! per-request label translation strategy or the promote/ignore headers.
22
23use 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
47/// Prefixed like the other engine-managed tables, so a user metric is
48/// unlikely to claim the name.
49pub const OTEL_RESOURCE_INFO_TABLE_NAME: &str = "greptime_otel_resource_info";
50
51/// Attributes projected under their raw OTel keys. `service.instance.id` is
52/// absent on purpose: it lands in `instance`.
53///
54/// Matched instead of scanned: this runs for every attribute of every
55/// resource, and the compiler turns it into a length-and-prefix dispatch.
56fn 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
73/// Upper bound of [`is_projected_attr`] plus the derived `job`/`instance`,
74/// used to size the per-row buffers.
75const MAX_PROJECTED_TAGS: usize = 13;
76
77/// Projected attributes (sorted `(name, value)` pairs) -> graph window ->
78/// the newest data-point time seen in that window, which is what the row for
79/// that window is stamped with.
80///
81/// Windows are an inner map so the attributes are stored, and moved, once per
82/// resource. They are keyed separately because one request may carry data for
83/// several of them: a single row per resource would describe only the newest
84/// window, leaving the earlier ones with metric rows but no entities.
85#[derive(Debug, Default)]
86pub struct ResourceInfoData {
87    rows: BTreeMap<Vec<(String, String)>, BTreeMap<i64, i64>>,
88}
89
90impl ResourceInfoData {
91    /// Takes the raw attributes, before the promote filter runs on them.
92    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        // Sorted so equal attribute sets share a key, and so the emitted
112        // columns keep a stable order.
113        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    /// Every projected attribute becomes a tag, so auto-create puts it in the
137    /// primary key where the conventions expect it.
138    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
170/// Visits the times of the data points the encoder writes rows for, so a
171/// resource is described exactly where it is measured rather than wherever
172/// its request happens to reach.
173fn 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    /// Earlier windows would keep their metric rows but lose their entities.
300    #[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    /// Describing a resource whose only data the encoder drops invents an
317    /// entity with no measurements.
318    #[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}