use crate::core::agent_context::AgentContextView;
use crate::core::ctx::Ctx;
use crate::core::engine::Engine;
use crate::types::{CapToken, StepId, WorkerId};
use crate::worker::adapter::{SpawnError, SpawnerAdapter, WorkerError, WorkerResult};
use crate::worker::output::{ContentRef, OutputEvent};
use crate::worker::{Worker, WorkerJoinHandler};
use async_trait::async_trait;
use mlua_swarm_schema::SubprocessOutput;
use serde_json::Value;
use std::collections::BTreeMap;
use std::process::Stdio;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::process::Command;
use tokio::sync::oneshot;
use tokio_util::sync::CancellationToken;
pub const EMBED_PLACEHOLDERS: [&str; 8] = [
"system",
"system_file",
"prompt",
"model",
"tools_csv",
"work_dir",
"task_id",
"attempt",
];
#[derive(Debug, Clone, Default)]
pub struct EmbedTemplate {
pub argv: Vec<String>,
pub stdin: Option<String>,
pub env: BTreeMap<String, String>,
pub cwd: Option<String>,
pub output: Option<SubprocessOutput>,
pub system_prompt: Option<String>,
pub model: Option<String>,
pub tools_csv: String,
}
#[derive(Debug, Clone)]
pub enum StreamMode {
NdjsonLines,
SseEvents,
LengthPrefixed,
}
#[derive(Debug)]
pub struct ProcessSpawner {
pub program: String,
pub args: Vec<String>,
pub use_stdin: bool,
pub stream_mode: Option<StreamMode>,
pub embed: Option<EmbedTemplate>,
}
impl ProcessSpawner {
pub fn new(program: impl Into<String>) -> Self {
Self {
program: program.into(),
args: Vec::new(),
use_stdin: true,
stream_mode: None,
embed: None,
}
}
pub fn arg(mut self, a: impl Into<String>) -> Self {
self.args.push(a.into());
self
}
pub fn args(mut self, args: impl IntoIterator<Item = impl Into<String>>) -> Self {
self.args.extend(args.into_iter().map(|a| a.into()));
self
}
pub fn use_stdin(mut self, v: bool) -> Self {
self.use_stdin = v;
self
}
pub fn stream_mode(mut self, mode: StreamMode) -> Self {
self.stream_mode = Some(mode);
self
}
pub fn plain(mut self) -> Self {
self.stream_mode = None;
self
}
pub fn ndjson(mut self, v: bool) -> Self {
self.stream_mode = if v {
Some(StreamMode::NdjsonLines)
} else {
None
};
self
}
pub fn run(cmd: impl Into<String>) -> Self {
Self {
program: "sh".into(),
args: vec!["-c".into(), cmd.into()],
use_stdin: true,
stream_mode: None,
embed: None,
}
}
pub fn cmd(program: impl Into<String>) -> Self {
Self {
program: program.into(),
args: Vec::new(),
use_stdin: true,
stream_mode: None,
embed: None,
}
}
pub fn embed(mut self, template: EmbedTemplate) -> Self {
self.embed = Some(template);
self
}
}
#[async_trait]
impl SpawnerAdapter for ProcessSpawner {
async fn spawn(
&self,
engine: &Engine,
ctx: &Ctx,
task_id: StepId,
attempt: u32,
token: CapToken,
) -> Result<Box<dyn Worker>, SpawnError> {
if let Some(embed) = &self.embed {
return self
.spawn_embed(embed, engine, ctx, task_id, attempt, token)
.await;
}
let directive = engine
.fetch_prompt(&token, &task_id)
.await
.map_err(|e| SpawnError::Internal(format!("fetch_prompt: {e}")))?;
let directive = crate::core::engine::render_directive_to_string(&directive);
let mut cmd = Command::new(&self.program);
cmd.args(&self.args)
.env("MSE_TOKEN_AGENT_ID", &token.agent_id)
.env("MSE_TOKEN_NONCE", &token.nonce)
.env("MSE_TASK_ID", task_id.as_str())
.env("MSE_ATTEMPT", attempt.to_string())
.env("MSE_CTX_AGENT", &ctx.agent)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
if !self.use_stdin {
cmd.arg(&directive);
}
let mut child = cmd
.spawn()
.map_err(|e| SpawnError::Internal(format!("spawn failed: {e}")))?;
if self.use_stdin {
if let Some(mut stdin) = child.stdin.take() {
stdin
.write_all(directive.as_bytes())
.await
.map_err(|e| SpawnError::Internal(format!("stdin write: {e}")))?;
drop(stdin); }
}
let cancel = CancellationToken::new();
let cancel_inner = cancel.clone();
let worker_id = WorkerId::new();
tracing::debug!(worker_id = %worker_id, step_id = %task_id, "worker spawned (subprocess)");
let (tx, rx) = oneshot::channel();
let engine_for_emit = engine.clone();
let token_for_emit = token.clone();
let task_id_for_emit = task_id.clone();
let stream_mode = self.stream_mode.clone();
tokio::spawn(async move {
let result: Result<WorkerResult, WorkerError> = if let Some(mode) = stream_mode {
run_streaming_mode(
mode,
child,
&engine_for_emit,
&token_for_emit,
&task_id_for_emit,
attempt,
cancel_inner,
)
.await
} else {
let result = tokio::select! {
output = child.wait_with_output() => {
match output {
Ok(out) => {
let stdout = String::from_utf8_lossy(&out.stdout).to_string();
let value: Value = serde_json::from_str(stdout.trim())
.unwrap_or_else(|_| serde_json::json!({
"raw": stdout.trim_end(),
"stderr": String::from_utf8_lossy(&out.stderr).to_string(),
}));
Ok(WorkerResult { value, ok: out.status.success() })
}
Err(e) => Err(WorkerError::Failed(format!("wait_with_output: {e}"))),
}
}
_ = cancel_inner.cancelled() => Err(WorkerError::Cancelled),
};
if let Ok(wr) = &result {
let ev = OutputEvent::Final {
content: ContentRef::Inline {
value: wr.value.clone(),
},
ok: wr.ok,
};
let _ = engine_for_emit
.submit_output(&token_for_emit, &task_id_for_emit, attempt, ev)
.await;
}
result
};
let signal: Result<(), WorkerError> = result.map(|_| ());
let _ = tx.send(signal);
});
Ok(Box::new(ProcessWorker {
handler: WorkerJoinHandler {
worker_id,
cancel,
completion: rx,
},
}))
}
}
struct EmbedVars {
system: String,
system_file: Option<String>,
prompt: String,
model: Option<String>,
tools_csv: String,
work_dir: Option<String>,
task_id: String,
attempt: String,
}
impl EmbedVars {
fn lookup(&self, token: &str) -> Result<Option<&str>, SpawnError> {
let (value, cause): (Option<&str>, &str) = match token {
"system" => (Some(self.system.as_str()), ""),
"system_file" => (
self.system_file.as_deref(),
"no system prompt was baked for this attempt (agent has no profile.system_prompt?)",
),
"prompt" => (Some(self.prompt.as_str()), ""),
"model" => (
self.model.as_deref(),
"no model declared (set profile.model or Runner::Subprocess overrides.model)",
),
"tools_csv" => (Some(self.tools_csv.as_str()), ""),
"work_dir" => (
self.work_dir.as_deref(),
"no work_dir/project_root in the agent context view and no overrides.cwd",
),
"task_id" => (Some(self.task_id.as_str()), ""),
"attempt" => (Some(self.attempt.as_str()), ""),
_ => return Ok(None),
};
match value {
Some(v) => Ok(Some(v)),
None => Err(SpawnError::Internal(format!(
"placeholder {{{token}}}: {cause}"
))),
}
}
fn render(&self, tmpl: &str) -> Result<String, SpawnError> {
let mut out = String::with_capacity(tmpl.len());
let mut rest = tmpl;
while let Some(start) = rest.find('{') {
out.push_str(&rest[..start]);
let after = &rest[start + 1..];
let Some(end) = after.find('}') else {
out.push_str(&rest[start..]);
return Ok(out);
};
let token = &after[..end];
match self.lookup(token)? {
Some(value) => {
out.push_str(value);
rest = &after[end + 1..];
}
None => {
out.push('{');
rest = after;
}
}
}
out.push_str(rest);
Ok(out)
}
}
fn embed_references(embed: &EmbedTemplate, token: &str) -> bool {
embed.argv.iter().any(|a| a.contains(token))
|| embed.stdin.as_deref().is_some_and(|s| s.contains(token))
|| embed.env.values().any(|v| v.contains(token))
|| embed.cwd.as_deref().is_some_and(|c| c.contains(token))
}
fn normalize_plain_output(
out: &std::process::Output,
decl: Option<&SubprocessOutput>,
) -> WorkerResult {
let stdout = String::from_utf8_lossy(&out.stdout).to_string();
let stderr = || String::from_utf8_lossy(&out.stderr).to_string();
let exit_ok = out.status.success();
let Some(decl) = decl else {
let value: Value = serde_json::from_str(stdout.trim()).unwrap_or_else(|_| {
serde_json::json!({
"raw": stdout.trim_end(),
"stderr": stderr(),
})
});
return WorkerResult { value, ok: exit_ok };
};
let parsed: Result<Value, _> = serde_json::from_str(stdout.trim());
let parsed = match parsed {
Ok(v) => v,
Err(e) => {
if decl.format.as_deref() == Some("json") {
return WorkerResult {
value: serde_json::json!({
"raw": stdout.trim_end(),
"stderr": stderr(),
"parse_error": e.to_string(),
}),
ok: false,
};
}
return WorkerResult {
value: serde_json::json!({
"raw": stdout.trim_end(),
"stderr": stderr(),
}),
ok: exit_ok,
};
}
};
let value = match decl.result_ptr.as_deref() {
Some(ptr) => match parsed.pointer(ptr) {
Some(v) => v.clone(),
None => {
return WorkerResult {
value: serde_json::json!({
"error": format!("result_ptr '{ptr}' not found in stdout JSON"),
"raw": parsed,
"stderr": stderr(),
}),
ok: false,
}
}
},
None => parsed.clone(),
};
let ok = match decl.ok_from.as_deref() {
None | Some("exit_code") => exit_ok,
Some(ptr) => parsed.pointer(ptr) == Some(&Value::Bool(true)),
};
WorkerResult { value, ok }
}
impl ProcessSpawner {
async fn spawn_embed(
&self,
embed: &EmbedTemplate,
engine: &Engine,
ctx: &Ctx,
task_id: StepId,
attempt: u32,
token: CapToken,
) -> Result<Box<dyn Worker>, SpawnError> {
let prompt_value = engine
.fetch_prompt(&token, &task_id)
.await
.map_err(|e| SpawnError::Internal(format!("fetch_prompt: {e}")))?;
let system = match embed.system_prompt.as_deref() {
Some(tmpl) => {
let slots = crate::operator::render::slots_from_prompt(&prompt_value);
let rendered = crate::operator::render::render_system(tmpl, &slots)
.map_err(|e| SpawnError::Internal(format!("render system_prompt: {e}")))?;
Some(rendered)
}
None => None,
};
engine
.bake_worker_system_prompt(&task_id, attempt, system.clone())
.await
.map_err(|e| SpawnError::Internal(format!("bake system_prompt: {e}")))?;
let payload = engine
.fetch_worker_payload(&token, &task_id)
.await
.map_err(|e| SpawnError::Internal(format!("fetch_worker_payload: {e}")))?;
let system_file = if embed_references(embed, "{system_file}") {
match engine
.materialize_system_file(&task_id, attempt)
.await
.map_err(|e| SpawnError::Internal(format!("materialize system_file: {e}")))?
{
Some(path) => Some(path.display().to_string()),
None => {
return Err(SpawnError::Internal(
"placeholder {system_file}: no system prompt was baked for this \
attempt (agent has no profile.system_prompt?)"
.into(),
))
}
}
} else {
None
};
let view = AgentContextView::materialized_or_from_ctx(ctx);
let work_dir = view.work_dir.clone().or_else(|| view.project_root.clone());
let vars = EmbedVars {
system: system.unwrap_or_default(),
system_file,
prompt: payload.prompt.clone(),
model: embed.model.clone(),
tools_csv: embed.tools_csv.clone(),
work_dir,
task_id: task_id.to_string(),
attempt: attempt.to_string(),
};
if embed.argv.is_empty() {
return Err(SpawnError::Internal(
"embed template: argv must not be empty".into(),
));
}
let mut rendered_argv = Vec::with_capacity(embed.argv.len());
for a in &embed.argv {
rendered_argv.push(vars.render(a)?);
}
let mut rendered_env = Vec::with_capacity(embed.env.len());
for (k, v) in &embed.env {
rendered_env.push((k.clone(), vars.render(v)?));
}
let rendered_cwd = match embed.cwd.as_deref() {
Some(c) => Some(vars.render(c)?),
None => None,
};
let rendered_stdin = match embed.stdin.as_deref() {
Some(s) => Some(vars.render(s)?),
None => None,
};
let mut cmd = Command::new(&rendered_argv[0]);
cmd.args(&rendered_argv[1..])
.env("MSE_TOKEN_AGENT_ID", &token.agent_id)
.env("MSE_TOKEN_NONCE", &token.nonce)
.env("MSE_TASK_ID", task_id.as_str())
.env("MSE_ATTEMPT", attempt.to_string())
.env("MSE_CTX_AGENT", &ctx.agent)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
for (k, v) in &rendered_env {
cmd.env(k, v);
}
if let Some(cwd) = &rendered_cwd {
cmd.current_dir(cwd);
}
let mut child = cmd
.spawn()
.map_err(|e| SpawnError::Internal(format!("spawn failed: {e}")))?;
if let Some(stdin_body) = rendered_stdin {
if let Some(mut stdin) = child.stdin.take() {
stdin
.write_all(stdin_body.as_bytes())
.await
.map_err(|e| SpawnError::Internal(format!("stdin write: {e}")))?;
drop(stdin); }
} else {
drop(child.stdin.take());
}
let cancel = CancellationToken::new();
let cancel_inner = cancel.clone();
let worker_id = WorkerId::new();
tracing::debug!(worker_id = %worker_id, step_id = %task_id, "worker spawned (subprocess embed)");
let (tx, rx) = oneshot::channel();
let engine_for_emit = engine.clone();
let token_for_emit = token.clone();
let task_id_for_emit = task_id.clone();
let stream_mode = self.stream_mode.clone();
let output_decl = embed.output.clone();
tokio::spawn(async move {
let result: Result<WorkerResult, WorkerError> = if let Some(mode) = stream_mode {
run_streaming_mode(
mode,
child,
&engine_for_emit,
&token_for_emit,
&task_id_for_emit,
attempt,
cancel_inner,
)
.await
} else {
let result = tokio::select! {
output = child.wait_with_output() => {
match output {
Ok(out) => Ok(normalize_plain_output(&out, output_decl.as_ref())),
Err(e) => Err(WorkerError::Failed(format!("wait_with_output: {e}"))),
}
}
_ = cancel_inner.cancelled() => Err(WorkerError::Cancelled),
};
if let Ok(wr) = &result {
let ev = OutputEvent::Final {
content: ContentRef::Inline {
value: wr.value.clone(),
},
ok: wr.ok,
};
let _ = engine_for_emit
.submit_output(&token_for_emit, &task_id_for_emit, attempt, ev)
.await;
}
result
};
let signal: Result<(), WorkerError> = result.map(|_| ());
let _ = tx.send(signal);
});
Ok(Box::new(ProcessWorker {
handler: WorkerJoinHandler {
worker_id,
cancel,
completion: rx,
},
}))
}
}
pub struct ProcessWorker {
pub handler: WorkerJoinHandler,
}
#[async_trait]
impl Worker for ProcessWorker {
fn id(&self) -> &WorkerId {
&self.handler.worker_id
}
fn cancel_token(&self) -> CancellationToken {
self.handler.cancel.clone()
}
async fn join(self: Box<Self>) -> Result<(), WorkerError> {
self.handler.await_completion().await
}
}
async fn run_streaming_mode(
mode: StreamMode,
mut child: tokio::process::Child,
engine: &Engine,
token: &CapToken,
task_id: &StepId,
attempt: u32,
cancel: CancellationToken,
) -> Result<WorkerResult, WorkerError> {
let stdout = child
.stdout
.take()
.ok_or_else(|| WorkerError::Failed("streaming: stdout pipe missing".into()))?;
let last_final = match mode {
StreamMode::NdjsonLines => {
read_ndjson(stdout, engine, token, task_id, attempt, cancel.clone()).await?
}
StreamMode::SseEvents => {
read_sse(stdout, engine, token, task_id, attempt, cancel.clone()).await?
}
StreamMode::LengthPrefixed => {
read_length_prefixed(stdout, engine, token, task_id, attempt, cancel.clone()).await?
}
};
let status = child
.wait()
.await
.map_err(|e| WorkerError::Failed(format!("streaming wait: {e}")))?;
match last_final {
Some((value, ok)) => Ok(WorkerResult {
value,
ok: ok && status.success(),
}),
None => {
let value = serde_json::json!({
"raw": "",
"note": "streaming mode: no Final event received",
"exit_success": status.success(),
});
let _ = engine
.submit_output(
token,
task_id,
attempt,
OutputEvent::Final {
content: ContentRef::Inline {
value: value.clone(),
},
ok: false,
},
)
.await;
Ok(WorkerResult { value, ok: false })
}
}
}
async fn forward_event(
engine: &Engine,
token: &CapToken,
task_id: &StepId,
attempt: u32,
ev: OutputEvent,
last_final: &mut Option<(Value, bool)>,
) {
if let OutputEvent::Final { content, ok } = &ev {
let value = match content {
ContentRef::Inline { value } => value.clone(),
ContentRef::FileRef {
path,
mime,
size_hint,
} => serde_json::json!({
"file_ref": path.to_string_lossy(),
"mime": mime,
"size_hint": size_hint,
}),
};
*last_final = Some((value, *ok));
}
let _ = engine.submit_output(token, task_id, attempt, ev).await;
}
async fn read_ndjson(
stdout: tokio::process::ChildStdout,
engine: &Engine,
token: &CapToken,
task_id: &StepId,
attempt: u32,
cancel: CancellationToken,
) -> Result<Option<(Value, bool)>, WorkerError> {
let mut reader = BufReader::new(stdout).lines();
let mut last_final = None;
loop {
tokio::select! {
line_res = reader.next_line() => match line_res {
Ok(Some(line)) => {
let trimmed = line.trim();
if trimmed.is_empty() { continue; }
if let Ok(ev) = serde_json::from_str::<OutputEvent>(trimmed) {
forward_event(engine, token, task_id, attempt, ev, &mut last_final).await;
}
}
Ok(None) => break,
Err(e) => return Err(WorkerError::Failed(format!("ndjson read: {e}"))),
},
_ = cancel.cancelled() => return Err(WorkerError::Cancelled),
}
}
Ok(last_final)
}
async fn read_sse(
stdout: tokio::process::ChildStdout,
engine: &Engine,
token: &CapToken,
task_id: &StepId,
attempt: u32,
cancel: CancellationToken,
) -> Result<Option<(Value, bool)>, WorkerError> {
let mut reader = BufReader::new(stdout).lines();
let mut last_final = None;
let mut data_buf = String::new();
loop {
tokio::select! {
line_res = reader.next_line() => match line_res {
Ok(Some(line)) => {
if line.is_empty() {
if !data_buf.is_empty() {
if let Ok(ev) = serde_json::from_str::<OutputEvent>(data_buf.trim()) {
forward_event(engine, token, task_id, attempt, ev, &mut last_final).await;
}
data_buf.clear();
}
} else if let Some(rest) = line.strip_prefix("data:") {
let payload = rest.strip_prefix(' ').unwrap_or(rest);
if !data_buf.is_empty() {
data_buf.push('\n');
}
data_buf.push_str(payload);
}
}
Ok(None) => {
if !data_buf.is_empty() {
if let Ok(ev) = serde_json::from_str::<OutputEvent>(data_buf.trim()) {
forward_event(engine, token, task_id, attempt, ev, &mut last_final).await;
}
}
break;
}
Err(e) => return Err(WorkerError::Failed(format!("sse read: {e}"))),
},
_ = cancel.cancelled() => return Err(WorkerError::Cancelled),
}
}
Ok(last_final)
}
async fn read_length_prefixed(
mut stdout: tokio::process::ChildStdout,
engine: &Engine,
token: &CapToken,
task_id: &StepId,
attempt: u32,
cancel: CancellationToken,
) -> Result<Option<(Value, bool)>, WorkerError> {
use tokio::io::AsyncReadExt;
let mut last_final = None;
loop {
let mut len_buf = [0u8; 4];
let read_fut = stdout.read_exact(&mut len_buf);
let read_res = tokio::select! {
r = read_fut => r,
_ = cancel.cancelled() => return Err(WorkerError::Cancelled),
};
match read_res {
Ok(_) => {}
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break, Err(e) => return Err(WorkerError::Failed(format!("len read: {e}"))),
}
let len = u32::from_be_bytes(len_buf) as usize;
if len == 0 || len > 16 * 1024 * 1024 {
break;
}
let mut payload = vec![0u8; len];
let read_fut = stdout.read_exact(&mut payload);
let read_res = tokio::select! {
r = read_fut => r,
_ = cancel.cancelled() => return Err(WorkerError::Cancelled),
};
if read_res.is_err() {
break;
}
if let Ok(ev) = serde_json::from_slice::<OutputEvent>(&payload) {
forward_event(engine, token, task_id, attempt, ev, &mut last_final).await;
}
}
Ok(last_final)
}
#[cfg(test)]
mod embed_tests {
use super::*;
fn vars() -> EmbedVars {
EmbedVars {
system: "SYS".to_string(),
system_file: Some("/tmp/sys.md".to_string()),
prompt: "do the task".to_string(),
model: Some("small".to_string()),
tools_csv: "Read,Grep".to_string(),
work_dir: Some("/tmp/wd".to_string()),
task_id: "ST-1".to_string(),
attempt: "1".to_string(),
}
}
#[test]
fn render_substitutes_every_closed_set_token() {
let v = vars();
let out = v
.render("{system}|{system_file}|{prompt}|{model}|{tools_csv}|{work_dir}|{task_id}|{attempt}")
.expect("render");
assert_eq!(
out,
"SYS|/tmp/sys.md|do the task|small|Read,Grep|/tmp/wd|ST-1|1"
);
}
#[test]
fn render_leaves_non_placeholder_braces_alone() {
let v = vars();
let out = v
.render(r#"echo '{"result": "{prompt}"}'"#)
.expect("render");
assert_eq!(out, r#"echo '{"result": "do the task"}'"#);
}
#[test]
fn render_fails_loud_when_model_referenced_but_absent() {
let mut v = vars();
v.model = None;
let err = v.render("--model {model}").unwrap_err();
let msg = format!("{err:?}");
assert!(msg.contains("{model}"), "actionable token in error: {msg}");
assert!(msg.contains("profile.model"), "actionable cause: {msg}");
}
#[test]
fn render_fails_loud_when_work_dir_referenced_but_absent() {
let mut v = vars();
v.work_dir = None;
let err = v.render("{work_dir}").unwrap_err();
let msg = format!("{err:?}");
assert!(msg.contains("{work_dir}"), "actionable token: {msg}");
}
#[test]
fn render_empty_system_is_legal() {
let mut v = vars();
v.system = String::new();
assert_eq!(v.render("[{system}]").expect("render"), "[]");
}
#[test]
fn render_never_substitutes_inside_substituted_values() {
let mut v = vars();
v.prompt = "please mention {model} and {work_dir} literally".to_string();
v.model = None; v.work_dir = None;
let out = v.render("task: {prompt}").expect("render");
assert_eq!(out, "task: please mention {model} and {work_dir} literally");
}
#[test]
fn render_substitutes_placeholder_nested_in_literal_braces() {
let v = vars();
let out = v.render(r#"{"task": "{prompt}", "n": 1}"#).expect("render");
assert_eq!(out, r#"{"task": "do the task", "n": 1}"#);
}
#[test]
fn embed_references_scans_argv_stdin_env_cwd() {
let mut t = EmbedTemplate {
argv: vec!["cat".to_string()],
..Default::default()
};
assert!(!embed_references(&t, "{system_file}"));
t.stdin = Some("{system_file}".to_string());
assert!(embed_references(&t, "{system_file}"));
t.stdin = None;
t.env.insert("X".to_string(), "{system_file}".to_string());
assert!(embed_references(&t, "{system_file}"));
t.env.clear();
t.cwd = Some("{system_file}".to_string());
assert!(embed_references(&t, "{system_file}"));
}
fn fake_output(stdout: &str, stderr: &str, code: i32) -> std::process::Output {
#[cfg(unix)]
use std::os::unix::process::ExitStatusExt;
#[cfg(windows)]
use std::os::windows::process::ExitStatusExt;
use std::process::ExitStatus;
#[cfg(unix)]
let status = ExitStatus::from_raw(code << 8);
#[cfg(windows)]
let status = ExitStatus::from_raw(code as u32);
std::process::Output {
status,
stdout: stdout.as_bytes().to_vec(),
stderr: stderr.as_bytes().to_vec(),
}
}
#[test]
fn normalize_without_decl_is_historical_json_or_raw() {
let out = fake_output(r#"{"a": 1}"#, "", 0);
let wr = normalize_plain_output(&out, None);
assert_eq!(wr.value, serde_json::json!({"a": 1}));
assert!(wr.ok);
let out = fake_output("plain text\n", "warned", 0);
let wr = normalize_plain_output(&out, None);
assert_eq!(
wr.value,
serde_json::json!({"raw": "plain text", "stderr": "warned"})
);
assert!(wr.ok);
let out = fake_output("boom", "", 1);
let wr = normalize_plain_output(&out, None);
assert!(!wr.ok, "non-zero exit is a failed step");
}
#[test]
fn normalize_result_ptr_extracts_declared_value() {
let decl = SubprocessOutput {
format: Some("json".to_string()),
result_ptr: Some("/result".to_string()),
ok_from: Some("exit_code".to_string()),
};
let out = fake_output(r#"{"result": {"answer": 42}, "noise": true}"#, "", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert_eq!(wr.value, serde_json::json!({"answer": 42}));
assert!(wr.ok);
}
#[test]
fn normalize_declared_json_unparsable_stdout_fails_loud() {
let decl = SubprocessOutput {
format: Some("json".to_string()),
result_ptr: None,
ok_from: None,
};
let out = fake_output("not json at all", "stderr text", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert!(!wr.ok, "declared-JSON unparsable stdout is a failed step");
assert_eq!(wr.value["raw"], "not json at all");
assert_eq!(wr.value["stderr"], "stderr text");
assert!(wr.value["parse_error"].is_string());
}
#[test]
fn normalize_missing_result_ptr_fails_loud_with_actionable_value() {
let decl = SubprocessOutput {
format: Some("json".to_string()),
result_ptr: Some("/missing".to_string()),
ok_from: None,
};
let out = fake_output(r#"{"present": 1}"#, "", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert!(!wr.ok);
assert!(wr.value["error"]
.as_str()
.expect("actionable error message")
.contains("/missing"));
}
#[test]
fn normalize_ok_from_pointer_reads_boolean() {
let decl = SubprocessOutput {
format: Some("json".to_string()),
result_ptr: Some("/result".to_string()),
ok_from: Some("/ok".to_string()),
};
let out = fake_output(r#"{"result": "r", "ok": true}"#, "", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert!(wr.ok);
let out = fake_output(r#"{"result": "r", "ok": false}"#, "", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert!(!wr.ok);
let out = fake_output(r#"{"result": "r", "ok": "yes"}"#, "", 0);
let wr = normalize_plain_output(&out, Some(&decl));
assert!(!wr.ok);
}
}