Initial commit
This commit is contained in:
184
src/executor/structured.rs
Normal file
184
src/executor/structured.rs
Normal file
@@ -0,0 +1,184 @@
|
||||
use crate::error::{AppError, Result};
|
||||
use serde_json::Value;
|
||||
use std::process::Stdio;
|
||||
use tokio::{
|
||||
io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
|
||||
process::Command,
|
||||
sync::mpsc::{self, UnboundedReceiver, UnboundedSender},
|
||||
};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum StructuredCommand {
|
||||
Input {
|
||||
prompt: String,
|
||||
default: Option<String>,
|
||||
secret: bool,
|
||||
},
|
||||
Menu {
|
||||
prompt: String,
|
||||
options: Vec<MenuOption>,
|
||||
},
|
||||
Confirm {
|
||||
prompt: String,
|
||||
},
|
||||
Message {
|
||||
level: MessageLevel,
|
||||
text: String,
|
||||
},
|
||||
Progress {
|
||||
percent: u8,
|
||||
message: Option<String>,
|
||||
},
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct MenuOption {
|
||||
pub id: String,
|
||||
pub label: String,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum MessageLevel {
|
||||
Info,
|
||||
Warn,
|
||||
Error,
|
||||
}
|
||||
|
||||
/// Запускает bash-скрипт в структурированном режиме.
|
||||
/// Возвращает:
|
||||
/// - `output_rx` – канал для строк вывода (обычный текст и stderr)
|
||||
/// - `command_rx` – канал для команд от скрипта
|
||||
/// - `finished_rx` – одноразовый канал с кодом завершения
|
||||
/// - `reply_tx` – канал для отправки ответов обратно в stdin скрипта
|
||||
pub fn spawn(
|
||||
script: &str,
|
||||
) -> Result<(
|
||||
UnboundedReceiver<String>,
|
||||
UnboundedReceiver<StructuredCommand>,
|
||||
tokio::sync::oneshot::Receiver<i32>,
|
||||
UnboundedSender<String>,
|
||||
)> {
|
||||
let mut child = Command::new("bash")
|
||||
.arg("-c")
|
||||
.arg(script)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()?;
|
||||
|
||||
let stdin = child.stdin.take().unwrap();
|
||||
let stdout = child.stdout.take().unwrap();
|
||||
let stderr = child.stderr.take().unwrap();
|
||||
|
||||
let (output_tx, output_rx) = mpsc::unbounded_channel();
|
||||
let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
|
||||
let (finished_tx, finished_rx) = tokio::sync::oneshot::channel();
|
||||
let (reply_tx, mut reply_rx) = mpsc::unbounded_channel::<String>();
|
||||
|
||||
// Чтение stdout
|
||||
let output_tx_clone = output_tx.clone();
|
||||
let cmd_tx_stdout = cmd_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut reader = BufReader::new(stdout).lines();
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
if line.starts_with("CMD:") {
|
||||
if parse_command(&line[4..], cmd_tx_stdout.clone()).is_err() {
|
||||
let _ = output_tx_clone.send(format!("[PROTOCOL ERROR] {}", line));
|
||||
}
|
||||
} else {
|
||||
let _ = output_tx_clone.send(line);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Чтение stderr
|
||||
let output_tx_stderr = output_tx.clone();
|
||||
let cmd_tx_stderr = cmd_tx.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut reader = BufReader::new(stderr).lines();
|
||||
while let Ok(Some(line)) = reader.next_line().await {
|
||||
if line.starts_with("CMD:") {
|
||||
if parse_command(&line[4..], cmd_tx_stderr.clone()).is_err() {
|
||||
let _ = output_tx_stderr.send(format!("[PROTOCOL ERROR] {}", line));
|
||||
}
|
||||
} else {
|
||||
let _ = output_tx_stderr.send(format!("stderr: {}", line));
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Запись в stdin
|
||||
tokio::spawn(async move {
|
||||
let mut stdin = stdin;
|
||||
while let Some(response) = reply_rx.recv().await {
|
||||
if stdin.write_all(response.as_bytes()).await.is_err() {
|
||||
break;
|
||||
}
|
||||
if stdin.write_all(b"\n").await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Ожидание завершения
|
||||
tokio::spawn(async move {
|
||||
let status = child.wait().await;
|
||||
let code = status.map(|s| s.code().unwrap_or(1)).unwrap_or(1);
|
||||
let _ = finished_tx.send(code);
|
||||
});
|
||||
|
||||
Ok((output_rx, cmd_rx, finished_rx, reply_tx))
|
||||
}
|
||||
|
||||
fn parse_command(json_str: &str, cmd_tx: UnboundedSender<StructuredCommand>) -> Result<()> {
|
||||
let v: Value = serde_json::from_str(json_str)?;
|
||||
let typ = v["type"].as_str().ok_or_else(|| AppError::Protocol("Missing type".into()))?;
|
||||
|
||||
let cmd = match typ {
|
||||
"input" => {
|
||||
let prompt = v["prompt"].as_str().unwrap_or("").to_string();
|
||||
let default = v["default"].as_str().map(String::from);
|
||||
let secret = v["secret"].as_bool().unwrap_or(false);
|
||||
StructuredCommand::Input { prompt, default, secret }
|
||||
}
|
||||
"menu" => {
|
||||
let prompt = v["prompt"].as_str().unwrap_or("").to_string();
|
||||
let options = v["options"]
|
||||
.as_array()
|
||||
.map(|arr| {
|
||||
arr.iter()
|
||||
.filter_map(|opt| {
|
||||
Some(MenuOption {
|
||||
id: opt["id"].as_str()?.to_string(),
|
||||
label: opt["label"].as_str()?.to_string(),
|
||||
})
|
||||
})
|
||||
.collect()
|
||||
})
|
||||
.unwrap_or_default();
|
||||
StructuredCommand::Menu { prompt, options }
|
||||
}
|
||||
"confirm" => {
|
||||
let prompt = v["prompt"].as_str().unwrap_or("").to_string();
|
||||
StructuredCommand::Confirm { prompt }
|
||||
}
|
||||
"message" => {
|
||||
let level = match v["level"].as_str() {
|
||||
Some("warn") => MessageLevel::Warn,
|
||||
Some("error") => MessageLevel::Error,
|
||||
_ => MessageLevel::Info,
|
||||
};
|
||||
let text = v["text"].as_str().unwrap_or("").to_string();
|
||||
StructuredCommand::Message { level, text }
|
||||
}
|
||||
"progress" => {
|
||||
let percent = v["percent"].as_u64().unwrap_or(0) as u8;
|
||||
let message = v["message"].as_str().map(String::from);
|
||||
StructuredCommand::Progress { percent, message }
|
||||
}
|
||||
_ => return Err(AppError::Protocol(format!("Unknown command type: {}", typ))),
|
||||
};
|
||||
|
||||
cmd_tx.send(cmd).map_err(|_| AppError::Protocol("Command channel closed".into()))?;
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user