diff --git a/neurosploit-rs/app/src/main.rs b/neurosploit-rs/app/src/main.rs index dbc1186..bb67ae8 100644 --- a/neurosploit-rs/app/src/main.rs +++ b/neurosploit-rs/app/src/main.rs @@ -1348,21 +1348,77 @@ pub(crate) fn finalize_run(mut out: RunOutput, workdir: &Path) -> RunOutput { async fn run_mode(base: &Path, cfg: RunConfig, mcp: bool, mode: Mode) -> anyhow::Result { subscription_preflight(&cfg).await; - let Spawned { mut task, mut rx, cancel, workdir, .. } = spawn_engagement(base, cfg, mcp, mode); + // Keep the pause/resume handles: a one-shot run that parks on quota/auth + // exhaustion must still be resumable. Without reading stdin here the pool's + // "type /continue" notice would be a dead end — nothing would accept it, and + // the process would hang forever on the parked task. + let Spawned { mut task, mut rx, cancel, paused, resume, fallback, workdir, .. } = + spawn_engagement(base, cfg, mcp, mode); let printer = tokio::spawn(async move { while let Some(line) = rx.recv().await { render_line(&line); } }); + // Read operator input ONLY at a real terminal: `/continue [provider:model]` + // to resume a parked run, `/model provider:model` to switch then resume. + // Over a pipe (CI, web) there is no one to type, so we skip it and rely on + // the auto-resume / Ctrl-C paths instead. + use std::io::IsTerminal as _; + let interactive = std::io::stdin().is_terminal(); + let mut stdin_lines = interactive.then(|| { + use tokio::io::AsyncBufReadExt as _; + tokio::io::BufReader::new(tokio::io::stdin()).lines() + }); + + let resume_run = |arg: &str| { + let arg = arg.trim(); + if !arg.is_empty() { + let m = ModelRef::parse(arg); + println!(" \x1b[1;35m▶ resuming with\x1b[0m {}:{}", m.provider, m.model); + if let Ok(mut fb) = fallback.lock() { fb.push(m); } + } else { + println!(" \x1b[1;35m▶ resuming\x1b[0m — retrying with the current model(s)."); + } + paused.store(false, std::sync::atomic::Ordering::Relaxed); + resume.notify_waiters(); + }; + let mut cancelled = false; - let out: RunOutput = tokio::select! { - r = &mut task => r.unwrap_or_default(), - _ = tokio::signal::ctrl_c() => { - cancelled = true; - cancel.store(true, std::sync::atomic::Ordering::Relaxed); - println!("\n \x1b[33m⏸ stopping — finishing in-flight work… (Ctrl-C again to abort now)\x1b[0m"); - tokio::select! { - r = &mut task => r.unwrap_or_default(), - _ = tokio::signal::ctrl_c() => { task.abort(); println!(" \x1b[31m✗ aborted.\x1b[0m"); RunOutput::default() } + let out: RunOutput = loop { + tokio::select! { + r = &mut task => break r.unwrap_or_default(), + _ = tokio::signal::ctrl_c() => { + cancelled = true; + cancel.store(true, std::sync::atomic::Ordering::Relaxed); + paused.store(false, std::sync::atomic::Ordering::Relaxed); // don't let a parked task hold the Ctrl-C + resume.notify_waiters(); + println!("\n \x1b[33m⏸ stopping — finishing in-flight work… (Ctrl-C again to abort now)\x1b[0m"); + break tokio::select! { + r = &mut task => r.unwrap_or_default(), + _ = tokio::signal::ctrl_c() => { task.abort(); println!(" \x1b[31m✗ aborted.\x1b[0m"); RunOutput::default() } + }; + } + line = async { match stdin_lines.as_mut() { Some(l) => l.next_line().await, None => Ok(None) } }, if stdin_lines.is_some() => { + match line { + Ok(Some(l)) => { + let t = l.trim(); + let (cmd, arg) = t.split_once(char::is_whitespace).unwrap_or((t, "")); + match cmd { + "/continue" | "/resume" => { + if paused.load(std::sync::atomic::Ordering::Relaxed) { resume_run(arg); } + else { println!(" run is not paused — it's still working."); } + } + "/model" => { + // Switch provider/model and resume if parked. + if arg.trim().is_empty() { println!(" usage: /model "); } + else { resume_run(arg); } + } + "" => {} + other => println!(" while a run is active only /continue [provider:model], /model and Ctrl-C are accepted (got {other})"), + } + } + Ok(None) => { stdin_lines = None; } // EOF — stop polling stdin + Err(_) => { stdin_lines = None; } + } } } };