vectordotdev/vector · error
validated by take_while
Error message
validated by take_while
What it means
Assertion comment in handle_s3_event_record describing why a downstream filter is trusted: the line stream is bounded by take_while, which stops consuming once the decoding/compression error condition is met, so the code after it may assume the error has already been surfaced and need not re-check it. The text documents a streaming invariant rather than reporting a runtime fault.
Solutions
- If the assumption is violated during debugging, verify the take_while predicate actually breaks on the error sentinel
- Add a trace log after the take_while boundary to confirm stream termination behavior
Defensive patterns
Strategy: validation
When it happens
Trigger: Thrown at src/sources/aws_s3/sqs.rs:743 when the library encounters an invalid state.
Common situations: See trigger scenarios.
AI-assisted analysis of vectordotdev/vector@bdb87aeaa4 (2026-09-16).
Data as JSON: /api/errors/19170930ac116799.
Report an issue: GitHub.
Appendix: source
Thrown at src/sources/aws_s3/sqs.rs:743
// prefer duplicate lines over message loss. Future work could include recording
// the offset of the object that has been read, but this would only be relevant in
// the case that the same vector instance processes the same message.
let mut read_error = None;
let bytes_received = self.bytes_received.clone();
let events_received = self.events_received.clone();
let lines: Box<dyn Stream<Item = Bytes> + Send + Unpin> = Box::new(
FramedRead::new(object_reader, self.state.decoder.framer.clone())
.map(|res| {
res.inspect(|bytes| {
bytes_received.emit(ByteSize(bytes.len()));
})
.map_err(|err| {
read_error = Some(err);
})
.ok()
})
.take_while(|res| ready(res.is_some()))
.map(|r| r.expect("validated by take_while")),
);
let lines: Box<dyn Stream<Item = Bytes> + Send + Unpin> = match &self.state.multiline {
Some(config) => Box::new(
LineAgg::new(
lines.map(|line| ((), line, ())),
line_agg::Logic::new(config.clone()),
)
.map(|(_src, line, _context, _lastline_context)| line),
),
None => lines,
};
let mut stream = lines.flat_map(|line| {
let events = match self.state.decoder.deserializer_parse(line) {
Ok((events, _events_size)) => events,
Err(_error) => {
// Error is handled by `codecs::Decoder`, no further handlingView on GitHub (pinned to bdb87aeaa4)