Skip to content

Commit 9e06d3f

Browse files
committed
clean up OTLP ingest
1 parent 1ae0397 commit 9e06d3f

3 files changed

Lines changed: 134 additions & 34 deletions

File tree

open-tsdb/src/db.rs

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,9 +23,10 @@ use std::sync::{Arc, atomic::AtomicU32};
2323

2424
use dashmap::DashMap;
2525
use opendata_common::{Storage, storage::StorageSnapshot};
26+
use opentelemetry_proto::tonic::metrics::v1::MetricsData;
2627
use tokio::sync::{Mutex, RwLock};
2728

28-
use crate::delta::TsdbDelta;
29+
use crate::delta::{TsdbDelta, TsdbDeltaBuilder};
2930
use crate::storage::OpenTsdbStorageReadExt;
3031
use crate::util::Result;
3132
use crate::{
@@ -76,12 +77,29 @@ impl Tsdb {
7677
})
7778
}
7879

79-
pub(crate) async fn ingest(&self, delta: TsdbDelta) -> Result<()> {
80+
/// Ingest an OLTP payload (this is the main entry point for ingestion)
81+
pub(crate) async fn ingest(&self, data: MetricsData) -> Result<()> {
82+
let delta =
83+
TsdbDeltaBuilder::new(self.bucket.clone(), &self.series_dict, &self.next_series_id)
84+
.ingest_metrics_data(data)?
85+
.build();
86+
87+
self.ingest_delta(&delta).await
88+
}
89+
90+
/// Ingest a delta directly (this can be used when ingesting data as a
91+
/// read-only replica to stay up to date with the main writer)
92+
pub(crate) async fn ingest_delta(&self, delta: &TsdbDelta) -> Result<()> {
93+
// TODO(agavra): log the delta to a WAL to avoid losing data
8094
let state = self.state.read().await;
8195
state.head.merge(&delta)?;
8296
Ok(())
8397
}
8498

99+
/// Flush the head to storage, making it durable. Eventually we will support
100+
/// a native WAL so that we can get durability without waiting until a flush,
101+
/// but for now the WAL is coupled with the storage layer so we accept the risk
102+
/// of losing a small amount of data in the event of a crash.
85103
pub(crate) async fn flush(&self, storage: Arc<dyn Storage>) -> Result<()> {
86104
let _flush_guard = self.flush_mutex.lock().await;
87105

open-tsdb/src/delta.rs

Lines changed: 62 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
use std::collections::HashMap;
2+
use std::sync::atomic::{AtomicU32, Ordering};
23

34
use dashmap::DashMap;
4-
use opentelemetry_proto::tonic::metrics::v1::Metric;
5+
use opentelemetry_proto::tonic::metrics::v1::{Metric, MetricsData};
56

67
use crate::{
78
index::{ForwardIndex, InvertedIndex},
8-
model::{MetricType, Sample, SeriesFingerprint, SeriesId, SeriesSpec, TimeBucket},
9-
otel::OtelUtil,
9+
model::{Attribute, MetricType, Sample, SeriesFingerprint, SeriesId, SeriesSpec, TimeBucket},
10+
otel::{OtelUtil, collect_resource_attributes, collect_scope_attributes},
1011
util::{Fingerprint, Result, normalize_str},
1112
};
1213

@@ -19,14 +20,14 @@ pub(crate) struct TsdbDeltaBuilder<'a> {
1920
pub(crate) series_dict: &'a DashMap<SeriesFingerprint, SeriesId>,
2021
pub(crate) series_dict_delta: HashMap<SeriesFingerprint, SeriesId>,
2122
pub(crate) samples: HashMap<SeriesId, Vec<Sample>>,
22-
pub(crate) next_series_id: u32,
23+
pub(crate) next_series_id: &'a AtomicU32,
2324
}
2425

2526
impl<'a> TsdbDeltaBuilder<'a> {
2627
pub(crate) fn new(
2728
bucket: TimeBucket,
2829
series_dict: &'a DashMap<SeriesFingerprint, SeriesId>,
29-
next_series_id: u32,
30+
next_series_id: &'a AtomicU32,
3031
) -> Self {
3132
Self {
3233
bucket,
@@ -39,10 +40,33 @@ impl<'a> TsdbDeltaBuilder<'a> {
3940
}
4041
}
4142

42-
pub(crate) fn ingest_metric(mut self, metric: &Metric) -> Result<()> {
43+
/// Ingest a full MetricsData payload, processing all resource_metrics, scope_metrics, and metrics.
44+
/// Resource and scope attributes are merged into each sample's attributes.
45+
pub(crate) fn ingest_metrics_data(mut self, data: MetricsData) -> Result<Self> {
46+
for resource_metrics in data.resource_metrics {
47+
let resource_attrs = collect_resource_attributes(resource_metrics.resource.as_ref());
48+
49+
for scope_metrics in resource_metrics.scope_metrics {
50+
let scope_attrs = collect_scope_attributes(scope_metrics.scope.as_ref());
51+
52+
for metric in scope_metrics.metrics {
53+
self.ingest_metric_internal(&metric, &resource_attrs, &scope_attrs)?;
54+
}
55+
}
56+
}
57+
Ok(self)
58+
}
59+
60+
/// Internal method to ingest a single metric with pre-collected resource and scope attributes
61+
fn ingest_metric_internal(
62+
&mut self,
63+
metric: &Metric,
64+
resource_attrs: &[Attribute],
65+
scope_attrs: &[Attribute],
66+
) -> Result<()> {
4367
let metric_unit = normalize_str(&metric.unit);
4468
let metric_type = MetricType::try_from(metric)?;
45-
let samples_with_attrs = OtelUtil::samples(metric);
69+
let samples_with_attrs = OtelUtil::samples(metric, resource_attrs, scope_attrs);
4670

4771
for sample_with_attrs in samples_with_attrs {
4872
self.ingest_sample(
@@ -56,9 +80,9 @@ impl<'a> TsdbDeltaBuilder<'a> {
5680
Ok(())
5781
}
5882

59-
pub(crate) fn ingest_sample(
83+
fn ingest_sample(
6084
&mut self,
61-
mut attributes: Vec<crate::model::Attribute>,
85+
mut attributes: Vec<Attribute>,
6286
metric_unit: Option<String>,
6387
metric_type: MetricType,
6488
sample: Sample,
@@ -78,8 +102,7 @@ impl<'a> TsdbDeltaBuilder<'a> {
78102
.map(|r| *r.value())
79103
.or_else(|| self.series_dict_delta.get(&fingerprint).copied())
80104
.unwrap_or_else(|| {
81-
let series_id = self.next_series_id;
82-
self.next_series_id = series_id + 1;
105+
let series_id = self.next_series_id.fetch_add(1, Ordering::SeqCst);
83106

84107
self.series_dict_delta.insert(fingerprint, series_id);
85108

@@ -131,6 +154,7 @@ mod tests {
131154
use super::*;
132155
use crate::model::{Attribute, MetricType, Temporality};
133156
use dashmap::DashMap;
157+
use std::sync::atomic::AtomicU32;
134158

135159
fn create_test_bucket() -> TimeBucket {
136160
TimeBucket::hour(1000)
@@ -161,7 +185,8 @@ mod tests {
161185
// given
162186
let bucket = create_test_bucket();
163187
let series_dict = DashMap::new();
164-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
188+
let next_series_id = AtomicU32::new(0);
189+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
165190
let attributes = create_test_attributes();
166191
let sample = create_test_sample();
167192
let metric_unit = Some("bytes".to_string());
@@ -176,7 +201,7 @@ mod tests {
176201
);
177202

178203
// then
179-
assert_eq!(builder.next_series_id, 1);
204+
assert_eq!(next_series_id.load(Ordering::SeqCst), 1);
180205
assert_eq!(builder.series_dict_delta.len(), 1);
181206
assert_eq!(builder.samples.len(), 1);
182207

@@ -209,7 +234,8 @@ mod tests {
209234
// given
210235
let bucket = create_test_bucket();
211236
let series_dict = DashMap::new();
212-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
237+
let next_series_id = AtomicU32::new(0);
238+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
213239
let attributes = create_test_attributes();
214240
let sample1 = Sample {
215241
timestamp: 1000,
@@ -236,7 +262,7 @@ mod tests {
236262
);
237263

238264
// then
239-
assert_eq!(builder.next_series_id, 1); // Only one series created
265+
assert_eq!(next_series_id.load(Ordering::SeqCst), 1); // Only one series created
240266
assert_eq!(builder.series_dict_delta.len(), 1);
241267
assert_eq!(builder.samples.len(), 1);
242268

@@ -252,7 +278,8 @@ mod tests {
252278
// given
253279
let bucket = create_test_bucket();
254280
let series_dict = DashMap::new();
255-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
281+
let next_series_id = AtomicU32::new(0);
282+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
256283
let attributes1 = vec![Attribute {
257284
key: "service".to_string(),
258285
value: "api".to_string(),
@@ -278,7 +305,7 @@ mod tests {
278305
);
279306

280307
// then
281-
assert_eq!(builder.next_series_id, 2); // Two series created
308+
assert_eq!(next_series_id.load(Ordering::SeqCst), 2); // Two series created
282309
assert_eq!(builder.series_dict_delta.len(), 2);
283310
assert_eq!(builder.samples.len(), 2);
284311
assert!(builder.samples.contains_key(&0));
@@ -295,7 +322,8 @@ mod tests {
295322
attributes.sort_by(|a, b| a.key.cmp(&b.key));
296323
let fingerprint = attributes.fingerprint();
297324
series_dict.insert(fingerprint, 42); // Existing series_id
298-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
325+
let next_series_id = AtomicU32::new(0);
326+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
299327
let metric_type = MetricType::Gauge;
300328

301329
// when
@@ -307,7 +335,7 @@ mod tests {
307335
);
308336

309337
// then
310-
assert_eq!(builder.next_series_id, 0); // No new series_id created
338+
assert_eq!(next_series_id.load(Ordering::SeqCst), 0); // No new series_id created
311339
assert_eq!(builder.series_dict_delta.len(), 0); // Not added to delta
312340
assert_eq!(builder.samples.len(), 1);
313341
assert!(builder.samples.contains_key(&42)); // Uses existing series_id
@@ -318,7 +346,8 @@ mod tests {
318346
// given
319347
let bucket = create_test_bucket();
320348
let series_dict = DashMap::new();
321-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
349+
let next_series_id = AtomicU32::new(0);
350+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
322351
let attributes = create_test_attributes();
323352
let metric_type = MetricType::Gauge;
324353

@@ -337,7 +366,7 @@ mod tests {
337366
);
338367

339368
// then
340-
assert_eq!(builder.next_series_id, 1); // Only one series_id created
369+
assert_eq!(next_series_id.load(Ordering::SeqCst), 1); // Only one series_id created
341370
assert_eq!(builder.series_dict_delta.len(), 1);
342371
assert_eq!(builder.samples.len(), 1);
343372
assert!(builder.samples.contains_key(&0)); // Reused series_id 0
@@ -349,7 +378,8 @@ mod tests {
349378
// given
350379
let bucket = create_test_bucket();
351380
let series_dict = DashMap::new();
352-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
381+
let next_series_id = AtomicU32::new(0);
382+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
353383
let attributes1 = vec![
354384
Attribute {
355385
key: "z_key".to_string(),
@@ -387,7 +417,7 @@ mod tests {
387417
);
388418

389419
// then
390-
assert_eq!(builder.next_series_id, 1); // Same series_id reused
420+
assert_eq!(next_series_id.load(Ordering::SeqCst), 1); // Same series_id reused
391421
assert_eq!(builder.series_dict_delta.len(), 1);
392422
assert_eq!(builder.samples.len(), 1);
393423
}
@@ -397,7 +427,8 @@ mod tests {
397427
// given
398428
let bucket = create_test_bucket();
399429
let series_dict = DashMap::new();
400-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
430+
let next_series_id = AtomicU32::new(0);
431+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
401432
let attributes = create_test_attributes();
402433
let metric_unit = Some("requests_per_second".to_string());
403434
let metric_type = MetricType::Sum {
@@ -433,7 +464,8 @@ mod tests {
433464
// given
434465
let bucket = create_test_bucket();
435466
let series_dict = DashMap::new();
436-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
467+
let next_series_id = AtomicU32::new(0);
468+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
437469
let attributes = vec![
438470
Attribute {
439471
key: "service".to_string(),
@@ -472,7 +504,8 @@ mod tests {
472504
// given
473505
let bucket = create_test_bucket();
474506
let series_dict = DashMap::new();
475-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
507+
let next_series_id = AtomicU32::new(0);
508+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
476509
let attributes = Vec::<Attribute>::new();
477510
let metric_type = MetricType::Gauge;
478511

@@ -485,7 +518,7 @@ mod tests {
485518
);
486519

487520
// then
488-
assert_eq!(builder.next_series_id, 1);
521+
assert_eq!(next_series_id.load(Ordering::SeqCst), 1);
489522
assert_eq!(builder.series_dict_delta.len(), 1);
490523
assert_eq!(builder.samples.len(), 1);
491524
assert_eq!(builder.inverted_index.postings.len(), 0); // No attributes to index
@@ -498,7 +531,8 @@ mod tests {
498531
// given
499532
let bucket = create_test_bucket();
500533
let series_dict = DashMap::new();
501-
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, 0);
534+
let next_series_id = AtomicU32::new(0);
535+
let mut builder = TsdbDeltaBuilder::new(bucket, &series_dict, &next_series_id);
502536
let attributes = create_test_attributes();
503537
let metric_type = MetricType::Gauge;
504538

open-tsdb/src/otel.rs

Lines changed: 52 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,40 @@
1-
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value};
1+
use opentelemetry_proto::tonic::common::v1::{AnyValue, InstrumentationScope, KeyValue, any_value};
22
use opentelemetry_proto::tonic::metrics::v1::{
33
AggregationTemporality, HistogramDataPoint, Metric, metric, number_data_point,
44
};
5+
use opentelemetry_proto::tonic::resource::v1::Resource;
56

67
use crate::model::{Attribute, MetricType, Sample, Temporality};
78
use crate::util::{OpenTsdbError, Result};
89

10+
/// Collect attributes from an OTLP Resource
11+
pub(crate) fn collect_resource_attributes(resource: Option<&Resource>) -> Vec<Attribute> {
12+
match resource {
13+
Some(resource) => key_values_to_attributes(&resource.attributes),
14+
None => Vec::new(),
15+
}
16+
}
17+
18+
/// Collect attributes from an OTLP InstrumentationScope
19+
pub(crate) fn collect_scope_attributes(scope: Option<&InstrumentationScope>) -> Vec<Attribute> {
20+
let mut attrs = Vec::new();
21+
if let Some(scope) = scope {
22+
if !scope.name.is_empty() {
23+
attrs.push(Attribute {
24+
key: "otel_scope_name".to_string(),
25+
value: scope.name.clone(),
26+
});
27+
}
28+
if !scope.version.is_empty() {
29+
attrs.push(Attribute {
30+
key: "otel_scope_version".to_string(),
31+
value: scope.version.clone(),
32+
});
33+
}
34+
}
35+
attrs
36+
}
37+
938
/// A sample with all its attributes, including the metric name
1039
#[derive(Clone, Debug)]
1140
pub(crate) struct SampleWithAttributes {
@@ -17,12 +46,31 @@ pub(crate) struct SampleWithAttributes {
1746
pub(crate) struct OtelUtil;
1847

1948
impl OtelUtil {
20-
/// Returns all samples with attributes for a given metric
21-
pub(crate) fn samples(metric: &Metric) -> Vec<SampleWithAttributes> {
22-
match &metric.data {
49+
/// Returns all samples with attributes for a given metric, merging in resource and scope attributes
50+
pub(crate) fn samples(
51+
metric: &Metric,
52+
resource_attrs: &[Attribute],
53+
scope_attrs: &[Attribute],
54+
) -> Vec<SampleWithAttributes> {
55+
let mut samples = match &metric.data {
2356
Some(data) => data.to_samples(&metric.name),
2457
None => Vec::new(),
58+
};
59+
60+
// Merge resource and scope attributes into each sample
61+
for sample in &mut samples {
62+
// Prepend resource attrs first, then scope attrs
63+
// This ensures consistent ordering: resource -> scope -> metric/datapoint attrs
64+
let mut merged_attrs = Vec::with_capacity(
65+
resource_attrs.len() + scope_attrs.len() + sample.attributes.len(),
66+
);
67+
merged_attrs.extend_from_slice(resource_attrs);
68+
merged_attrs.extend_from_slice(scope_attrs);
69+
merged_attrs.append(&mut sample.attributes);
70+
sample.attributes = merged_attrs;
2571
}
72+
73+
samples
2674
}
2775
}
2876

0 commit comments

Comments
 (0)