From 049ca16aba83a2558d9f4311ee8a7a0c003ee25d Mon Sep 17 00:00:00 2001 From: Kris Hicks Date: Wed, 30 Sep 2026 13:57:22 -0700 Subject: [PATCH] feat(sandbox): write agent output to the container log Since #2726 the canonical main process's stdout and stderr are captured in pipes that feed only the in-memory replay buffer used by sandbox connect. Agent output therefore never reaches the container's own stdout and stderr, so it is missing from kubectl logs, docker logs, and podman logs and from anything that collects container logs. Before #2726 the entrypoint inherited the container's descriptors and its output appeared there. Copy the main process's output to the launcher's stdout and stderr in addition to the replay buffer, restoring the earlier behavior: - Output is copied byte for byte to the matching stream from a forwarder thread per stream, after it is published to the replay buffer. When the container runtime falls behind on a stream, that stream's reader waits instead of dropping output, so backpressure reaches the agent as it did with inherited descriptors, while the other stream and attachments keep receiving output. - Before the main process's exit is published, the output readers finish and queued output is drained to the container log, so an agent's final lines are not lost at shutdown. A 30 second deadline covers both; when it expires, readers waiting on the container log are released and drain the pipes into the replay buffer only, so a stalled container log cannot block exit reporting. - PTY-mode processes are not copied. The terminal stream carries escape sequences and echoed input, and terminal commands never reached the container log before #2726. - Exec, SSH, and SFTP sessions are not copied. Launcher log lines keep their existing format and remain in the container's stderr. They are written as whole lines, and a newline is inserted first when the agent left stderr mid-line, so launcher and agent lines do not merge. The Docker and VM drivers appended the tail of the workload's output to failure messages: Docker the workload container's log, and the VM driver the guest console, which carries the launcher's stdout and stderr. Those messages land in the sandbox's Ready condition and in platform events that the gateway republishes to the sandbox event stream. With agent output in that log, those messages would carry arbitrary agent output, including anything sensitive the agent prints, into gateway status and events. The supervisor starts its health endpoint only after the agent starts, so every Docker failure path could include agent output, and the VM driver reports one whenever the VM or host supervisor exits. Forward only the supervisor's log tail, matching the Podman driver, which reads the workload log solely to match fixed launcher markers and never forwards raw workload output. The workload's output remains available through docker logs and the VM's rootfs-console.log. Document where main process output appears in the logging docs and the cluster debugging skill. Closes #3928 Signed-off-by: Kris Hicks --- crates/openshell-driver-docker/src/lib.rs | 36 +- crates/openshell-driver-docker/src/tests.rs | 9 +- crates/openshell-driver-vm/src/driver.rs | 77 +++- crates/openshell-sandbox/src/container_log.rs | 399 ++++++++++++++++++ crates/openshell-sandbox/src/delegated.rs | 6 +- crates/openshell-sandbox/src/lib.rs | 1 + crates/openshell-sandbox/src/main.rs | 8 +- crates/openshell-sandbox/src/main_session.rs | 255 ++++++++++- docs/observability/accessing-logs.mdx | 14 + skills/debug-openshell-cluster/SKILL.md | 4 + 10 files changed, 747 insertions(+), 62 deletions(-) create mode 100644 crates/openshell-sandbox/src/container_log.rs diff --git a/crates/openshell-driver-docker/src/lib.rs b/crates/openshell-driver-docker/src/lib.rs index 3f839876e9..950b2742e8 100644 --- a/crates/openshell-driver-docker/src/lib.rs +++ b/crates/openshell-driver-docker/src/lib.rs @@ -5288,12 +5288,6 @@ async fn spawn_docker_control_process( if !log_tail.is_empty() { write!(message, "; log tail: {log_tail}").ok(); } - let sandbox_log_tail = - docker_container_log_tail(&monitored_docker, &failure_context.container_id) - .await; - if !sandbox_log_tail.is_empty() { - write!(message, "; sandbox log tail: {sandbox_log_tail}").ok(); - } let _ = monitored_docker.remove_container( &monitored_supervisor_id, Some(RemoveContainerOptionsBuilder::default().force(true).build()), @@ -5344,11 +5338,9 @@ async fn wait_for_docker_supervisor_ready( Status::internal(format!("inspect Docker sandbox container: {error}")) })?; if sandbox.state.unwrap_or_default().running == Some(false) { - let sandbox_log_tail = docker_container_log_tail(docker, sandbox_id).await; - return Err(Status::unavailable(format!( - "Docker sandbox exited before supervisor became ready{}", - format_named_log_tail("sandbox log tail", &sandbox_log_tail) - ))); + return Err(Status::unavailable( + "Docker sandbox exited before supervisor became ready", + )); } let inspected = docker .inspect_container(supervisor_id, None) @@ -5361,12 +5353,10 @@ async fn wait_for_docker_supervisor_ready( Some(HealthStatusEnum::HEALTHY) => return Ok(()), _ if state.running == Some(false) => { let log_tail = docker_container_log_tail(docker, supervisor_id).await; - let sandbox_log_tail = docker_container_log_tail(docker, sandbox_id).await; - warn!(sandbox_id, supervisor_id, supervisor_logs = %log_tail, sandbox_logs = %sandbox_log_tail, "Docker supervisor exited before becoming ready"); + warn!(sandbox_id, supervisor_id, supervisor_logs = %log_tail, "Docker supervisor exited before becoming ready"); return Err(Status::unavailable(format!( - "Docker supervisor exited before becoming ready{}{}", - format_log_tail(&log_tail), - format_named_log_tail("sandbox log tail", &sandbox_log_tail) + "Docker supervisor exited before becoming ready{}", + format_log_tail(&log_tail) ))); } _ => tokio::time::sleep(Duration::from_millis(100)).await, @@ -5375,13 +5365,9 @@ async fn wait_for_docker_supervisor_ready( } fn format_log_tail(log_tail: &str) -> String { - format_named_log_tail("log tail", log_tail) -} - -fn format_named_log_tail(label: &str, log_tail: &str) -> String { - // gRPC status messages travel in HTTP/2 headers. Two 16 KiB container - // tails exceed the client's 16 KiB header budget and hide the real error - // behind PROTOCOL_ERROR. Allow for up to 3x percent-encoding expansion. + // gRPC status messages travel in HTTP/2 headers. A full 16 KiB container + // tail can exceed the client's 16 KiB header budget and hide the real + // error behind PROTOCOL_ERROR. Allow for up to 3x percent-encoding expansion. const MAX_STATUS_LOG_TAIL_BYTES: usize = 1024; if log_tail.is_empty() { String::new() @@ -5390,9 +5376,9 @@ fn format_named_log_tail(label: &str, log_tail: &str) -> String { while !log_tail.is_char_boundary(start) { start += 1; } - format!("; {label}: [truncated] {}", &log_tail[start..]) + format!("; log tail: [truncated] {}", &log_tail[start..]) } else { - format!("; {label}: {log_tail}") + format!("; log tail: {log_tail}") } } diff --git a/crates/openshell-driver-docker/src/tests.rs b/crates/openshell-driver-docker/src/tests.rs index f09b2c6f91..08576f422a 100644 --- a/crates/openshell-driver-docker/src/tests.rs +++ b/crates/openshell-driver-docker/src/tests.rs @@ -27,15 +27,14 @@ use tempfile::TempDir; #[test] fn startup_error_log_tails_fit_grpc_header_budget() { // Multibyte text exercises both the UTF-8 cut and worst-case gRPC message - // percent encoding. Preserve the final diagnostic from each container. + // percent encoding. Preserve the supervisor's final diagnostic. let logs = format!("{}\nstartup timed out", "🦀".repeat(8192)); let message = format!( - "Docker supervisor exited before becoming ready{}{}", + "Docker supervisor exited before becoming ready{}", format_log_tail(&logs), - format_named_log_tail("sandbox log tail", &logs), ); - assert_eq!(message.matches("[truncated]").count(), 2); - assert_eq!(message.matches("startup timed out").count(), 2); + assert_eq!(message.matches("[truncated]").count(), 1); + assert_eq!(message.matches("startup timed out").count(), 1); let response = Status::unavailable(message).into_http::<()>(); let header_bytes: usize = response .headers() diff --git a/crates/openshell-driver-vm/src/driver.rs b/crates/openshell-driver-vm/src/driver.rs index 3fceaf0e80..30eec2cd5e 100644 --- a/crates/openshell-driver-vm/src/driver.rs +++ b/crates/openshell-driver-vm/src/driver.rs @@ -4171,16 +4171,6 @@ impl VmDriver { || format!("{component} process exited"), |code| format!("{component} process exited with status {code}"), ); - if component == "VM" - && let Some(state_dir) = state_dir.as_deref() - && let Some(console) = read_vm_console_tail( - &state_dir.join("rootfs-console.log"), - VM_CONSOLE_DIAGNOSTIC_BYTES, - ) - { - write!(message, "; guest console tail:\n{console}") - .expect("writing to String cannot fail"); - } if component == "host supervisor" && let Some(state_dir) = state_dir.as_deref() && let Some(stderr) = read_vm_console_tail( @@ -4191,16 +4181,6 @@ impl VmDriver { write!(message, "; supervisor stderr tail:\n{stderr}") .expect("writing to String cannot fail"); } - if component == "host supervisor" - && let Some(state_dir) = state_dir.as_deref() - && let Some(console) = read_vm_console_tail( - &state_dir.join("rootfs-console.log"), - VM_CONSOLE_DIAGNOSTIC_BYTES, - ) - { - write!(message, "; guest console tail:\n{console}") - .expect("writing to String cannot fail"); - } if let Some(snapshot) = self .set_snapshot_condition( &sandbox_id, @@ -7517,6 +7497,63 @@ mod tests { })); } + #[tokio::test] + async fn process_exit_status_omits_guest_console_output() { + let running_child = || { + Command::new("sleep") + .arg("30") + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .kill_on_drop(true) + .spawn() + .unwrap() + }; + for (component, child) in [ + ("VM", spawn_exited_child()), + ("host supervisor", running_child()), + ] { + let temp = tempfile::tempdir().unwrap(); + std::fs::write(temp.path().join("rootfs-console.log"), "agent output\n").unwrap(); + std::fs::write(temp.path().join("supervisor.err.log"), "supervisor error\n").unwrap(); + let driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + let mut events = driver.events.subscribe(); + insert_test_record(&driver, "sb-exit", temp.path().to_path_buf(), child).await; + + driver.monitor_sandbox("sb-exit".to_string()).await; + + let condition = driver.registry.lock().await["sb-exit"] + .snapshot + .status + .as_ref() + .and_then(|status| { + status + .conditions + .iter() + .find(|condition| condition.reason == "ProcessExited") + .cloned() + }) + .expect("ProcessExited condition"); + let mut event_message = None; + while let Ok(event) = events.try_recv() { + if let Some(watch_sandboxes_event::Payload::PlatformEvent(platform)) = event.payload + && let Some(event) = platform.event + && event.reason == "ProcessExited" + { + event_message = Some(event.message); + } + } + let event_message = event_message.expect("ProcessExited platform event"); + for message in [&condition.message, &event_message] { + assert!(message.starts_with(&format!("{component} process exited"))); + assert!(!message.contains("agent output"), "{component}: {message}"); + } + if component == "host supervisor" { + assert!(condition.message.contains("supervisor error")); + } + } + } + #[tokio::test] async fn background_provisioning_does_not_extend_the_rpc_span_lifetime() { let traced = TestTracing::new(); diff --git a/crates/openshell-sandbox/src/container_log.rs b/crates/openshell-sandbox/src/container_log.rs new file mode 100644 index 0000000000..86264d1255 --- /dev/null +++ b/crates/openshell-sandbox/src/container_log.rs @@ -0,0 +1,399 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! Container stdout/stderr shared by the agent's output and launcher log lines. + +use std::io::Write; +use std::sync::{Arc, Mutex, OnceLock}; + +use bytes::Bytes; +use tokio::sync::{mpsc, oneshot}; + +const AGENT_OUTPUT_QUEUE_CHUNKS: usize = 16; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AgentStream { + Stdout, + Stderr, +} + +struct StderrState { + writer: Box, + mid_line: bool, +} + +pub struct ContainerLog { + stdout: Mutex>, + stderr: Mutex, +} + +impl ContainerLog { + #[must_use] + pub fn new(stdout: Box, stderr: Box) -> Arc { + Arc::new(Self { + stdout: Mutex::new(stdout), + stderr: Mutex::new(StderrState { + writer: stderr, + mid_line: false, + }), + }) + } + + /// The process's own stdout and stderr. + pub fn process() -> &'static Arc { + static LOG: OnceLock> = OnceLock::new(); + LOG.get_or_init(|| Self::new(Box::new(std::io::stdout()), Box::new(std::io::stderr()))) + } + + /// The agent output forwarder for the process's own stdout and stderr. + pub fn process_agent_output() -> AgentOutputSink { + static SINK: OnceLock = OnceLock::new(); + SINK.get_or_init(|| Self::process().agent_output()).clone() + } + + fn write_agent(&self, stream: AgentStream, data: &[u8]) { + match stream { + AgentStream::Stdout => { + let mut stdout = self.stdout.lock().expect("container stdout lock poisoned"); + let _ = stdout.write_all(data); + let _ = stdout.flush(); + } + AgentStream::Stderr => { + let mut stderr = self.stderr.lock().expect("container stderr lock poisoned"); + if let Some(last) = data.last() { + stderr.mid_line = *last != b'\n'; + } + let _ = stderr.writer.write_all(data); + let _ = stderr.writer.flush(); + } + } + } + + fn write_launcher_line(&self, line: &[u8]) { + let mut stderr = self.stderr.lock().expect("container stderr lock poisoned"); + let mut record = Vec::with_capacity(line.len() + 1); + if stderr.mid_line { + record.push(b'\n'); + stderr.mid_line = false; + } + record.extend_from_slice(line); + let _ = stderr.writer.write_all(&record); + let _ = stderr.writer.flush(); + } + + /// A writer for launcher log output that emits whole lines. + #[must_use] + pub fn launcher_writer(self: &Arc) -> LauncherLogWriter { + LauncherLogWriter { + log: Arc::clone(self), + pending: Vec::new(), + } + } + + /// Forward agent output with a dedicated thread per stream. + #[must_use] + pub fn agent_output(self: &Arc) -> AgentOutputSink { + AgentOutputSink { + stdout: self.forwarder(AgentStream::Stdout, "openshell-agent-stdout"), + stderr: self.forwarder(AgentStream::Stderr, "openshell-agent-stderr"), + } + } + + fn forwarder(self: &Arc, stream: AgentStream, name: &str) -> mpsc::Sender { + let (sender, mut receiver) = mpsc::channel(AGENT_OUTPUT_QUEUE_CHUNKS); + let log = Arc::clone(self); + std::thread::Builder::new() + .name(name.to_string()) + .spawn(move || { + while let Some(message) = receiver.blocking_recv() { + match message { + AgentOutput::Data(data) => log.write_agent(stream, &data), + AgentOutput::Drained(done) => { + let _ = done.send(()); + } + } + } + }) + .expect("spawn agent output forwarder"); + sender + } +} + +pub struct LauncherLogWriter { + log: Arc, + pending: Vec, +} + +impl Write for LauncherLogWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.pending.extend_from_slice(buf); + while let Some(end) = self.pending.iter().position(|byte| *byte == b'\n') { + let line: Vec = self.pending.drain(..=end).collect(); + self.log.write_launcher_line(&line); + } + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +enum AgentOutput { + Data(Bytes), + Drained(oneshot::Sender<()>), +} + +#[derive(Clone)] +pub struct AgentOutputSink { + stdout: mpsc::Sender, + stderr: mpsc::Sender, +} + +impl AgentOutputSink { + /// Waits while the container log is not keeping up with this stream. + pub async fn send(&self, stream: AgentStream, data: Bytes) { + let _ = self.sender(stream).send(AgentOutput::Data(data)).await; + } + + /// Waits until previously sent output has been written. + pub async fn drain(&self) { + tokio::join!(drain(&self.stdout), drain(&self.stderr)); + } + + const fn sender(&self, stream: AgentStream) -> &mpsc::Sender { + match stream { + AgentStream::Stdout => &self.stdout, + AgentStream::Stderr => &self.stderr, + } + } +} + +async fn drain(sender: &mpsc::Sender) { + let (done, drained) = oneshot::channel(); + if sender.send(AgentOutput::Drained(done)).await.is_ok() { + let _ = drained.await; + } +} + +#[cfg(test)] +pub(crate) mod test_support { + use super::*; + + #[derive(Clone, Default)] + pub struct Capture(Arc>>); + + impl Capture { + pub fn contents(&self) -> String { + String::from_utf8(self.0.lock().unwrap().clone()).unwrap() + } + + pub async fn wait_for(&self, expected: &str) -> String { + let _ = tokio::time::timeout(std::time::Duration::from_secs(2), async { + while self.contents() != expected { + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + } + }) + .await; + self.contents() + } + } + + impl Write for Capture { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + pub fn capture_log() -> (Arc, Capture, Capture) { + let stdout = Capture::default(); + let stderr = Capture::default(); + let log = ContainerLog::new(Box::new(stdout.clone()), Box::new(stderr.clone())); + (log, stdout, stderr) + } + + /// Holds container stdout writes until opened. + #[derive(Clone, Default)] + pub struct Gate(Arc<(Mutex, std::sync::Condvar)>); + + impl Gate { + pub fn open(&self) { + *self.0.0.lock().unwrap() = true; + self.0.1.notify_all(); + } + + fn wait(&self) { + let mut open = self.0.0.lock().unwrap(); + while !*open { + open = self.0.1.wait(open).unwrap(); + } + } + } + + struct Gated(Gate, Capture); + + impl Write for Gated { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.wait(); + self.1.write(buf) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + pub fn gated_capture_log(gate: &Gate) -> (Arc, Capture, Capture) { + let stdout = Capture::default(); + let stderr = Capture::default(); + let log = ContainerLog::new( + Box::new(Gated(gate.clone(), stdout.clone())), + Box::new(stderr.clone()), + ); + (log, stdout, stderr) + } +} + +#[cfg(test)] +mod tests { + use super::test_support::{Gate, capture_log, gated_capture_log}; + use super::*; + + #[test] + fn agent_output_is_written_unmodified() { + let (log, stdout, stderr) = capture_log(); + log.write_agent(AgentStream::Stdout, b"{\"msg\":\"out\"}\n"); + log.write_agent(AgentStream::Stderr, b"err\n"); + assert_eq!(stdout.contents(), "{\"msg\":\"out\"}\n"); + assert_eq!(stderr.contents(), "err\n"); + } + + #[test] + fn launcher_lines_follow_complete_agent_lines() { + let (log, stdout, stderr) = capture_log(); + log.write_agent(AgentStream::Stderr, b"agent line\n"); + log.write_launcher_line(b"WARN denied\n"); + assert_eq!(stderr.contents(), "agent line\nWARN denied\n"); + assert_eq!(stdout.contents(), ""); + } + + #[test] + fn launcher_line_starts_on_a_new_line_after_partial_agent_stderr() { + let (log, _stdout, stderr) = capture_log(); + log.write_agent(AgentStream::Stderr, b"partial"); + log.write_launcher_line(b"WARN denied\n"); + log.write_agent(AgentStream::Stderr, b" rest\n"); + assert_eq!(stderr.contents(), "partial\nWARN denied\n rest\n"); + } + + #[test] + fn partial_agent_stdout_does_not_split_launcher_lines() { + let (log, stdout, stderr) = capture_log(); + log.write_agent(AgentStream::Stdout, b"prompt> "); + log.write_launcher_line(b"WARN denied\n"); + assert_eq!(stdout.contents(), "prompt> "); + assert_eq!(stderr.contents(), "WARN denied\n"); + } + + #[test] + fn launcher_writer_joins_fragmented_writes_into_lines() { + let (log, _stdout, stderr) = capture_log(); + let mut writer = log.launcher_writer(); + write!(writer, "2026-01-01T00:00:00.000Z ").unwrap(); + write!(writer, "WARN target: one\nWARN target: ").unwrap(); + assert_eq!( + stderr.contents(), + "2026-01-01T00:00:00.000Z WARN target: one\n" + ); + writeln!(writer, "two").unwrap(); + assert_eq!( + stderr.contents(), + "2026-01-01T00:00:00.000Z WARN target: one\nWARN target: two\n" + ); + } + + #[tokio::test] + async fn agent_output_sink_waits_for_a_stalled_log_without_dropping() { + let gate = Gate::default(); + let (log, stdout, _stderr) = gated_capture_log(&gate); + let sink = log.agent_output(); + let chunks = AGENT_OUTPUT_QUEUE_CHUNKS * 4; + let sender = tokio::spawn(async move { + for _ in 0..chunks { + sink.send(AgentStream::Stdout, Bytes::from_static(b"x\n")) + .await; + } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + assert!(!sender.is_finished()); + + gate.open(); + sender.await.unwrap(); + let expected = "x\n".repeat(chunks); + assert_eq!(stdout.wait_for(&expected).await, expected); + } + + #[tokio::test] + async fn stalled_stdout_does_not_block_stderr() { + let gate = Gate::default(); + let (log, _stdout, stderr) = gated_capture_log(&gate); + let sink = log.agent_output(); + let stdout_sink = sink.clone(); + let stdout_sender = tokio::spawn(async move { + for _ in 0..AGENT_OUTPUT_QUEUE_CHUNKS * 4 { + stdout_sink + .send(AgentStream::Stdout, Bytes::from_static(b"x\n")) + .await; + } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + assert!(!stdout_sender.is_finished()); + + tokio::time::timeout( + std::time::Duration::from_secs(1), + sink.send(AgentStream::Stderr, Bytes::from_static(b"still flowing\n")), + ) + .await + .expect("stderr send waited on stalled stdout"); + assert_eq!(stderr.wait_for("still flowing\n").await, "still flowing\n"); + gate.open(); + } + + #[tokio::test] + async fn drain_waits_for_queued_output_to_be_written() { + let gate = Gate::default(); + let (log, stdout, _stderr) = gated_capture_log(&gate); + let sink = log.agent_output(); + sink.send(AgentStream::Stdout, Bytes::from_static(b"last words\n")) + .await; + let drain = tokio::spawn({ + let sink = sink.clone(); + async move { sink.drain().await } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + assert!(!drain.is_finished()); + + gate.open(); + drain.await.unwrap(); + assert_eq!(stdout.contents(), "last words\n"); + } + + #[tokio::test] + async fn agent_output_sink_forwards_in_order() { + let (log, stdout, stderr) = capture_log(); + let sink = log.agent_output(); + sink.send(AgentStream::Stdout, Bytes::from_static(b"one\n")) + .await; + sink.send(AgentStream::Stderr, Bytes::from_static(b"two\n")) + .await; + sink.send(AgentStream::Stdout, Bytes::from_static(b"three\n")) + .await; + assert_eq!(stdout.wait_for("one\nthree\n").await, "one\nthree\n"); + assert_eq!(stderr.wait_for("two\n").await, "two\n"); + } +} diff --git a/crates/openshell-sandbox/src/delegated.rs b/crates/openshell-sandbox/src/delegated.rs index e294f81ac0..f594fa35f2 100644 --- a/crates/openshell-sandbox/src/delegated.rs +++ b/crates/openshell-sandbox/src/delegated.rs @@ -113,7 +113,11 @@ pub async fn spawn_workload( )?; entrypoint_pid.store(handle.pid(), Ordering::Release); - let main_session = crate::main_session::MainSession::new(handle.take_io(), handle.pid()); + let main_session = crate::main_session::MainSession::new( + handle.take_io(), + handle.pid(), + Some(crate::container_log::ContainerLog::process_agent_output()), + ); let (terminal, signal_lock) = handle.signaling_state(); boundary_runtime .register_process_group(handle.pid(), terminal.clone(), signal_lock.clone()) diff --git a/crates/openshell-sandbox/src/lib.rs b/crates/openshell-sandbox/src/lib.rs index 6c3a9829f9..a8d31fbfe9 100644 --- a/crates/openshell-sandbox/src/lib.rs +++ b/crates/openshell-sandbox/src/lib.rs @@ -9,6 +9,7 @@ pub mod boundary_exec; pub mod boundary_io; mod boundary_server; pub mod child_env; +pub mod container_log; #[cfg(target_os = "linux")] pub(crate) mod delegated; #[cfg(target_os = "linux")] diff --git a/crates/openshell-sandbox/src/main.rs b/crates/openshell-sandbox/src/main.rs index 8ef97be1a3..99f37b202e 100644 --- a/crates/openshell-sandbox/src/main.rs +++ b/crates/openshell-sandbox/src/main.rs @@ -1840,9 +1840,11 @@ fn run_boundary(bootstrap: &Path, log_level: &str) -> Result<()> { EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(log_level)); let _ = tracing_subscriber::registry() .with( - OcsfShorthandLayer::new(std::io::stderr()) - .with_non_ocsf(true) - .with_filter(console_filter), + OcsfShorthandLayer::new( + openshell_sandbox::container_log::ContainerLog::process().launcher_writer(), + ) + .with_non_ocsf(true) + .with_filter(console_filter), ) .try_init(); let (qualification, _) = qualify_runtime()?; diff --git a/crates/openshell-sandbox/src/main_session.rs b/crates/openshell-sandbox/src/main_session.rs index 334c6e7223..ff77f7d464 100644 --- a/crates/openshell-sandbox/src/main_session.rs +++ b/crates/openshell-sandbox/src/main_session.rs @@ -21,6 +21,7 @@ use openshell_isolation_interface::contract::{ BoundaryProcess, BoundarySignal, BoundaryTerminal, ProcessAttachment, }; +use crate::container_log::{AgentOutputSink, AgentStream}; use crate::process::ProcessIo; const OUTPUT_BUFFER_BYTES: usize = 1024 * 1024; @@ -352,10 +353,13 @@ pub struct MainSession { finished: AtomicBool, terminal_attachments: Mutex, terminal_attachments_done: Notify, + container_log: Option, + container_log_released: watch::Sender, } impl MainSession { const REMOTE_OUTPUT_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + const CONTAINER_LOG_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); #[cfg(test)] pub fn inert() -> Arc { let (input, _input_rx) = tokio::sync::mpsc::channel(64); @@ -379,6 +383,8 @@ impl MainSession { expectation: AttachmentExpectation::None, }), terminal_attachments_done: Notify::new(), + container_log: None, + container_log_released: watch::channel(false).0, }) } @@ -387,7 +393,7 @@ impl MainSession { let pty = nix::pty::openpty(None, None).expect("open test PTY"); let slave = std::fs::File::from(pty.slave); ( - Self::new(ProcessIo::Pty(std::fs::File::from(pty.master)), 1), + Self::new(ProcessIo::Pty(std::fs::File::from(pty.master)), 1, None), slave, ) } @@ -403,7 +409,7 @@ impl MainSession { } #[must_use] - pub fn new(io: ProcessIo, pid: u32) -> Arc { + pub fn new(io: ProcessIo, pid: u32, container_log: Option) -> Arc { let terminal = matches!(io, ProcessIo::Pty(_)); let (input, input_rx) = tokio::sync::mpsc::channel(64); let pty_master = match &io { @@ -433,6 +439,8 @@ impl MainSession { expectation: AttachmentExpectation::None, }), terminal_attachments_done: Notify::new(), + container_log: if terminal { None } else { container_log }, + container_log_released: watch::channel(false).0, }); Self::start_io(&session, io, input_rx); session @@ -474,6 +482,8 @@ impl MainSession { expectation: AttachmentExpectation::None, }), terminal_attachments_done: Notify::new(), + container_log: None, + container_log_released: watch::channel(false).0, }); let stdout_session = Arc::clone(&session); tokio::spawn(async move { @@ -608,9 +618,12 @@ impl MainSession { match stdout.read(&mut buffer).await { Ok(0) | Err(_) => break, Ok(read) => { - stdout_session.publish(MainOutput::Stdout(Bytes::copy_from_slice( - &buffer[..read], - ))); + stdout_session + .publish_agent( + AgentStream::Stdout, + Bytes::copy_from_slice(&buffer[..read]), + ) + .await; } } } @@ -623,9 +636,12 @@ impl MainSession { match stderr.read(&mut buffer).await { Ok(0) | Err(_) => break, Ok(read) => { - stderr_session.publish(MainOutput::Stderr(Bytes::copy_from_slice( - &buffer[..read], - ))); + stderr_session + .publish_agent( + AgentStream::Stderr, + Bytes::copy_from_slice(&buffer[..read]), + ) + .await; } } } @@ -652,6 +668,20 @@ impl MainSession { self.output.publish(event); } + async fn publish_agent(&self, stream: AgentStream, data: Bytes) { + self.publish(match stream { + AgentStream::Stdout => MainOutput::Stdout(data.clone()), + AgentStream::Stderr => MainOutput::Stderr(data.clone()), + }); + if let Some(container_log) = &self.container_log { + let mut released = self.container_log_released.subscribe(); + tokio::select! { + () = container_log.send(stream, data) => {} + _ = released.wait_for(|released| *released) => {} + } + } + } + fn reader_finished(&self) { if self.readers_remaining.fetch_sub(1, Ordering::AcqRel) == 1 { self.readers_done.notify_waiters(); @@ -663,6 +693,30 @@ impl MainSession { /// /// Returns whether terminal delivery must complete before shutdown. pub async fn finish(&self, exit_code: i32, attachment_expected: bool) -> bool { + self.finish_with_timeout( + exit_code, + attachment_expected, + Self::CONTAINER_LOG_DRAIN_TIMEOUT, + ) + .await + } + + async fn finish_with_timeout( + &self, + exit_code: i32, + attachment_expected: bool, + container_log_timeout: std::time::Duration, + ) -> bool { + if let Some(container_log) = &self.container_log { + let delivered = tokio::time::timeout(container_log_timeout, async { + self.wait_for_output_readers().await; + container_log.drain().await; + }) + .await; + if delivered.is_err() { + self.container_log_released.send_replace(true); + } + } self.wait_for_output_readers().await; self.complete_finish(exit_code, attachment_expected) } @@ -1477,4 +1531,189 @@ mod tests { assert_eq!(&received[..read], b"client input\n"); session.release_input(owner); } + + #[tokio::test] + async fn terminal_output_is_not_copied_to_the_container_log() { + let (log, stdout, stderr) = crate::container_log::test_support::capture_log(); + let pty = nix::pty::openpty(None, None).expect("open test PTY"); + let mut slave = std::fs::File::from(pty.slave); + let session = MainSession::new( + ProcessIo::Pty(std::fs::File::from(pty.master)), + 1, + Some(log.agent_output()), + ); + let mut output = session.subscribe(); + + slave.write_all(b"agent output").expect("write PTY output"); + let event = tokio::time::timeout(std::time::Duration::from_secs(1), output.recv()) + .await + .expect("PTY output timed out") + .expect("PTY output was retained"); + assert!(matches!(event, MainOutput::Stdout(data) if data == b"agent output"[..])); + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + assert_eq!(stdout.contents(), ""); + assert_eq!(stderr.contents(), ""); + } + + #[tokio::test] + async fn piped_output_is_copied_to_the_matching_container_stream() { + let (log, stdout, stderr) = crate::container_log::test_support::capture_log(); + let mut child = tokio::process::Command::new("sh") + .args(["-c", "printf 'to stdout\\n'; printf 'to stderr\\n' >&2"]) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn() + .expect("spawn test process"); + let io = ProcessIo::Pipes { + stdin: child.stdin.take().expect("stdin"), + stdout: child.stdout.take().expect("stdout"), + stderr: child.stderr.take().expect("stderr"), + }; + let session = MainSession::new(io, child.id().unwrap_or(0), Some(log.agent_output())); + let _ = child.wait().await; + session.wait_for_output_readers().await; + + assert_eq!(stdout.wait_for("to stdout\n").await, "to stdout\n"); + assert_eq!(stderr.wait_for("to stderr\n").await, "to stderr\n"); + } + + #[tokio::test] + async fn finish_waits_for_output_to_reach_the_container_log() { + let gate = crate::container_log::test_support::Gate::default(); + let (log, stdout, _stderr) = crate::container_log::test_support::gated_capture_log(&gate); + let mut child = tokio::process::Command::new("sh") + .args(["-c", "printf 'last words\\n'"]) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .spawn() + .expect("spawn test process"); + let io = ProcessIo::Pipes { + stdin: child.stdin.take().expect("stdin"), + stdout: child.stdout.take().expect("stdout"), + stderr: child.stderr.take().expect("stderr"), + }; + let session = MainSession::new(io, child.id().unwrap_or(0), Some(log.agent_output())); + let _ = child.wait().await; + let finish = tokio::spawn({ + let session = session.clone(); + async move { session.finish(1, false).await } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + assert!(!finish.is_finished()); + + gate.open(); + finish.await.expect("finish task"); + assert_eq!(stdout.contents(), "last words\n"); + } + + fn spawn_piped(script: &str) -> (tokio::process::Child, ProcessIo) { + let mut child = tokio::process::Command::new("sh") + .args(["-c", script]) + .stdin(std::process::Stdio::piped()) + .stdout(std::process::Stdio::piped()) + .stderr(std::process::Stdio::piped()) + .kill_on_drop(true) + .spawn() + .expect("spawn test process"); + let io = ProcessIo::Pipes { + stdin: child.stdin.take().expect("stdin"), + stdout: child.stdout.take().expect("stdout"), + stderr: child.stderr.take().expect("stderr"), + }; + (child, io) + } + + #[tokio::test] + async fn finish_releases_readers_blocked_on_a_stalled_container_log() { + let gate = crate::container_log::test_support::Gate::default(); + let (log, _stdout, _stderr) = crate::container_log::test_support::gated_capture_log(&gate); + let (child, io) = spawn_piped("head -c 1048576 /dev/zero"); + let session = MainSession::new(io, child.id().unwrap_or(0), Some(log.agent_output())); + + let finished = tokio::time::timeout( + std::time::Duration::from_secs(5), + session.finish_with_timeout(1, false, std::time::Duration::from_millis(50)), + ) + .await; + assert!( + finished.is_ok(), + "finish waited on the stalled container log" + ); + gate.open(); + } + + #[tokio::test] + async fn stderr_and_attachments_flow_while_stdout_log_is_stalled() { + let gate = crate::container_log::test_support::Gate::default(); + let (log, _stdout, stderr) = crate::container_log::test_support::gated_capture_log(&gate); + let (child, io) = spawn_piped( + "head -c 1048576 /dev/zero & for i in $(seq 1 100); do echo tick >&2; sleep 0.05; done", + ); + let session = MainSession::new(io, child.id().unwrap_or(0), Some(log.agent_output())); + let mut output = session.subscribe(); + + let replayed = tokio::time::timeout(std::time::Duration::from_secs(3), async { + let mut ticks = 0; + while ticks < 20 { + if let Ok(MainOutput::Stderr(data)) = output.recv().await { + ticks += data.windows(4).filter(|window| window == b"tick").count(); + } + } + }) + .await; + assert!(replayed.is_ok(), "attachment stopped receiving stderr"); + let logged = tokio::time::timeout(std::time::Duration::from_secs(2), async { + while stderr.contents().matches("tick").count() < 20 { + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + } + }) + .await; + assert!(logged.is_ok(), "container stderr stopped receiving output"); + gate.open(); + } + + #[tokio::test] + async fn attachments_receive_output_waiting_on_a_stalled_container_log() { + let gate = crate::container_log::test_support::Gate::default(); + let (log, _stdout, _stderr) = crate::container_log::test_support::gated_capture_log(&gate); + let sink = log.agent_output(); + let (child, io) = spawn_piped("sleep 5"); + let session = MainSession::new(io, child.id().unwrap_or(0), Some(sink.clone())); + let saturate = tokio::spawn(async move { + loop { + sink.send(AgentStream::Stdout, Bytes::from_static(b"x")) + .await; + } + }); + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + + let mut output = session.subscribe(); + let publisher = tokio::spawn({ + let session = session.clone(); + async move { + session + .publish_agent(AgentStream::Stdout, Bytes::from_static(b"marker")) + .await; + } + }); + let received = tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + if let Ok(MainOutput::Stdout(data)) = output.recv().await + && data == b"marker"[..] + { + break; + } + } + }) + .await; + assert!( + received.is_ok(), + "attachment waited on the stalled container log" + ); + saturate.abort(); + publisher.abort(); + gate.open(); + } } diff --git a/docs/observability/accessing-logs.mdx b/docs/observability/accessing-logs.mdx index e5981c9583..f1197cff99 100644 --- a/docs/observability/accessing-logs.mdx +++ b/docs/observability/accessing-logs.mdx @@ -33,6 +33,20 @@ Gateway-originated policy mutations also appear in this stream. When the gateway The TUI dashboard displays sandbox logs in real time. Logs appear in the log panel with the same format as the CLI. +## Main Process Output + +The sandbox's main process writes its stdout and stderr to the sandbox container's stdout and stderr, alongside warnings from the sandbox runtime. Read it with the container tooling for the compute driver: + +```shell +kubectl -n logs -c agent +docker logs +podman logs +``` + +A main process started with a TTY writes only to its attachment; use `openshell sandbox connect` to view that output. Output from `sandbox exec` and SSH sessions does not reach the container log, and `openshell logs` does not include main process output. + +When the container runtime falls behind reading the log, the main process blocks on writes, as with any container. + ## Gateway Log Storage The sandbox pushes logs to the gateway over gRPC in real time. The gateway stores a bounded buffer of recent log lines per sandbox. This buffer is not persisted to disk and is lost when the gateway restarts. diff --git a/skills/debug-openshell-cluster/SKILL.md b/skills/debug-openshell-cluster/SKILL.md index c8b763e399..46d5abc254 100644 --- a/skills/debug-openshell-cluster/SKILL.md +++ b/skills/debug-openshell-cluster/SKILL.md @@ -240,6 +240,10 @@ errors as connectivity, authorization, or lifecycle failures. ### Step 4: Check Docker-Backed Gateways +The sandbox container's log holds the sandbox runtime's warnings and the +main process's stdout and stderr when it runs without a TTY. The supervisor +container's log holds supervisor diagnostics. + ```bash docker info docker ps --filter name=openshell