use anyhow::Result; use common::config::{ChainStep, SynspecInput, TlustyInput}; use common::embedded::{ensure_specific_data_files, RuntimePaths}; use common::models::{ModelSummary, TaskSpec}; use common::result_filter::is_result_worthy; use common::runner::ExecutionRunner; use reqwest::Client; use std::path::{Path, PathBuf}; use tracing::{debug, info, warn}; /// 进程级互斥锁:序列化谱线表按需下载,防止多 Slot 并发首次派发时对同一 238MB 文件 /// 重复发起下载请求(POSIX rename 保证写安全,但冗余下载浪费带宽与 IO)。 static LINELIST_DOWNLOAD_MUTEX: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); pub async fn execute_task( client: &Client, server_url: &str, runtime: &RuntimePaths, work_dir: &Path, result_dir: &Path, task: &TaskSpec, shutdown: Option>, ) -> Result<(ModelSummary, Option>)> { info!( "开始执行计算任务 {} (网格点: {})", task.task_id, task.point_name ); // 1. Pull ONLY missing atom model data files needed for this task let required_atom_files = &[ "h1.dat", "he1.dat", "he2.dat", "c1.dat", "c2.dat", "c3_34+12lev.dat", "c4.dat", "n1.dat", "n2_32+10lev.dat", "n3.dat", "n4_34+14lev.dat", "n5.dat", "o1_23+10lev.dat", "o2_36+12lev.dat", "o3_28+13lev.dat", "o4.dat", "o5.dat", ]; if let Err(e) = ensure_specific_data_files(&runtime.data_dir, server_url, client, required_atom_files).await { warn!("拉取缺失原子数据文件失败: {}", e); } let mut seed_atmos_path: Option = None; // 提前创建 per-slot 隔离沙盒目录:种子下载后需复制一份私有副本进沙盒(见下方), // 故沙盒必须先于种子下载就绪。 let slot_work_dir = work_dir.join(format!("task_{}", task.task_id)); tokio::fs::create_dir_all(&slot_work_dir).await?; // 2. 种子大气获取(见 docs/task_engine_decoupling_design.md §5): // - TLUSTY 启用 + 策略为 seed_step:下载近邻种子 .7 作热启动种子(既有逻辑)。 // - TLUSTY 关闭(仅 SYNSPEC 场景):需拉取目标点 .7 大气作光谱合成输入。 // 顺序:本地 result 归档 → server 拉取。 // 注:不查沙盒本地——TLUSTY 与 SYNSPEC 在同一任务内串行,种子获取先于 runner, // 全新 slot 内不可能已有目标大气;重试任务的 slot 亦全新(task_id 唯一)。 let tlusty_enabled = task.tlusty_config.enabled; let current_strategy = task.tlusty_config.current_strategy("cold_run"); let needs_seed_download = (tlusty_enabled && current_strategy == "seed_step") || (!tlusty_enabled && task.synspec_config.enabled); if needs_seed_download { let seed_name = if !tlusty_enabled { // 仅 SYNSPEC 场景:大气来自目标点本身(atmosphere_ref 或 point_name)。 task.atmosphere_ref .clone() .unwrap_or_else(|| task.point_name.clone()) } else { // SeedStep 热启动:大气来自近邻种子点。 task.seed_point_name.clone().ok_or_else(|| { anyhow::anyhow!( "SeedStep 任务 {} 缺少 seed_point_name,无法热启动", task.point_name ) })? }; // 仅 SYNSPEC 场景先查本地 result 归档目录(节点此前算过同点大气,避免 server 拉取)。 if !tlusty_enabled { let archived = result_dir .join(&task.point_name) .join(format!("{}.7", task.point_name)); if archived.is_file() { info!( "SYNSPEC-only:在本地 result 归档找到大气 {},复用避免 server 拉取", archived.display() ); seed_atmos_path = Some(archived); } } // 归档无 → 向 server 拉取。 if seed_atmos_path.is_none() { let seed_url = format!("{}/api/seed/{}", server_url, seed_name); info!( "正在从服务端下载大气文件 ({}, 用途: {}): {}", seed_name, if tlusty_enabled { "TLUSTY 热启动种子" } else { "SYNSPEC 输入大气" }, seed_url ); let resp = client.get(&seed_url).send().await.map_err(|e| { anyhow::anyhow!( "任务 {} 下载大气文件 {} 失败: {}", task.point_name, seed_url, e ) })?; if !resp.status().is_success() { anyhow::bail!( "任务 {} 下载大气文件 {} 失败: HTTP {}", task.point_name, seed_url, resp.status() ); } let bytes = resp.bytes().await.map_err(|e| { anyhow::anyhow!( "任务 {} 读取大气文件 {} 响应体失败: {}", task.point_name, seed_url, e ) })?; let temp_seed_dir = work_dir.join(".seed_cache"); tokio::fs::create_dir_all(&temp_seed_dir).await?; cleanup_seed_cache(&temp_seed_dir).await; let tmp_path = temp_seed_dir.join(format!( "{}.{}.tmp", seed_name, uuid::Uuid::new_v4().simple() )); let final_seed_path = temp_seed_dir.join(format!("{}.seed.7", seed_name)); tokio::fs::write(&tmp_path, bytes).await?; tokio::fs::rename(&tmp_path, &final_seed_path).await?; let private_seed = slot_work_dir.join("seed_atmos.seed.7"); tokio::fs::copy(&final_seed_path, &private_seed).await?; seed_atmos_path = Some(private_seed); } } // 3. (slot_work_dir 已在种子下载前提前创建,种子私有副本亦已落盘于沙盒内。) // 反序列化工作流携带的 SYNSPEC 数值参数(波长范围等)。None → runner 用硬编码默认。 let synspec_cfg: Option = task .synspec_params .as_ref() .and_then(|v| serde_json::from_value::(v.clone()).ok()); // 谱线表覆盖:TaskSpec.linelist 指定时(如 YAML 配 linelist: gfATO.dat), // 按需从服务端下载到 runtime_dir 根目录,并构造覆盖了 linelist 路径的 RuntimePaths。 // None → 用 node 启动时下载的默认线表(ensure_runtime 的 default_linelist 参数)。 let runtime_override: Option = if let Some(ref ll_name) = task.linelist { // 线表放在 runtime_dir 根目录(与默认线表同级,runner symlink 为 fort.19)。 let ll_path = runtime .data_dir .parent() .unwrap_or(&runtime.data_dir) .join(ll_name); // 互斥保护:多 Slot 并发首次派发时,仅第一个进入临界区的 Slot 执行下载, // 后续 Slot 在获得锁后 double-check 发现文件已存在直接跳过,避免冗余下载。 { let _guard = LINELIST_DOWNLOAD_MUTEX.lock().await; if !ll_path.exists() { info!("任务指定谱线表 {} 本地缺失,从服务端按需下载...", ll_name); let url = format!("{}/api/data/file/{}", server_url, ll_name); let resp = client.get(&url).send().await?; if resp.status().is_success() { let bytes = resp.bytes().await?; let ll_dir = ll_path.parent().unwrap_or(std::path::Path::new(".")); let tmp = ll_dir.join(format!("{}.{}.tmp", ll_name, uuid::Uuid::new_v4().simple())); tokio::fs::write(&tmp, &bytes).await?; tokio::fs::rename(&tmp, &ll_path).await?; info!("成功下载谱线表 {} ({} bytes)", ll_name, bytes.len()); } else { anyhow::bail!( "下载任务指定谱线表 {} 失败,HTTP {}", ll_name, resp.status() ); } } } let mut rt = runtime.clone(); rt.linelist = ll_path; Some(rt) } else { None }; let effective_runtime = runtime_override.as_ref().unwrap_or(runtime); let runner = ExecutionRunner::new(effective_runtime, slot_work_dir.clone()); // 执行链来源(优先级): // 1. TaskSpec.tlusty_chain_params(用户在 YAML `tlusty_chain:` 配置的多阶段 ChainStep // 数组,由 scheduler 序列化注入)——非空时优先使用,使用户能细粒度控制 niter/chmax/ // ilvlin 等阶段参数。 // 2. default_chain_for_strategy(current_strategy) 兜底——按策略名(cold_run/seed_step) // 选预设默认链(runner.rs 的 default_cold_chain / default_seed_chain)。 // 历史:Phase 6 起仅用 default 链(用户 config.chain 被忽略,是死字段);本次接通后 // 用户配置真正生效,default 链降级为兜底。旧 MQ payload(无 tlusty_chain_params 字段) // 反序列化为 None → 回退 default 链,行为与旧版完全一致(向后兼容)。 let chain = resolve_execution_chain( current_strategy, &task.tlusty_chain_params, &task.seed_chain_params, &task.task_id.to_string(), ); // TLUSTY 输入文件全局参数(NFREAD/ions 表/nst extra_keys 等)。 // None → runner 用代码内硬编码默认(向后兼容)。 let tlusty_input = task.tlusty_input_params.as_ref().and_then(|v| { match serde_json::from_value::(v.clone()) { Ok(t) => Some(t), Err(e) => { warn!( "任务 {} 的 tlusty_input_params 反序列化失败,回退默认输入: {}", task.task_id, e ); None } } }); let summary = runner .run_model_with_timeout( &task.params, // 用权威的 point_name(DB grid_points.name 列,源精度正确)作为模型名, // 而非 task.params.model_name()(后者经 DB REAL 列回读已丢精度 "5.0"→"5")。 &task.point_name, current_strategy, Some(chain), seed_atmos_path.as_deref(), synspec_cfg.as_ref(), // 阶段独立配置开关(见 docs/task_engine_decoupling_design.md §5)。 task.tlusty_config.enabled, task.synspec_config.enabled, task.timeout_sec, shutdown, tlusty_input.as_ref(), task.energy_tolerance, task.temp_max_factor, task.temp_floor, task.temp_ceiling, task.emflux_tolerance, task.convergence_min_ratio, task.bfac_max, task.bfac_min, ) .await?; info!( "完成计算任务 {} (网格点: {}, 结果可用: {})", task.task_id, task.point_name, summary.result_valid ); // Read seed bytes if result usable and clean let mut seed_bytes: Option> = None; if summary.result_valid && !summary.atmosphere_has_nan { let model_sub_dir = slot_work_dir.join(&summary.name); let candidates = [ model_sub_dir.join(format!("{}.7", summary.name)), model_sub_dir.join(format!("{}.nl.7", summary.name)), model_sub_dir.join(format!("{}.nc.7", summary.name)), model_sub_dir.join("fort.7"), slot_work_dir.join(format!("{}.7", summary.name)), ]; for cand in &candidates { if cand.is_file() { if let Ok(bytes) = tokio::fs::read(cand).await { info!( "找到网格点 {} 的种子二进制文件: {}", summary.name, cand.display() ); seed_bytes = Some(bytes); break; } } } } info!( "任务 {} 计算完成,沙盒目录: {}", task.task_id, slot_work_dir.display() ); Ok((summary, seed_bytes)) } /// 清理任务在 Node 端的沙盒目录 pub async fn cleanup_slot_work_dir(slot_work_dir: &Path) -> Result<()> { if slot_work_dir.exists() { tokio::fs::remove_dir_all(slot_work_dir).await?; info!("已清理 Node 端沙盒目录: {}", slot_work_dir.display()); } Ok(()) } /// 把单个任务沙盒内的产物拷贝到持久归档目录。 /// /// 采用**白名单**策略([`is_result_worthy`])而非「拷贝所有普通文件」的 catch-all: /// 只保留有语义价值的产物,丢弃 Tlusty/Synspec 运行时产生的中间工作单元 /// (`fort.1/2/3/13/14/18/22/42/44/50/57/69/82/95` 等,旧版 catch-all 会把它们 /// 一并搬进归档,每个模型浪费约 2MB / 4.8MB)。 /// /// 保留内容(详见 [`is_result_worthy`]): /// - 裸名:`conv.json`、`fort.8`(synspec 输入大气)、`fort.55`(synspec 控制卡) /// - 科学核心:`.7/.spec/.cont/.iden/.log`,以及 runner 在 SYNSPEC 覆盖前快照的 /// TLUSTY 最终产物:`.bfac`(b 因子/非 LTE 偏离因子)、`.emflux`(出射谱 λ–Fλ) /// - 阶段快照:`.