Skip to main content

vtcode_commons/
line_framing.rs

1//! Bounded byte framing for newline-delimited subprocess streams.
2
3use tokio::io::{AsyncBufRead, AsyncBufReadExt};
4
5/// Whether LF is retained and charged against the byte limit. CR is always content.
6#[derive(Debug, Clone, Copy, PartialEq, Eq)]
7pub enum LineEnding {
8    IncludeLf,
9    ExcludeLf,
10}
11
12/// Read and drain one physical line, retaining at most `max_bytes` in `line`.
13///
14/// Clears and reuses the caller's buffer. Returns `None` only at EOF without a
15/// line, and `Some(truncated)` for a complete or unterminated final line. Even
16/// with a zero-byte cap, a physical line is returned and drained. Oversized
17/// content is discarded through LF (or EOF), preserving the next frame boundary.
18/// This is byte framing: retained data need not end at a UTF-8 character boundary.
19///
20/// # Errors
21/// Returns the underlying read error. After error or cancellation, callers must
22/// not assume the stream is aligned at a new line; this operation is not
23/// cancellation-safe and transports should terminate the reader on failure.
24pub async fn read_bounded_line<R: AsyncBufRead + Unpin>(
25    reader: &mut R,
26    line: &mut Vec<u8>,
27    max_bytes: usize,
28    ending: LineEnding,
29) -> std::io::Result<Option<bool>> {
30    line.clear();
31    let mut truncated = false;
32    loop {
33        let available = reader.fill_buf().await?;
34        if available.is_empty() {
35            return Ok(if line.is_empty() && !truncated {
36                None
37            } else {
38                Some(truncated)
39            });
40        }
41        let newline = available.iter().position(|byte| *byte == b'\n');
42        let consumed = newline.map_or(available.len(), |position| position + 1);
43        let content_len = consumed - usize::from(newline.is_some() && ending == LineEnding::ExcludeLf);
44        let copy_len = content_len.min(max_bytes.saturating_sub(line.len()));
45        line.extend(available.iter().copied().take(copy_len));
46        truncated |= copy_len < content_len;
47        reader.consume(consumed);
48        if newline.is_some() {
49            return Ok(Some(truncated));
50        }
51    }
52}
53
54#[cfg(test)]
55mod tests {
56    use super::*;
57    use std::pin::Pin;
58    use std::task::{Context, Poll};
59    use tokio::io::{AsyncRead, BufReader, ReadBuf};
60
61    #[tokio::test]
62    async fn delimiter_policy_preserves_exact_limit_and_crlf() -> std::io::Result<()> {
63        for (ending, expected, truncated) in [
64            (LineEnding::ExcludeLf, b"abc".as_slice(), false),
65            (LineEnding::IncludeLf, b"abc".as_slice(), true),
66        ] {
67            let mut reader = BufReader::with_capacity(1, b"abc\n\r\nnext".as_slice());
68            let mut line = Vec::new();
69            assert_eq!(read_bounded_line(&mut reader, &mut line, 3, ending).await?, Some(truncated));
70            assert_eq!(line, expected);
71            assert_eq!(read_bounded_line(&mut reader, &mut line, 3, ending).await?, Some(false));
72            assert_eq!(
73                line,
74                if ending == LineEnding::ExcludeLf {
75                    b"\r".as_slice()
76                } else {
77                    b"\r\n".as_slice()
78                }
79            );
80        }
81        Ok(())
82    }
83
84    #[tokio::test]
85    async fn framing_matches_physical_lines_across_caps_and_chunk_boundaries() -> std::io::Result<()> {
86        // Independent oracle splits complete input first; production frames incremental chunks.
87        for input in [
88            b"".as_slice(),
89            b"\n",
90            b"abc\nnext\n",
91            b"\nabc\r\nz\nlast",
92            "đ\n終".as_bytes(),
93        ] {
94            for ending in [LineEnding::IncludeLf, LineEnding::ExcludeLf] {
95                for cap in 0..=8 {
96                    for chunk in 1..=9 {
97                        let mut reader = BufReader::with_capacity(chunk, input);
98                        let mut line = vec![b'!'; 20];
99                        for physical_line in input.split_inclusive(|byte| *byte == b'\n') {
100                            let content = if ending == LineEnding::ExcludeLf {
101                                physical_line.strip_suffix(b"\n").unwrap_or(physical_line)
102                            } else {
103                                physical_line
104                            };
105                            assert_eq!(
106                                read_bounded_line(&mut reader, &mut line, cap, ending).await?,
107                                Some(content.len() > cap)
108                            );
109                            assert_eq!(
110                                line,
111                                &content[..content.len().min(cap)],
112                                "input={input:?}, ending={ending:?}, cap={cap}, chunk={chunk}"
113                            );
114                            assert!(line.len() <= cap);
115                        }
116                        assert_eq!(read_bounded_line(&mut reader, &mut line, cap, ending).await?, None);
117                        assert!(line.is_empty());
118                    }
119                }
120            }
121        }
122        Ok(())
123    }
124
125    #[tokio::test]
126    async fn oversized_unterminated_line_is_drained_and_bounded() -> std::io::Result<()> {
127        for ending in [LineEnding::IncludeLf, LineEnding::ExcludeLf] {
128            let input = vec![b'x'; 4096];
129            let mut reader = BufReader::with_capacity(7, input.as_slice());
130            let mut line = Vec::new();
131            assert_eq!(read_bounded_line(&mut reader, &mut line, 5, ending).await?, Some(true));
132            assert_eq!(line, b"xxxxx");
133            assert_eq!(read_bounded_line(&mut reader, &mut line, 5, ending).await?, None);
134        }
135        Ok(())
136    }
137
138    struct FailingReader {
139        sent_prefix: bool,
140    }
141
142    impl AsyncRead for FailingReader {
143        fn poll_read(
144            mut self: Pin<&mut Self>,
145            _: &mut Context<'_>,
146            buffer: &mut ReadBuf<'_>,
147        ) -> Poll<std::io::Result<()>> {
148            if self.sent_prefix {
149                Poll::Ready(Err(std::io::Error::new(std::io::ErrorKind::BrokenPipe, "test read failure")))
150            } else {
151                self.sent_prefix = true;
152                buffer.put_slice(b"x");
153                Poll::Ready(Ok(()))
154            }
155        }
156    }
157
158    #[tokio::test]
159    async fn read_errors_propagate_after_a_partial_line() {
160        for ending in [LineEnding::IncludeLf, LineEnding::ExcludeLf] {
161            let mut reader = BufReader::new(FailingReader { sent_prefix: false });
162            let mut line = Vec::new();
163            let error = read_bounded_line(&mut reader, &mut line, 3, ending).await.unwrap_err();
164            assert_eq!(error.kind(), std::io::ErrorKind::BrokenPipe);
165            assert_eq!(error.to_string(), "test read failure");
166            assert_eq!(line, b"x");
167        }
168    }
169}