{"record":{"id":"19170930ac116799","repo":"vectordotdev/vector","slug":"validated-by-take-while","errorCode":null,"errorMessage":"validated by take_while","messagePattern":"validated by take_while","errorType":"panic","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/sources/aws_s3/sqs.rs","lineNumber":743,"sourceCode":"        // prefer duplicate lines over message loss. Future work could include recording\n        // the offset of the object that has been read, but this would only be relevant in\n        // the case that the same vector instance processes the same message.\n        let mut read_error = None;\n        let bytes_received = self.bytes_received.clone();\n        let events_received = self.events_received.clone();\n        let lines: Box<dyn Stream<Item = Bytes> + Send + Unpin> = Box::new(\n            FramedRead::new(object_reader, self.state.decoder.framer.clone())\n                .map(|res| {\n                    res.inspect(|bytes| {\n                        bytes_received.emit(ByteSize(bytes.len()));\n                    })\n                    .map_err(|err| {\n                        read_error = Some(err);\n                    })\n                    .ok()\n                })\n                .take_while(|res| ready(res.is_some()))\n                .map(|r| r.expect(\"validated by take_while\")),\n        );\n\n        let lines: Box<dyn Stream<Item = Bytes> + Send + Unpin> = match &self.state.multiline {\n            Some(config) => Box::new(\n                LineAgg::new(\n                    lines.map(|line| ((), line, ())),\n                    line_agg::Logic::new(config.clone()),\n                )\n                .map(|(_src, line, _context, _lastline_context)| line),\n            ),\n            None => lines,\n        };\n\n        let mut stream = lines.flat_map(|line| {\n            let events = match self.state.decoder.deserializer_parse(line) {\n                Ok((events, _events_size)) => events,\n                Err(_error) => {\n                    // Error is handled by `codecs::Decoder`, no further handling","sourceCodeStart":725,"sourceCodeEnd":761,"githubUrl":"https://github.com/vectordotdev/vector/blob/bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013/src/sources/aws_s3/sqs.rs#L725-L761","documentation":"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.","triggerScenarios":"Thrown at src/sources/aws_s3/sqs.rs:743 when the library encounters an invalid state.","commonSituations":"See trigger scenarios.","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"],"exampleFix":null,"handlingStrategy":"validation","validationCode":null,"typeGuard":null,"tryCatchPattern":null,"preventionTips":[],"tags":[],"backgroundTag":null,"analyzedSha":"bdb87aeaa4c4ff27c0ba643c1c77b21bf2ef4013","analyzedAt":"2026-09-16T02:53:35.741Z","contentChangedAt":"2026-09-16T02:53:35.741Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}