diff --git a/src-tauri/src/commands/nodejs.rs b/src-tauri/src/commands/nodejs.rs index d735433..7f56cb7 100644 --- a/src-tauri/src/commands/nodejs.rs +++ b/src-tauri/src/commands/nodejs.rs @@ -23,8 +23,9 @@ use serde::Serialize; use serde_json::Value; +use std::collections::HashMap; use tauri::{AppHandle, Emitter, Window}; -use tokio::sync::mpsc; +use tokio::sync::{mpsc, watch, Mutex}; use crate::app_config::AppConfig; use crate::js_runtime::PipelineEvent; @@ -32,6 +33,11 @@ use crate::nodejs; const EVENT_NAME: &str = "nodejs:event"; +#[derive(Default)] +pub struct NodeTaskRegistry { + inner: Mutex>>, +} + #[derive(Debug, Serialize, Clone)] #[serde(tag = "type", rename_all = "camelCase")] pub enum FrontendEvent { @@ -102,16 +108,33 @@ pub async fn run_nodejs_script( app: AppHandle, window: Window, app_config: tauri::State<'_, AppConfig>, + registry: tauri::State<'_, NodeTaskRegistry>, script_name: String, params: Value, ) -> Result { + let (cancel_tx, cancel_rx) = watch::channel(false); + { + let mut map = registry.inner.lock().await; + if map.contains_key(&script_name) { + return Err(format!("脚本正在运行: {script_name}")); + } + map.insert(script_name.clone(), cancel_tx); + } + let (tx, rx) = mpsc::unbounded_channel::(); let pump = spawn_event_pump(app, window, script_name.clone(), rx); let config_env = app_config.env_for_node().await; - let result = nodejs::run_node_script(script_name, params, tx, config_env) + let result = nodejs::run_node_script(script_name.clone(), params, tx, config_env, cancel_rx) .await - .map_err(|e| e.to_string())?; + .map_err(|e| e.to_string()); + + { + let mut map = registry.inner.lock().await; + map.remove(&script_name); + } + + let result = result?; pump.await.ok(); Ok(result) @@ -123,16 +146,34 @@ pub async fn run_nodejs_script_source( app: AppHandle, window: Window, app_config: tauri::State<'_, AppConfig>, + registry: tauri::State<'_, NodeTaskRegistry>, script_source: String, params: Value, ) -> Result { + let script_key = "".to_string(); + let (cancel_tx, cancel_rx) = watch::channel(false); + { + let mut map = registry.inner.lock().await; + if map.contains_key(&script_key) { + return Err("调试脚本正在运行".to_string()); + } + map.insert(script_key.clone(), cancel_tx); + } + let (tx, rx) = mpsc::unbounded_channel::(); - let pump = spawn_event_pump(app, window, "".to_string(), rx); + let pump = spawn_event_pump(app, window, script_key.clone(), rx); let config_env = app_config.env_for_node().await; - let result = nodejs::run_node_script_source(script_source, params, tx, config_env) + let result = nodejs::run_node_script_source(script_source, params, tx, config_env, cancel_rx) .await - .map_err(|e| e.to_string())?; + .map_err(|e| e.to_string()); + + { + let mut map = registry.inner.lock().await; + map.remove(&script_key); + } + + let result = result?; pump.await.ok(); Ok(result) @@ -143,3 +184,19 @@ pub async fn run_nodejs_script_source( pub async fn list_nodejs_scripts() -> Result, String> { nodejs::list_scripts().await.map_err(|e| e.to_string()) } + +#[tauri::command] +pub async fn stop_nodejs_script( + registry: tauri::State<'_, NodeTaskRegistry>, + script_name: String, +) -> Result { + let sender = { + let map = registry.inner.lock().await; + map.get(&script_name).cloned() + }; + if let Some(tx) = sender { + tx.send(true).map_err(|e| e.to_string())?; + return Ok(true); + } + Ok(false) +} diff --git a/src-tauri/src/commands/publish.rs b/src-tauri/src/commands/publish.rs index 1eee062..62594a8 100644 --- a/src-tauri/src/commands/publish.rs +++ b/src-tauri/src/commands/publish.rs @@ -3,6 +3,7 @@ use serde::Deserialize; use serde_json::{json, Value}; use tauri::{AppHandle, Manager, State, Window}; +use tokio::sync::watch; use crate::app_config::AppConfig; use crate::publish_db::PlatformAccountRecord; @@ -64,6 +65,7 @@ async fn run_publish_script( params, tx, config_env, + watch::channel(false).1, ) .await .map_err(|e| e.to_string())?; diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index f2fc535..734511a 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -44,6 +44,7 @@ pub fn run() { .manage(audio_store::AudioStore::new()) .manage(video_store::VideoStore::new()) .manage(publish_store::PublishStore::new()) + .manage(commands::nodejs::NodeTaskRegistry::default()) .setup(|app| { let handle = app.handle().clone(); @@ -129,6 +130,7 @@ pub fn run() { commands::nodejs::run_nodejs_script, commands::nodejs::run_nodejs_script_source, commands::nodejs::list_nodejs_scripts, + commands::nodejs::stop_nodejs_script, commands::fs_util::read_local_file_base64, commands::fs_util::get_media_duration_seconds, commands::fs_util::copy_local_file, diff --git a/src-tauri/src/nodejs.rs b/src-tauri/src/nodejs.rs index f9ccc2d..66705b2 100644 --- a/src-tauri/src/nodejs.rs +++ b/src-tauri/src/nodejs.rs @@ -28,6 +28,7 @@ use thiserror::Error; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::process::Command; use tokio::sync::mpsc::UnboundedSender; +use tokio::sync::watch; use uuid::Uuid; use crate::js_runtime::PipelineEvent; @@ -66,6 +67,8 @@ pub enum NodeError { NoResult, #[error("json: {0}")] Json(#[from] serde_json::Error), + #[error("任务已取消")] + Cancelled, } // --------------------------------------------------------------------------- @@ -346,11 +349,12 @@ pub async fn run_node_script( params: Value, events: UnboundedSender, config_env: HashMap, + mut cancel_rx: watch::Receiver, ) -> Result { let entry = normalize_script_name(&script_name) .ok_or_else(|| NodeError::FetchScript(format!("无效的脚本名称: {script_name}")))?; let modules = resolve_script_bundle(&entry).await?; - run_script_bundle(&entry, modules, params, events, config_env).await + run_script_bundle(&entry, modules, params, events, config_env, &mut cancel_rx).await } /// 调试入口:直接传源码;其中 `require('*.js')` 仍从服务端拉取依赖。 @@ -359,9 +363,18 @@ pub async fn run_node_script_source( params: Value, events: UnboundedSender, config_env: HashMap, + mut cancel_rx: watch::Receiver, ) -> Result { let modules = resolve_script_bundle_for_inline(&source).await?; - run_script_bundle(INLINE_ENTRY, modules, params, events, config_env).await + run_script_bundle( + INLINE_ENTRY, + modules, + params, + events, + config_env, + &mut cancel_rx, + ) + .await } async fn run_script_bundle( @@ -370,6 +383,7 @@ async fn run_script_bundle( params: Value, events: UnboundedSender, config_env: HashMap, + cancel_rx: &mut watch::Receiver, ) -> Result { let node = locate_node().ok_or(NodeError::NotFound)?; let node_modules_cwd = nodejs_dir().unwrap_or_else(std::env::temp_dir); @@ -378,6 +392,7 @@ async fn run_script_bundle( let mut cmd = Command::new(&node); cmd.arg(&script_path) + .kill_on_drop(true) .current_dir(&bundle_dir) .env( "NODE_PATH", @@ -491,18 +506,26 @@ async fn run_script_bundle( result }); - // 与 stdout/stderr 读取并行等待,避免子进程在管道未排空时退出(Windows libuv) - let (status, stderr_tail, result_line) = tokio::join!( - child.wait(), - stderr_task, - stdout_task, - ); - let status = status?; + // 与取消信号竞争等待;取消时立即杀掉子进程 + let (status, cancelled) = tokio::select! { + status = child.wait() => (status?, false), + _ = wait_for_cancel(cancel_rx) => { + let _ = child.start_kill(); + let status = child.wait().await?; + (status, true) + } + }; + + let (stderr_tail, result_line) = tokio::join!(stderr_task, stdout_task); let stderr_tail = stderr_tail.unwrap_or_default(); let result_line = result_line.unwrap_or(None); let _ = tokio::fs::remove_dir_all(&bundle_dir).await; + if cancelled { + return Err(NodeError::Cancelled); + } + if !status.success() { return Err(NodeError::NonZero { status: status.code(), @@ -514,3 +537,14 @@ async fn run_script_bundle( let v: Value = serde_json::from_str(&result_str)?; Ok(v) } + +async fn wait_for_cancel(cancel_rx: &mut watch::Receiver) { + if *cancel_rx.borrow() { + return; + } + while cancel_rx.changed().await.is_ok() { + if *cancel_rx.borrow() { + return; + } + } +} diff --git a/src/components/workflow/VideoLinkDialog.vue b/src/components/workflow/VideoLinkDialog.vue index 9535e20..2f04f6f 100644 --- a/src/components/workflow/VideoLinkDialog.vue +++ b/src/components/workflow/VideoLinkDialog.vue @@ -1,6 +1,7 @@