Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions changelog.d/otlp_native_log_encoding.feature.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
The `otlp` codec now converts native Vector log events to OTLP log records, so logs from any source can
be sent with the `opentelemetry` sink without a `remap` transform that builds the `resourceLogs` structure.
For example, this log event from the `file` source:

```yaml
message: disk full
host: web-1
timestamp: 2026-10-09T12:00:00Z
source_type: file
```

is sent as this OTLP log record (shown as YAML, without default fields):

```yaml
resourceLogs:
- scopeLogs:
- logRecords:
- timeUnixNano: "1791547200000000000"
body:
stringValue: disk full
attributes:
- key: host
value:
stringValue: web-1
```

The conversion is the inverse of the `opentelemetry` source decoding. A log event with a `resourceLogs`,
`resourceMetrics`, or `resourceSpans` root field is treated as an OTLP request that is already built,
and is sent as it is.

The source type marker identifies the Vector source component type that produced the event, for example
`source_type: file`.

With `log_namespace: false` (Legacy namespace), fields with no OTLP equivalent are sent as log record
attributes. A timestamp inside the configured message field still sets `timeUnixNano` and remains in
the body. The marker's location is configured with `log_schema.source_type_key` (default:
`.source_type`). The codec omits this internal field. If the configured path points to metadata and
that field exists (for example, `%source_type`), the matching payload field (`.source_type`) is kept as
an attribute. Otherwise, the matching event field is omitted, because some Legacy sources still write
the marker there.

With `log_namespace: true` (Vector namespace), the entire event payload becomes the OTLP `body`.
A payload field named `.source_type` is preserved in the body, not sent as an attribute. The internal
marker is stored separately in `%vector.source_type` metadata and is not sent.

