diff --git a/lib/propolis/src/hw/virtio/p9fs.rs b/lib/propolis/src/hw/virtio/p9fs.rs index 00f887958..f346d07f4 100644 --- a/lib/propolis/src/hw/virtio/p9fs.rs +++ b/lib/propolis/src/hw/virtio/p9fs.rs @@ -845,8 +845,7 @@ impl P9Handler for HostFSHandler { let mut entries: Vec = Vec::new(); - let mut offset = 1; - for de in &dir[msg.offset as usize..] { + for (offset, de) in (1..).zip(dir[msg.offset as usize..].iter()) { let metadata = match de.metadata() { Ok(m) => m, Err(e) => { @@ -885,7 +884,6 @@ impl P9Handler for HostFSHandler { space_left -= dirent.wire_size(); entries.push(dirent); - offset += 1; } let response = Rreaddir::new(entries); diff --git a/lib/propolis/src/hw/virtio/softnpu.rs b/lib/propolis/src/hw/virtio/softnpu.rs index 66546c543..c5921b48d 100644 --- a/lib/propolis/src/hw/virtio/softnpu.rs +++ b/lib/propolis/src/hw/virtio/softnpu.rs @@ -8,7 +8,7 @@ use std::{ io::{Result, Write}, sync::{Arc, Mutex}, thread::{sleep, spawn}, - time::Duration, + time::{Duration, Instant}, }; use crate::{ @@ -244,15 +244,23 @@ impl SoftNpu { log: Logger, ) { info!(log, "management handler thread started"); + let mut needs_resync = false; loop { let r = ManagementMessageReader::new(uart.clone(), log.clone()); - let msg = r.read(); + let msg = r.read(&mut needs_resync); info!(log, "received management message: {:#?}", msg); let pipeline = pipeline.clone(); let uart = uart.clone(); let log = log.clone(); - handle_management_message(msg, pipeline, uart, radix, log.clone()); + handle_management_message( + msg, + pipeline, + uart, + radix, + &mut needs_resync, + log.clone(), + ); info!(log, "handled management message"); } } @@ -664,18 +672,80 @@ fn read_buf(mem: &MemCtx, chain: &mut Chain, buf: &mut [u8]) -> usize { }) } +/// Write each byte of `buf` to the uart, yielding while the one-byte FIFO +/// is full. Gives up once `deadline` passes. +/// +/// Returns the number of bytes written. +fn write_with_deadline(uart: &LpcUart, buf: &[u8], deadline: Instant) -> usize { + for (i, b) in buf.iter().enumerate() { + if Instant::now() >= deadline { + return i; + } + while !uart.write(*b) { + if Instant::now() >= deadline { + return i; + } + sleep(Duration::from_millis(1)); + } + } + buf.len() +} + +/// Write a response buffer to the management uart, yielding while the guest +/// drains the FIFO. This gives up once a deadline passes. So, a guest that +/// stopped reading the management tty cannot block this thread. +/// +/// A timed out write can leave a partial, unterminated frame in the tty; +/// `needs_resync` makes the next call write a newline first to terminate it. +/// +/// Returns true if the full buffer was written. +fn write_management_response( + uart: &LpcUart, + buf: &[u8], + needs_resync: &mut bool, + log: &Logger, +) -> bool { + // The management protocol has no client side timeout to inherit from. + // + // `scadm` performs one blocking `read()` with a 1 KiB buffer, and only for + // the radix query. Every other command writes the tty and never reads a + // response. A guest that hits this deadline stopped reading. + // + // See . + const WRITE_TIMEOUT: Duration = Duration::from_secs(30); + let deadline = Instant::now() + WRITE_TIMEOUT; + + if *needs_resync { + if write_with_deadline(uart, b"\n", deadline) != 1 { + warn!(log, "management uart write timed out, dropping response"); + return false; + } + *needs_resync = false; + } + + let written = write_with_deadline(uart, buf, deadline); + if written == buf.len() { + return true; + } + if written > 0 { + *needs_resync = true; + } + warn!(log, "management uart write timed out, dropping response"); + false +} + /// Handle ASIC management messages from the guest using the loaded program. fn handle_management_message( msg: ManagementRequest, pipeline: Arc>>, uart: Arc, radix: usize, + needs_resync: &mut bool, log: Logger, ) { - let mut pl_opt = pipeline.lock().unwrap(); - match msg { ManagementRequest::TableAdd(tm) => { + let mut pl_opt = pipeline.lock().unwrap(); let pl = match &mut *pl_opt { Some(pl) => pl, None => return, @@ -689,6 +759,7 @@ fn handle_management_message( ); } ManagementRequest::TableRemove(tm) => { + let mut pl_opt = pipeline.lock().unwrap(); let pl = match &mut *pl_opt { Some(pl) => pl, None => return, @@ -705,16 +776,21 @@ fn handle_management_message( let mut buf: Vec = Vec::new(); buf.extend_from_slice(radix.to_string().as_bytes()); buf.push(b'\n'); - for b in &buf { - while !uart.write(*b) { - std::thread::yield_now(); - } + if write_management_response(&uart, &buf, needs_resync, &log) { + info!(log, "wrote: {} bytes", buf.len()); } - info!(log, "wrote: {:?}", buf.len()); } ManagementRequest::DumpRequest => { info!(log, "dumping state"); + // Collect the table state under the pipeline lock, then serialize + // and write the response after releasing it. Holding the lock + // across the uart write loop below can deadlock the whole guest. + // For example, if the guest stops draining the management tty, this + // thread spins with the lock held while a vcpu servicing a queue + // notify blocks on the same lock in process_guest_packet. That + // vcpu is stuck in its exit and nothing ever drains the tty. let result = { + let mut pl_opt = pipeline.lock().unwrap(); let pl = match &mut *pl_opt { Some(pl) => &pl.1, None => return, @@ -726,7 +802,9 @@ fn handle_management_message( for id in pl.get_table_ids() { let entries = pl.get_table_entries(id); - result.insert(id, entries); + // The table ids borrow from the pipeline, so own them to + // let the map outlive the lock. + result.insert(id.to_owned(), entries); } result }; @@ -734,26 +812,20 @@ fn handle_management_message( let buf = match serde_json::to_string(&result) { Ok(j) => { let mut buf = j.as_bytes().to_vec(); - info!(log, "writing: {}", j); + info!(log, "writing table dump: {} bytes", j.len()); // Add trailing newline for proper tty handling. buf.push(b'\n'); buf } Err(e) => { - warn!(log, "failed to serialize table state: {}", e); + warn!(log, "failed to serialize table state: {e}"); b"{}\n".to_vec() } }; - for b in &buf { - while !uart.write(*b) { - // If we cannot write to the uart, yield and come back once - // scheduled again. - std::thread::yield_now(); - } + if write_management_response(&uart, &buf, needs_resync, &log) { + info!(log, "management wrote: {}", buf.len()); } - - info!(log, "management wrote: {}", buf.len()); } } } @@ -773,12 +845,15 @@ impl ManagementMessageReader { Self { uart, log } } - fn read(&self) -> ManagementRequest { + fn read(&self, needs_resync: &mut bool) -> ManagementRequest { loop { let mut buf = vec![0; 10240]; let mut i = 0; let mut in_message = false; loop { + if *needs_resync && self.uart.write(b'\n') { + *needs_resync = false; + } let x = match self.uart.read() { Some(b) => b, None => {