Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,9 @@ name: CI

on:
push:
branches: [main]
branches: [main, develop]
pull_request:
branches: [main]
branches: [main, develop]

concurrency:
group: ci-${{ github.workflow }}-${{ github.ref }}
Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/docs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,10 @@ name: Documentation

on:
push:
branches: [main]
branches: [main, develop]
paths: ['docs/**', 'scripts/*docs*.py', 'justfile', '.github/workflows/docs.yml']
pull_request:
branches: [main]
branches: [main, develop]
paths: ['docs/**', 'scripts/*docs*.py', 'justfile', '.github/workflows/docs.yml']
workflow_dispatch:

Expand Down
4 changes: 2 additions & 2 deletions .github/workflows/pvisor-benchmark.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ name: pVisor Benchmark

on:
push:
branches: [main]
branches: [main, develop]
paths:
- 'Cargo.*'
- '.cargo/**'
Expand All @@ -14,7 +14,7 @@ on:
- '.github/actions/setup-build-env/**'
- '.github/workflows/pvisor-benchmark.yml'
pull_request:
branches: [main]
branches: [main, develop]
paths:
- 'Cargo.*'
- '.cargo/**'
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ tokio = { version = "1", default-features = false }
tokio-stream = "0.1"
tokio-util = { version = "0.7", default-features = false }
toml = "0.8"
toml_edit = "0.22"
tracing = "0.1"
url = "2"
uuid = "1"
Expand Down
100 changes: 82 additions & 18 deletions crates/persisting-overlaynet/src/vm.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use std::collections::{HashMap, VecDeque};
use std::io;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::os::unix::net::UnixStream as StdUnixStream;
use std::sync::Arc;
use std::thread;
use std::time::Duration as StdDuration;

Expand All @@ -26,7 +27,7 @@ use smoltcp::wire::{
};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpStream, UnixStream};
use tokio::sync::{mpsc, oneshot};
use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc, oneshot};
use tokio::task::JoinHandle;

