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

  1. If the assumption is violated during debugging, verify the take_while predicate actually breaks on the error sentinel
  2. 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 handling

View on GitHub (pinned to bdb87aeaa4)