authors: thomasqueirozb
53 changes: 45 additions & 8 deletions lib/codecs/src/encoding/format/otlp.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
use crate::encoding::ProtobufSerializer;
use bytes::BytesMut;
use opentelemetry_proto::logs::log_event_to_export_request;
use opentelemetry_proto::metrics::metric_event_to_export_request;
use opentelemetry_proto::proto::{
DESCRIPTOR_BYTES, LOGS_REQUEST_MESSAGE_TYPE, METRICS_REQUEST_MESSAGE_TYPE,
Expand Down Expand Up @@ -51,14 +52,18 @@ impl OtlpSerializerConfig {
///
/// # Implementation approach
///
/// This serializer converts Vector's internal event representation to the appropriate OTLP message type
/// based on the top-level field in the event:
/// This serializer converts Vector's internal event representation to the appropriate OTLP message type.
/// Events that already have the OTLP structure are encoded as they are, based on the top-level field:
/// - `resourceLogs` → `ExportLogsServiceRequest`
/// - `resourceMetrics` → `ExportMetricsServiceRequest`
/// - `resourceSpans` → `ExportTraceServiceRequest`
///
/// The implementation is the inverse of what the `opentelemetry` source does when decoding,
/// ensuring round-trip compatibility.
/// The first field in this list that the event has is used, and all other event fields are
/// not encoded.
///
/// Native Vector logs and metrics are converted to `ExportLogsServiceRequest` and
/// `ExportMetricsServiceRequest`. The conversions are the inverse of what the `opentelemetry`
/// source does when decoding, ensuring round-trip compatibility.
#[derive(Debug, Clone)]
#[allow(dead_code)] // Fields will be used once encoding is implemented
pub struct OtlpSerializer {
Expand Down Expand Up @@ -119,11 +124,12 @@ impl Encoder<Event> for OtlpSerializer {
} else if log.contains(event_path!(RESOURCE_METRICS_JSON_FIELD)) {
// Currently the OTLP metrics are Vector logs (not metrics).
self.metrics_descriptor.encode(Event::Log(log), buffer)
} else if log.contains(event_path!(RESOURCE_SPANS_JSON_FIELD)) {
// OTLP-structured traces can be log events, for example when read as JSON.
self.traces_descriptor.encode(Event::Log(log), buffer)
} else {
Err(format!(
"Log event does not contain OTLP top-level fields ({RESOURCE_LOGS_JSON_FIELD} or {RESOURCE_METRICS_JSON_FIELD})",
)
.into())
let request = log_event_to_export_request(log);
request.encode(buffer).map_err(Into::into)
}
}
Event::Trace(trace) => {
Expand Down Expand Up @@ -194,6 +200,37 @@ mod tests {
assert_eq!(metric.clone(), round_trip_metric(metric));
}

#[test]
fn native_log_round_trips_through_otlp_source_decoding() {
use opentelemetry_proto::proto::collector::logs::v1::ExportLogsServiceRequest;
use vector_core::{config::LogNamespace, event::LogEvent};

let mut log = LogEvent::from("disk full");
log.insert(event_path!("host"), "web-1");
log.insert(event_path!("timestamp"), Utc.timestamp_nanos(1_000_000_000));

let mut buffer = BytesMut::new();
OtlpSerializer::new()
.unwrap()
.encode(Event::Log(log), &mut buffer)
.expect("native log must be converted, not rejected");

let request = ExportLogsServiceRequest::decode(buffer.freeze()).unwrap();
let mut events: Vec<Event> = request
.resource_logs
.into_iter()
.flat_map(|logs| logs.into_event_iter(LogNamespace::Legacy))
.collect();
assert_eq!(events.len(), 1);
let decoded = events.remove(0).into_log();
assert_eq!(decoded["message"], "disk full".into());
assert_eq!(decoded["attributes.host"], "web-1".into());
assert_eq!(
decoded["timestamp"],
Utc.timestamp_nanos(1_000_000_000).into()
);
}

#[test]
fn round_trip_gauge() {
let metric = with_empty_tags(
Expand Down
3 changes: 3 additions & 0 deletions lib/opentelemetry-proto/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,3 +23,6 @@ tonic.workspace = true
vrl.workspace = true
vector-common = { path = "../vector-common", default-features = false }
vector-core = { path = "../vector-core", default-features = false }

[dev-dependencies]
vector-core = { path = "../vector-core", default-features = false, features = ["test"] }
49 changes: 48 additions & 1 deletion lib/opentelemetry-proto/src/common.rs
Original file line number Diff line number Diff line change
@@ -1,9 +1,12 @@
use bytes::Bytes;
use chrono::SecondsFormat;
use ordered_float::NotNan;
use vector_core::event::metric::{TagValue, TagValueSet};
use vrl::value::{ObjectMap, Value};

use super::proto::common::v1::{AnyValue, ArrayValue, KeyValue, any_value::Value as PBValue};
use super::proto::common::v1::{
AnyValue, ArrayValue, KeyValue, KeyValueList, any_value::Value as PBValue,
};

impl From<PBValue> for Value {
fn from(av: PBValue) -> Self {
Expand All @@ -24,6 +27,38 @@ impl From<PBValue> for Value {
}
}

/// Inverse of `From<PBValue> for Value`. Strings become `string_value`; byte strings become
/// `string_value` with a lossy UTF-8 decode. Timestamps and regexes have no
/// OTLP equivalent and are encoded as strings. `Null` becomes an empty `AnyValue`.
impl From<Value> for AnyValue {
fn from(value: Value) -> Self {
let value = match value {
Value::Bytes(_) => PBValue::StringValue(
value
.to_str_lossy()
.expect("`Value::Bytes` always converts to a string")
.into_owned(),
),
Value::String(string) => PBValue::StringValue(string.into()),
Value::Regex(regex) => PBValue::StringValue(regex.as_str().to_owned()),
Value::Integer(int) => PBValue::IntValue(int),
Value::Float(float) => PBValue::DoubleValue(float.into_inner()),
Value::Boolean(boolean) => PBValue::BoolValue(boolean),
Value::Timestamp(timestamp) => {
PBValue::StringValue(timestamp.to_rfc3339_opts(SecondsFormat::AutoSi, true))
}
Value::Object(object) => PBValue::KvlistValue(KeyValueList {
values: object_into_kv_list(object),
}),
Value::Array(array) => PBValue::ArrayValue(ArrayValue {
values: array.into_iter().map(Into::into).collect(),
}),
Value::Null => return Self { value: None },
};
Self { value: Some(value) }
}
}

impl From<PBValue> for TagValue {
fn from(pb: PBValue) -> Self {
match pb {
Expand Down Expand Up @@ -80,6 +115,18 @@ pub fn kv_list_into_value(arr: Vec<KeyValue>) -> Value {
)
}

/// Inverse of [`kv_list_into_value`].
#[must_use]
pub fn object_into_kv_list(object: ObjectMap) -> Vec<KeyValue> {
object
.into_iter()
.map(|(key, value)| KeyValue {
key: key.into(),
value: Some(value.into()),
})
.collect()
}

#[must_use]
pub fn to_hex(d: &[u8]) -> String {
if d.is_empty() {
Expand Down
Loading
Loading