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
1 change: 1 addition & 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 crates/agentos-actor-plugin/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ bytes = "1"
http = "1"
uuid = { version = "1", features = ["v4"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
futures = "0.3"

[dev-dependencies]
Expand Down
103 changes: 93 additions & 10 deletions crates/agentos-actor-plugin/src/actions/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ pub mod network;
pub mod preview;
pub mod process;
pub mod session;
pub mod shell;

use std::collections::HashMap;

Expand All @@ -41,6 +42,11 @@ pub struct Vars {
pub capture_tasks: HashMap<String, JoinHandle<()>>,
/// `live_session_id -> permission-request pump task`.
pub permission_tasks: HashMap<String, JoinHandle<()>>,
/// Shell data/stderr/exit broadcast pump tasks (one triple per `openShell`).
/// The pumps end on their own when the shell exits (stream close); this
/// list exists so VM teardown aborts any still-live pumps. Bounded by the
/// client's shell registries, not here.
pub shell_tasks: Vec<JoinHandle<()>>,
}

impl Vars {
Expand All @@ -62,6 +68,9 @@ impl Vars {
for (_, task) in self.permission_tasks.drain() {
task.abort();
}
for task in self.shell_tasks.drain(..) {
task.abort();
}
self.live_sessions.clear();
}
}
Expand Down Expand Up @@ -207,18 +216,45 @@ pub(crate) async fn dispatch(
},
Err(error) => reply_err(host, token, error),
},
"spawn" => match decode_as::<(String, Vec<String>)>(args) {
Ok((command, spawn_args)) => match process::spawn(vm, &command, spawn_args) {
Ok(handle) => reply_ok(host, token, &handle),
"spawn" => {
// The trailing options object is optional on the TS side, so a
// 2-arg call decodes via the fallback.
let decoded =
decode_as::<(String, Vec<String>, Option<process::SpawnActionOptions>)>(args)
.or_else(|_| {
decode_as::<(String, Vec<String>)>(args)
.map(|(command, spawn_args)| (command, spawn_args, None))
});
match decoded {
Ok((command, spawn_args, options)) => match process::spawn(
host,
vm,
vars,
&command,
spawn_args,
options.unwrap_or_default(),
) {
Ok(handle) => reply_ok(host, token, &handle),
Err(error) => reply_err(host, token, error),
},
Err(error) => reply_err(host, token, error),
},
Err(error) => reply_err(host, token, error),
},
}
}
// Long-running wait: replies from a spawned task so it does not occupy
// the serial action worker (a waitProcess held for the process lifetime
// would starve every later action, including the stdin writes the
// process needs to make progress).
"waitProcess" => match decode_as::<(u32,)>(args) {
Ok((pid,)) => match process::wait_process(vm, pid).await {
Ok(code) => reply_ok(host, token, &code),
Err(error) => reply_err(host, token, error),
},
Ok((pid,)) => {
let host = host.clone();
let vm = vm.clone();
vars.shell_tasks.push(tokio::spawn(async move {
match process::wait_process(&vm, pid).await {
Ok(code) => reply_ok(&host, token, &code),
Err(error) => reply_err(&host, token, error),
}
}));
}
Err(error) => reply_err(host, token, error),
},
"killProcess" => match decode_as::<(u32,)>(args) {
Expand Down Expand Up @@ -373,6 +409,53 @@ pub(crate) async fn dispatch(
},
Err(error) => reply_err(host, token, error),
},
"openShell" => match decode_as::<(Option<shell::OpenShellActionOptions>,)>(args) {
Ok((options,)) => {
match shell::open_shell(host, vm, vars, options.unwrap_or_default()) {
Ok(dto) => reply_ok(host, token, &dto),
Err(error) => reply_err(host, token, error),
}
}
Err(error) => reply_err(host, token, error),
},
"writeShell" => match decode_as::<(String, WriteFileContent)>(args) {
Ok((shell_id, data)) => match shell::write_shell(vm, &shell_id, data).await {
Ok(()) => reply_ok(host, token, &()),
Err(error) => reply_err(host, token, error),
},
Err(error) => reply_err(host, token, error),
},
"resizeShell" => match decode_as::<(String, u16, u16)>(args) {
Ok((shell_id, cols, rows)) => match shell::resize_shell(vm, &shell_id, cols, rows) {
Ok(()) => reply_ok(host, token, &()),
Err(error) => reply_err(host, token, error),
},
Err(error) => reply_err(host, token, error),
},
"closeShell" => match decode_as::<(String,)>(args) {
Ok((shell_id,)) => match shell::close_shell(vm, &shell_id) {
Ok(()) => reply_ok(host, token, &()),
Err(error) => reply_err(host, token, error),
},
Err(error) => reply_err(host, token, error),
},
// Long-running wait: replies from a spawned task so it does not occupy
// the serial action worker (the shell CLI calls waitShell up front and
// streams writeShell input afterwards; holding the worker here would
// deadlock the shell — input can never arrive to end the wait).
"waitShell" => match decode_as::<(String,)>(args) {
Ok((shell_id,)) => {
let host = host.clone();
let vm = vm.clone();
vars.shell_tasks.push(tokio::spawn(async move {
match shell::wait_shell(&vm, &shell_id).await {
Ok(exit_code) => reply_ok(&host, token, &exit_code),
Err(error) => reply_err(&host, token, error),
}
}));
}
Err(error) => reply_err(host, token, error),
},
other => {
host.reply_err(
token,
Expand Down
44 changes: 39 additions & 5 deletions crates/agentos-actor-plugin/src/actions/process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,45 @@ pub async fn exec(vm: &AgentOs, command: &str) -> Result<ExecResultDto> {
.map(ExecResultDto::from)
}

/// `spawn(command, args)` — port of [`AgentOs::spawn`]. Returns the
/// [`SpawnHandle`] `{ pid }` directly; the underlying type already
/// derives `Serialize`.
pub fn spawn(vm: &AgentOs, command: &str, args: Vec<String>) -> Result<SpawnHandle> {
vm.spawn(command, args, SpawnOptions::default())
/// JSON options for the `spawn` action — the serializable subset of the TS
/// `SpawnOptions` (output callbacks are replaced by the `processOutput` /
/// `processExit` broadcasts).
#[derive(Debug, Default, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SpawnActionOptions {
#[serde(default)]
pub env: std::collections::BTreeMap<String, String>,
pub cwd: Option<String>,
pub stream_stdin: Option<bool>,
}

/// `spawn(command, args, options?)` — port of [`AgentOs::spawn`]. Returns the
/// [`SpawnHandle`] `{ pid }`; stdout/stderr chunks stream to connected clients
/// as `processOutput` events and the exit code broadcasts as `processExit`.
pub fn spawn(
host: &crate::host_ctx::HostCtx,
vm: &AgentOs,
vars: &mut super::Vars,
command: &str,
args: Vec<String>,
options: SpawnActionOptions,
) -> Result<SpawnHandle> {
let mut base = ExecOptions::default();
base.env = options.env;
if options.cwd.is_some() {
base.cwd = options.cwd;
}
let handle = vm.spawn(
command,
args,
SpawnOptions {
base,
stream_stdin: options.stream_stdin,
..SpawnOptions::default()
},
)?;
super::shell::spawn_process_output_pumps(host, vm, vars, handle.pid);
Ok(handle)
}

/// `waitProcess(pid)` — port of [`AgentOs::wait_process`]. Returns the
Expand Down
Loading
Loading