| 1 | use std::collections::VecDeque; |
| 2 | use std::io::{self, Write}; |
| 3 | use std::path::PathBuf; |
| 4 | use std::sync::{Arc, Mutex}; |
| 5 | |
| 6 | use encoding_rs::{CoderResult, Decoder, UTF_8}; |
| 7 | |
| 8 | const BOUNDED_OUTPUT_MAX_LINES: usize = 2_000; |
| 9 | const BOUNDED_OUTPUT_MAX_BYTES: usize = 50 * 1024; |
| 10 | const BOUNDED_OUTPUT_RETAIN_BYTES: usize = BOUNDED_OUTPUT_MAX_BYTES + 4; |
| 11 | |
| 12 | #[derive(Debug)] |
| 13 | pub(super) struct BoundedOutputSnapshot { |
| 14 | pub(super) content: String, |
| 15 | pub(super) total_bytes: usize, |
| 16 | pub(super) retained_bytes: usize, |
| 17 | pub(super) truncated: bool, |
| 18 | } |
| 19 | |
| 20 | /// One decoded, arrival-ordered stream: complete output goes to disk while |
| 21 | /// memory retains only enough tail bytes for the 2,000-line/50KiB result bound. |
| 22 | pub(super) struct BoundedOutputAccumulator { |
| 23 | tail: VecDeque<u8>, |
| 24 | tail_newlines: usize, |
| 25 | total_bytes: usize, |
| 26 | total_newlines: usize, |
| 27 | current_line_bytes: usize, |
| 28 | last_line_bytes: usize, |
| 29 | front_clipped: bool, |
| 30 | last_byte: Option<u8>, |
| 31 | decoder: Decoder, |
| 32 | stream_finished: bool, |
| 33 | stream_error: Option<String>, |
| 34 | temp: Option<tempfile::NamedTempFile>, |
| 35 | full_output_path: Option<PathBuf>, |
| 36 | /// Why the on-disk spill file could not be created (disk full, descriptor |
| 37 | /// exhaustion, unwritable temp dir). The stream still runs and the bounded |
| 38 | /// tail is still delivered; only "Full output: <path>" is unavailable. |
| 39 | spill_unavailable: Option<String>, |
| 40 | } |
| 41 | |
| 42 | impl BoundedOutputAccumulator { |
| 43 | /// Build an accumulator whose complete-output spill file lives in |
| 44 | /// `spill_dir` (`None` = process temp dir). Never fails: when the spill |
| 45 | /// file cannot be created (disk full, `EMFILE`, missing temp dir) the |
| 46 | /// command still runs and the bounded tail is still returned — the spill |
| 47 | /// is a convenience, not a precondition for executing `echo ok`. Tests |
| 48 | /// pass a nonexistent dir to fault-inject the failure. |
| 49 | pub(super) fn new_in(spill_dir: Option<&std::path::Path>) -> Self { |
| 50 | let mut builder = tempfile::Builder::new(); |
| 51 | builder.prefix("codewhale-bash-"); |
| 52 | let temp = match spill_dir { |
| 53 | Some(dir) => builder.tempfile_in(dir), |
| 54 | None => builder.tempfile(), |
| 55 | }; |
| 56 | let (temp, spill_unavailable) = match temp { |
| 57 | Ok(temp) => (Some(temp), None), |
| 58 | Err(error) => { |
| 59 | tracing::warn!( |
| 60 | error = %error, |
| 61 | "shell output spill file unavailable; continuing with the in-memory tail only" |
| 62 | ); |
| 63 | (None, Some(spill_unavailable_reason(&error))) |
| 64 | } |
| 65 | }; |
| 66 | Self { |
| 67 | tail: VecDeque::with_capacity(BOUNDED_OUTPUT_RETAIN_BYTES), |
| 68 | tail_newlines: 0, |
| 69 | total_bytes: 0, |
| 70 | total_newlines: 0, |
| 71 | current_line_bytes: 0, |
| 72 | last_line_bytes: 0, |
| 73 | front_clipped: false, |
| 74 | last_byte: None, |
| 75 | decoder: UTF_8.new_decoder_without_bom_handling(), |
| 76 | stream_finished: false, |
| 77 | stream_error: None, |
| 78 | temp, |
| 79 | full_output_path: None, |
| 80 | spill_unavailable, |
| 81 | } |
| 82 | } |
| 83 | |
| 84 | /// Why the complete output is not being persisted, if it is not. |
| 85 | #[cfg(test)] |
| 86 | pub(super) fn spill_unavailable(&self) -> Option<&str> { |
| 87 | self.spill_unavailable.as_deref() |
| 88 | } |
| 89 | |
| 90 | fn decode(&mut self, bytes: &[u8], last: bool) -> String { |
| 91 | let capacity = self |
| 92 | .decoder |
| 93 | .max_utf8_buffer_length(bytes.len()) |
| 94 | .unwrap_or(bytes.len().saturating_mul(3).saturating_add(3)); |
| 95 | let mut decoded = String::with_capacity(capacity); |
| 96 | let mut offset = 0; |
| 97 | loop { |
| 98 | let (result, read, _) = |
| 99 | self.decoder |
| 100 | .decode_to_string(&bytes[offset..], &mut decoded, last); |
| 101 | offset += read; |
| 102 | if result == CoderResult::InputEmpty { |
| 103 | return decoded; |
| 104 | } |
| 105 | decoded.reserve(capacity.max(4)); |
| 106 | } |
| 107 | } |
| 108 | |
| 109 | pub(super) fn append(&mut self, raw: &[u8]) -> io::Result<()> { |
| 110 | if self.stream_finished { |
| 111 | return Err(io::Error::other( |
| 112 | "shell output arrived after the stream closed", |
| 113 | )); |
| 114 | } |
| 115 | if let Some(temp) = self.temp.as_mut() { |
| 116 | temp.write_all(raw)?; |
| 117 | } |
| 118 | let decoded = self.decode(raw, false); |
| 119 | self.append_decoded(decoded.as_bytes()); |
| 120 | Ok(()) |
| 121 | } |
| 122 | |
| 123 | pub(super) fn finish(&mut self) -> io::Result<()> { |
| 124 | if !self.stream_finished { |
| 125 | let decoded = self.decode(&[], true); |
| 126 | self.append_decoded(decoded.as_bytes()); |
| 127 | if let Some(temp) = self.temp.as_mut() { |
| 128 | temp.flush()?; |
| 129 | } |
| 130 | self.stream_finished = true; |
| 131 | } |
| 132 | Ok(()) |
| 133 | } |
| 134 | |
| 135 | pub(super) fn record_error(&mut self, error: &io::Error) { |
| 136 | self.stream_error = Some(error.to_string()); |
| 137 | } |
| 138 | |
| 139 | fn append_decoded(&mut self, bytes: &[u8]) { |
| 140 | self.total_bytes = self.total_bytes.saturating_add(bytes.len()); |
| 141 | for &byte in bytes { |
| 142 | self.tail.push_back(byte); |
| 143 | if byte == b'\n' { |
| 144 | self.tail_newlines += 1; |
| 145 | self.total_newlines += 1; |
| 146 | self.last_line_bytes = self.current_line_bytes; |
| 147 | self.current_line_bytes = 0; |
| 148 | } else { |
| 149 | self.current_line_bytes += 1; |
| 150 | } |
| 151 | self.last_byte = Some(byte); |
| 152 | } |
| 153 | while self.tail.len() > BOUNDED_OUTPUT_RETAIN_BYTES { |
| 154 | self.pop_front(); |
| 155 | self.front_clipped = true; |
| 156 | } |
| 157 | while self.tail_lines() > BOUNDED_OUTPUT_MAX_LINES { |
| 158 | while let Some(byte) = self.tail.pop_front() { |
| 159 | if byte == b'\n' { |
| 160 | self.tail_newlines -= 1; |
| 161 | break; |
| 162 | } |
| 163 | } |
| 164 | self.front_clipped = false; |
| 165 | } |
| 166 | } |
| 167 | |
| 168 | fn pop_front(&mut self) { |
| 169 | if self.tail.pop_front() == Some(b'\n') { |
| 170 | self.tail_newlines -= 1; |
| 171 | } |
| 172 | } |
| 173 | |
| 174 | fn tail_lines(&self) -> usize { |
| 175 | self.tail_newlines + usize::from(self.tail.back().is_some_and(|byte| *byte != b'\n')) |
| 176 | } |
| 177 | |
| 178 | fn total_lines(&self) -> usize { |
| 179 | self.total_newlines + usize::from(self.last_byte.is_some_and(|byte| byte != b'\n')) |
| 180 | } |
| 181 | |
| 182 | fn selected(&self) -> (Vec<u8>, bool) { |
| 183 | let mut bytes = self.tail.iter().copied().collect::<Vec<_>>(); |
| 184 | let recent_line_bytes = if self.last_byte == Some(b'\n') { |
| 185 | self.last_line_bytes |
| 186 | } else { |
| 187 | self.current_line_bytes |
| 188 | }; |
| 189 | let partial_line = recent_line_bytes > BOUNDED_OUTPUT_MAX_BYTES; |
| 190 | if partial_line { |
| 191 | if bytes.last() == Some(&b'\n') { |
| 192 | bytes.pop(); |
| 193 | } |
| 194 | let floor = bytes.len().saturating_sub(BOUNDED_OUTPUT_MAX_BYTES); |
| 195 | let start = (floor..bytes.len()) |
| 196 | .find(|index| std::str::from_utf8(&bytes[*index..]).is_ok()) |
| 197 | .unwrap_or(bytes.len()); |
| 198 | bytes.drain(..start); |
| 199 | } else if self.front_clipped |
| 200 | && let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') |
| 201 | { |
| 202 | bytes.drain(..=newline); |
| 203 | } |
| 204 | (bytes, partial_line) |
| 205 | } |
| 206 | |
| 207 | fn format_size(bytes: usize) -> String { |
| 208 | if bytes < 1024 { |
| 209 | format!("{bytes}B") |
| 210 | } else if bytes < 1024 * 1024 { |
| 211 | format!("{:.1}KB", bytes as f64 / 1024.0) |
| 212 | } else { |
| 213 | format!("{:.1}MB", bytes as f64 / (1024.0 * 1024.0)) |
| 214 | } |
| 215 | } |
| 216 | |
| 217 | pub(super) fn total_bytes(&self) -> usize { |
| 218 | self.total_bytes |
| 219 | } |
| 220 | |
| 221 | /// Decoded output appended after `cursor` (a previous `total_bytes`), for |
| 222 | /// incremental readers such as `wait`. The second value counts bytes of |
| 223 | /// that range the memory bound already discarded; the complete output is |
| 224 | /// still in the spill file. |
| 225 | pub(super) fn delta_since(&self, cursor: usize) -> io::Result<(String, usize)> { |
| 226 | if let Some(error) = self.stream_error.as_ref() { |
| 227 | return Err(io::Error::other(error.clone())); |
| 228 | } |
| 229 | let new_bytes = self.total_bytes.saturating_sub(cursor); |
| 230 | let available = new_bytes.min(self.tail.len()); |
| 231 | let mut bytes = self |
| 232 | .tail |
| 233 | .range(self.tail.len() - available..) |
| 234 | .copied() |
| 235 | .collect::<Vec<_>>(); |
| 236 | // A cursor is always a character boundary; a clipped tail's front may |
| 237 | // not be, so drop any continuation bytes it starts with. |
| 238 | let partial = bytes |
| 239 | .iter() |
| 240 | .take_while(|byte| (**byte & 0xC0) == 0x80) |
| 241 | .count(); |
| 242 | bytes.drain(..partial); |
| 243 | let omitted = new_bytes - available + partial; |
| 244 | Ok((String::from_utf8_lossy(&bytes).into_owned(), omitted)) |
| 245 | } |
| 246 | |
| 247 | /// First line of a delta that lost `omitted` bytes to the memory bound: |
| 248 | /// how many, and where the complete output is (or why it is not kept). |
| 249 | pub(super) fn omitted_notice(&self, omitted: usize) -> String { |
| 250 | let location = self |
| 251 | .full_output_path |
| 252 | .as_deref() |
| 253 | .or_else(|| self.temp.as_ref().map(tempfile::NamedTempFile::path)); |
| 254 | match (location, self.spill_unavailable.as_deref()) { |
| 255 | (Some(path), _) => format!( |
| 256 | "[{omitted} bytes of earlier output not retained in memory. Full output: {}]\n", |
| 257 | path.display() |
| 258 | ), |
| 259 | (None, Some(reason)) => format!( |
| 260 | "[{omitted} bytes of earlier output not retained in memory. Full output was not persisted: {reason}]\n" |
| 261 | ), |
| 262 | (None, None) => { |
| 263 | format!("[{omitted} bytes of earlier output not retained in memory]\n") |
| 264 | } |
| 265 | } |
| 266 | } |
| 267 | |
| 268 | pub(super) fn snapshot(&mut self, finalize: bool) -> io::Result<BoundedOutputSnapshot> { |
| 269 | if let Some(error) = self.stream_error.as_ref() { |
| 270 | return Err(io::Error::other(error.clone())); |
| 271 | } |
| 272 | let (selected, partial_line) = self.selected(); |
| 273 | let retained_bytes = selected.len(); |
| 274 | let truncated = retained_bytes < self.total_bytes; |
| 275 | let total_lines = self.total_lines(); |
| 276 | let kept_lines = selected.iter().filter(|byte| **byte == b'\n').count() |
| 277 | + usize::from(selected.last().is_some_and(|byte| *byte != b'\n')); |
| 278 | let mut content = String::from_utf8(selected).expect("stream decoder emits valid UTF-8"); |
| 279 | |
| 280 | if finalize && self.stream_finished && self.full_output_path.is_none() { |
| 281 | if truncated { |
| 282 | if let Some(mut temp) = self.temp.take() { |
| 283 | temp.flush()?; |
| 284 | let (_, path) = temp.keep().map_err(|error| error.error)?; |
| 285 | self.full_output_path = Some(path); |
| 286 | } |
| 287 | } else { |
| 288 | self.temp.take(); |
| 289 | } |
| 290 | } |
| 291 | if truncated && finalize && self.full_output_path.is_none() { |
| 292 | let reason = self.spill_unavailable.as_deref().unwrap_or( |
| 293 | "the output stream did not close cleanly, so the spill file was not kept", |
| 294 | ); |
| 295 | content.push_str(&format!( |
| 296 | "\n\n[Showing the last {} of {} lines ({} limit). Full output was not persisted: {reason}]", |
| 297 | Self::format_size(retained_bytes), |
| 298 | total_lines, |
| 299 | Self::format_size(BOUNDED_OUTPUT_MAX_BYTES), |
| 300 | )); |
| 301 | } else if truncated |
| 302 | && finalize |
| 303 | && let Some(path) = self.full_output_path.as_ref() |
| 304 | { |
| 305 | if partial_line { |
| 306 | content.push_str(&format!( |
| 307 | "\n\n[Showing last {} of line {} (line is {}). Full output: {}]", |
| 308 | Self::format_size(retained_bytes), |
| 309 | total_lines, |
| 310 | Self::format_size(self.current_line_bytes), |
| 311 | path.display() |
| 312 | )); |
| 313 | } else { |
| 314 | let start = total_lines.saturating_sub(kept_lines) + 1; |
| 315 | let limit = if self.front_clipped { |
| 316 | format!(" ({} limit)", Self::format_size(BOUNDED_OUTPUT_MAX_BYTES)) |
| 317 | } else { |
| 318 | String::new() |
| 319 | }; |
| 320 | content.push_str(&format!( |
| 321 | "\n\n[Showing lines {start}-{total_lines} of {total_lines}{limit}. Full output: {}]", |
| 322 | path.display() |
| 323 | )); |
| 324 | } |
| 325 | } |
| 326 | Ok(BoundedOutputSnapshot { |
| 327 | content, |
| 328 | total_bytes: self.total_bytes, |
| 329 | retained_bytes, |
| 330 | truncated, |
| 331 | }) |
| 332 | } |
| 333 | |
| 334 | #[cfg(test)] |
| 335 | pub(super) fn retained_memory_bytes(&self) -> usize { |
| 336 | self.tail.len() |
| 337 | } |
| 338 | |
| 339 | #[cfg(test)] |
| 340 | pub(super) fn full_output_path(&self) -> Option<&std::path::Path> { |
| 341 | self.full_output_path.as_deref() |
| 342 | } |
| 343 | } |
| 344 | |
| 345 | /// Human-readable, actionable reason for a failed spill-file creation. |
| 346 | pub(super) fn spill_unavailable_reason(error: &io::Error) -> String { |
| 347 | match resource_exhaustion_hint(error) { |
| 348 | Some(hint) => format!("{error} ({hint})"), |
| 349 | None => error.to_string(), |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | /// When an I/O error looks like host resource exhaustion, name the likely |
| 354 | /// cause and the remedy. Returns `None` for ordinary errors. |
| 355 | pub(super) fn resource_exhaustion_hint(error: &io::Error) -> Option<&'static str> { |
| 356 | use io::ErrorKind; |
| 357 | match error.kind() { |
| 358 | ErrorKind::StorageFull | ErrorKind::QuotaExceeded => { |
| 359 | return Some("the disk holding the temp dir is full; free space and retry"); |
| 360 | } |
| 361 | ErrorKind::OutOfMemory => { |
| 362 | return Some("the host is out of memory; close heavy processes and retry"); |
| 363 | } |
| 364 | _ => {} |
| 365 | } |
| 366 | let code = error.raw_os_error()?; |
| 367 | // ENOSPC / EDQUOT / EMFILE / ENFILE / ENOMEM / EAGAIN — the codes fork(2), |
| 368 | // pipe(2), and open(2) return when the machine is thrashing. |
| 369 | #[cfg(unix)] |
| 370 | { |
| 371 | if code == libc::ENOSPC || code == libc::EDQUOT { |
| 372 | return Some("the disk holding the temp dir is full; free space and retry"); |
| 373 | } |
| 374 | if code == libc::EMFILE || code == libc::ENFILE { |
| 375 | return Some( |
| 376 | "the process or host has run out of file descriptors; close background jobs or raise `ulimit -n` and retry", |
| 377 | ); |
| 378 | } |
| 379 | if code == libc::ENOMEM { |
| 380 | return Some("the host is out of memory; close heavy processes and retry"); |
| 381 | } |
| 382 | if code == libc::EAGAIN { |
| 383 | return Some( |
| 384 | "the host refused to create a process, thread, or pipe (resource limit reached); close heavy processes and retry", |
| 385 | ); |
| 386 | } |
| 387 | } |
| 388 | #[cfg(windows)] |
| 389 | { |
| 390 | // ERROR_DISK_FULL, ERROR_HANDLE_DISK_FULL, ERROR_NOT_ENOUGH_MEMORY, ERROR_TOO_MANY_OPEN_FILES |
| 391 | if code == 112 || code == 39 { |
| 392 | return Some("the disk holding the temp dir is full; free space and retry"); |
| 393 | } |
| 394 | if code == 8 { |
| 395 | return Some("the host is out of memory; close heavy processes and retry"); |
| 396 | } |
| 397 | if code == 4 { |
| 398 | return Some( |
| 399 | "the process has run out of file handles; close background jobs and retry", |
| 400 | ); |
| 401 | } |
| 402 | } |
| 403 | let _ = code; |
| 404 | None |
| 405 | } |
| 406 | |
| 407 | /// Hard in-flight ceiling for one raw shell stream held in memory (#5472). |
| 408 | /// Past this the oldest bytes are dropped — counted, never silently lost — so |
| 409 | /// one chatty command (`cargo build -v`, `git log -p`) cannot grow the process |
| 410 | /// by its entire output. Deliberately far above every consumer of these bytes: |
| 411 | /// the 30 KB tool-result truncation (`shell_output::MAX_OUTPUT_SIZE`), the |
| 412 | /// 1,200-char job-panel tail and the 1 KiB completion tail all fit with three |
| 413 | /// orders of magnitude to spare. The only surface a clip can reach is the |
| 414 | /// durable completion artifact, which records the omission explicitly. |
| 415 | pub(super) const RAW_STREAM_MAX_BYTES: usize = 16 * 1024 * 1024; |
| 416 | |
| 417 | /// Extra headroom before a front-drop, so the O(len) compaction runs once per |
| 418 | /// `cap / 4` bytes appended instead of once per chunk. |
| 419 | const RAW_STREAM_DROP_SLACK: usize = RAW_STREAM_MAX_BYTES / 4; |
| 420 | |
| 421 | /// Tail retained once a job's output has been *delivered* — the foreground |
| 422 | /// result is already the tool result, or the completion evidence is already |
| 423 | /// written to its session artifact. Everything past this is dead weight for |
| 424 | /// the up-to-1 h the finished record stays listed (#5472 finding 1). |
| 425 | pub(super) const RAW_STREAM_SETTLED_TAIL_BYTES: usize = 64 * 1024; |
| 426 | |
| 427 | /// One raw (undecoded) shell stream retained in memory for a live job. |
| 428 | /// |
| 429 | /// Bounded two independent ways, which is the whole point of the type: |
| 430 | /// `append` enforces `cap` while the command runs, and `release_to_tail` |
| 431 | /// collapses the buffer the moment its bytes have been delivered. Both record |
| 432 | /// how many leading bytes were discarded so `total_len` — and therefore every |
| 433 | /// `stdout_len` / `byte_length` the model and the artifact see — stays honest. |
| 434 | pub(super) struct RawOutputBuffer { |
| 435 | data: Vec<u8>, |
| 436 | dropped: usize, |
| 437 | cap: usize, |
| 438 | abandoned: bool, |
| 439 | } |
| 440 | |
| 441 | impl RawOutputBuffer { |
| 442 | pub(super) fn new() -> Self { |
| 443 | Self::with_cap(RAW_STREAM_MAX_BYTES) |
| 444 | } |
| 445 | |
| 446 | pub(super) fn with_cap(cap: usize) -> Self { |
| 447 | Self { |
| 448 | data: Vec::new(), |
| 449 | dropped: 0, |
| 450 | cap: cap.max(1), |
| 451 | abandoned: false, |
| 452 | } |
| 453 | } |
| 454 | |
| 455 | /// Append, returning `false` once nobody will ever read this stream again. |
| 456 | /// |
| 457 | /// The reader thread uses that as its exit condition, which is the only way |
| 458 | /// out when a descendant has escaped the process group and holds the pipe |
| 459 | /// write-end open: `read()` will never see EOF, so without this the thread |
| 460 | /// runs — and retains its buffer — for the life of the process (#5472 |
| 461 | /// finding 2). |
| 462 | pub(super) fn append(&mut self, bytes: &[u8]) -> bool { |
| 463 | if self.abandoned { |
| 464 | // Keep the total honest even though the bytes are discarded. |
| 465 | self.dropped = self.dropped.saturating_add(bytes.len()); |
| 466 | return false; |
| 467 | } |
| 468 | self.data.extend_from_slice(bytes); |
| 469 | if self.data.len() > self.cap.saturating_add(RAW_STREAM_DROP_SLACK.min(self.cap)) { |
| 470 | self.drop_front_to(self.cap); |
| 471 | } |
| 472 | true |
| 473 | } |
| 474 | |
| 475 | /// Give up on this stream: release everything held and stop accepting more. |
| 476 | /// |
| 477 | /// Called when the bounded reader join times out. The shell is already |
| 478 | /// terminal and its result already delivered, so nothing can consume these |
| 479 | /// bytes; holding them until the writer eventually closes is pure residency. |
| 480 | pub(super) fn abandon(&mut self) { |
| 481 | self.abandoned = true; |
| 482 | self.dropped = self.dropped.saturating_add(self.data.len()); |
| 483 | self.data = Vec::new(); |
| 484 | } |
| 485 | |
| 486 | /// Total bytes this stream has produced, including bytes no longer held. |
| 487 | pub(super) fn total_len(&self) -> usize { |
| 488 | self.dropped.saturating_add(self.data.len()) |
| 489 | } |
| 490 | |
| 491 | /// Leading bytes discarded by the in-flight cap or by `release_to_tail`. |
| 492 | pub(super) fn dropped(&self) -> usize { |
| 493 | self.dropped |
| 494 | } |
| 495 | |
| 496 | pub(super) fn retained(&self) -> &[u8] { |
| 497 | &self.data |
| 498 | } |
| 499 | |
| 500 | /// Collapse to at most `keep` trailing bytes and give the allocation back. |
| 501 | /// Called once a job is terminal *and* its output has been delivered. |
| 502 | pub(super) fn release_to_tail(&mut self, keep: usize) { |
| 503 | if self.data.len() <= keep { |
| 504 | return; |
| 505 | } |
| 506 | self.drop_front_to(keep); |
| 507 | self.data.shrink_to_fit(); |
| 508 | } |
| 509 | |
| 510 | fn drop_front_to(&mut self, keep: usize) { |
| 511 | let mut start = self.data.len().saturating_sub(keep); |
| 512 | // Snap forward off a UTF-8 continuation byte so the retained slice |
| 513 | // never begins mid-character (the leading-U+FFFD bug guarded against |
| 514 | // in `tail_from_buffer`). |
| 515 | while start < self.data.len() && (self.data[start] & 0xC0) == 0x80 { |
| 516 | start += 1; |
| 517 | } |
| 518 | self.data.drain(..start); |
| 519 | self.dropped = self.dropped.saturating_add(start); |
| 520 | } |
| 521 | } |
| 522 | |
| 523 | impl Default for RawOutputBuffer { |
| 524 | fn default() -> Self { |
| 525 | Self::new() |
| 526 | } |
| 527 | } |
| 528 | |
| 529 | pub(super) type SharedRawOutput = Arc<Mutex<RawOutputBuffer>>; |
| 530 | |
| 531 | pub(super) fn new_shared_raw_output() -> SharedRawOutput { |
| 532 | Arc::new(Mutex::new(RawOutputBuffer::new())) |
| 533 | } |
| 534 | |
| 535 | pub(super) fn take_delta_from_buffer( |
| 536 | buffer: &SharedRawOutput, |
| 537 | cursor: &mut usize, |
| 538 | ) -> (Vec<u8>, usize) { |
| 539 | let guard = buffer.lock().unwrap_or_else(|e| e.into_inner()); |
| 540 | let total = guard.total_len(); |
| 541 | // The cursor is an absolute offset into the stream. Bytes the bound already |
| 542 | // discarded can never be delivered as a delta, so skip forward over them |
| 543 | // rather than re-sending the retained tail as if it were new. |
| 544 | let start_abs = (*cursor).max(guard.dropped()).min(total); |
| 545 | let start = start_abs - guard.dropped(); |
| 546 | let retained = guard.retained(); |
| 547 | // Clone only the unread portion (the delta), not the entire accumulated buffer. |
| 548 | // Long-running processes can produce megabytes of output; cloning the full |
| 549 | // buffer on every poll held the ShellManager mutex for O(total_bytes) time. |
| 550 | let unread = &retained[start..]; |
| 551 | // A poll can land mid-character: the caller decodes this delta as UTF-8, so |
| 552 | // handing back a truncated multibyte sequence renders it as replacement |
| 553 | // glyphs and corrupts the next delta's leading byte too (the streaming-client |
| 554 | // bug from #1675, in the shell preview path). Leave an incomplete trailing |
| 555 | // sequence in the buffer for the next poll. Bytes that are genuinely invalid |
| 556 | // rather than merely unfinished still pass through, so binary output cannot |
| 557 | // stall the cursor, and the final result is read from the whole buffer. |
| 558 | let consumed = match std::str::from_utf8(unread) { |
| 559 | Ok(_) => unread.len(), |
| 560 | Err(error) if error.error_len().is_none() => error.valid_up_to(), |
| 561 | Err(_) => unread.len(), |
| 562 | }; |
| 563 | let delta = unread[..consumed].to_vec(); |
| 564 | *cursor = start_abs + consumed; |
| 565 | (delta, total) |
| 566 | } |
| 567 | |
| 568 | /// Read only the tail of a byte buffer and return (total_len, tail_string). |
| 569 | /// |
| 570 | /// Avoids cloning the full buffer when only a trailing excerpt is needed |
| 571 | /// (e.g. for the job-panel display). `max_tail_chars` is in Unicode scalar |
| 572 | /// values; we read at most `max_tail_chars * 4` bytes from the end to account |
| 573 | /// for multi-byte UTF-8 sequences. |
| 574 | pub(super) fn tail_from_buffer(buffer: &SharedRawOutput, max_tail_chars: usize) -> (usize, String) { |
| 575 | let guard = buffer.lock().unwrap_or_else(|e| e.into_inner()); |
| 576 | // The reported length is the stream's total, not what is still held: a |
| 577 | // released or clipped buffer must not make the model believe the command |
| 578 | // printed less than it did. |
| 579 | let total = guard.total_len(); |
| 580 | let retained = guard.retained(); |
| 581 | let retained_len = retained.len(); |
| 582 | // Over-estimate byte count (4 bytes per char worst case for UTF-8). |
| 583 | let mut tail_start = retained_len.saturating_sub(max_tail_chars.saturating_mul(4)); |
| 584 | // Snap forward to the next valid UTF-8 codepoint boundary so we don't |
| 585 | // pass a slice beginning with continuation bytes (0x80-0xBF) to |
| 586 | // from_utf8_lossy, which would emit a leading U+FFFD replacement char. |
| 587 | while tail_start < retained_len && (retained[tail_start] & 0xC0) == 0x80 { |
| 588 | tail_start += 1; |
| 589 | } |
| 590 | let tail_str = String::from_utf8_lossy(&retained[tail_start..]).into_owned(); |
| 591 | (total, tail_text(&tail_str, max_tail_chars)) |
| 592 | } |
| 593 | |
| 594 | pub(super) fn tail_text(text: &str, max_chars: usize) -> String { |
| 595 | if text.chars().count() <= max_chars { |
| 596 | return text.to_string(); |
| 597 | } |
| 598 | let tail = text |
| 599 | .chars() |
| 600 | .rev() |
| 601 | .take(max_chars) |
| 602 | .collect::<Vec<_>>() |
| 603 | .into_iter() |
| 604 | .rev() |
| 605 | .collect::<String>(); |
| 606 | format!("...{tail}") |
| 607 | } |
| 608 | |
| 609 | #[cfg(test)] |
| 610 | mod tests { |
| 611 | use super::{ |
| 612 | BOUNDED_OUTPUT_MAX_BYTES, BOUNDED_OUTPUT_MAX_LINES, BOUNDED_OUTPUT_RETAIN_BYTES, |
| 613 | BoundedOutputAccumulator, RAW_STREAM_MAX_BYTES, RawOutputBuffer, SharedRawOutput, |
| 614 | tail_from_buffer, take_delta_from_buffer, |
| 615 | }; |
| 616 | use std::sync::{Arc, Mutex}; |
| 617 | |
| 618 | fn raw(bytes: &[u8]) -> SharedRawOutput { |
| 619 | let mut buffer = RawOutputBuffer::new(); |
| 620 | buffer.append(bytes); |
| 621 | Arc::new(Mutex::new(buffer)) |
| 622 | } |
| 623 | |
| 624 | fn append(buffer: &SharedRawOutput, bytes: &[u8]) { |
| 625 | buffer.lock().unwrap().append(bytes); |
| 626 | } |
| 627 | |
| 628 | #[test] |
| 629 | fn delta_holds_back_an_incomplete_trailing_utf8_sequence() { |
| 630 | // "宽" is three bytes; deliver two of them, then the rest. |
| 631 | let wide = "宽".as_bytes(); |
| 632 | let buffer = raw(b"ok "); |
| 633 | append(&buffer, &wide[..2]); |
| 634 | let mut cursor = 0usize; |
| 635 | |
| 636 | let (delta, total) = take_delta_from_buffer(&buffer, &mut cursor); |
| 637 | assert_eq!( |
| 638 | String::from_utf8(delta).expect("delta must be whole characters"), |
| 639 | "ok " |
| 640 | ); |
| 641 | assert_eq!(total, 5, "total still reports every buffered byte"); |
| 642 | assert_eq!(cursor, 3, "the split character stays unread"); |
| 643 | |
| 644 | append(&buffer, &wide[2..]); |
| 645 | let (delta, _) = take_delta_from_buffer(&buffer, &mut cursor); |
| 646 | assert_eq!( |
| 647 | String::from_utf8(delta).expect("delta must be whole characters"), |
| 648 | "宽" |
| 649 | ); |
| 650 | } |
| 651 | |
| 652 | #[test] |
| 653 | fn delta_does_not_stall_on_genuinely_invalid_bytes() { |
| 654 | // A lone 0xFF is never a valid start byte: passing it through keeps |
| 655 | // binary output flowing instead of parking the cursor forever. |
| 656 | let buffer = raw(&[b'a', 0xFF, b'b']); |
| 657 | let mut cursor = 0usize; |
| 658 | let (delta, total) = take_delta_from_buffer(&buffer, &mut cursor); |
| 659 | assert_eq!(delta, vec![b'a', 0xFF, b'b']); |
| 660 | assert_eq!(cursor, total); |
| 661 | } |
| 662 | |
| 663 | // === #5472: in-memory retention bounds for the raw `Bash` streams === |
| 664 | |
| 665 | #[test] |
| 666 | fn raw_buffer_caps_in_flight_bytes_and_keeps_the_total_honest() { |
| 667 | let mut buffer = RawOutputBuffer::with_cap(1_024); |
| 668 | // 4 MiB through a 1 KiB cap: the analogue of `cargo build -v` through |
| 669 | // the 16 MiB production ceiling. |
| 670 | for _ in 0..1_024 { |
| 671 | buffer.append(&[b'x'; 4_096]); |
| 672 | } |
| 673 | let produced = 1_024 * 4_096; |
| 674 | assert_eq!( |
| 675 | buffer.total_len(), |
| 676 | produced, |
| 677 | "the stream's length must survive the bound" |
| 678 | ); |
| 679 | assert_eq!(buffer.dropped(), produced - buffer.retained().len()); |
| 680 | assert!( |
| 681 | buffer.retained().len() <= 1_024 + 1_024 / 4, |
| 682 | "retained {} exceeded cap + slack", |
| 683 | buffer.retained().len() |
| 684 | ); |
| 685 | } |
| 686 | |
| 687 | #[test] |
| 688 | fn raw_buffer_release_collapses_to_a_tail_and_reports_the_omission() { |
| 689 | let mut buffer = RawOutputBuffer::new(); |
| 690 | buffer.append(&[b'y'; 200_000]); |
| 691 | assert_eq!(buffer.dropped(), 0, "200 KB is under the in-flight ceiling"); |
| 692 | |
| 693 | buffer.release_to_tail(1_000); |
| 694 | assert_eq!(buffer.retained().len(), 1_000); |
| 695 | assert_eq!(buffer.dropped(), 199_000); |
| 696 | assert_eq!( |
| 697 | buffer.total_len(), |
| 698 | 200_000, |
| 699 | "releasing memory must not rewrite how much the command printed" |
| 700 | ); |
| 701 | } |
| 702 | |
| 703 | #[test] |
| 704 | fn raw_buffer_never_retains_a_split_character() { |
| 705 | let mut buffer = RawOutputBuffer::with_cap(8); |
| 706 | // Each "宽" is 3 bytes, so a byte-exact tail would land mid-character. |
| 707 | for _ in 0..64 { |
| 708 | buffer.append("宽".as_bytes()); |
| 709 | } |
| 710 | assert!( |
| 711 | std::str::from_utf8(buffer.retained()).is_ok(), |
| 712 | "front-drop must snap off continuation bytes" |
| 713 | ); |
| 714 | |
| 715 | let mut released = RawOutputBuffer::new(); |
| 716 | for _ in 0..64 { |
| 717 | released.append("宽".as_bytes()); |
| 718 | } |
| 719 | released.release_to_tail(10); |
| 720 | assert!(std::str::from_utf8(released.retained()).is_ok()); |
| 721 | } |
| 722 | |
| 723 | #[test] |
| 724 | fn delta_skips_bytes_the_bound_already_discarded() { |
| 725 | // A consumer that stops reading while output keeps arriving must be |
| 726 | // moved forward, not handed the retained tail as if it were new bytes. |
| 727 | let buffer = Arc::new(Mutex::new(RawOutputBuffer::with_cap(16))); |
| 728 | append(&buffer, b"first-chunk-that-will-be-dropped-entirely"); |
| 729 | let mut cursor = 0usize; |
| 730 | let (delta, total) = take_delta_from_buffer(&buffer, &mut cursor); |
| 731 | let dropped = buffer.lock().unwrap().dropped(); |
| 732 | assert!(dropped > 0, "the cap must have clipped the front"); |
| 733 | assert_eq!(cursor, total, "cursor lands at the stream's true position"); |
| 734 | assert_eq!( |
| 735 | delta.len(), |
| 736 | total - dropped, |
| 737 | "only bytes still held can be delivered" |
| 738 | ); |
| 739 | |
| 740 | append(&buffer, b"tail"); |
| 741 | let (delta, _) = take_delta_from_buffer(&buffer, &mut cursor); |
| 742 | assert_eq!( |
| 743 | delta, |
| 744 | b"tail".to_vec(), |
| 745 | "subsequent deltas continue from the corrected cursor" |
| 746 | ); |
| 747 | } |
| 748 | |
| 749 | #[test] |
| 750 | fn tail_reports_the_stream_total_not_the_retained_length() { |
| 751 | let buffer = Arc::new(Mutex::new(RawOutputBuffer::new())); |
| 752 | append(&buffer, b"abcdefghij"); |
| 753 | buffer.lock().unwrap().release_to_tail(4); |
| 754 | let (total, tail) = tail_from_buffer(&buffer, 100); |
| 755 | assert_eq!(total, 10, "stdout_len must not shrink when memory is freed"); |
| 756 | assert_eq!(tail, "ghij"); |
| 757 | } |
| 758 | |
| 759 | #[test] |
| 760 | fn abandoning_a_stream_releases_it_and_stops_the_reader() { |
| 761 | let mut buffer = RawOutputBuffer::new(); |
| 762 | assert!(buffer.append(&[b'a'; 5_000]), "a live stream keeps reading"); |
| 763 | buffer.abandon(); |
| 764 | |
| 765 | assert_eq!(buffer.retained().len(), 0, "held bytes are released"); |
| 766 | assert_eq!( |
| 767 | buffer.total_len(), |
| 768 | 5_000, |
| 769 | "the stream's length survives the release" |
| 770 | ); |
| 771 | assert!( |
| 772 | !buffer.append(&[b'b'; 100]), |
| 773 | "an abandoned stream tells the reader thread to exit" |
| 774 | ); |
| 775 | assert_eq!(buffer.retained().len(), 0, "and retains nothing further"); |
| 776 | assert_eq!( |
| 777 | buffer.total_len(), |
| 778 | 5_100, |
| 779 | "bytes that arrive after the give-up are still counted, not hidden" |
| 780 | ); |
| 781 | } |
| 782 | |
| 783 | #[test] |
| 784 | fn raw_stream_ceiling_clears_every_downstream_bound() { |
| 785 | // The clip must be unreachable by the model-visible surfaces: the 30 KB |
| 786 | // result truncation, the 1,200-char job tail, the 1 KiB completion tail. |
| 787 | const { assert!(RAW_STREAM_MAX_BYTES > 30_000 * 100) }; |
| 788 | const { assert!(super::RAW_STREAM_SETTLED_TAIL_BYTES > 30_000) }; |
| 789 | } |
| 790 | |
| 791 | #[test] |
| 792 | fn bounded_output_keeps_last_two_thousand_complete_lines() { |
| 793 | let source = (0..=BOUNDED_OUTPUT_MAX_LINES) |
| 794 | .map(|index| format!("line-{index}")) |
| 795 | .collect::<Vec<_>>() |
| 796 | .join("\n"); |
| 797 | let mut output = BoundedOutputAccumulator::new_in(None); |
| 798 | output.append(source.as_bytes()).expect("append"); |
| 799 | output.finish().expect("finish"); |
| 800 | let snapshot = output.snapshot(true).expect("snapshot"); |
| 801 | assert!(snapshot.truncated); |
| 802 | assert!(snapshot.content.starts_with("line-1\n")); |
| 803 | assert!(snapshot.content.contains("Showing lines 2-2001 of 2001")); |
| 804 | } |
| 805 | |
| 806 | #[test] |
| 807 | fn bounded_output_streams_raw_full_output_and_bounds_decoded_tail() { |
| 808 | let raw = vec![0xFF; 2 * 1024 * 1024]; |
| 809 | let mut output = BoundedOutputAccumulator::new_in(None); |
| 810 | for chunk in raw.chunks(4_096) { |
| 811 | output.append(chunk).expect("append"); |
| 812 | assert!(output.retained_memory_bytes() <= BOUNDED_OUTPUT_MAX_BYTES + 4); |
| 813 | } |
| 814 | output.finish().expect("finish"); |
| 815 | let snapshot = output.snapshot(true).expect("snapshot"); |
| 816 | assert!(snapshot.truncated); |
| 817 | assert!(snapshot.retained_bytes <= BOUNDED_OUTPUT_MAX_BYTES); |
| 818 | assert!(snapshot.content.contains('\u{FFFD}')); |
| 819 | let path = output |
| 820 | .full_output_path() |
| 821 | .expect("full output") |
| 822 | .to_path_buf(); |
| 823 | assert_eq!(std::fs::read(&path).expect("read full output"), raw); |
| 824 | drop(output); |
| 825 | std::fs::remove_file(path).expect("remove full output"); |
| 826 | } |
| 827 | |
| 828 | #[test] |
| 829 | fn bounded_output_delta_since_returns_only_new_bytes() { |
| 830 | let mut output = BoundedOutputAccumulator::new_in(None); |
| 831 | output.append(b"first\n").expect("append"); |
| 832 | let cursor = output.total_bytes(); |
| 833 | assert_eq!(output.delta_since(0).expect("delta").0, "first\n"); |
| 834 | output.append(b"second\n").expect("append"); |
| 835 | assert_eq!( |
| 836 | output.delta_since(cursor).expect("delta"), |
| 837 | ("second\n".to_string(), 0) |
| 838 | ); |
| 839 | assert_eq!( |
| 840 | output.delta_since(output.total_bytes()).expect("delta"), |
| 841 | (String::new(), 0) |
| 842 | ); |
| 843 | |
| 844 | // Past the memory bound: the dropped span is counted, and a clipped |
| 845 | // multi-byte character at the front is not emitted as garbage. |
| 846 | let mut output = BoundedOutputAccumulator::new_in(None); |
| 847 | // One ASCII byte first puts the retained tail's front mid-character. |
| 848 | output.append(b"a").expect("append"); |
| 849 | output |
| 850 | .append("é".repeat(BOUNDED_OUTPUT_MAX_BYTES).as_bytes()) |
| 851 | .expect("append"); |
| 852 | let total = output.total_bytes(); |
| 853 | let (delta, omitted) = output.delta_since(0).expect("delta"); |
| 854 | assert!(!delta.contains('\u{FFFD}')); |
| 855 | assert_eq!(delta.len() + omitted, total); |
| 856 | assert!(delta.len() <= BOUNDED_OUTPUT_RETAIN_BYTES); |
| 857 | } |
| 858 | |
| 859 | /// The notice for dropped delta bytes says where the full output is. |
| 860 | #[test] |
| 861 | fn bounded_output_omitted_notice_points_at_the_spill_file() { |
| 862 | let output = BoundedOutputAccumulator::new_in(None); |
| 863 | let spill = output |
| 864 | .temp |
| 865 | .as_ref() |
| 866 | .expect("spill file") |
| 867 | .path() |
| 868 | .display() |
| 869 | .to_string(); |
| 870 | let notice = output.omitted_notice(42); |
| 871 | assert!(notice.starts_with("[42 bytes"), "{notice}"); |
| 872 | assert!(notice.contains(&spill), "{notice}"); |
| 873 | |
| 874 | let missing = tempfile::tempdir().expect("tempdir").path().join("gone"); |
| 875 | let unspilled = BoundedOutputAccumulator::new_in(Some(&missing)); |
| 876 | let notice = unspilled.omitted_notice(7); |
| 877 | assert!(notice.contains("not persisted"), "{notice}"); |
| 878 | } |
| 879 | |
| 880 | #[test] |
| 881 | fn bounded_output_huge_terminal_line_matches_upstream_notice() { |
| 882 | let mut source = vec![b'x'; BOUNDED_OUTPUT_MAX_BYTES + 1_024]; |
| 883 | source.push(b'\n'); |
| 884 | let mut output = BoundedOutputAccumulator::new_in(None); |
| 885 | output.append(&source).expect("append"); |
| 886 | output.finish().expect("finish"); |
| 887 | let snapshot = output.snapshot(true).expect("snapshot"); |
| 888 | assert!(snapshot.content.contains("Showing last 50.0KB of line 1")); |
| 889 | assert!(snapshot.content.contains("line is 0B")); |
| 890 | let path = output |
| 891 | .full_output_path() |
| 892 | .expect("full output") |
| 893 | .to_path_buf(); |
| 894 | drop(output); |
| 895 | std::fs::remove_file(path).expect("remove full output"); |
| 896 | } |
| 897 | |
| 898 | #[test] |
| 899 | fn spill_failure_is_soft_and_names_the_reason() { |
| 900 | // A missing spill dir simulates a full or broken temp volume: the |
| 901 | // stream still runs, the tail is still delivered, and the notice says |
| 902 | // why "Full output: <path>" is absent instead of failing the command. |
| 903 | let missing = |
| 904 | std::env::temp_dir().join(format!("codewhale-missing-spill-{}", std::process::id())); |
| 905 | let mut output = BoundedOutputAccumulator::new_in(Some(&missing)); |
| 906 | let reason = output.spill_unavailable().expect("spill unavailable"); |
| 907 | assert!(!reason.is_empty(), "reason must name the io error"); |
| 908 | |
| 909 | output.append(b"ok\n").expect("append works without spill"); |
| 910 | output.finish().expect("finish works without spill"); |
| 911 | let short = output.snapshot(true).expect("snapshot"); |
| 912 | assert_eq!(short.content, "ok\n"); |
| 913 | assert!(!short.truncated); |
| 914 | assert!(output.full_output_path().is_none()); |
| 915 | |
| 916 | let source = (0..=BOUNDED_OUTPUT_MAX_LINES) |
| 917 | .map(|index| format!("line-{index}")) |
| 918 | .collect::<Vec<_>>() |
| 919 | .join("\n"); |
| 920 | let mut output = BoundedOutputAccumulator::new_in(Some(&missing)); |
| 921 | output.append(source.as_bytes()).expect("append"); |
| 922 | output.finish().expect("finish"); |
| 923 | let snapshot = output.snapshot(true).expect("snapshot"); |
| 924 | assert!(snapshot.truncated); |
| 925 | assert!(snapshot.content.starts_with("line-1\n")); |
| 926 | assert!( |
| 927 | snapshot.content.contains("Full output was not persisted:"), |
| 928 | "{}", |
| 929 | snapshot.content |
| 930 | ); |
| 931 | assert!(!snapshot.content.contains("Full output: ")); |
| 932 | assert!(output.full_output_path().is_none()); |
| 933 | } |
| 934 | |
| 935 | #[test] |
| 936 | fn resource_exhaustion_hint_names_disk_descriptors_and_memory() { |
| 937 | use std::io::{Error, ErrorKind}; |
| 938 | assert!( |
| 939 | super::resource_exhaustion_hint(&Error::from(ErrorKind::StorageFull)) |
| 940 | .expect("storage full") |
| 941 | .contains("disk") |
| 942 | ); |
| 943 | assert!( |
| 944 | super::resource_exhaustion_hint(&Error::from(ErrorKind::OutOfMemory)) |
| 945 | .expect("oom") |
| 946 | .contains("memory") |
| 947 | ); |
| 948 | #[cfg(unix)] |
| 949 | { |
| 950 | assert!( |
| 951 | super::resource_exhaustion_hint(&Error::from_raw_os_error(libc::ENOSPC)) |
| 952 | .expect("enospc") |
| 953 | .contains("disk") |
| 954 | ); |
| 955 | assert!( |
| 956 | super::resource_exhaustion_hint(&Error::from_raw_os_error(libc::EMFILE)) |
| 957 | .expect("emfile") |
| 958 | .contains("file descriptors") |
| 959 | ); |
| 960 | assert!( |
| 961 | super::resource_exhaustion_hint(&Error::from_raw_os_error(libc::EAGAIN)) |
| 962 | .expect("eagain") |
| 963 | .contains("retry") |
| 964 | ); |
| 965 | } |
| 966 | assert!(super::resource_exhaustion_hint(&Error::from(ErrorKind::NotFound)).is_none()); |
| 967 | assert!( |
| 968 | super::resource_exhaustion_hint(&Error::from(ErrorKind::PermissionDenied)).is_none() |
| 969 | ); |
| 970 | } |
| 971 | } |
| 972 |