Skip to content
Merged
Show file tree
Hide file tree
Changes from 62 commits
Commits
Show all changes
63 commits
Select commit Hold shift + click to select a range
f469b3f
fix(buffers): add nesting depth limit to prevent protobuf decode corr…
connoryy Mar 26, 2026
b05d98a
fix: also check event metadata value for nesting depth
connoryy Mar 26, 2026
87d25cd
docs: replace approximate nesting formula with verified exact derivation
connoryy Mar 26, 2026
eff1c70
fix: lower MAX_NESTING_DEPTH from 33 to 32 to cover all proto paths
connoryy Mar 26, 2026
59de3d4
test: add EventWrapper path boundary test for vector sink gRPC
connoryy Mar 26, 2026
da1e2eb
docs: simplify MAX_NESTING_DEPTH comment to state only verified facts
connoryy Mar 26, 2026
fbd398c
docs: remove proto field names from MAX_NESTING_DEPTH comment
connoryy Mar 26, 2026
2191d9b
docs: simplify test comments to remove proto-specific terminology
connoryy Mar 26, 2026
5046ec5
make error slightly more informative
connoryy Apr 1, 2026
390c30b
move imports
connoryy Apr 1, 2026
3eb9fb4
add test
connoryy Apr 1, 2026
39a9924
fix conditional to new function signature
connoryy Apr 6, 2026
d835bc8
style: auto-fix lint/format errors
connoryy Apr 6, 2026
220be47
Merge branch 'master' into connor/protobuf-nesting-depth-limit
connoryy Apr 6, 2026
573dafb
fix metrics case and test failure
connoryy Apr 6, 2026
4d53caf
Merge branch 'connor/protobuf-nesting-depth-limit' of github.com:conn…
connoryy Apr 6, 2026
7936f12
Merge branch 'master' into connor/protobuf-nesting-depth-limit
connoryy Apr 15, 2026
6ff03c5
Add roundtrip tests for depth-32 metadata via prost
connoryy Apr 15, 2026
13bc237
Replace individual nesting tests with exhaustive path coverage
connoryy Apr 15, 2026
79e46b8
Simplify nesting tests: saturate all Value fields instead of enumerat…
connoryy Apr 15, 2026
6800e85
style: auto-fix lint/format errors
connoryy Apr 15, 2026
9ea18cd
Fix clippy doc_markdown lints in nesting depth tests
connoryy Apr 15, 2026
30e6a86
Merge branch 'connor/protobuf-nesting-depth-limit' of github.com:conn…
connoryy Apr 15, 2026
7b42478
Add tests proving metadata_full is the tightest path
connoryy Apr 16, 2026
bf5119c
style: auto-fix lint/format errors
connoryy Apr 16, 2026
5ebc616
Test full boundaries for loosest and tightest encoding paths
connoryy Apr 16, 2026
1021e1e
style: auto-fix lint/format errors
connoryy Apr 16, 2026
fbd2ad2
Add Trace to flat events test for uniformity
connoryy Apr 16, 2026
8c81391
Use per-path depth limits: 33 for event data, 32 for metadata
connoryy Apr 16, 2026
f0fc05a
style: auto-fix lint/format errors
connoryy Apr 16, 2026
ef06620
Fix clippy doc_markdown lint
connoryy Apr 16, 2026
2a58740
Fix doc comment arithmetic and add native codec metadata tests
connoryy Apr 16, 2026
5f87a9f
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
oganel May 12, 2026
87038f5
Fix compile
oganel May 12, 2026
57c9c91
fix(buffers): drop oversized events gracefully on disk-buffer write
oganel May 18, 2026
9d44c36
fix(buffers): account for array vs object cost in nesting check
oganel May 18, 2026
2447fda
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
ganelo May 19, 2026
4993bb5
fix(buffers): charge Value::Timestamp for one nesting frame
oganel May 19, 2026
b5d2e9d
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
pront Jun 12, 2026
0b639d4
Update changelog.d/protobuf_nesting_depth_limit.fix.md
ganelo Jun 16, 2026
eec68b2
refactor(buffers): move pre-encode filtering to Bufferable::filter_un…
oganel Jun 16, 2026
3598122
fix(buffers): account for filter drops in buffer usage instrumentation
oganel Jun 16, 2026
860f40a
fix(console sink): keep per-event encoder failures from terminating t…
oganel Jun 16, 2026
f28b2a3
fix(buffers): skip filter_unencodable when disk is already full
oganel Jun 16, 2026
8a22d86
fix(codecs): silent-drop over-budget events from the native encoder
oganel Jun 17, 2026
95ba0d0
test(vector sink): pin PushEventsRequest decode boundary at the same …
oganel Jun 17, 2026
806371d
fix(buffers): report disk-v2 filter drops via the ledger's usage handle
oganel Jun 17, 2026
13999c5
fix(sinks): finalize socket-sink encoder drops with explicit status
oganel Jun 17, 2026
e90dd47
fix(codecs): preserve framing for legitimate empty payloads
oganel Jun 18, 2026
b172e6c
refactor(codecs): drop native-encoder nesting guard, scope fix to dec…
oganel Jun 18, 2026
f57988f
docs(buffers): document over-budget overflow-routing limitation in tr…
oganel Jun 18, 2026
e1d219d
chore(buffers): merge master into protobuf nesting depth limit branch.
EricaJ6 Aug 12, 2026
20c5cbd
fix(buffers): route unencodable items to overflow regardless of buffe…
EricaJ6 Aug 12, 2026
261ddcd
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 12, 2026
4426682
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 13, 2026
23da462
Update changelog.d/protobuf_nesting_depth_limit.fix.md
EricaJ6 Aug 13, 2026
e47d598
fix(buffers): gate overflow diversion on the base stage's encoding co…
EricaJ6 Aug 13, 2026
4e997b0
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 13, 2026
a426b6b
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 13, 2026
9e165b1
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 13, 2026
f7ca90a
fix(buffers): use a common protobuf nesting budget
pront Aug 17, 2026
d40de33
Merge branch 'master' into og/continue-protobuf-nesting-depth-limit
EricaJ6 Aug 17, 2026
39b74b4
Update changelog.d/protobuf_nesting_depth_limit.fix.md
pront Aug 17, 2026
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
3 changes: 3 additions & 0 deletions changelog.d/protobuf_nesting_depth_limit.fix.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
Fixed unrecoverable disk buffer corruption and vector-to-vector retry loops caused by event data or metadata that protobuf could encode but prost could not decode. Vector now drops only protobuf-unsafe nested payloads before disk buffer or `vector` sink gRPC encoding, while preserving nested shapes that prost can safely decode. Under `when_full = "overflow"`, an event the buffer cannot encode is passed to the overflow stage intact rather than dropped, and does so regardless of how full the buffer currently is.
Comment thread
pront marked this conversation as resolved.
Outdated

authors: connoryy ganelo EricaJ6 jonodera97
4 changes: 3 additions & 1 deletion lib/vector-buffers/benches/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use metrics_util::{debugging::DebuggingRecorder, layers::Layer};
use tracing::Span;
use tracing_subscriber::prelude::__tracing_subscriber_SubscriberExt;
use vector_buffers::{
BufferType, EventCount,
BufferType, Bufferable, EventCount,
encoding::FixedEncodable,
topology::{
builder::TopologyBuilder,
Expand Down Expand Up @@ -57,6 +57,8 @@ impl<const N: usize> EventCount for Message<N> {
}
}

impl<const N: usize> Bufferable for Message<N> {}

impl<const N: usize> Finalizable for Message<N> {
fn take_finalizers(&mut self) -> EventFinalizers {
Default::default() // This benchmark doesn't need finalization
Expand Down
2 changes: 2 additions & 0 deletions lib/vector-buffers/examples/buffer_perf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,8 @@ impl EventCount for VariableMessage {
}
}

impl Bufferable for VariableMessage {}

impl Finalizable for VariableMessage {
fn take_finalizers(&mut self) -> EventFinalizers {
std::mem::take(&mut self.finalizers)
Expand Down
58 changes: 55 additions & 3 deletions lib/vector-buffers/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -129,10 +129,62 @@ impl<T> InMemoryBufferable for T where
/// An item that can be buffered.
///
/// This supertrait serves as the base trait for any item that can be pushed into a buffer.
pub trait Bufferable: InMemoryBufferable + Encodable + GroupedFinalizable {}
pub trait Bufferable: InMemoryBufferable + Encodable + GroupedFinalizable {
/// Drops any sub-items that cannot be persisted by the calling backend (e.g. due to
/// format-imposed nesting depth limits), reporting them as dropped via the appropriate
/// telemetry. Returns `None` if nothing remains worth writing.
///
/// # Who calls this
///
/// Only persistent backends with wire-format constraints invoke this — today that's
/// the disk-v2 sender (`SenderAdapter::send`/`try_send`). In-memory channels skip it
/// entirely because they hold the in-memory representation and have no nesting-limit
/// risk. A new backend with similar constraints should call this in the same place
/// and surface the resulting `FilterDrops` to `BufferSender` so that buffer-usage
/// instrumentation stays consistent with what actually lands in the buffer.
///
/// # Default behaviour
///
/// The default returns `Some(self)` if the item carries any events, and `None` if
/// it is already empty. This means an item that arrives empty (`event_count() == 0`)
/// is silently *not* persisted — preserving the pre-existing
/// "don't write empty records to disk" behaviour the call site used to enforce.
/// Types whose owners want empty items to be persisted must override this.
///
/// # Skipping this call
///
/// If a persistent backend writes an item without first calling `filter_unencodable`,
/// any sub-item that exceeds the format's limits will surface as a hard
/// [`Encodable::encode`] error and the *entire* item is rejected — including any
/// sibling sub-items that would otherwise have encoded fine. The filter is the only
/// path that produces graceful per-item drop with telemetry and a `Rejected` event
/// status; the encode-level check exists purely as defense-in-depth to ensure a
/// corrupt record cannot reach disk if a future caller forgets to filter.
fn filter_unencodable(self) -> Option<Self> {
if self.event_count() > 0 {
Some(self)
} else {
None
}
}

// Blanket implementation for anything that is already bufferable.
impl<T> Bufferable for T where T: InMemoryBufferable + Encodable + GroupedFinalizable {}
/// Returns whether every sub-item can be persisted by a backend with wire-format
/// constraints, without consuming or modifying the item.
///
/// This is the non-destructive counterpart to [`Bufferable::filter_unencodable`], and
/// exists so routing policy can be decided *before* any filtering happens. In
/// particular `WhenFull::Overflow` needs to know that an item can never reach disk, so
/// it can hand the item to the overflow stage intact rather than pruning sub-items for
/// a write that would not have succeeded at any buffer occupancy.
///
/// The default returns `true`, which is correct for any type without format limits.
/// Implementors overriding [`Bufferable::filter_unencodable`] must override this too,
/// and the two must agree: this returns `false` exactly when `filter_unencodable` would
/// drop at least one sub-item.
fn is_fully_encodable(&self) -> bool {
true
}
}

/// Hook for observing items as they are sent into a `BufferSender`.
pub trait BufferInstrumentation<T: Bufferable>: Send + Sync + 'static {
Expand Down
6 changes: 5 additions & 1 deletion lib/vector-buffers/src/test/messages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,11 @@ use vector_common::{
},
};

use crate::{EventCount, encoding::FixedEncodable};
use crate::{Bufferable, EventCount, encoding::FixedEncodable};

impl Bufferable for SizedRecord {}
impl Bufferable for UndecodableRecord {}
impl Bufferable for MultiEventRecord {}

macro_rules! message_wrapper {
($id:ident: $ty:ty, $event_count:expr) => {
Expand Down
90 changes: 85 additions & 5 deletions lib/vector-buffers/src/topology/channel/sender.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,22 @@ impl<T> SenderAdapter<T>
where
T: Bufferable,
{
/// Whether this backend can only persist items satisfying [`Bufferable::is_fully_encodable`].
///
/// In-memory stages hold the in-memory representation and have no wire format, so they can
/// accept any item regardless of its nesting depth. Disk stages encode to protobuf on write
/// and cannot.
///
/// Callers use this to avoid assuming a stage is constrained: an item that one stage cannot
/// encode may be perfectly storable by another, so the check must be asked of the specific
/// stage rather than applied to every topology.
pub(crate) fn requires_encodable_items(&self) -> bool {
match self {
Self::InMemory(_) => false,
Self::DiskV2(_) => true,
}
}

pub(crate) async fn send(&mut self, item: T) -> crate::Result<TryWriteOutcome<T>> {
match self {
Self::InMemory(tx) => tx
Expand All @@ -51,8 +67,25 @@ where
.map(|()| TryWriteOutcome::Written)
.map_err(Into::into),
Self::DiskV2(writer) => {
let pre_count = item.event_count() as u64;
let pre_size = item.size_of() as u64;
let mut writer = writer.lock().await;

let Some(item) = item.filter_unencodable() else {
// The whole item was filtered out (e.g. every sub-item over the
// protobuf nesting budget). Report the drop directly via the
// ledger's usage handle so it shows up in the disk-v2 stage's
// `received` / `dropped` metrics — `BufferSender` does not carry
// its own handle for backends that `provides_instrumentation()`.
writer.track_dropped(pre_count, pre_size);
return Ok(TryWriteOutcome::Dropped);
};
if item.event_count() as u64 != pre_count {
let dropped_events = pre_count - item.event_count() as u64;
let dropped_bytes = pre_size.saturating_sub(item.size_of() as u64);
writer.track_dropped(dropped_events, dropped_bytes);
}

writer.write_record_outcome(item).await.map_err(|e| {
// Record-level failures that can never succeed (a record too large to encode
// within the max record size) are handled inside the writer and surfaced as a
Expand All @@ -76,6 +109,26 @@ where
Self::DiskV2(writer) => {
let mut writer = writer.lock().await;

// Filtering here is unconditional and independent of current occupancy.
// Whether an unencodable item should be dropped or handed to an overflow
// stage is a `WhenFull` policy decision, so it is made in `BufferSender`
// before the item ever reaches this backend: `WhenFull::Overflow` diverts
// items failing `is_fully_encodable` straight to the overflow stage, and
// anything arriving here is therefore expected to be persistable. Keeping
// the filter unconditional means a given item is treated the same at 99%
// full as at 100% full.
let pre_count = item.event_count() as u64;
let pre_size = item.size_of() as u64;
let Some(item) = item.filter_unencodable() else {
writer.track_dropped(pre_count, pre_size);
return Ok(TryWriteOutcome::Dropped);
};
if item.event_count() as u64 != pre_count {
let dropped_events = pre_count - item.event_count() as u64;
let dropped_bytes = pre_size.saturating_sub(item.size_of() as u64);
writer.track_dropped(dropped_events, dropped_bytes);
}

writer.try_write_record(item).await.map_err(|e| {
// Record-level failures that can never succeed (a record too large to encode
// within the max record size) are handled inside the writer and surfaced as a
Expand Down Expand Up @@ -263,20 +316,47 @@ impl<T: Bufferable> BufferSender<T> {
TryWriteOutcome::Full(_) => UsageAccounting::DroppedNewest,
TryWriteOutcome::Dropped => UsageAccounting::NotAccepted,
},
WhenFull::Overflow => match self.base.try_send(item).await? {
TryWriteOutcome::Written => UsageAccounting::Accepted,
TryWriteOutcome::Full(item) => {
WhenFull::Overflow => {
// An item the base stage can never encode is routed to the overflow stage
// intact, whatever the current occupancy. Deciding this here, rather than
// letting the backend filter it, is what makes the behaviour
// state-independent: previously an over-nested item was pruned while the
// disk had room and forwarded whole once the disk reported full, so the
// same item took different paths at 99% and 100%.
//
// The check is gated on the base stage actually having a wire-format
// constraint. A memory stage overflowing to disk can store an over-nested
// item perfectly well, so diverting it past memory would send an item the
// base could have kept to a stage that must drop it.
if self.base.requires_encodable_items() && !item.is_fully_encodable() {
Comment thread
pront marked this conversation as resolved.
self.overflow
.as_mut()
.unwrap_or_else(|| unreachable!("overflow must exist"))
.send(item, send_reference)
.await?;
UsageAccounting::NotAccepted
} else {
match self.base.try_send(item).await? {
TryWriteOutcome::Written => UsageAccounting::Accepted,
TryWriteOutcome::Full(item) => {
self.overflow
.as_mut()
.unwrap_or_else(|| unreachable!("overflow must exist"))
.send(item, send_reference)
.await?;
UsageAccounting::NotAccepted
}
TryWriteOutcome::Dropped => UsageAccounting::NotAccepted,
}
}
TryWriteOutcome::Dropped => UsageAccounting::NotAccepted,
},
}
};

// Backend filter drops are accounted directly through the backend's own
// usage handle (e.g. disk-v2's ledger), so they show up in the buffer
// stage's `received` / `dropped` metrics even when the `BufferSender`
// does not carry instrumentation. This block only reports fullness-driven
// drops captured via `was_dropped`.
if let Some(instrumentation) = self.usage_instrumentation.as_ref()
&& let Some((item_count, item_size)) = item_sizing
{
Expand Down
2 changes: 2 additions & 0 deletions lib/vector-buffers/src/topology/test_util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,8 @@ impl EventCount for Sample {
}
}

impl Bufferable for Sample {}

#[derive(Debug)]
#[allow(dead_code)] // The inner _is_ read by the `Debug` impl, but that's ignored
pub struct BasicError(pub(crate) String);
Expand Down
14 changes: 14 additions & 0 deletions lib/vector-buffers/src/variants/disk_v2/ledger.rs
Original file line number Diff line number Diff line change
Expand Up @@ -534,6 +534,20 @@ where
next_record_id
}

/// Tracks events that arrived at the buffer but were rejected before being
/// persisted (e.g. `Bufferable::filter_unencodable` dropping over-budget
/// sub-items). Bumps both `received` and the unintentional-`dropped` counter
/// on the usage handle so `buffer_size = received - sent - dropped` stays
/// consistent and operators can see the rejection in buffer-usage metrics.
/// `total_buffer_size` is intentionally left alone — these events never
/// reached disk.
pub fn track_dropped(&self, event_count: u64, byte_size: u64) {
self.usage_handle
.increment_received_event_count_and_byte_size(event_count, byte_size);
self.usage_handle
.increment_dropped_event_count_and_byte_size(event_count, byte_size, false);
}

/// Tracks the statistics of multiple successful reads.
pub fn track_reads(&self, event_count: u64, total_record_size: u64) {
self.decrement_total_buffer_size(total_record_size);
Expand Down
Loading
Loading