From c986c0b68d53812f891951ed1be57357ca881bce Mon Sep 17 00:00:00 2001 From: dawn Date: Sun, 23 Aug 2026 04:04:27 +0900 Subject: [PATCH] shuttle,spindle/{agentproto,engines/microvm}: consolidate CI and debug ssh env construction into a single code path Signed-off-by: dawn --- shuttle/src/exec.rs | 94 +++++------ shuttle/src/gen/file_descriptor_set.bin | Bin 315846 -> 315907 bytes .../gen/spindle/agent/v1/spindle.agent.v1.rs | 9 +- shuttle/src/pty.rs | 146 +++++++++++------- spindle/agentproto/gen/agent.pb.go | 79 ++++++---- .../agentproto/spindle/agent/v1/agent.proto | 5 + spindle/engines/microvm/README.md | 16 +- spindle/engines/microvm/debug.go | 11 +- spindle/engines/microvm/debug_test.go | 51 +++++- spindle/engines/microvm/engine.go | 71 +++++---- 10 files changed, 314 insertions(+), 168 deletions(-) diff --git a/shuttle/src/exec.rs b/shuttle/src/exec.rs index 20f7ac43..942f8296 100644 --- a/shuttle/src/exec.rs +++ b/shuttle/src/exec.rs @@ -23,49 +23,14 @@ pub async fn run(id: String, req: v1::ExecStart, out: Sender) { let _ = out.send(msg).await; }; - if req.argv.is_empty() { - send_exit(127, Some("missing argv".to_owned()), false).await; - return; - } - - let user = if req.user.is_empty() { - DEFAULT_USER - } else { - req.user.as_str() - }; - let run_as = match resolve_user(user) { - Ok(run_as) => run_as, + let (spec, run_as) = match request_spec(&req) { + Ok(result) => result, Err(err) => { send_exit(127, Some(err), false).await; return; } }; - let mut env = run_as.login_env(); - let runtime_dir = run_as.runtime_dir(); - match runtime_dir.try_exists() { - Ok(true) => env.push(( - OsString::from("XDG_RUNTIME_DIR"), - runtime_dir.into_os_string(), - )), - Ok(false) => {} - Err(err) => warn!(error = %err, "could not stat XDG_RUNTIME_DIR for workflow user"), - } - - let mut spec = Spec::new(req.argv[0].clone()) - .args(req.argv[1..].iter().cloned()) - .envs(env) - .envs(parse_env(&req.env)) - .run_as(run_as.uid, run_as.gid); - if !req.cwd.is_empty() { - spec = spec.cwd(req.cwd.clone()); - } - let timeout = - (req.timeout_seconds > 0).then(|| Duration::from_secs(u64::from(req.timeout_seconds))); - if let Some(timeout) = timeout { - spec = spec.timeout(timeout); - } - info!( %id, user = %run_as.name, @@ -114,13 +79,50 @@ pub async fn run(id: String, req: v1::ExecStart, out: Sender) { send_exit(exit.exit_code, exit.error, exit.timed_out).await } +pub(crate) fn request_spec(req: &v1::ExecStart) -> Result<(Spec, ResolvedUser), String> { + if req.argv.is_empty() { + return Err("missing argv".to_owned()); + } + let run_as = resolve_user(&req.user)?; + let spec = request_spec_for_user(req, &run_as); + Ok((spec, run_as)) +} + +pub(crate) fn request_spec_for_user(req: &v1::ExecStart, run_as: &ResolvedUser) -> Spec { + let mut env = run_as.login_env(); + let runtime_dir = run_as.runtime_dir(); + match runtime_dir.try_exists() { + Ok(true) => env.push(( + OsString::from("XDG_RUNTIME_DIR"), + runtime_dir.into_os_string(), + )), + Ok(false) => {} + Err(err) => warn!(error = %err, "could not stat XDG_RUNTIME_DIR for workflow user"), + } + + let mut spec = Spec::new(req.argv[0].clone()) + .args(req.argv[1..].iter().cloned()) + .envs(env) + .envs(parse_env(&req.env)) + .run_as(run_as.uid, run_as.gid); + if !req.cwd.is_empty() { + spec = spec.cwd(req.cwd.clone()); + } + let timeout = + (req.timeout_seconds > 0).then(|| Duration::from_secs(u64::from(req.timeout_seconds))); + if let Some(timeout) = timeout { + spec = spec.timeout(timeout); + } + spec +} + #[derive(Clone, Debug)] -struct ResolvedUser { - name: String, - uid: u32, - gid: u32, - home: OsString, - shell: OsString, +pub(crate) struct ResolvedUser { + pub(crate) name: String, + pub(crate) uid: u32, + pub(crate) gid: u32, + pub(crate) home: OsString, + pub(crate) shell: OsString, } impl ResolvedUser { @@ -141,9 +143,13 @@ impl ResolvedUser { fn runtime_dir(&self) -> PathBuf { PathBuf::from(format!("/run/user/{}", self.uid)) } + + pub(crate) fn shell(&self) -> &OsString { + &self.shell + } } -fn resolve_user(spec: &str) -> Result { +pub(crate) fn resolve_user(spec: &str) -> Result { let spec = spec.trim(); if spec.is_empty() { return resolve_user(DEFAULT_USER); diff --git a/shuttle/src/gen/file_descriptor_set.bin b/shuttle/src/gen/file_descriptor_set.bin index efb93df701b5a2be88358e15ad5d2490beb5f2c6..1519daeadb97ae6e318984475e10ca2b7cffd2a3 100644 GIT binary patch delta 120 zcmX@MSh#tKa6=1Y3sVd87M9LUj6b$_ZDOfoWjr^Xa|cVVfDIRST4H8SYD#=@Nos)v z>+}zMSR|$g>|$Xt3oc14Dhc9(NC%gs76>VG@fK&K=H$c|6s6{rrld-+DKKhq2JwIe O10Z7CXYXLqmIMGzjw&kv delta 76 zcmZo(B7AJIa6=1Y3sVd87M9LUjJvmYZDOfoWt=#ja|cVVst^}TX>n?i1giq024@hM YEhNCjoLrtF!3GurF}L5?!J;h*0Kz#I)Bpeg diff --git a/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs b/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs index e3eb309a..a5f954e0 100644 --- a/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs +++ b/shuttle/src/gen/spindle/agent/v1/spindle.agent.v1.rs @@ -120,10 +120,15 @@ pub struct OpenDebugShell { pub term: ::prost::alloc::string::String, #[prost(uint32, tag = "3")] pub rows: u32, - /// debug shells always run as the spindle-workflow user, with that user's - /// login shell, starting in its home dir. nothing here is client-specifiable. #[prost(uint32, tag = "4")] pub cols: u32, + /// the exact process context used to launch the failed workflow step. the + /// guest ignores its argv and timeout when constructing the debug shell. + #[prost(message, optional, tag = "5")] + pub failed_step: ::core::option::Option, + /// the same dependency-environment prelude prepended to the failed step. + #[prost(string, tag = "6")] + pub shell_prelude: ::prost::alloc::string::String, } /// changes meaning based on who sends this: /// guest->host is shell output, host->guest is keyboard input diff --git a/shuttle/src/pty.rs b/shuttle/src/pty.rs index 52ca433d..2ae4d811 100644 --- a/shuttle/src/pty.rs +++ b/shuttle/src/pty.rs @@ -1,16 +1,14 @@ -use std::path::Path; - use crate::command::{self, Spec}; +use crate::exec; use crate::protocol::{self, Message, v1}; -use anyhow::{Context, Result, bail}; +use anyhow::{Context, Result}; use nix::sys::signal::{Signal, kill}; -use nix::unistd::{Pid, User}; +use nix::unistd::Pid; use pty_process::Size; use tokio::io::{AsyncReadExt, BufReader}; use tokio_vsock::{VsockAddr, VsockStream}; use tracing::{info, warn}; -const WF_USER: &str = "spindle-workflow"; const READ_CHUNK: usize = 32 * 1024; pub async fn run(host_cid: u32, open: v1::OpenDebugShell) { @@ -25,16 +23,10 @@ async fn serve(host_cid: u32, open: v1::OpenDebugShell) -> Result<()> { .with_context(|| format!("dial host debug vsock port {}", open.vsock_port))?; info!(port = open.vsock_port, "debug shell connected"); - let user = resolve_user(WF_USER)?; - let rows = clamp_tty_dim(open.rows); let cols = clamp_tty_dim(open.cols); - let spec = Spec::new(&user.shell) - .arg("-l") - .envs(user.env(&open.term)) - .run_as(user.uid, user.gid) - .cwd(Path::new(&user.home).join("repo")); // /workflow/repo + let spec = debug_shell_spec(&open)?; let (mut pty_reader, mut pty_writer, mut child) = command::spawn_pty(spec, rows, cols).context("spawn pty shell")?; @@ -128,55 +120,95 @@ async fn serve(host_cid: u32, open: v1::OpenDebugShell) -> Result<()> { Ok(()) } -struct ResolvedUser { - uid: u32, - gid: u32, - name: String, - home: String, - shell: String, -} - -impl ResolvedUser { - fn env(&self, term: &str) -> Vec<(String, String)> { - let term = if term.is_empty() { - "xterm-256color" - } else { - term - }; - vec![ - ("TERM".to_owned(), term.to_owned()), - ("HOME".to_owned(), self.home.clone()), - ("USER".to_owned(), self.name.clone()), - ("LOGNAME".to_owned(), self.name.clone()), - ("SHELL".to_owned(), self.shell.clone()), - ( - "PATH".to_owned(), - "/run/current-system/sw/bin:/usr/bin:/bin".to_owned(), - ), - ] - } +fn debug_shell_spec(open: &v1::OpenDebugShell) -> Result { + let failed_step = open + .failed_step + .as_ref() + .context("debug shell request is missing failed step context")?; + let user = exec::resolve_user(&failed_step.user).map_err(anyhow::Error::msg)?; + Ok(debug_shell_spec_for_user(open, failed_step, &user)) } -fn resolve_user(name: &str) -> Result { - let user = User::from_name(name) - .with_context(|| format!("lookup user {name:?}"))? - .with_context(|| format!("debug shell user {name:?} not found"))?; - if user.uid.as_raw() == 0 || user.gid.as_raw() == 0 { - bail!("refusing to open a debug shell as privileged user {name:?}"); - } - let shell = user.shell.to_string_lossy().into_owned(); - if shell.is_empty() { - bail!("debug shell user {name:?} has no login shell set in the image"); - } - Ok(ResolvedUser { - uid: user.uid.as_raw(), - gid: user.gid.as_raw(), - name: user.name, - home: user.dir.to_string_lossy().into_owned(), - shell, - }) +fn debug_shell_spec_for_user( + open: &v1::OpenDebugShell, + failed_step: &v1::ExecStart, + user: &exec::ResolvedUser, +) -> Spec { + let shell = user.shell().to_string_lossy().into_owned(); + let command = format!("{}exec \"$0\" -l", open.shell_prelude); + + let mut request = failed_step.clone(); + request.argv = vec![shell.clone(), "-lc".to_owned(), command, shell]; + request.timeout_seconds = 0; + let mut spec = exec::request_spec_for_user(&request, user); + let term = if open.term.is_empty() { + "xterm-256color" + } else { + &open.term + }; + spec.env.push(("TERM".into(), term.into())); + spec } fn clamp_tty_dim(value: u32) -> u16 { value.clamp(1, u16::MAX as u32) as u16 } + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::ffi::{OsStr, OsString}; + use std::path::Path; + + use super::*; + + #[test] + fn debug_shell_uses_failed_step_environment_and_devshell() { + let open = v1::OpenDebugShell { + term: "xterm".to_owned(), + failed_step: Some(v1::ExecStart { + argv: vec!["/run/current-system/sw/bin/bash".to_owned()], + env: vec![ + "HOME=/workspace".to_owned(), + "CI=true".to_owned(), + "SECRET=value".to_owned(), + ], + cwd: "/workspace/repo".to_owned(), + user: "65534".to_owned(), + timeout_seconds: 60, + }), + shell_prelude: ". /run/spindle/devshell-env.sh; ".to_owned(), + ..Default::default() + }; + + let user = exec::ResolvedUser { + name: "spindle-workflow".to_owned(), + uid: 1000, + gid: 1000, + home: OsString::from("/workspace"), + shell: OsString::from("/bin/bash"), + }; + + let failed_step = open.failed_step.as_ref().unwrap(); + let spec = debug_shell_spec_for_user(&open, failed_step, &user); + let env: HashMap<_, _> = spec.env.into_iter().collect(); + + assert_eq!(spec.program, "/bin/bash"); + assert_eq!(spec.cwd, Some(Path::new("/workspace/repo").to_path_buf())); + assert_eq!( + env.get(OsStr::new("XDG_CACHE_HOME")).unwrap(), + "/workspace/.cache" + ); + assert_eq!(env.get(OsStr::new("CI")).unwrap(), "true"); + assert_eq!(env.get(OsStr::new("SECRET")).unwrap(), "value"); + assert_eq!(env.get(OsStr::new("TERM")).unwrap(), "xterm"); + assert_eq!(spec.timeout, None); + assert_eq!(spec.args[0], "-lc"); + assert_eq!(spec.args[2], "/bin/bash"); + assert!( + spec.args[1] + .to_string_lossy() + .contains(". /run/spindle/devshell-env.sh") + ); + } +} diff --git a/spindle/agentproto/gen/agent.pb.go b/spindle/agentproto/gen/agent.pb.go index c0ac3750..d37a049a 100644 --- a/spindle/agentproto/gen/agent.pb.go +++ b/spindle/agentproto/gen/agent.pb.go @@ -781,11 +781,16 @@ func (x *PoweroffResult) GetError() string { } type OpenDebugShell struct { - state protoimpl.MessageState `protogen:"open.v1"` - VsockPort uint32 `protobuf:"varint,1,opt,name=vsock_port,json=vsockPort,proto3" json:"vsock_port,omitempty"` - Term string `protobuf:"bytes,2,opt,name=term,proto3" json:"term,omitempty"` - Rows uint32 `protobuf:"varint,3,opt,name=rows,proto3" json:"rows,omitempty"` - Cols uint32 `protobuf:"varint,4,opt,name=cols,proto3" json:"cols,omitempty"` + state protoimpl.MessageState `protogen:"open.v1"` + VsockPort uint32 `protobuf:"varint,1,opt,name=vsock_port,json=vsockPort,proto3" json:"vsock_port,omitempty"` + Term string `protobuf:"bytes,2,opt,name=term,proto3" json:"term,omitempty"` + Rows uint32 `protobuf:"varint,3,opt,name=rows,proto3" json:"rows,omitempty"` + Cols uint32 `protobuf:"varint,4,opt,name=cols,proto3" json:"cols,omitempty"` + // the exact process context used to launch the failed workflow step. the + // guest ignores its argv and timeout when constructing the debug shell. + FailedStep *ExecStart `protobuf:"bytes,5,opt,name=failed_step,json=failedStep,proto3" json:"failed_step,omitempty"` + // the same dependency-environment prelude prepended to the failed step. + ShellPrelude string `protobuf:"bytes,6,opt,name=shell_prelude,json=shellPrelude,proto3" json:"shell_prelude,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -848,6 +853,20 @@ func (x *OpenDebugShell) GetCols() uint32 { return 0 } +func (x *OpenDebugShell) GetFailedStep() *ExecStart { + if x != nil { + return x.FailedStep + } + return nil +} + +func (x *OpenDebugShell) GetShellPrelude() string { + if x != nil { + return x.ShellPrelude + } + return "" +} + // changes meaning based on who sends this: // guest->host is shell output, host->guest is keyboard input type PtyData struct { @@ -1180,13 +1199,16 @@ const file_spindle_agent_v1_agent_proto_rawDesc = "" + "\n" + "\bPoweroff\"&\n" + "\x0ePoweroffResult\x12\x14\n" + - "\x05error\x18\x01 \x01(\tR\x05error\"k\n" + + "\x05error\x18\x01 \x01(\tR\x05error\"\xce\x01\n" + "\x0eOpenDebugShell\x12\x1d\n" + "\n" + "vsock_port\x18\x01 \x01(\rR\tvsockPort\x12\x12\n" + "\x04term\x18\x02 \x01(\tR\x04term\x12\x12\n" + "\x04rows\x18\x03 \x01(\rR\x04rows\x12\x12\n" + - "\x04cols\x18\x04 \x01(\rR\x04cols\"\x1d\n" + + "\x04cols\x18\x04 \x01(\rR\x04cols\x12<\n" + + "\vfailed_step\x18\x05 \x01(\v2\x1b.spindle.agent.v1.ExecStartR\n" + + "failedStep\x12#\n" + + "\rshell_prelude\x18\x06 \x01(\tR\fshellPrelude\"\x1d\n" + "\aPtyData\x12\x12\n" + "\x04data\x18\x01 \x01(\fR\x04data\"3\n" + "\tPtyResize\x12\x12\n" + @@ -1270,27 +1292,28 @@ var file_spindle_agent_v1_agent_proto_goTypes = []any{ (*Message)(nil), // 16: spindle.agent.v1.Message } var file_spindle_agent_v1_agent_proto_depIdxs = []int32{ - 0, // 0: spindle.agent.v1.Message.hello:type_name -> spindle.agent.v1.Hello - 1, // 1: spindle.agent.v1.Message.init:type_name -> spindle.agent.v1.Init - 2, // 2: spindle.agent.v1.Message.exec_start:type_name -> spindle.agent.v1.ExecStart - 3, // 3: spindle.agent.v1.Message.exec_stdout:type_name -> spindle.agent.v1.ExecStdout - 4, // 4: spindle.agent.v1.Message.exec_stderr:type_name -> spindle.agent.v1.ExecStderr - 5, // 5: spindle.agent.v1.Message.exec_exit:type_name -> spindle.agent.v1.ExecExit - 6, // 6: spindle.agent.v1.Message.activate_config:type_name -> spindle.agent.v1.ActivateConfig - 7, // 7: spindle.agent.v1.Message.activate_config_result:type_name -> spindle.agent.v1.ActivateConfigResult - 8, // 8: spindle.agent.v1.Message.built_paths:type_name -> spindle.agent.v1.BuiltPaths - 9, // 9: spindle.agent.v1.Message.cache_drain:type_name -> spindle.agent.v1.CacheDrain - 10, // 10: spindle.agent.v1.Message.cache_drain_result:type_name -> spindle.agent.v1.CacheDrainResult - 11, // 11: spindle.agent.v1.Message.poweroff:type_name -> spindle.agent.v1.Poweroff - 12, // 12: spindle.agent.v1.Message.poweroff_result:type_name -> spindle.agent.v1.PoweroffResult - 13, // 13: spindle.agent.v1.Message.open_debug_shell:type_name -> spindle.agent.v1.OpenDebugShell - 14, // 14: spindle.agent.v1.Message.pty_data:type_name -> spindle.agent.v1.PtyData - 15, // 15: spindle.agent.v1.Message.pty_resize:type_name -> spindle.agent.v1.PtyResize - 16, // [16:16] is the sub-list for method output_type - 16, // [16:16] is the sub-list for method input_type - 16, // [16:16] is the sub-list for extension type_name - 16, // [16:16] is the sub-list for extension extendee - 0, // [0:16] is the sub-list for field type_name + 2, // 0: spindle.agent.v1.OpenDebugShell.failed_step:type_name -> spindle.agent.v1.ExecStart + 0, // 1: spindle.agent.v1.Message.hello:type_name -> spindle.agent.v1.Hello + 1, // 2: spindle.agent.v1.Message.init:type_name -> spindle.agent.v1.Init + 2, // 3: spindle.agent.v1.Message.exec_start:type_name -> spindle.agent.v1.ExecStart + 3, // 4: spindle.agent.v1.Message.exec_stdout:type_name -> spindle.agent.v1.ExecStdout + 4, // 5: spindle.agent.v1.Message.exec_stderr:type_name -> spindle.agent.v1.ExecStderr + 5, // 6: spindle.agent.v1.Message.exec_exit:type_name -> spindle.agent.v1.ExecExit + 6, // 7: spindle.agent.v1.Message.activate_config:type_name -> spindle.agent.v1.ActivateConfig + 7, // 8: spindle.agent.v1.Message.activate_config_result:type_name -> spindle.agent.v1.ActivateConfigResult + 8, // 9: spindle.agent.v1.Message.built_paths:type_name -> spindle.agent.v1.BuiltPaths + 9, // 10: spindle.agent.v1.Message.cache_drain:type_name -> spindle.agent.v1.CacheDrain + 10, // 11: spindle.agent.v1.Message.cache_drain_result:type_name -> spindle.agent.v1.CacheDrainResult + 11, // 12: spindle.agent.v1.Message.poweroff:type_name -> spindle.agent.v1.Poweroff + 12, // 13: spindle.agent.v1.Message.poweroff_result:type_name -> spindle.agent.v1.PoweroffResult + 13, // 14: spindle.agent.v1.Message.open_debug_shell:type_name -> spindle.agent.v1.OpenDebugShell + 14, // 15: spindle.agent.v1.Message.pty_data:type_name -> spindle.agent.v1.PtyData + 15, // 16: spindle.agent.v1.Message.pty_resize:type_name -> spindle.agent.v1.PtyResize + 17, // [17:17] is the sub-list for method output_type + 17, // [17:17] is the sub-list for method input_type + 17, // [17:17] is the sub-list for extension type_name + 17, // [17:17] is the sub-list for extension extendee + 0, // [0:17] is the sub-list for field type_name } func init() { file_spindle_agent_v1_agent_proto_init() } diff --git a/spindle/agentproto/spindle/agent/v1/agent.proto b/spindle/agentproto/spindle/agent/v1/agent.proto index 2c25e2aa..0880df8b 100644 --- a/spindle/agentproto/spindle/agent/v1/agent.proto +++ b/spindle/agentproto/spindle/agent/v1/agent.proto @@ -87,6 +87,11 @@ message OpenDebugShell { string term = 2; uint32 rows = 3; uint32 cols = 4; + // the exact process context used to launch the failed workflow step. the + // guest ignores its argv and timeout when constructing the debug shell. + ExecStart failed_step = 5; + // the same dependency-environment prelude prepended to the failed step. + string shell_prelude = 6; } // changes meaning based on who sends this: diff --git a/spindle/engines/microvm/README.md b/spindle/engines/microvm/README.md index 6c516818..a4a8b374 100644 --- a/spindle/engines/microvm/README.md +++ b/spindle/engines/microvm/README.md @@ -268,8 +268,20 @@ The shell is deliberately not configurable from either end. It always: - runs as the `spindle-workflow` user (the ssh username selects the *job*, not a unix user), - uses that user's login shell from the image's passwd, launched as a login - shell (`-l`), and -- starts in the dir where the repo was cloned to. + shell (`-l`), +- starts in the dir where the repo was cloned to, +- receives the failed step's initial workflow/step environment and unlocked + secrets, +- uses the workflow user's cache directory instead of inheriting shuttle's + service cache, +- and sources the activated dependency devshell before starting the interactive + login shell. + +The executor retains the actual `ExecStart` sent for the failed step instead of +reconstructing its environment later. Shuttle uses the same process-spec builder +for CI execs and debug PTYs; the debug path replaces only the argv and timeout, +then adds the SSH terminal. The dependency prelude is also generated once and +shared by both paths. The only things the client influences are the terminal type and window size (forwarded from the ssh pty request, and on resize). This relies on the image diff --git a/spindle/engines/microvm/debug.go b/spindle/engines/microvm/debug.go index d8c92384..28d5ddda 100644 --- a/spindle/engines/microvm/debug.go +++ b/spindle/engines/microvm/debug.go @@ -26,6 +26,7 @@ const debugAcceptTimeout = 15 * time.Second type debugTarget struct { cid uint32 agent *AgentSession + failedStep *agentv1.ExecStart knot string repoDid string maxAliveAt time.Time @@ -110,10 +111,12 @@ func (e *Engine) OpenDebugSession(ctx context.Context, jobID, term string, rows, filtered := &cidFilteredVsockListener{Listener: ln, cid: target.cid, logger: e.l} if err := target.agent.OpenDebugShell(&agentv1.OpenDebugShell{ - VsockPort: port, - Term: term, - Rows: clampDim(rows), - Cols: clampDim(cols), + VsockPort: port, + Term: term, + Rows: clampDim(rows), + Cols: clampDim(cols), + FailedStep: target.failedStep, + ShellPrelude: depsSourcePrelude(), }); err != nil { _ = ln.Close() return nil, fmt.Errorf("ask guest to open debug shell: %w", err) diff --git a/spindle/engines/microvm/debug_test.go b/spindle/engines/microvm/debug_test.go index c5a2d0e5..c7f9c201 100644 --- a/spindle/engines/microvm/debug_test.go +++ b/spindle/engines/microvm/debug_test.go @@ -2,7 +2,13 @@ package microvm -import "testing" +import ( + "strings" + "testing" + + "tangled.org/core/spindle/models" + "tangled.org/core/spindle/secrets" +) func TestDebugSSHCommandUsesJumpHost(t *testing.T) { got := debugSSHCommand("0.0.0.0:2224", "executor-a.internal", "executor-a", "spindle.example", "job-1") @@ -19,3 +25,46 @@ func TestDebugSSHCommandUsesExecutorHostWithoutJump(t *testing.T) { t.Fatalf("debug ssh command = %q, want %q", got, want) } } + +func TestStepEnvironmentPreservesFailedStepContext(t *testing.T) { + workflow := &models.Workflow{ + Environment: map[string]string{ + "CI": "true", + "OVERRIDE": "workflow", + }, + } + step := Step{ + environment: map[string]string{ + "OVERRIDE": "step", + "STEP_ONLY": "yes", + }, + } + unlocked := []secrets.UnlockedSecret{{ + Key: "SECRET", + Value: "value", + }} + + got := make(map[string]string) + for _, value := range stepEnvironment(workflow, step, unlocked) { + key, value, ok := strings.Cut(value, "=") + if ok { + got[key] = value + } + } + + want := map[string]string{ + "HOME": "/workspace", + "LOGNAME": guestWorkflowUser, + "PATH": guestBasePATH, + "USER": guestWorkflowUser, + "CI": "true", + "OVERRIDE": "step", + "STEP_ONLY": "yes", + "SECRET": "value", + } + for key, value := range want { + if got[key] != value { + t.Errorf("%s = %q, want %q", key, got[key], value) + } + } +} diff --git a/spindle/engines/microvm/engine.go b/spindle/engines/microvm/engine.go index 60a79744..95937b6a 100644 --- a/spindle/engines/microvm/engine.go +++ b/spindle/engines/microvm/engine.go @@ -362,30 +362,19 @@ func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *m return nil } -func applyDepsSource(command string) string { +func depsSourcePrelude() string { return fmt.Sprintf( // check if it exists because not all images have this - `if [ -f %s ]; then . %s; export PATH="$PATH:%s"; fi; %s`, - guestDevShellEnvPath, guestDevShellEnvPath, guestBasePATH, command, + `if [ -f %s ]; then . %s; export PATH="$PATH:%s"; fi; `, + guestDevShellEnvPath, guestDevShellEnvPath, guestBasePATH, ) } -func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { - state, ok := w.Data.(*workflowState) - if !ok || state == nil || state.Agent == nil { - return fmt.Errorf("microVM workflow is not connected to agent") - } - - stderr := wfLogger.DataWriter(idx, "stderr") - - execCtx, vmExited, cancelWatch := watchVMExit(ctx, state.VM) - defer cancelWatch() +func applyDepsSource(command string) string { + return depsSourcePrelude() + command +} - step := w.Steps[idx] - if s, ok := step.(Step); ok && s.action == activationStepAction { - err := e.activateConfig(execCtx, wid, state, s, wfLogger.DataWriter(idx, "stdout")) - return e.classifyStepError(ctx, wid, step, state, stderr, vmExited, "Failed to activate config", err) - } +func stepEnvironment(w *models.Workflow, step models.Step, unlocked []secrets.UnlockedSecret) []string { env := []string{ "HOME=/workspace", "LOGNAME=" + guestWorkflowUser, @@ -395,27 +384,48 @@ func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.W for k, v := range w.Environment { env = append(env, k+"="+v) } - for _, s := range secrets { - env = append(env, s.Key+"="+s.Value) + for _, secret := range unlocked { + env = append(env, secret.Key+"="+secret.Value) } if s, ok := step.(Step); ok { for k, v := range s.environment { env = append(env, k+"="+v) } } + return env +} + +func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, secrets []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { + state, ok := w.Data.(*workflowState) + if !ok || state == nil || state.Agent == nil { + return fmt.Errorf("microVM workflow is not connected to agent") + } + + stderr := wfLogger.DataWriter(idx, "stderr") + + execCtx, vmExited, cancelWatch := watchVMExit(ctx, state.VM) + defer cancelWatch() + + step := w.Steps[idx] + if s, ok := step.(Step); ok && s.action == activationStepAction { + err := e.activateConfig(execCtx, wid, state, s, wfLogger.DataWriter(idx, "stdout")) + return e.classifyStepError(ctx, wid, step, state, stderr, vmExited, "Failed to activate config", err) + } + env := stepEnvironment(w, step, secrets) stdout := wfLogger.DataWriter(idx, "stdout") + execStart := &agentv1.ExecStart{ + Argv: []string{state.ImageSpec.Shell, "-lc", applyDepsSource(step.Command())}, + Env: env, + Cwd: guestWorkDir, + User: guestWorkflowUser, + // timeout not set here, Exec will fill it + } exitCode, err := state.Agent.Exec(execCtx, AgentExec{ - ID: fmt.Sprintf("%s-%d", wid.String(), idx), - ExecStart: &agentv1.ExecStart{ - Argv: []string{state.ImageSpec.Shell, "-lc", applyDepsSource(step.Command())}, - Env: env, - Cwd: guestWorkDir, - User: guestWorkflowUser, - // timeout not set here, Exec will fill it - }, - Stdout: stdout, - Stderr: stderr, + ID: fmt.Sprintf("%s-%d", wid.String(), idx), + ExecStart: execStart, + Stdout: stdout, + Stderr: stderr, }) if err != nil { return e.classifyStepError(ctx, wid, step, state, stderr, vmExited, "User step error", err) @@ -428,6 +438,7 @@ func (e *Engine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.W cid: state.CID, agent: state.Agent, knot: w.Environment["TANGLED_REPO_KNOT"], + failedStep: execStart, repoDid: w.Environment["TANGLED_REPO_REPO_DID"], maxAliveAt: state.StartedAt.Add(e.WorkflowTimeout()), connected: make(chan struct{}), -- 2.51.2