use crate::egress::{
Expand Down Expand Up @@ -188,7 +189,7 @@ struct Flow {
phase: FlowPhase,
upstream: Option<mpsc::Sender<Vec<u8>>>,
upstream_task: Option<JoinHandle<()>>,
inbound: VecDeque<Vec<u8>>,
inbound: VecDeque<(Vec<u8>, Option<OwnedSemaphorePermit>)>,
inbound_offset: usize,
dns_input: Vec<u8>,
remote_eof: bool,
Expand All @@ -203,6 +204,7 @@ enum FlowEvent {
Data {
key: FlowKey,
bytes: Vec<u8>,
permit: OwnedSemaphorePermit,
},
Uploaded {
key: FlowKey,
Expand Down Expand Up @@ -561,7 +563,7 @@ fn drive_dns_tcp(
let mut framed = Vec::with_capacity(response.len() + 2);
framed.extend_from_slice(&(response.len() as u16).to_be_bytes());
framed.extend_from_slice(&response);
flow.inbound.push_back(framed);
flow.inbound.push_back((framed, None));
}
flush_inbound(socket, flow);
if !socket.may_recv() && flow.inbound.is_empty() && socket.may_send() {
Expand All @@ -571,7 +573,7 @@ fn drive_dns_tcp(

fn flush_inbound(socket: &mut tcp::Socket<'_>, flow: &mut Flow) {
while socket.can_send() {
let Some(front) = flow.inbound.front() else {
let Some((front, _permit)) = flow.inbound.front() else {
break;
};
match socket.send_slice(&front[flow.inbound_offset..]) {
Expand Down Expand Up @@ -702,7 +704,11 @@ async fn bridge_upstream(
};
let download = async {
let mut buffer = vec![0; 16 * 1024];
let credits = Arc::new(Semaphore::new(FLOW_BUFFER_CHUNKS));
loop {
// Bound each flow across both the shared event channel and inbound
// queue. Reading resumes only after the guest socket accepts data.
let permit = credits.clone().acquire_owned().await.unwrap();
let length = read.read(&mut buffer).await?;
if length == 0 {
let _ = events.send(FlowEvent::RemoteEof(key)).await;
Expand All @@ -713,6 +719,7 @@ async fn bridge_upstream(
.send(FlowEvent::Data {
key,
bytes: buffer[..length].to_vec(),
permit,
})
.await
.is_err()
Expand Down Expand Up @@ -758,21 +765,10 @@ fn apply_flow_event(
}
}
}
FlowEvent::Data { key, bytes } => {
FlowEvent::Data { key, bytes, permit } => {
if let Some(flow) = flows.get_mut(&key) {
let queued = flow
.inbound
.iter()
.map(Vec::len)
.sum::<usize>()
.saturating_sub(flow.inbound_offset);
if queued.saturating_add(bytes.len()) > TCP_BUFFER_BYTES {
sockets.get_mut::<tcp::Socket>(flow.handle).abort();
metrics.tcp_connect_failure();
} else {
metrics.host_to_guest(bytes.len());
flow.inbound.push_back(bytes);
}
metrics.host_to_guest(bytes.len());
flow.inbound.push_back((bytes, Some(permit)));
}
}
FlowEvent::Uploaded { key, bytes } => {
Expand Down Expand Up @@ -1244,6 +1240,74 @@ impl TxToken for FrameTxToken<'_> {
mod tests {
use super::*;

#[tokio::test]
async fn slow_guest_backpressures_large_download_without_losing_data() {
tokio::time::timeout(StdDuration::from_secs(10), async {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let stream = TcpStream::connect(listener.local_addr().unwrap())
.await
.unwrap();
let (mut server, _) = listener.accept().await.unwrap();
let payload: Vec<u8> = (0..1024 * 1024).map(|i| (i % 251) as u8).collect();
let expected = payload.clone();
let server_task = tokio::spawn(async move {
server.write_all(&payload).await.unwrap();
server.shutdown().await.unwrap();
});
let key = FlowKey {
guest_port: 49152,
destination_addr: Ipv4Addr::new(198, 18, 0, 1),
destination_port: 80,
};
let (sender, outbound) = mpsc::channel(FLOW_BUFFER_CHUNKS);
let (events, mut received) = mpsc::channel(TCP_CHANNEL_DEPTH);
let bridge = tokio::spawn(bridge_upstream(
key,
stream,
outbound,
Default::default(),
events,
));
drop(sender);
let mut bytes = Vec::new();
let mut held = Vec::new();
for _ in 0..FLOW_BUFFER_CHUNKS {
match received.recv().await.unwrap() {
FlowEvent::Data {
bytes: chunk,
permit,
..
} => {
bytes.extend(chunk);
held.push(permit);
}
_ => panic!("expected download data"),
}
}
assert!(bytes.len() <= TCP_BUFFER_BYTES);
// A stalled guest must stop this flow's reader, not abort the TCP
// connection or keep filling the shared event channel.
assert!(
tokio::time::timeout(StdDuration::from_millis(50), received.recv())
.await
.is_err()
);
drop(held);
loop {
match received.recv().await.unwrap() {
FlowEvent::Data { bytes: chunk, .. } => bytes.extend(chunk),
FlowEvent::RemoteEof(_) => break,
_ => panic!("download failed while guest resumed"),
}
}
assert_eq!(bytes, expected);
bridge.await.unwrap();
server_task.await.unwrap();
})
.await
.expect("download stalled after guest resumed");
}

#[test]
fn denied_vm_flow_records_its_destination() {
let key = FlowKey {
Expand Down
8 changes: 8 additions & 0 deletions crates/persisting-pvisor/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ flate2.workspace = true
fs2.workspace = true
globset.workspace = true
libc.workspace = true
lru = "0.12.5"
persisting-control.workspace = true
persisting-gateway = { workspace = true, default-features = false }
persisting-overlaynet.workspace = true
Expand All @@ -34,6 +35,7 @@ sha2.workspace = true
tar.workspace = true
thiserror.workspace = true
toml.workspace = true
toml_edit.workspace = true
tempfile.workspace = true
tokio = { workspace = true, features = ["io-util", "macros", "process", "rt-multi-thread", "signal", "sync", "time", "net"] }
tokio-util = { workspace = true, features = ["rt"] }
Expand All @@ -56,3 +58,9 @@ proptest.workspace = true

[build-dependencies]
libloading.workspace = true

[target.'cfg(target_os = "macos")'.dependencies]
fuser = { workspace = true, features = ["abi-7-19", "libfuse", "macfuse-5"] }

[target.'cfg(target_os = "linux")'.dependencies]
fuser = { workspace = true, default-features = false, features = ["abi-7-31"] }
20 changes: 18 additions & 2 deletions crates/persisting-pvisor/src/cli/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,9 @@ enum Command {
long_about = run::RUN_COMMAND_LONG_ABOUT
)]
Run(Box<run::RunArgs>),
/// Serve or query the shared OCI file cache.
#[cfg(unix)]
Cache(crate::image::cache::CacheArgs),
/// Apply selected staged changes from a stopped Job.
Apply(runtime::ApplyArgs),
/// Discard staged changes from a stopped Job.
Expand Down Expand Up @@ -78,11 +81,11 @@ pub fn main() -> anyhow::Result<()> {
if run.tui_requested() || audit {
anyhow::ensure!(
run.wants_tui(audit),
"--tui/--audit requires inherited stdio and a normal Job"
"--tui/--ask requires inherited stdio and a normal Job"
);
anyhow::ensure!(
tui::available(),
"--tui/--audit requires an interactive terminal"
"--tui/--ask requires an interactive terminal"
);
let code = tui::run(args, audit)?;
if code != 0 {
Expand All @@ -92,6 +95,8 @@ pub fn main() -> anyhow::Result<()> {
}
}
match parsed.command {
#[cfg(unix)]
Command::Cache(args) => crate::image::cache::run(args)?,
Command::Run(args) => {
let code = tokio::runtime::Runtime::new()?.block_on(run::run(*args))?;
if code != 0 {
Expand Down Expand Up @@ -138,6 +143,7 @@ fn normalize_default_run(mut args: Vec<std::ffi::OsString>) -> Vec<std::ffi::OsS
"run",
"replay",
"env",
"cache",
"status",
"kill",
"inspect",
Expand All @@ -164,6 +170,9 @@ mod tests {
fn standalone_cli_is_small_and_run_can_be_explicit() {
for args in [
vec!["pvisor", "status"],
vec!["pvisor", "cache", "serve"],
vec!["pvisor", "cache", "prepare", "alpine:latest"],
vec!["pvisor", "cache", "list", "sha256:example"],
vec!["pvisor", "inspect", "run-1", "--", "rg", "TODO"],
vec!["pvisor", "status", "run-1", "--review"],
vec!["pvisor", "kill", "run-1"],
Expand Down Expand Up @@ -269,6 +278,13 @@ mod tests {
assert!(help.contains("including the replayed prefix and any live continuation"));
}

#[test]
fn cache_is_not_rewritten_to_run() {
let args = normalize_default_run(vec!["pvisor".into(), "cache".into(), "serve".into()]);
assert_eq!(args[1], "cache");
Cli::try_parse_from(args).unwrap();
}

#[test]
fn unknown_first_token_becomes_default_run() {
let args = normalize_default_run(vec!["pvisor".into(), "/bin/true".into()]);
Expand Down
Loading
Loading