diff --git a/relay-config/src/config.rs b/relay-config/src/config.rs index e57e2e1f9f2..20a58f806ed 100644 --- a/relay-config/src/config.rs +++ b/relay-config/src/config.rs @@ -639,6 +639,8 @@ pub struct Limits { pub max_trace_metric_size: ByteSize, /// The maximum payload size for a log. pub max_log_size: ByteSize, + /// The maximum number of logs that can result from a log expansion. + pub max_expanded_log_count: usize, /// The maximum payload size for a span. pub max_span_size: ByteSize, /// The maximum payload size for an item container. @@ -728,6 +730,7 @@ impl Default for Limits { max_profile_size: ByteSize::mebibytes(50), max_trace_metric_size: ByteSize::mebibytes(1), max_log_size: ByteSize::mebibytes(1), + max_expanded_log_count: 1000, max_span_size: ByteSize::mebibytes(10), max_container_size: ByteSize::mebibytes(12), max_statsd_size: ByteSize::mebibytes(1), @@ -2412,6 +2415,11 @@ impl Config { self.values.limits.max_log_size.as_bytes() } + /// Returns the maximum number of logs that can result from a log expansion. + pub fn max_expanded_log_count(&self) -> usize { + self.values.limits.max_expanded_log_count + } + /// Returns the maximum payload size of a span in bytes. pub fn max_span_size(&self) -> usize { self.values.limits.max_span_size.as_bytes() diff --git a/relay-server/src/processing/logs/integrations/mod.rs b/relay-server/src/processing/logs/integrations/mod.rs index 29c2b070c03..7aa67e82c53 100644 --- a/relay-server/src/processing/logs/integrations/mod.rs +++ b/relay-server/src/processing/logs/integrations/mod.rs @@ -8,6 +8,8 @@ use crate::processing::logs::Settings; mod nel; mod otel; +mod otel_json_deserializer; +mod otel_proto_deserializer; mod vercel; /// Expands a log [`Integration`] into a list of logs. @@ -17,6 +19,7 @@ pub fn expand( item: Item, records: &mut RecordKeeper<'_>, headers: &EnvelopeHeaders, + max_expanded_log_count: usize, ) -> Option<(Settings, ContainerItems)> { let integration = match item.integration() { Some(Integration::Logs(integration)) => integration, @@ -46,7 +49,9 @@ pub fn expand( let settings = match integration { LogsIntegration::Nel => nel::expand(&payload, headers, produce), - LogsIntegration::OtelV1 { format } => otel::expand(format, &payload, produce), + LogsIntegration::OtelV1 { format } => { + otel::expand(format, &payload, max_expanded_log_count, produce) + } LogsIntegration::VercelDrainLog { format } => vercel::expand(format, &payload, produce), }; let settings = match settings { diff --git a/relay-server/src/processing/logs/integrations/otel.rs b/relay-server/src/processing/logs/integrations/otel.rs index 3d51c5b0cb0..d81cb393421 100644 --- a/relay-server/src/processing/logs/integrations/otel.rs +++ b/relay-server/src/processing/logs/integrations/otel.rs @@ -1,17 +1,23 @@ -use opentelemetry_proto::tonic::logs::v1::LogsData; -use prost::Message as _; -use relay_event_schema::protocol::OurLog; - use crate::integrations::OtelFormat; -use crate::processing::logs::{Error, Result, Settings}; +use crate::processing::logs::{ + Error, Result, Settings, + integrations::{otel_json_deserializer, otel_proto_deserializer}, +}; use crate::services::outcome::DiscardReason; +use opentelemetry_proto::tonic::logs::v1::LogsData; +use relay_event_schema::protocol::OurLog; /// Expands OTeL logs into the [`OurLog`] format. -pub fn expand(format: OtelFormat, payload: &[u8], mut produce: F) -> Result +pub fn expand( + format: OtelFormat, + payload: &[u8], + max_logs: usize, + mut produce: F, +) -> Result where F: FnMut(OurLog), { - let logs = parse_logs_data(format, payload)?; + let logs = parse_logs_data(format, payload, max_logs)?; for resource_logs in logs.resource_logs { let resource = resource_logs.resource.as_ref(); @@ -27,21 +33,129 @@ where Ok(Settings::default()) } -fn parse_logs_data(format: OtelFormat, payload: &[u8]) -> Result { +fn parse_logs_data(format: OtelFormat, payload: &[u8], max_logs: usize) -> Result { match format { - OtelFormat::Json => serde_json::from_slice(payload).map_err(|e| { + OtelFormat::Json => otel_json_deserializer::deserialize(payload, max_logs).map_err(|e| { relay_log::debug!( error = &e as &dyn std::error::Error, "Failed to parse logs data as JSON" ); Error::Invalid(DiscardReason::InvalidJson) }), - OtelFormat::Protobuf => LogsData::decode(payload).map_err(|e| { - relay_log::debug!( - error = &e as &dyn std::error::Error, - "Failed to parse logs data as protobuf" - ); - Error::Invalid(DiscardReason::InvalidProtobuf) - }), + OtelFormat::Protobuf => { + otel_proto_deserializer::deserialize(payload, max_logs).map_err(|e| { + relay_log::debug!( + error = &e as &dyn std::error::Error, + "Failed to parse logs data as protobuf" + ); + Error::Invalid(DiscardReason::InvalidProtobuf) + }) + } + } +} +#[cfg(test)] +mod tests { + + use opentelemetry_proto::tonic::common::v1::any_value::Value; + use opentelemetry_proto::tonic::common::v1::{ + AnyValue, ArrayValue, InstrumentationScope, KeyValue, + }; + use opentelemetry_proto::tonic::resource::v1::Resource; + use relay_ourlogs::otel_logs::{LogRecord, LogsData, ResourceLogs, ScopeLogs}; + + use crate::processing::logs::integrations::otel::parse_logs_data; + + #[test] + fn test_basic_json() { + let log_data = LogsData { + resource_logs: vec![ResourceLogs { + resource: Some(Resource { + attributes: vec![KeyValue { + key: "service.name".to_owned(), + value: Some(AnyValue { + value: Some(Value::StringValue("test-service".to_owned())), + }), + }], + dropped_attributes_count: 0, + entity_refs: vec![], + }), + scope_logs: vec![ScopeLogs { + scope: Some(InstrumentationScope { + name: "test-library".to_owned(), + version: "".to_owned(), + attributes: vec![], + dropped_attributes_count: 0, + }), + log_records: vec![LogRecord { + time_unix_nano: 123, + observed_time_unix_nano: 123, + severity_number: 2, + severity_text: "Information".to_owned(), + body: Some(AnyValue { + value: Some(Value::StringValue("a body".to_owned())), + }), + attributes: vec![ + KeyValue { + key: "attribute".to_owned(), + value: Some(AnyValue { + value: Some(Value::StringValue("value".to_owned())), + }), + }, + KeyValue { + key: "nested attribute".to_owned(), + value: Some(AnyValue { + value: Some(Value::ArrayValue(ArrayValue { + values: vec![AnyValue { + value: Some(Value::StringValue("value".to_owned())), + }], + })), + }), + }, + ], + dropped_attributes_count: 0, + flags: 0, + trace_id: "5B8EFFF798038103D269B633813FC60C".into(), + span_id: "EEE19B7EC3C1B174".into(), + event_name: "".to_owned(), + }], + schema_url: "".to_owned(), + }], + schema_url: "http://example.com".to_owned(), + }], + }; + + let json = serde_json::to_string(&log_data).unwrap(); + + let unjson: LogsData = serde_json::from_str(&json).unwrap(); + + assert_eq!(unjson, log_data); + } + + #[test] + fn test_abusive_json() { + let mut abusive_log = "{},".repeat(1_001); + abusive_log.pop(); + + let json = r#"{ + "resourceLogs": [ + { + "resource": { + "attributes": [ + { + "key": "service.name", + "value": {"stringValue": "test-service"} + } + ] + }, + "scopeLogs": [ + { + "scope": {"name": "test-library"}, + "logRecords": ["# + .to_owned(); + let json = json + &abusive_log + "]}]}]}"; + + assert!( + parse_logs_data(crate::integrations::OtelFormat::Json, json.as_bytes(), 1000).is_err() + ); } } diff --git a/relay-server/src/processing/logs/integrations/otel_json_deserializer.rs b/relay-server/src/processing/logs/integrations/otel_json_deserializer.rs new file mode 100644 index 00000000000..1101f98d9d3 --- /dev/null +++ b/relay-server/src/processing/logs/integrations/otel_json_deserializer.rs @@ -0,0 +1,260 @@ +use std::cell::Cell; +use std::fmt::{self, Display}; + +use opentelemetry_proto::tonic::common::v1::InstrumentationScope; +use opentelemetry_proto::tonic::logs::*; +use opentelemetry_proto::tonic::resource::v1::Resource; +use serde::Deserialize; +use serde::de::{self, DeserializeSeed, Deserializer, IgnoredAny, MapAccess, SeqAccess, Visitor}; + +use crate::processing::logs; +use crate::services::outcome::DiscardReason; + +struct Meter { + remaining: Cell, +} + +struct MeterError {} + +impl Display for MeterError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "ran out of meter") + } +} + +impl Meter { + fn new(max: usize) -> Self { + Self { + remaining: Cell::new(max), + } + } + + fn spend(&self) -> Result<(), E> { + match self.remaining.get() { + 0 => Err(de::Error::custom(MeterError {})), + n => { + self.remaining.set(n - 1); + Ok(()) + } + } + } + + fn is_empty(&self) -> bool { + self.remaining.get() == 0 + } +} + +// A deserializer for generic arrays that can pass deserializer state ("seed") to the +// children elements. +struct Array(S); + +impl<'de, S: DeserializeSeed<'de> + Copy> DeserializeSeed<'de> for Array { + type Value = Vec; + + fn deserialize>(self, d: D) -> Result { + d.deserialize_seq(self) + } +} + +impl<'de, S: DeserializeSeed<'de> + Copy> Visitor<'de> for Array { + type Value = Vec; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("an array") + } + + fn visit_seq>(self, mut seq: A) -> Result { + let mut out = Vec::new(); + while let Some(item) = seq.next_element_seed(self.0)? { + out.push(item); + } + Ok(out) + } +} + +/// Mapping for OTEL v1::LogsData +#[derive(Copy, Clone)] +struct LogsDataSeed<'b>(&'b Meter); + +#[derive(Deserialize)] +#[serde(field_identifier, rename_all = "camelCase")] +enum LogsDataField { + ResourceLogs, + #[serde(other)] + Other, +} + +impl<'de> DeserializeSeed<'de> for LogsDataSeed<'_> { + type Value = v1::LogsData; + + fn deserialize>(self, d: D) -> Result { + d.deserialize_map(self) + } +} + +impl<'de> Visitor<'de> for LogsDataSeed<'_> { + type Value = v1::LogsData; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("logs data") + } + + fn visit_map>(self, mut map: A) -> Result { + let mut resource_logs = Vec::new(); + while let Some(field) = map.next_key::()? { + match field { + LogsDataField::ResourceLogs => { + resource_logs = map.next_value_seed(Array(ResourceLogsSeed(self.0)))?; + } + LogsDataField::Other => drop(map.next_value::()?), + } + } + Ok(v1::LogsData { resource_logs }) + } +} + +/// Mapping for OTEL v1::ResourceLogs +#[derive(Copy, Clone)] +struct ResourceLogsSeed<'b>(&'b Meter); + +#[derive(Deserialize)] +#[serde(field_identifier, rename_all = "camelCase")] +enum ResourceLogsField { + Resource, + ScopeLogs, + SchemaUrl, + #[serde(other)] + Other, +} + +impl<'de> DeserializeSeed<'de> for ResourceLogsSeed<'_> { + type Value = v1::ResourceLogs; + + fn deserialize>(self, d: D) -> Result { + self.0.spend()?; + d.deserialize_map(self) + } +} + +impl<'de> Visitor<'de> for ResourceLogsSeed<'_> { + type Value = v1::ResourceLogs; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("resource logs") + } + + fn visit_map>(self, mut map: A) -> Result { + let mut out = v1::ResourceLogs::default(); + while let Some(field) = map.next_key::()? { + match field { + // Below this point the derived impls take over unchanged. + ResourceLogsField::Resource => { + out.resource = map.next_value::>()? + } + ResourceLogsField::SchemaUrl => out.schema_url = map.next_value()?, + ResourceLogsField::ScopeLogs => { + out.scope_logs = map.next_value_seed(Array(ScopeLogsSeed(self.0)))?; + } + ResourceLogsField::Other => drop(map.next_value::()?), + } + } + Ok(out) + } +} + +/// Mapping for OTEL v1::ScopeLogs +#[derive(Copy, Clone)] +struct ScopeLogsSeed<'b>(&'b Meter); + +#[derive(Deserialize)] +#[serde(field_identifier, rename_all = "camelCase")] +enum ScopeLogsField { + Scope, + LogRecords, + SchemaUrl, + #[serde(other)] + Other, +} + +impl<'de> DeserializeSeed<'de> for ScopeLogsSeed<'_> { + type Value = v1::ScopeLogs; + + fn deserialize>(self, d: D) -> Result { + self.0.spend()?; + d.deserialize_map(self) + } +} + +impl<'de> Visitor<'de> for ScopeLogsSeed<'_> { + type Value = v1::ScopeLogs; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("scope logs") + } + + fn visit_map>(self, mut map: A) -> Result { + let mut out = v1::ScopeLogs::default(); + while let Some(field) = map.next_key::()? { + match field { + ScopeLogsField::Scope => { + out.scope = map.next_value::>()? + } + ScopeLogsField::SchemaUrl => out.schema_url = map.next_value()?, + ScopeLogsField::LogRecords => { + out.log_records = map.next_value_seed(LogRecordsSeed(self.0))?; + } + ScopeLogsField::Other => drop(map.next_value::()?), + } + } + Ok(out) + } +} + +/// Mapping for OTEL v1::LogRecord +#[derive(Copy, Clone)] +struct LogRecordsSeed<'b>(&'b Meter); + +impl<'de> DeserializeSeed<'de> for LogRecordsSeed<'_> { + type Value = Vec; + + fn deserialize>(self, d: D) -> Result { + d.deserialize_seq(self) + } +} + +impl<'de> Visitor<'de> for LogRecordsSeed<'_> { + type Value = Vec; + + fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str("log records") + } + + fn visit_seq>(self, mut seq: A) -> Result { + let mut out = Vec::new(); + while let Some(record) = seq.next_element::()? { + self.0.spend()?; + out.push(record); + } + Ok(out) + } +} + +/// Deserialize the supplied JSON into v1::LogsData, returning an error if the number of logs +/// elements in the payload exceeds the supplied maximum. +pub fn deserialize(payload: &[u8], max: usize) -> Result { + let budget = Meter::new(max); + let mut de = serde_json::Deserializer::from_slice(payload); + let res = LogsDataSeed(&budget).deserialize(&mut de); + + match res { + Ok(logs) => { + // Only call 'end' if we're on the Ok path (end expects EOF/trailing whitespace, will + // error otherwise.) + de.end() + .map_err(|_| logs::Error::Invalid(DiscardReason::InvalidJson))?; + Ok(logs) + } + Err(_) if budget.is_empty() => Err(logs::Error::TooManyExpandedLogs), + Err(_) => Err(logs::Error::Invalid(DiscardReason::InvalidJson)), + } +} diff --git a/relay-server/src/processing/logs/integrations/otel_proto_deserializer.rs b/relay-server/src/processing/logs/integrations/otel_proto_deserializer.rs new file mode 100644 index 00000000000..14d67af2fb8 --- /dev/null +++ b/relay-server/src/processing/logs/integrations/otel_proto_deserializer.rs @@ -0,0 +1,318 @@ +use bytes::Buf; +use opentelemetry_proto::tonic::logs::*; +use prost::Message as _; +use prost::encoding::{DecodeContext, WireType, check_wire_type, decode_key, decode_varint}; +use std::fmt::Display; + +use crate::processing::logs; +use crate::services::outcome::DiscardReason; + +/// Field tag of `LogsData::resource_logs`. +const TAG_RESOURCE_LOGS: u32 = 1; +/// Field tag of `ResourceLogs::scope_logs`. +const TAG_SCOPE_LOGS: u32 = 2; +/// Field tag of `ScopeLogs::log_records`. +const TAG_LOG_RECORDS: u32 = 2; + +pub enum Error { + MeterExhausted, + BufferUnderflow, + DelimitedLengthExceeded, + ProstError(prost::DecodeError), +} + +struct Meter { + remaining: usize, +} + +impl Display for Error { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Error::MeterExhausted => write!(f, "meter exhausted"), + Error::BufferUnderflow => write!(f, "buffer underflow "), + Error::DelimitedLengthExceeded => write!(f, "delimited length exceeded "), + Error::ProstError(decode_error) => { + write!(f, "prost decoding error: {}", decode_error) + } + } + } +} + +impl From for Error { + fn from(value: prost::DecodeError) -> Self { + Self::ProstError(value) + } +} + +impl Meter { + fn new(max: usize) -> Self { + Self { remaining: max } + } + + fn spend(&mut self) -> Result<(), Error> { + if self.remaining == 0 { + Err(Error::MeterExhausted) + } else { + self.remaining -= 1; + Ok(()) + } + } + + fn is_empty(&self) -> bool { + self.remaining == 0 + } +} + +/// Merges a top-level `v1::LogsData` message, which spans the entire buffer. +fn merge_logs_data( + meter: &mut Meter, + out: &mut v1::LogsData, + buf: &mut B, + ctx: DecodeContext, +) -> Result<(), Error> { + while buf.has_remaining() { + let (tag, wire_type) = decode_key(buf)?; + merge_logs_data_field(meter, out, tag, wire_type, buf, ctx.clone())?; + } + Ok(()) +} + +fn merge_logs_data_field( + meter: &mut Meter, + out: &mut v1::LogsData, + tag: u32, + wire_type: WireType, + buf: &mut B, + ctx: DecodeContext, +) -> Result<(), Error> { + match tag { + TAG_RESOURCE_LOGS => { + check_wire_type(WireType::LengthDelimited, wire_type)?; + meter.spend()?; + + let mut resource_logs = v1::ResourceLogs::default(); + merge_loop(&mut resource_logs, buf, ctx, |value, buf, ctx| { + let (tag, wire_type) = decode_key(buf)?; + merge_resource_logs_field(meter, value, tag, wire_type, buf, ctx) + })?; + out.resource_logs.push(resource_logs); + Ok(()) + } + _ => Ok(out.merge_field(tag, wire_type, buf, ctx)?), + } +} + +fn merge_resource_logs_field( + meter: &mut Meter, + out: &mut v1::ResourceLogs, + tag: u32, + wire_type: WireType, + buf: &mut B, + ctx: DecodeContext, +) -> Result<(), Error> { + match tag { + TAG_SCOPE_LOGS => { + check_wire_type(WireType::LengthDelimited, wire_type)?; + meter.spend()?; + + let mut scope_logs = v1::ScopeLogs::default(); + merge_loop(&mut scope_logs, buf, ctx, |value, buf, ctx| { + let (tag, wire_type) = decode_key(buf)?; + merge_scope_logs_field(meter, value, tag, wire_type, buf, ctx) + })?; + out.scope_logs.push(scope_logs); + Ok(()) + } + _ => Ok(out.merge_field(tag, wire_type, buf, ctx)?), + } +} + +fn merge_scope_logs_field( + meter: &mut Meter, + out: &mut v1::ScopeLogs, + tag: u32, + wire_type: WireType, + buf: &mut B, + ctx: DecodeContext, +) -> Result<(), Error> { + if tag == TAG_LOG_RECORDS { + meter.spend()?; + } + Ok(out.merge_field(tag, wire_type, buf, ctx)?) +} + +fn merge_loop( + value: &mut T, + buf: &mut B, + ctx: DecodeContext, + mut merge: M, +) -> Result<(), Error> +where + M: FnMut(&mut T, &mut B, DecodeContext) -> Result<(), Error>, + B: Buf, +{ + let len = decode_varint(buf)?; + let remaining = buf.remaining(); + if len > remaining as u64 { + return Err(Error::BufferUnderflow); + } + + let limit = remaining - len as usize; + while buf.remaining() > limit { + merge(value, buf, ctx.clone())?; + } + + if buf.remaining() != limit { + return Err(Error::DelimitedLengthExceeded); + } + Ok(()) +} + +/// Deserialize the supplied protobuf into v1::LogsData, returning an error if the number of logs +/// elements in the payload exceeds the supplied maximum. +pub fn deserialize(mut payload: &[u8], max: usize) -> Result { + let mut meter = Meter::new(max); + let mut out = v1::LogsData::default(); + + match merge_logs_data(&mut meter, &mut out, &mut payload, DecodeContext::default()) { + Ok(()) => Ok(out), + Err(_) if meter.is_empty() => Err(logs::Error::TooManyExpandedLogs), + Err(_) => Err(logs::Error::Invalid(DiscardReason::InvalidProtobuf)), + } +} + +#[cfg(test)] +mod tests { + use opentelemetry_proto::tonic::common::v1::any_value::Value; + use opentelemetry_proto::tonic::common::v1::{AnyValue, InstrumentationScope, KeyValue}; + use opentelemetry_proto::tonic::logs::v1::*; + use opentelemetry_proto::tonic::resource::v1::Resource; + + use super::*; + + fn log_record(body: &str) -> LogRecord { + LogRecord { + time_unix_nano: 123, + observed_time_unix_nano: 123, + severity_number: 2, + severity_text: "Information".to_owned(), + body: Some(AnyValue { + value: Some(Value::StringValue(body.to_owned())), + }), + attributes: vec![KeyValue { + key: "attribute".to_owned(), + value: Some(AnyValue { + value: Some(Value::StringValue("value".to_owned())), + }), + }], + dropped_attributes_count: 0, + flags: 0, + trace_id: "5B8EFFF798038103D269B633813FC60C".into(), + span_id: "EEE19B7EC3C1B174".into(), + event_name: "".to_owned(), + } + } + + fn scope_logs(num_records: usize) -> ScopeLogs { + ScopeLogs { + scope: Some(InstrumentationScope { + name: "test-library".to_owned(), + version: "".to_owned(), + attributes: vec![], + dropped_attributes_count: 0, + }), + log_records: (0..num_records) + .map(|i| log_record(&i.to_string())) + .collect(), + schema_url: "".to_owned(), + } + } + + fn resource_logs(num_scope_logs: usize, records: usize) -> ResourceLogs { + ResourceLogs { + resource: Some(Resource { + attributes: vec![KeyValue { + key: "service.name".to_owned(), + value: Some(AnyValue { + value: Some(Value::StringValue("test-service".to_owned())), + }), + }], + dropped_attributes_count: 0, + entity_refs: vec![], + }), + scope_logs: (0..num_scope_logs).map(|_| scope_logs(records)).collect(), + schema_url: "http://example.com".to_owned(), + } + } + + fn logs_data(num_resource_logs: usize, num_scope_logs: usize, num_records: usize) -> LogsData { + LogsData { + resource_logs: (0..num_resource_logs) + .map(|_| resource_logs(num_scope_logs, num_records)) + .collect(), + } + } + + #[test] + fn test_basic_protobuf() { + let expected = logs_data(2, 2, 3); + let payload = expected.encode_to_vec(); + + assert_eq!(deserialize(&payload, 1000).unwrap(), expected); + } + + #[test] + fn test_empty_payload() { + assert_eq!(deserialize(&[], 1000).unwrap(), LogsData::default()); + } + + #[test] + fn test_invalid_protobuf() { + let err = deserialize(b"this is not protobuf", 1000).unwrap_err(); + assert!(matches!( + err, + logs::Error::Invalid(DiscardReason::InvalidProtobuf) + )); + } + + #[test] + fn test_exact_limit() { + let payload = logs_data(1, 1, 8).encode_to_vec(); + + assert!(deserialize(&payload, 10).is_ok()); + assert!(matches!( + deserialize(&payload, 9).unwrap_err(), + logs::Error::TooManyExpandedLogs + )); + } + + #[test] + fn test_too_many_log_records() { + let payload = logs_data(1, 1, 1_001).encode_to_vec(); + + assert!(matches!( + deserialize(&payload, 1000).unwrap_err(), + logs::Error::TooManyExpandedLogs + )); + } + + #[test] + fn test_too_many_scope_logs() { + let payload = logs_data(1, 1_001, 0).encode_to_vec(); + + assert!(matches!( + deserialize(&payload, 1000).unwrap_err(), + logs::Error::TooManyExpandedLogs + )); + } + + #[test] + fn test_too_many_resource_logs() { + let payload = logs_data(1_001, 0, 0).encode_to_vec(); + + assert!(matches!( + deserialize(&payload, 1000).unwrap_err(), + logs::Error::TooManyExpandedLogs + )); + } +} diff --git a/relay-server/src/processing/logs/mod.rs b/relay-server/src/processing/logs/mod.rs index fcc7634b0cd..fc382300fc1 100644 --- a/relay-server/src/processing/logs/mod.rs +++ b/relay-server/src/processing/logs/mod.rs @@ -55,6 +55,9 @@ pub enum Error { /// The log is invalid. #[error("invalid: {0}")] Invalid(DiscardReason), + /// The expanded logs exceed the maximum number allowed. + #[error("expanded logs exeeds limit")] + TooManyExpandedLogs, } impl OutcomeError for Error { @@ -74,6 +77,8 @@ impl OutcomeError for Error { } Self::ProcessingFailed(_) => Some(Outcome::Invalid(DiscardReason::Internal)), Self::Invalid(reason) => Some(Outcome::Invalid(*reason)), + // TODO: Or should this be abuse? Or should this be filtered, or rate-limited? + Self::TooManyExpandedLogs => Some(Outcome::Invalid(DiscardReason::InvalidLog)), }; (outcome, self) @@ -153,7 +158,7 @@ impl processing::Processor for LogsProcessor { // Fast filters, which do not need expanded logs. filter::feature_flag(ctx).reject(&logs)?; - let mut logs = process::expand(logs)?; + let mut logs = process::expand(logs, ctx.config.max_expanded_log_count())?; validate::size(&mut logs, ctx); diff --git a/relay-server/src/processing/logs/process.rs b/relay-server/src/processing/logs/process.rs index 88c891439ef..584c9bac9a3 100644 --- a/relay-server/src/processing/logs/process.rs +++ b/relay-server/src/processing/logs/process.rs @@ -17,7 +17,10 @@ use crate::services::outcome::DiscardReason; /// Parses all serialized logs into their [`ExpandedLogs`] representation. /// /// Individual, invalid logs will be discarded. -pub fn expand(logs: Managed) -> Result, Rejected> { +pub fn expand( + logs: Managed, + max_expanded_log_count: usize, +) -> Result, Rejected> { let trust = logs.headers.meta().request_trust(); logs.try_map(|logs, records| { @@ -44,7 +47,8 @@ pub fn expand(logs: Managed) -> Result, Re let (settings, logs) = match items { LogItems::Container(item) => expand_log_container(&item, trust)?, LogItems::Integration(item) => { - logs::integrations::expand(item, records, &headers).unwrap_or_default() + logs::integrations::expand(item, records, &headers, max_expanded_log_count) + .unwrap_or_default() } };