Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
22 changes: 10 additions & 12 deletions lib/codecs/src/decoding/format/gelf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -316,38 +316,36 @@ mod tests {

assert_eq!(
log.get(event_path!(VERSION)),
Some(&Value::Bytes(Bytes::from_static(b"1.1")))
Some(&Value::from_static_str("1.1"))
);
assert_eq!(
log.get(event_path!(HOST)),
Some(&Value::Bytes(Bytes::from_static(b"example.org")))
Some(&Value::from_static_str("example.org"))
);
assert_eq!(
log.get(log_schema().message_key_target_path().unwrap()),
Some(&Value::Bytes(Bytes::from_static(
b"A short message that helps you identify what is going on"
)))
Some(&Value::from_static_str(
"A short message that helps you identify what is going on"
))
);
assert_eq!(
log.get(event_path!(FULL_MESSAGE)),
Some(&Value::Bytes(Bytes::from_static(
b"Backtrace here\n\nmore stuff"
)))
Some(&Value::from_static_str("Backtrace here\n\nmore stuff"))
);
let dt = DateTime::from_timestamp(1_385_053_862, 307_200_000).expect("invalid timestamp");
assert_eq!(log.get(event_path!(TIMESTAMP)), Some(&Value::Timestamp(dt)));
assert_eq!(log.get(event_path!(LEVEL)), Some(&Value::Integer(1)));
assert_eq!(
log.get(event_path!(FACILITY)),
Some(&Value::Bytes(Bytes::from_static(b"foo")))
Some(&Value::from_static_str("foo"))
);
assert_eq!(
log.get(event_path!(LINE)),
Some(&Value::Float(ordered_float::NotNan::new(42.0).unwrap()))
);
assert_eq!(
log.get(event_path!(FILE)),
Some(&Value::Bytes(Bytes::from_static(b"/tmp/bar")))
Some(&Value::from_static_str("/tmp/bar"))
);
assert_eq!(
log.get(event_path!(add_on_int_in)),
Expand All @@ -357,7 +355,7 @@ mod tests {
);
assert_eq!(
log.get(event_path!(add_on_str_in)),
Some(&Value::Bytes(Bytes::from_static(b"A Space Odyssey")))
Some(&Value::from_static_str("A Space Odyssey"))
);
}

Expand Down Expand Up @@ -486,7 +484,7 @@ mod tests {

assert_eq!(
log.get(event_path!(VERSION)),
Some(&Value::Bytes(Bytes::from_static(b"1.0")))
Some(&Value::from_static_str("1.0"))
);

assert_eq!(
Expand Down
2 changes: 1 addition & 1 deletion lib/codecs/src/decoding/format/syslog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -297,7 +297,7 @@ impl Deserializer for SyslogDeserializer {
syslog_loose::parse_message_with_year_exact(line, resolve_year, Variant::Either)?;

let log = if let (Some(source), LogNamespace::Vector) = (self.source, log_namespace) {
let mut log = LogEvent::from(Value::Bytes(Bytes::from(parsed.msg.to_string())));
let mut log = LogEvent::from(Value::from(parsed.msg));
insert_metadata_fields_from_syslog(&mut log, source, parsed, log_namespace);
log
} else {
Expand Down
2 changes: 1 addition & 1 deletion lib/opentelemetry-proto/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ use super::proto::common::v1::{AnyValue, ArrayValue, KeyValue, any_value::Value
impl From<PBValue> for Value {
fn from(av: PBValue) -> Self {
match av {
PBValue::StringValue(v) => Value::Bytes(Bytes::from(v)),
PBValue::StringValue(v) => Value::from(v),
PBValue::BoolValue(v) => Value::Boolean(v),
PBValue::IntValue(v) => Value::Integer(v),
PBValue::DoubleValue(v) => NotNan::new(v).map_or(Value::Null, Value::Float),
Expand Down
7 changes: 3 additions & 4 deletions lib/opentelemetry-proto/src/logs.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,3 @@
use bytes::Bytes;
use chrono::{DateTime, TimeZone, Utc};
use vector_core::{
config::{LegacyKey, LogNamespace, log_schema},
Expand Down Expand Up @@ -147,7 +146,7 @@ impl ResourceLog {
&mut log,
Some(LegacyKey::Overwrite(path!(TRACE_ID_KEY))),
path!(TRACE_ID_KEY),
Bytes::from(to_hex(&self.log_record.trace_id)),
Value::from(to_hex(&self.log_record.trace_id)),
);
}
if !self.log_record.span_id.is_empty() {
Expand All @@ -156,7 +155,7 @@ impl ResourceLog {
&mut log,
Some(LegacyKey::Overwrite(path!(SPAN_ID_KEY))),
path!(SPAN_ID_KEY),
Bytes::from(to_hex(&self.log_record.span_id)),
Value::from(to_hex(&self.log_record.span_id)),
);
}
if !self.log_record.severity_text.is_empty() {
Expand Down Expand Up @@ -240,7 +239,7 @@ impl ResourceLog {
&mut log,
log_schema().source_type_key(),
path!("source_type"),
Bytes::from_static(SOURCE_NAME.as_bytes()),
Value::from_static_str(SOURCE_NAME),
);
if log_namespace == LogNamespace::Vector {
log.metadata_mut()
Expand Down
3 changes: 1 addition & 2 deletions lib/vector-core/src/config/mod.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
use std::{collections::HashMap, fmt, num::NonZeroUsize, sync::Arc};

use bitmask_enum::bitmask;
use bytes::Bytes;
use chrono::{DateTime, Utc};

mod global_options;
Expand Down Expand Up @@ -501,7 +500,7 @@ impl LogNamespace {
log,
log_schema().source_type_key(),
path!("source_type"),
Bytes::from_static(source_name.as_bytes()),
Value::from_static_str(source_name),
Comment thread
pront marked this conversation as resolved.
);
self.insert_vector_metadata(
log,
Expand Down
6 changes: 3 additions & 3 deletions lib/vector-core/src/event/log_event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -792,9 +792,9 @@ impl From<&tracing::Event<'_>> for LogEvent {
log.insert(
&TRACING_TARGET_PATHS.kind,
if meta.is_event() {
Value::Bytes("event".to_string().into())
Value::from_static_str("event")
} else if meta.is_span() {
Value::Bytes("span".to_string().into())
Value::from_static_str("span")
} else {
Value::Null
},
Expand All @@ -803,7 +803,7 @@ impl From<&tracing::Event<'_>> for LogEvent {
log.insert(
&TRACING_TARGET_PATHS.module_path,
meta.module_path()
.map_or(Value::Null, |mp| Value::Bytes(mp.to_string().into())),
.map_or(Value::Null, |mp| Value::from(mp.to_string())),
);
log.insert(&TRACING_TARGET_PATHS.target, meta.target().to_string());
log
Expand Down
2 changes: 1 addition & 1 deletion lib/vector-vrl/metrics/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,7 @@ pub(crate) fn metric_into_vrl(value: &Metric) -> Value {
Value::Array(
v.iter()
.filter_map(|v| {
v.map(ToString::to_string).map(Into::into).map(Value::Bytes)
v.map(ToString::to_string).map(Value::from)
})
.collect(),
),
Expand Down
33 changes: 25 additions & 8 deletions src/conditions/datadog_search.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
use std::{borrow::Cow, str::FromStr};

use bytes::Bytes;
use vector_lib::{
configurable::configurable_component,
event::{Event, EventRef, LogEvent, Value},
Expand Down Expand Up @@ -140,17 +139,13 @@ impl Filter<LogEvent> for EventFilter {
Field::Reserved(field) if field == "tags" => {
let to_match = to_match.to_owned();

array_match_multiple(vec!["ddtags", "tags"], move |values| {
values.contains(&Value::Bytes(Bytes::copy_from_slice(to_match.as_bytes())))
})
any_string_match_multiple(vec!["ddtags", "tags"], move |value| value == to_match)
}
// Individual tags are compared by element key:value.
Field::Tag(tag) => {
let value_bytes = Value::Bytes(format!("{tag}:{to_match}").into());
let to_match = format!("{tag}:{to_match}");

array_match_multiple(vec!["ddtags", "tags"], move |values| {
values.contains(&value_bytes)
})
any_string_match_multiple(vec!["ddtags", "tags"], move |value| value == to_match)
}
// A literal "source" field should string match in "source" and "ddsource" fields (OR condition).
Field::Reserved(field) if field == "source" => {
Expand Down Expand Up @@ -1648,6 +1643,28 @@ mod test {
test_filter(EventFilter, vector_lib::event::Event::into_log);
}

#[test]
fn tag_equality_matches_byte_values() {
for query in ["tags:foo", "env:prod"] {
let config: DatadogSearchConfig = query.parse().unwrap();
let runner = DatadogSearchRunner::try_from(&config).unwrap();
let mut log = LogEvent::default();
log.insert(
vrl::event_path!("tags"),
Value::Array(vec![Value::Bytes(
if query == "tags:foo" {
"foo"
} else {
"env:prod"
}
.into(),
)]),
);

assert!(runner.matches(&Event::Log(log)), "query: {query}");
}
}

#[test]
fn generate_config() {
crate::test_util::test_generate_config::<DatadogSearchConfig>();
Expand Down
6 changes: 1 addition & 5 deletions src/enrichment_tables/memory/bloom_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ use std::{

use async_trait::async_trait;
use bloomy::{BloomFilter, bloom};
use bytes::Bytes;
use futures::{
Stream, StreamExt,
stream::{self, BoxStream},
Expand Down Expand Up @@ -164,10 +163,7 @@ impl Table for BloomMemoryTable {
include_key_metric_tag: self.config.internal_metrics.include_key_tag
});
let result = ObjectMap::from([
(
KeyString::from("key"),
Value::Bytes(Bytes::copy_from_slice(key.as_bytes())),
),
(KeyString::from("key"), Value::from(key)),
(KeyString::from("value"), Value::Null),
]);
Ok(vec![result])
Expand Down
11 changes: 2 additions & 9 deletions src/enrichment_tables/memory/cuckoo_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,6 @@ use std::{
};

use async_trait::async_trait;
use bytes::Bytes;
use cuckoo_clock::{
CuckooFilter, ExportableRandomState, InsertValues, LookupValues,
config::{CounterConfig, CuckooConfiguration, LruAgingStrategy, LruConfig, TtlConfig},
Expand Down Expand Up @@ -646,16 +645,10 @@ impl Table for CuckooMemoryTable {
include_key_metric_tag: self.config.internal_metrics.include_key_tag
});
let mut result = ObjectMap::from([
(
KeyString::from("key"),
Value::Bytes(Bytes::copy_from_slice(key.as_bytes())),
),
(KeyString::from("key"), Value::from(key)),
(
KeyString::from("fingerprint"),
Value::Bytes(Bytes::from(format!(
"{:X}",
associated_data.get_fingerprint()
))),
Value::from(format!("{:X}", associated_data.get_fingerprint())),
),
(KeyString::from("value"), Value::Null),
]);
Expand Down
6 changes: 1 addition & 5 deletions src/enrichment_tables/memory/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ use std::{
};

use async_trait::async_trait;
use bytes::Bytes;
use evmap::{
shallow_copy::CopyValue,
{self},
Expand Down Expand Up @@ -70,10 +69,7 @@ impl MemoryEntry {
.ttl
.saturating_sub(now.duration_since(*self.update_time).as_secs());
Ok(ObjectMap::from([
(
KeyString::from("key"),
Value::Bytes(Bytes::copy_from_slice(key.as_bytes())),
),
(KeyString::from("key"), Value::from(key)),
(
KeyString::from("value"),
// Unreachable in normal operation: `value` was serialized by `handle_value`.
Expand Down
5 changes: 2 additions & 3 deletions src/sinks/websocket_server/buffering.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
use std::{collections::VecDeque, net::SocketAddr, num::NonZeroUsize};

use bytes::Bytes;
use derivative::Derivative;
use tokio_tungstenite::tungstenite::{Message, handshake::server::Request};
use url::Url;
Expand All @@ -11,7 +10,7 @@ use vector_lib::{
event::{Event, MaybeAsLogMut},
lookup::lookup_v2::ConfigValuePath,
};
use vrl::prelude::VrlValueConvert;
use vrl::{prelude::VrlValueConvert, value::Value};

use crate::serde::default_decoding;

Expand Down Expand Up @@ -239,7 +238,7 @@ impl WsMessageBufferConfig for Option<MessageBufferingConfig> {
let mut buffer = [0; 36];
let uuid = message_id.hyphenated().encode_lower(&mut buffer);
log.value_mut()
.insert(message_id_path, Bytes::copy_from_slice(uuid.as_bytes()));
.insert(message_id_path, Value::from(uuid as &str));
}
message_id
}
Expand Down
4 changes: 2 additions & 2 deletions src/sources/amqp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ use vector_lib::{
internal_event::{CountByteSize, EventsReceived, InternalEventHandle as _},
lookup::{lookup_v2::OptionalValuePath, metadata_path, owned_value_path, path},
};
use vrl::value::Kind;
use vrl::value::{Kind, Value};

use crate::{
SourceSender,
Expand Down Expand Up @@ -286,7 +286,7 @@ fn populate_log_event(
log,
log_schema().source_type_key(),
path!("source_type"),
Bytes::from_static(AmqpSourceConfig::NAME.as_bytes()),
Value::from_static_str(AmqpSourceConfig::NAME),
);

// This handles the transition from the original timestamp logic. Originally the
Expand Down
2 changes: 1 addition & 1 deletion src/sources/aws_kinesis_firehose/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ pub(super) async fn firehose(
log,
log_schema().source_type_key(),
path!("source_type"),
Bytes::from_static(AwsKinesisFirehoseConfig::NAME.as_bytes()),
Value::from_static_str(AwsKinesisFirehoseConfig::NAME),
Comment thread
bruceg marked this conversation as resolved.
);
// This handles the transition from the original timestamp logic. Originally the
// `timestamp_key` was always populated by the `request.timestamp` time.
Expand Down
10 changes: 5 additions & 5 deletions src/sources/aws_kinesis_firehose/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -607,7 +607,7 @@ mod tests {
assert_event_data_eq!(
events,
vec![log_event! {
"source_type" => Bytes::from("aws_kinesis_firehose"),
"source_type" => "aws_kinesis_firehose",
"timestamp" => timestamp.trunc_subsecs(3), // AWS sends timestamps as ms
"message" => Bytes::from(expected),
"request_id" => REQUEST_ID,
Expand Down Expand Up @@ -792,7 +792,7 @@ mod tests {
assert_event_data_eq!(
events,
vec![log_event! {
"source_type" => Bytes::from("aws_kinesis_firehose"),
"source_type" => "aws_kinesis_firehose",
"timestamp" => timestamp.trunc_subsecs(3), // AWS sends timestamps as ms
"message"=> RECORD,
"request_id" => REQUEST_ID,
Expand Down Expand Up @@ -839,7 +839,7 @@ mod tests {
assert_event_data_eq!(
events,
vec![log_event! {
"source_type" => Bytes::from("aws_kinesis_firehose"),
"source_type" => "aws_kinesis_firehose",
"timestamp" => timestamp.trunc_subsecs(3), // AWS sends timestamps as ms
"message"=> RECORD,
"request_id" => REQUEST_ID,
Expand Down Expand Up @@ -975,7 +975,7 @@ mod tests {
assert_event_data_eq!(
events,
vec![log_event! {
"source_type" => Bytes::from("aws_kinesis_firehose"),
"source_type" => "aws_kinesis_firehose",
"timestamp" => timestamp.trunc_subsecs(3), // AWS sends timestamps as ms
"message"=> RECORD,
"request_id" => REQUEST_ID,
Expand Down Expand Up @@ -1309,7 +1309,7 @@ mod tests {
assert_event_data_eq!(
events,
vec![log_event! {
"source_type" => Bytes::from("aws_kinesis_firehose"),
"source_type" => "aws_kinesis_firehose",
"timestamp" => timestamp.trunc_subsecs(3), // AWS sends timestamps as ms
"message"=> Bytes::from(expected),
"request_id" => REQUEST_ID,
Expand Down
Loading
Loading