mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-31 19:16:29 +00:00
fix(bootstrap): make read_decoded_line cancel-safe under tokio::select!
The salvaged helper cleared its line buffer on entry. Inside run_script's tokio::select! loop, a stdout line arrival cancels the in-flight stderr read (and vice versa); read_until had already consumed bytes into the buffer, and the next call's clear() silently dropped that partial line. Keep partially-read bytes across cancellation (clear only after a full line is decoded) and emit an unterminated final line at EOF instead of swallowing it. Adds a cancellation regression test (fails against the clear-on-entry version) and an EOF-tail test.
This commit is contained in:
parent
acee4f25c7
commit
8fe9706da8
1 changed files with 57 additions and 6 deletions
|
|
@ -87,18 +87,27 @@ pub(crate) async fn read_decoded_line<R>(
|
|||
where
|
||||
R: AsyncBufReadExt + Unpin,
|
||||
{
|
||||
buf.clear();
|
||||
// Cancel-safety: `buf` is NOT cleared on entry. When this future is
|
||||
// dropped mid-read inside `tokio::select!` (the other stream produced a
|
||||
// line first), `read_until` has already appended any consumed bytes to
|
||||
// `buf`; the next call resumes and appends the rest of the line. Clearing
|
||||
// on entry would silently drop those bytes. We clear only after a full
|
||||
// line has been decoded.
|
||||
let n = reader.read_until(b'\n', buf).await?;
|
||||
if n == 0 {
|
||||
if n == 0 && buf.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
// n == 0 with a non-empty buf means EOF cut off an unterminated line
|
||||
// (possibly accumulated across cancelled reads) -- emit it.
|
||||
if buf.last() == Some(&b'\n') {
|
||||
buf.pop();
|
||||
if buf.last() == Some(&b'\r') {
|
||||
buf.pop();
|
||||
}
|
||||
}
|
||||
if buf.last() == Some(&b'\r') {
|
||||
buf.pop();
|
||||
}
|
||||
Ok(Some(decode_console_bytes(buf)))
|
||||
let line = decode_console_bytes(buf);
|
||||
buf.clear();
|
||||
Ok(Some(line))
|
||||
}
|
||||
|
||||
/// Hooks the caller installs to receive output.
|
||||
|
|
@ -499,4 +508,46 @@ info line
|
|||
.unwrap()
|
||||
.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read_decoded_line_preserves_partial_line_across_cancellation() {
|
||||
use std::time::Duration;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
let (mut tx, rx) = tokio::io::duplex(64);
|
||||
let mut reader = BufReader::new(rx);
|
||||
let mut buf = Vec::new();
|
||||
|
||||
tx.write_all(b"partial").await.unwrap();
|
||||
// Poll once, then cancel (drop) the future -- exactly what
|
||||
// tokio::select! does in run_script when the other stream produces
|
||||
// a line first. The consumed bytes must survive in `buf`.
|
||||
let _ = tokio::time::timeout(
|
||||
Duration::from_millis(0),
|
||||
read_decoded_line(&mut reader, &mut buf),
|
||||
)
|
||||
.await;
|
||||
|
||||
tx.write_all(b" line\n").await.unwrap();
|
||||
let line = read_decoded_line(&mut reader, &mut buf).await.unwrap();
|
||||
assert_eq!(line.as_deref(), Some("partial line"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read_decoded_line_emits_unterminated_final_line_at_eof() {
|
||||
let data: &[u8] = b"no trailing newline";
|
||||
let mut reader = BufReader::new(data);
|
||||
let mut buf = Vec::new();
|
||||
assert_eq!(
|
||||
read_decoded_line(&mut reader, &mut buf)
|
||||
.await
|
||||
.unwrap()
|
||||
.as_deref(),
|
||||
Some("no trailing newline")
|
||||
);
|
||||
assert!(read_decoded_line(&mut reader, &mut buf)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue