use anyhow::Result; use common::config::GridConfig; use common::models::{GridPointParams, PhaseConfig, ResumePolicy, TaskSpec}; use mq::sqlite_queue::SqliteTaskQueue; use std::sync::Arc; use tracing::{info, warn}; use uuid::Uuid; use crate::db::{Database, FallbackSnapshot}; pub struct GridScheduler { db: Database, queue: Arc, /// 调度互斥锁:防止 start_workflow 的即时调度与后台 30s 循环并发进入 /// schedule_pending_tasks,消除 TOCTOU 竞态导致的重复派发(#5 修复)。 schedule_lock: tokio::sync::Mutex<()>, } /// 解析工作流 config_yaml 为 GridConfig;失败时记 warn 日志(避免静默吞错)。 /// `field_hint` 标识调用方(如 "tlusty_chain"/"tlusty_input"),便于日志定位。 /// YAML 格式错误(如 tlusty_input 块缩进错、字段拼错导致反序列化失败)时, /// 旧实现静默返回 None → 用户配置完全不生效且无任何告警。现记 warn 让问题可见。 fn parse_grid_config_or_warn( config_yaml: &str, workflow_name: &str, field_hint: &str, ) -> Option { match serde_yaml::from_str::(config_yaml) { Ok(cfg) => Some(cfg), Err(e) => { warn!( "工作流 {} 的 config_yaml 解析失败(读取 {} 时),回退默认配置: {}", workflow_name, field_hint, e ); None } } } impl GridScheduler { pub fn new(db: Database, queue: Arc) -> Self { Self { db, queue, schedule_lock: tokio::sync::Mutex::new(()), } } /// Expands grid points from config and registers them into the database. /// /// 多工作流分区(#3 修复): /// - 仅清理**本工作流**的排队任务(clear_queue_by_workflow),不再 clear_queue() 全局清空, /// 避免启动工作流 B 时误删工作流 A 的在队任务。 /// - 仅重置**本工作流**的 queued 点为 pending(reset_queued_grid_points_to_pending 带 wf), /// 避免误伤其他工作流。 /// - upsert 带 workflow_name,使同一物理点可属于多个工作流。 pub async fn initialize_grid(&self, cfg: &GridConfig, workflow_name: &str) -> Result<()> { // stop/重启卫生(2026-08-02 涡旋事故修复):clear_queue_by_workflow 删除的 // pending 队列行须同步删除主库 tasks 表对应行,否则遗留的 pending 僵尸行会成为 // 重复派发涡旋的燃料(旧实现只删队列行,三次重启累积近万条僵尸)。 // claimed 队列行保留(#6 设计:在途任务领用凭证不可丢),其 tasks 行一并保留。 match self.queue.clear_queue_by_workflow(workflow_name).await { Ok(cleared_ids) if !cleared_ids.is_empty() => { if let Err(e) = self.db.delete_tasks_by_ids(&cleared_ids).await { tracing::warn!( "初始化工作流 {} 时同步清理 {} 条被删队列任务的 tasks 历史行失败: {}", workflow_name, cleared_ids.len(), e ); } } Err(e) => { tracing::warn!( "初始化工作流 {} 网格时清理该流闲散排队记录发生警告: {}", workflow_name, e ); } _ => {} } if let Err(e) = self .db .reset_queued_grid_points_to_pending(workflow_name) .await { tracing::warn!( "重置工作流 {} 网格状态到 pending 处理过程遇到异常: {}", workflow_name, e ); } let mut points = Vec::new(); for teff in &cfg.grid.teff { for logg in &cfg.grid.logg { for loghe in &cfg.grid.loghe { for logc in &cfg.grid.logc { for logn in &cfg.grid.logn { for logo in &cfg.grid.logo { points.push(GridPointParams { teff: teff.clone(), logg: logg.clone(), loghe: loghe.clone(), logc: logc.clone(), logn: logn.clone(), logo: logo.clone(), }); } } } } } } // Sort by difficulty: cno_sum ASC -> teff ASC -> -logg -> loghe ASC points.sort_by(|a, b| { a.cno_sum() .partial_cmp(&b.cno_sum()) .unwrap_or(std::cmp::Ordering::Equal) .then_with(|| { a.teff .partial_cmp(&b.teff) .unwrap_or(std::cmp::Ordering::Equal) }) .then_with(|| { b.logg .partial_cmp(&a.logg) .unwrap_or(std::cmp::Ordering::Equal) }) .then_with(|| { a.loghe .partial_cmp(&b.loghe) .unwrap_or(std::cmp::Ordering::Equal) }) }); // Group into Waves by cno_sum let mut current_cno: Option = None; let mut wave_idx = 0; for pt in &points { let cno = pt.cno_sum(); if let Some(cur) = current_cno { if (cur - cno).abs() > 1e-5 { wave_idx += 1; current_cno = Some(cno); } } else { current_cno = Some(cno); } // upsert 是幂等的 ON CONFLICT DO NOTHING:若 initialize_grid 中途失败, // 重新 start 该工作流会自然补齐(#4 半初始化回退由幂等性消解)。 self.db .upsert_grid_point(pt, wave_idx, workflow_name) .await?; } // 执行策略(见 docs/task_engine_decoupling_design.md §2.1):策略只决定**启动工作流时** // 对历史终态点的处理;失败后的策略链回退由启动时的策略链(回退优先级排序)独立驱动, // 不受策略门控(2026-08-04 语义修正,对齐 §4.2)。 // // - force_recompute(场景 B/C):任一步骤为强制重算 → 已收敛 + 已失败全部打回 pending, // 无视历史状态与产物全量重算。 // - skip_converged(默认,场景 A):跳过收敛、重试失败 → 仅把**已失败**的点打回 // pending 重试,收敛点保留(修复:此前 SkipConverged 与 SkipFailed 在启动时行为 // 相同,无法表达"只重试失败"这一常用增量语义)。 // - skip_failed:收敛 + 失败都保留(最保守增量,只算从未计算过的点)。 let tlusty_policy = cfg.resolve_tlusty_config().policy; let synspec_policy = cfg.resolve_synspec_config().policy; let tlusty_force = tlusty_policy == ResumePolicy::ForceRecompute; let synspec_force = synspec_policy == ResumePolicy::ForceRecompute; let tlusty_retry = tlusty_policy == ResumePolicy::SkipConverged; let synspec_retry = synspec_policy == ResumePolicy::SkipConverged; if tlusty_force || synspec_force { let n = self .db .reset_terminal_points_for_recompute(workflow_name) .await?; if n > 0 { let mut stages = Vec::new(); if tlusty_force { stages.push("TLUSTY"); } if synspec_force { stages.push("SYNSPEC"); } info!( "工作流 {} 阶段 [{}] 策略为 force_recompute:重置 {} 个终态点回 pending 强制重算", workflow_name, stages.join(" + "), n ); } } else if tlusty_retry || synspec_retry { let n = self.db.reset_failed_points_for_retry(workflow_name).await?; if n > 0 { info!( "工作流 {} 策略为 skip_converged:重置 {} 个失败点回 pending 重试(收敛点保留)", workflow_name, n ); } } // 全部为 skip_failed → 收敛/失败都保留,不重置。 info!( "已在数据库中成功初始化并记录工作流 {} 的 {} 个恒星大气网格点", workflow_name, points.len() ); Ok(()) } /// 读取指定工作流的 timeout_sec(按工作流分区:多工作流各有自己的超时配置)。 async fn get_workflow_timeout_sec(&self, workflow_name: &str) -> u64 { // L5 修复:与其它新 helper 统一用 parse_grid_config_or_warn——YAML 解析失败时 // 记 warn 而非静默回退默认(原 serde_yaml::from_str(...).ok() 吞错,运维无从知晓 // 某工作流 timeout 为何回落到默认 7200)。 if let Some(wf) = self.db.get_workflow(workflow_name).await.ok().flatten() { if let Some(cfg) = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "timeout_sec") { return cfg.timeout_sec; } } 7200 } /// 解析指定工作流的 TLUSTY / SYNSPEC 阶段独立配置(见设计文档 §3)。 /// /// 优先用新版 `tlusty:` / `synspec_stage:` 顶层块;缺省回退到旧字段推断 /// (`seed_step_fallback` / 旧 `synspec` 块存在性),由 `GridConfig::resolve_*` /// 统一口径。解析失败兜底为阶段默认配置,不阻断调度。 async fn get_workflow_stage_configs(&self, workflow_name: &str) -> (PhaseConfig, PhaseConfig) { // L5 修复:与其它新 helper 统一用 parse_grid_config_or_warn——YAML 解析失败时 // 记 warn 而非静默回退默认(原 serde_yaml::from_str(...).ok() 吞错)。 if let Some(wf) = self.db.get_workflow(workflow_name).await.ok().flatten() { if let Some(cfg) = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "stage_configs") { return (cfg.resolve_tlusty_config(), cfg.resolve_synspec_config()); } } ( PhaseConfig::default_tlusty(), PhaseConfig::default_synspec(), ) } /// 读取指定工作流的 SYNSPEC 数值参数(波长范围等),序列化为 JSON Value 供 /// TaskSpec 携带(节点 executor 反序列化为 SynspecInput 透传给 runner)。 /// 工作流未配置 synspec_input 块 → None(runner 用硬编码默认)。 async fn get_workflow_synspec_params(&self, workflow_name: &str) -> Option { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "synspec_params")?; cfg.synspec_input .as_ref() .and_then(|s| serde_json::to_value(s).ok()) } /// 读取指定工作流的 TLUSTY 物理迭代步进链(`config::GridConfig.tlusty_chain`), /// 序列化为 JSON Value 供 TaskSpec 携带。节点 executor 反序列化为 `Vec` /// 后透传给 runner 的 custom_chain 参数,使用户在 YAML 配置的 niter/chmax /// 等阶段参数真正生效(此前 executor 硬编码用 default 链,忽略用户配置)。 /// 工作流未配置 tlusty_chain(空数组)→ None(executor 用 default 链兜底)。 async fn get_workflow_tlusty_chain(&self, workflow_name: &str) -> Option { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "tlusty_chain")?; if cfg.tlusty_chain.is_empty() { return None; } serde_json::to_value(&cfg.tlusty_chain).ok() } /// 读取指定工作流的种子热启动链(`config::GridConfig.seed_chain`),序列化为 JSON /// Value 供 TaskSpec 携带。仅 seed_step 策略下由 executor 读取。 /// 与 `get_workflow_tlusty_chain` 对称。工作流未配置 seed_chain(空数组)→ None ///(executor 用 `default_seed_chain()` 兜底)。 async fn get_workflow_seed_chain(&self, workflow_name: &str) -> Option { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "seed_chain")?; if cfg.seed_chain.is_empty() { return None; } serde_json::to_value(&cfg.seed_chain).ok() } /// 读取指定工作流的 TLUSTY 输入文件全局参数(`config::GridConfig.tlusty_input`), /// 序列化为 JSON Value 供 TaskSpec 携带。包含 NFREAD 频率网格、ions 能级表、 /// nst extra_keys 等不随阶段变化的参数。节点 executor 反序列化为 `TlustyInput` /// 后透传给 runner → make_input5 / generate_nst_content。 /// 工作流未配置 tlusty_input → None(runner 用代码内硬编码默认)。 async fn get_workflow_tlusty_input(&self, workflow_name: &str) -> Option { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "tlusty_input")?; cfg.tlusty_input .as_ref() .and_then(|t| serde_json::to_value(t).ok()) } /// 读取工作流 YAML 的 `linelist` 字段(SYNSPEC 谱线表文件名,如 "gfATO.dat")。 /// None → node 端用默认线表(embedded.rs 启动时下载的 default_linelist)。 async fn get_workflow_linelist(&self, workflow_name: &str) -> Option { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "linelist")?; cfg.linelist.clone() } /// 一次性读取工作流的全部物理校验阈值(能量守恒 / 温度结构 / emflux)。 /// 统一读取避免对同一 YAML 多次解析。返回 8 元组,对应 TaskSpec 的 8 个标量字段: /// (energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, /// convergence_min_ratio, bfac_max, bfac_min)。 async fn get_workflow_validation_thresholds( &self, workflow_name: &str, ) -> Option<( Option, Option, Option, Option, Option, Option, Option, Option, )> { let wf = self.db.get_workflow(workflow_name).await.ok()??; let cfg = parse_grid_config_or_warn(&wf.config_yaml, workflow_name, "validation_thresholds")?; Some(( cfg.energy_tolerance, cfg.temp_max_factor, cfg.temp_floor, cfg.temp_ceiling, cfg.emflux_tolerance, cfg.convergence_min_ratio, cfg.bfac_max, cfg.bfac_min, )) } /// 从策略链解析出「首个可派发」的顺位(见 docs/task_engine_decoupling_design.md §4.2)。 /// /// 判定: /// - `seed_step`:必须存在近邻收敛种子才能热启动;无种子则跳过该顺位,继续看下一项 /// (与失败回退的种子门控同一口径,统一种子注入逻辑)。 /// - 其它策略(`cold_run` 等):无需种子,直接可派发,链保持原样。 /// - 链全被跳过 → 返回 `(空, None)`,调用方应跳过派发(保持 pending 等种子出现后自愈)。 /// /// 修复(审查 #3):初始派发此前直接取 `strategies[0]`——用户把 seed_step 排到链头时, /// 派发的任务带 `seed_point_name=None`,节点 executor 对 seed_step 任务强校验 /// seed_point_name 并 bail → 任务必失败。现初始派发与回退共用本函数,种子注入一致。 /// /// 修复(审查 #6 后续):无种子的 seed_step 顺位**暂存后追加到链尾而非丢弃**——若直接 /// 弹出(旧行为),持久化的剩余链丢失该顺位:cold_run 失败后即使近邻种子已出现 /// (他点收敛),回退也永远用不上 seed_step。追加到链尾把 seed_step 保留为后续回退 /// 顺位,语义为「种子暂不可用 → 先跑当前可派发策略,失败后再试种子步进」。链仅剩 /// seed_step 且无种子时无其它顺位可派,才判空(保持 failed / 打回 pending)。 /// 暂存式处理同时保证终止:重复 seed_step(如 `[seed_step, seed_step]`)不会因 /// 「弹出再压回」保持原序而陷入死循环。 async fn resolve_dispatchable_chain( db: &Database, workflow_name: &str, params: &GridPointParams, mut chain: Vec, ) -> Result<(Vec, Option)> { // 暂存被跳过(无种子)的 seed_step 顺位,确定可派发首项后追加到链尾保留。 let mut skipped: Vec = Vec::new(); loop { match chain.first().map(|s| s.as_str()) { Some("seed_step") => { // 种子查找(审查 #2 加固):`find_best_seed_from_db` 读内存缓存(seed_cache/ // seed_index),当前恒 Ok;但若未来改为 DB 直查,**不会**静默吞错——Err 直接 // 上抛(`?`)中断本轮派发,由下轮调度重试,避免误把「瞬时故障」当「无种子」 // 而降级冷启动。真「无种子」(Ok(None))才跳过 seed_step 顺位。 let found_seed = db.find_best_seed_from_db(params).await?; match found_seed { Some(seed) => { // seed_step 可派发:先前跳过的 seed_step 顺位一并保留到链尾。 if !skipped.is_empty() { chain.extend(skipped); } return Ok((chain, Some(seed.name.clone()))); } None if chain.len() > 1 => { // 无种子 → 暂存该顺位,继续解析下一项。 let head = chain.remove(0); skipped.push(head); info!("策略解析:seed_step 无近邻种子,暂存顺位(待追加回退链尾)"); } // 仅剩 seed_step 一个顺位且无种子 → 无可派发策略。 None => { chain.clear(); return Ok((chain, None)); } } } // 稳定化种子步进(2026-08-18,docs/failed81_cno_seed_popzer_dpsilg_2026_08_18.md): // 仅在 exact_family(同 Teff/logg/logHe、仅 CNO 不同)内找种子——稳定化配方 // (POPZER+DPSILG)的验证前提是种子与目标仅差丰度微扰,跨 Teff/logg 种子 // 不在适用域。排除目标点自身名(防止把自己旧产物当种子)。 // 2026-08-20 种子轮换:配方对种子逐点敏感(同距离换 N/换 O 邻居收敛性不同), // 排除本点历史任务已用过的种子,重试时轮换到下一个未试过的同族邻居。 Some("seed_step_stab") => { let mut exclude = vec![params.model_name()]; exclude.extend(db.list_used_seed_names(¶ms.model_name(), workflow_name).await?); let found_seed = db .find_exact_family_seed_from_db(params, &exclude) .await?; match found_seed { Some(seed) => { if !skipped.is_empty() { chain.extend(skipped); } return Ok((chain, Some(seed.name.clone()))); } None if chain.len() > 1 => { let head = chain.remove(0); skipped.push(head); info!("策略解析:seed_step_stab 无同物理族 CNO 邻居种子,暂存顺位"); } None => { chain.clear(); return Ok((chain, None)); } } } // 空链或非 seed_step 顺位:可直接派发。先前跳过的 seed_step 追加到链尾, // 保留为后续回退顺位(空链由调用方判空处理)。 _ => { if !skipped.is_empty() { chain.extend(skipped); } return Ok((chain, None)); } } } } /// Enqueues pending grid points into MQ with batching(冷启动优先). /// /// 多工作流分区(#3 修复):对**每个** running/initializing 工作流分别派发任务, /// 替代原来「全局只一个 running workflow」的 LIMIT 1 假设。各工作流独立 batch; /// 正常路径一律派发 ColdRun 冷启动,种子匹配只发生在失败后的 /// trigger_tlusty_fallback(seeds 仍是全局共享的物理资源池)。 pub async fn schedule_pending_tasks(&self) -> Result { // 互斥锁:start_workflow 的即时调度与后台 30s 循环可能并发调用本方法, // 两者各自 SELECT 同一批 pending 点会产生重复任务(#5 修复)。 let _guard = self.schedule_lock.lock().await; let workflows = self.db.get_running_workflow_names().await?; if workflows.is_empty() { return Ok(0); } let batch_limit: usize = std::env::var("DCTS_BATCH_LIMIT") .ok() .and_then(|v| v.parse().ok()) .unwrap_or(100); let mut total_dispatched = 0; for wf in &workflows { let dispatched = self .schedule_pending_tasks_for_workflow(wf, batch_limit) .await?; total_dispatched += dispatched; } if total_dispatched > 0 { info!( "已成功将 {} 个待计算网格点推进任务队列(跨 {} 个工作流)", total_dispatched, workflows.len() ); } Ok(total_dispatched) } /// 为单个工作流派发 pending 点。 async fn schedule_pending_tasks_for_workflow( &self, workflow_name: &str, batch_limit: usize, ) -> Result { let timeout_sec = self.get_workflow_timeout_sec(workflow_name).await; let (tlusty_cfg, synspec_cfg) = self.get_workflow_stage_configs(workflow_name).await; let synspec_params = self.get_workflow_synspec_params(workflow_name).await; let tlusty_chain = self.get_workflow_tlusty_chain(workflow_name).await; let seed_chain = self.get_workflow_seed_chain(workflow_name).await; let tlusty_input = self.get_workflow_tlusty_input(workflow_name).await; let linelist = self.get_workflow_linelist(workflow_name).await; let ( energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, convergence_min_ratio, bfac_max, bfac_min, ) = self .get_workflow_validation_thresholds(workflow_name) .await .unwrap_or((None, None, None, None, None, None, None, None)); // 双阶段全关是退化配置(save_workflow 已拦截,此处兜底防御):无可执行阶段, // 整工作流跳过派发(修复审查 #5)。 if !tlusty_cfg.enabled && !synspec_cfg.enabled { tracing::warn!( "工作流 {} 的 TLUSTY 与 SYNSPEC 阶段均关闭,无任务可派发", workflow_name ); return Ok(0); } // 原子选点(#5 修复):IMMEDIATE 事务内完成 SELECT + UPDATE status='queued', // 替代原来 get_pending_grid_points_limit(SELECT)+ update_grid_status(UPDATE) // 的分离操作,杜绝两个并发调度调用 SELECT 到同一批 pending 点的 TOCTOU 竞态。 let pending = self .db .claim_pending_grid_points(batch_limit, workflow_name) .await?; let mut dispatched = 0; for (name, params, wave) in pending { // 派发去重 + 僵尸自愈(2026-08-02 涡旋事故修复): // 若该点在 tasks 表仍有 pending 行,逐个做 MQ 活性交叉校验—— // - 任一活(队列行 pending/claimed):任务真在途(如 requeue_stale_tasks // 同 task_id 重投后点被重置回 pending 的情形),跳过本次派发,点保持 // queued,由在途任务的上报自然结算; // - 全死:僵尸行(典型为 insert_task 后、push_task 前崩溃遗留,点卡 queued // 且永无上报),清除后落回正常派发,点由此自愈。 // 不变量:任意时刻每个点至多一个活任务。 let zombies = self .db .has_pending_tasks_for_point(&name, workflow_name, None) .await?; if !zombies.is_empty() { let mut any_alive = false; for id in &zombies { if self.queue.task_row_exists(id).await? { any_alive = true; break; } } if any_alive { tracing::info!( "网格点 {} 已有在途任务({} 条活队列行关联),跳过重复派发", name, zombies.len() ); continue; } let removed = self.db.delete_tasks_by_ids(&zombies).await?; tracing::info!( "网格点 {} 清除 {} 条无队列凭证的僵尸任务行,重新派发", name, removed ); } // 阶段独立配置派发(见 docs/task_engine_decoupling_design.md §4.2): // 节点读取 tlusty_config.strategies[0] 决定 TLUSTY 执行链。正常路径派发 // 完整策略链(如 [cold_run, seed_step]),失败后由服务端弹出首项回退。 //(Phase 6 起无 task_type 兼容字段,策略链是唯一权威。) // // 修复(审查 #3):用 resolve_dispatchable_chain 解析「首个可派发」顺位—— // 链头为 seed_step 时须注入近邻种子;无种子则跳过该顺位。此前直接取 // strategies[0] 导致「seed_step 链头 + 无种子」的任务在节点端必失败。 // // H1 活锁修复:若该点被运行时回退(H1)打回 pending 并记录了剩余策略链标记 // (pending_strategies,见 trigger_strategy_fallback),则用该剩余链作为基础链 // 重新解析——跳过已失败的策略(如 cold_run),等种子出现才派 seed_step,避免 // 重跑已失败策略导致的无界失败重试活锁。force_recompute/skip_converged 启动 // 重置不设标记 → 走完整 YAML 链(保持各自语义)。 let base_chain: Vec = match self .db .get_pending_strategies(&name, workflow_name) .await? .and_then(|j| serde_json::from_str::>(&j).ok()) { Some(c) if !c.is_empty() => { tracing::info!( "网格点 {}(工作流 {})带剩余策略链标记重派 {:?}(跳过已失败策略,H1 活锁规避)", name, workflow_name, c ); c } _ => tlusty_cfg.strategies.clone(), }; let (dispatch_chain, seed_point_name) = if tlusty_cfg.enabled { let (chain, seed) = Self::resolve_dispatchable_chain(&self.db, workflow_name, ¶ms, base_chain).await?; if chain.is_empty() { // enabled 阶段无可派发顺位(基础链全为无种子的 seed_step)→ 打回 pending, // 待近邻种子出现后由下轮调度自愈派发(不是终态,不能被卡死在 queued)。 // **不**清除 pending_strategies 标记:若基础链来自 H1 标记且仍无种子, // 保留标记,下轮继续等种子(不清除则不会重跑已失败策略)。 tracing::warn!( "网格点 {}(工作流 {})的 TLUSTY 策略链无可派发顺位,打回 pending 待种子出现", name, workflow_name ); self.db .update_grid_status( &name, common::models::GridPointStatus::Pending, workflow_name, ) .await?; continue; } // 实际派发时消费并清除 H1 剩余链标记(下轮若再失败由回退重新设置)。 self.db .clear_pending_strategies(&name, workflow_name) .await?; (chain, seed) } else { // TLUSTY 关闭(仅 SYNSPEC 场景):策略链不参与门控,大气来自既有产物。 // 链归一化为 ["cold_run"]——tlusty 关闭时策略链无执行语义,若保留用户 // 完整链(如链头 seed_step),DB 行会误导「将执行大气链」。归一化后 // 归因/统计自洽(归因取 synspec 链首项,见 record_task_report tlusty/synspec 阶段归因)。 (vec!["cold_run".to_string()], None) }; let mut next_tlusty_cfg = tlusty_cfg.clone(); next_tlusty_cfg.strategies = dispatch_chain.clone(); let task_spec = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params, seed_point_name, timeout_sec, workflow_name: Some(workflow_name.to_string()), wave, tlusty_config: next_tlusty_cfg, synspec_config: synspec_cfg.clone(), synspec_params: synspec_params.clone(), tlusty_chain_params: tlusty_chain.clone(), seed_chain_params: seed_chain.clone(), tlusty_input_params: tlusty_input.clone(), // 显式绑定大气来源(设计 §5.2,修复审查 #3):仅 SYNSPEC-only(TLUSTY 关闭) // 场景需要外部大气——节点凭 atmosphere_ref(或 point_name 兜底)从本地归档 // /服务端拉取目标点 .7。TLUSTY 启用的任务大气在沙盒内自产,无需外部引用。 atmosphere_ref: if tlusty_cfg.enabled { None } else { Some(name.clone()) }, energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, convergence_min_ratio, bfac_max, bfac_min, linelist: linelist.clone(), }; self.db.insert_task(&task_spec).await?; // 点已在 claim_pending_grid_points 的 IMMEDIATE 事务中原子标记为 queued, // 无需再单独 update_grid_status(Queued)。push 失败时回滚为 pending 即可。 match self.queue.push_task(&task_spec).await { Ok(_) => { dispatched += 1; } Err(e) => { tracing::warn!( "将任务 {} 推入 MQ 队列失败,执行严格状态回滚以避免脏数据: {}", name, e ); if let Err(db_e) = self .db .update_grid_status( &name, common::models::GridPointStatus::Pending, workflow_name, ) .await { tracing::error!( "关键性回滚异常:任务 {} 无法重置回 Pending: {}", name, db_e ); } let _ = self.queue.remove_task(&task_spec.task_id.to_string()).await; // 同步清理先于 push 插入的 tasks 历史行,避免遗留 pending 历史记录 // 污染每点尝试计数统计(attempt_count 依赖 tasks 表聚合)。 let _ = self.db.delete_task(&task_spec.task_id).await; } } } Ok(dispatched) } /// 回收孤儿网格点(#6 修复兜底,2026-08-02 涡旋事故重构版)。 /// /// 场景:网格点卡在 running/queued,但其任务已无任何队列凭证(queue 行被误删、 /// server 崩溃丢队列、insert_task 后 push_task 前崩溃等),节点无法上报,点永远 /// 到不了终态。与 requeue_stale_tasks 互补:后者处理 queue 行仍存在的超时 claimed。 /// /// 判定流程(跨库活性交叉校验,取代旧版"存在老 pending tasks 行即孤儿"的危险 /// 代理判据——僵尸行使旧判据恒真,是 2026-08-02 重复派发涡旋的根源): /// 1. `find_stale_pending_points`:列出 running/queued 点中创建老于 stale_sec 的 /// pending tasks 行(候选,非孤儿); /// 2. 对每个点的每条候选行查 `queue.task_row_exists`: /// - 死行(无队列凭证)一律清除(僵尸卫生); /// - 只要还有活行(requeue 重投的同 task_id 行、正常长任务),该点**不重置**, /// 由在途任务自然结算; /// - 全部为死行 → 真孤儿,`rescue_orphaned_point` 置回 pending,下轮派发接管。 /// /// 返回被重置回 pending 的孤儿点数。 pub async fn reclaim_orphaned_points(&self, stale_sec: u64) -> Result { // 互斥锁:与 schedule_pending_tasks / trigger_strategy_fallback 共享调度锁, // 三者都会修改「点→任务」映射(insert_task + update_grid_status + push_task), // 必须互斥以消除 check-then-act TOCTOU 竞态(审查修复:此前仅 schedule 持锁, // reclaim 与 fallback 都裸跑,可派发重复任务违反「每点至多一个活任务」不变量)。 let _guard = self.schedule_lock.lock().await; let candidates = self.db.find_stale_pending_points(stale_sec).await?; if candidates.is_empty() { return Ok(0); } // 按 (name, wf) 分组候选 task_id(find 返回扁平行)。 let mut groups: std::collections::BTreeMap<(String, String), Vec> = std::collections::BTreeMap::new(); for (name, wf, task_id) in candidates { groups.entry((name, wf)).or_default().push(task_id); } let mut rescued = 0usize; for ((name, wf), task_ids) in groups { let mut dead = Vec::new(); let mut any_alive = false; for id in &task_ids { if self.queue.task_row_exists(id).await? { any_alive = true; } else { dead.push(id.clone()); } } // 死行一律清除(即使点仍有活任务,僵尸行也该清)。 if !dead.is_empty() { match self.db.delete_tasks_by_ids(&dead).await { Ok(n) if n > 0 => { tracing::info!( "孤儿回收:清除网格点 {}(工作流 {})的 {} 条无队列凭证僵尸任务行", name, wf, n ); } Err(e) => { tracing::warn!("孤儿回收:清理网格点 {} 僵尸任务行失败: {}", name, e); } _ => {} } } // 仅当全部候选行都无队列凭证时才判定为真孤儿并重置。 if !any_alive { match self.db.rescue_orphaned_point(&name, &wf).await { Ok(true) => { rescued += 1; tracing::info!( "孤儿回收:网格点 {}(工作流 {})领用凭证全部丢失,已重置为 pending 待重派", name, wf ); } Ok(false) => {} // 并发上报已置终态,无需处理 Err(e) => { tracing::warn!("孤儿回收:重置网格点 {} 失败: {}", name, e); } } } } Ok(rescued) } /// 策略链自动回退(见 docs/task_engine_decoupling_design.md §4.2)。 /// /// 替代旧版「种子回退仅一次」的硬编码逻辑:读取该网格点最近 tasks 行的策略链数组 /// (TLUSTY 或 SYNSPEC,由 `failed_stage` 决定),弹出已失败的首项策略,若剩余链 /// 非空则重置网格点为 queued 并派发下一顺位策略的新任务;剩余链为空则策略耗尽, /// 保持 failed。 /// /// 工作流: /// 1. 该工作流须仍 running 才考虑回退; /// 2. 仅 failed 点回退(状态守卫,沿用 2026-08-02 涡旋事故修复); /// 3. 清理无队列凭证的 pending 僵尸行(沿用结构性清僵尸); /// 4. `SkipFailed` 策略(见 §2.1):忽略历史失败 → 不回退重试(修复审查 #1, /// 此前 policy 仅存储不被消费); /// 5. 按失败阶段弹对应策略链: /// - `failed_stage="synspec"`:弹 `synspec_strategies`,重试光谱合成。无种子门控 /// ——SYNSPEC 输入大气来自目标点自身(修复审查 #2:旧实现误弹 TLUSTY 链并做 /// 无意义的邻居种子搜索,可能把可重试的 synspec 瞬态失败误判为不可回退)。 /// 半失败点(大气已收敛 + 光谱失败,修复审查 #2 后续):重试任务强制关闭 /// TLUSTY(大气产物有效,不重复计算),凭 `atmosphere_ref` 复用既有 .7; /// - 其余(含旧节点缺省):弹 `tlusty_strategies`。新首项为 seed_step 时在全局 /// 网格找近邻种子注入 seed_point_name(找不到则跳过该顺位),cold_run 直接派发。 /// 6. 不修改原 policy(保持用户初始配置,避免状态污染)。 /// /// 已知竞态(审查 #5,设计取舍,仅文档化):回退任务的**运行参数**(`synspec_cfg` / /// `timeout_sec` / `synspec_params`)取自**当前**工作流 YAML,而**策略链与 policy** /// 取自**派发时**落库的 DB 快照([`FallbackSnapshot`])。若用户在派发与回退之间修改 /// 了 YAML,重试任务以「新运行参数 + 旧链/旧 policy」执行——策略链与 policy 是任务级 /// 审计/回退语义记录必须用派发时的(对齐 §4.2「保持用户初始配置」),运行参数允许 /// 跟随最新配置(用户编辑的意图优先)。 /// /// Phase 7b 改名:原名 `trigger_seed_step_fallback`(TLUSTY-first 残留,易误读为"只处理 /// 种子步进");实际是 TLUSTY 阶段失败回退(等价于传 `failed_stage="tlusty"`)。api/task.rs /// 已直调 `trigger_strategy_fallback` 做精确归因,本方法仅测试用(所有调用点在 mod tests 内)。 #[cfg(test)] pub async fn trigger_tlusty_fallback( &self, params: &GridPointParams, name: &str, workflow_name: &str, ) -> Result { self.trigger_strategy_fallback(params, name, workflow_name, "tlusty") .await } /// 策略链回退的真正实现(见上方法文档)。`failed_stage` 为 `"tlusty"` / `"synspec"` /// (旧节点上报缺省 → 调用方兜底传 `"tlusty"`,行为与旧版一致)。 pub async fn trigger_strategy_fallback( &self, params: &GridPointParams, name: &str, workflow_name: &str, failed_stage: &str, ) -> Result { // 互斥锁:与 schedule_pending_tasks / reclaim_orphaned_points 共享调度锁。 // 失败上报(report_task API 在 HTTP 线程)可能并发到达同一网格点的多条报告, // 若不互斥,两次 fallback 各自读到相同的完整策略链、各自弹首项、各自 insert_task, // 派发两个回退任务(审查修复:非原子弹栈 + 无锁 TOCTOU 的叠加效应)。 let _guard = self.schedule_lock.lock().await; // 该工作流须仍处于 running 态才回退(避免 stop 后继续派发) let still_running = self .db .get_running_workflow_names() .await? .iter() .any(|w| w == workflow_name); if !still_running { return Ok(false); } // 状态守卫(2026-08-02 涡旋事故修复):仅 failed 点才回退。 match self.db.get_grid_point_status(name, workflow_name).await? { Some((status, _)) if status == "failed" => {} _ => return Ok(false), } // 僵尸卫生:清理该点无队列凭证的 pending 行(不论策略), // 避免遗留僵尸行干扰下面的「在途去重」。沿用结构性清僵尸口径。 let pending_all = self .db .has_pending_tasks_for_point(name, workflow_name, None) .await?; if !pending_all.is_empty() { let mut dead = Vec::new(); let mut any_alive = false; for id in &pending_all { if self.queue.task_row_exists(id).await? { any_alive = true; } else { dead.push(id.clone()); } } if !dead.is_empty() { let removed = self.db.delete_tasks_by_ids(&dead).await?; if removed > 0 { info!( "策略回退:网格点 {} 清除 {} 条无队列凭证的僵尸任务行", name, removed ); } } // 任一活行 → 任务真在途(如 requeue 重投),跳过本次回退。 if any_alive { info!( "策略回退:网格点 {} 已有在途任务(活队列行),跳过本次回退", name ); return Ok(false); } } // 阶段配置(enabled / 其它阶段 / 数值参数)取自当前工作流 YAML;策略链与 policy // 取自 DB 派发时快照(见 db.rs FallbackSnapshot 文档)——回退决策用派发时配置, // 不随运行期 YAML 编辑漂移(对齐设计 §4.2)。 let (tlusty_cfg, synspec_cfg) = self.get_workflow_stage_configs(workflow_name).await; let stage_label = if failed_stage == "synspec" { "SYNSPEC" } else { "TLUSTY" }; // 读取最新已上报行的对应策略链并弹出首项(只读 pop,不改写旧行——见 db.rs 方法文档), // 同时带回该行落库的 policy。若无已上报行或策略链耗尽(仅 1 项或为空)→ 保持 failed。 let pop_res = if failed_stage == "synspec" { self.db .pop_stage_strategy_for_fallback(name, workflow_name, "synspec") .await? } else { self.db .pop_tlusty_strategy_for_fallback(name, workflow_name) .await? }; let Some(snap) = pop_res else { info!( "策略回退:网格点 {} 无已上报 tasks 行可弹({} 链),保持 failed", name, stage_label ); return Ok(false); }; let FallbackSnapshot { rest_strategies, popped, policy: row_policy, } = snap; // 2026-08-04 语义修正:策略链回退**只由策略链驱动**(弹出失败首项、派发下一顺位, // 链耗尽保持 failed),不再受执行策略门控——策略只决定启动工作流时对历史终态点的 // 处理(见 initialize_grid 三分支)。`row_policy` 仍随重试任务落库供审计/启动时消费。 // 若用户不想重试失败点,用单策略链(如 ["cold_run"])+ SkipFailed 表达。 let timeout_sec = self.get_workflow_timeout_sec(workflow_name).await; let synspec_params = self.get_workflow_synspec_params(workflow_name).await; let tlusty_chain = self.get_workflow_tlusty_chain(workflow_name).await; let seed_chain = self.get_workflow_seed_chain(workflow_name).await; let tlusty_input = self.get_workflow_tlusty_input(workflow_name).await; let linelist = self.get_workflow_linelist(workflow_name).await; let ( energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, convergence_min_ratio, bfac_max, bfac_min, ) = self .get_workflow_validation_thresholds(workflow_name) .await .unwrap_or((None, None, None, None, None, None, None, None)); // 2026-08-20 种子轮换:seed_step_stab 常为链上最后一项,失败弹出后链空。 // 稳定化配方对种子逐点敏感(同 CNO 距离、换不同元素方向的邻居收敛性不同, // 见 docs/failed81_cno_seed_popzer_dpsilg_2026_08_18.md §九)——首个选中 // 的邻居失败不代表配方无效。排除本点历史已用种子后若还有未试过的同族 // 邻居,则以 seed_step_stab 单顺位重链重派。轮换次数天然受同族邻居数约束: // 全部试过后此处不再命中,落入下方 failed 终态,无死循环。 if rest_strategies.is_empty() && failed_stage != "synspec" && popped == "seed_step_stab" { let mut exclude = vec![params.model_name()]; exclude.extend(self.db.list_used_seed_names(name, workflow_name).await?); if let Some(fresh_seed) = self .db .find_exact_family_seed_from_db(params, &exclude) .await? { info!( "策略回退:工作流 {} 网格点 {} seed_step_stab 失败但存在未试过的同族邻居种子 {},轮换种子重派(已试过 {:?})", workflow_name, name, fresh_seed.name, &exclude[1..] ); // 轮换任务显式注入新种子(绕过 resolve 的种子查找,避免其再次 // 命中同一邻居)。 let mut stab_cfg = tlusty_cfg.clone(); stab_cfg.strategies = vec!["seed_step_stab".to_string()]; stab_cfg.policy = row_policy.clone(); let task_spec = TaskSpec { task_id: Uuid::new_v4(), point_name: name.to_string(), params: params.clone(), seed_point_name: Some(fresh_seed.name), timeout_sec, workflow_name: Some(workflow_name.to_string()), wave: 0, tlusty_config: stab_cfg, synspec_config: synspec_cfg.clone(), synspec_params, tlusty_chain_params: tlusty_chain.clone(), seed_chain_params: seed_chain.clone(), tlusty_input_params: tlusty_input.clone(), atmosphere_ref: None, energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, convergence_min_ratio, bfac_max, bfac_min, linelist: linelist.clone(), }; self.db.insert_task(&task_spec).await?; self.db .update_grid_status(name, common::models::GridPointStatus::Queued, workflow_name) .await?; if let Err(e) = self.queue.push_task(&task_spec).await { let _ = self .db .update_grid_status( name, common::models::GridPointStatus::Pending, workflow_name, ) .await; let _ = self.queue.remove_task(&task_spec.task_id.to_string()).await; let _ = self.db.delete_task(&task_spec.task_id).await; return Err(e); } return Ok(true); } } if rest_strategies.is_empty() { info!( "策略回退:网格点 {} 的 {} 策略链已耗尽(弹出最后项 {}),保持 failed 终态", name, stage_label, popped ); return Ok(false); } // SYNSPEC 链回退:重试光谱合成。无邻居种子门控(大气来自目标点自身既有产物, // 见 docs/task_engine_decoupling_design.md §5)——旧实现把 synspec 失败误归因到 // TLUSTY 链并做无意义的种子搜索(修复审查 #2)。 if failed_stage == "synspec" { let mut next_synspec_cfg = synspec_cfg.clone(); next_synspec_cfg.strategies = rest_strategies.clone(); // 回退任务继承派发时 policy(§4.2 保持用户初始配置,不随 YAML 漂移)。 next_synspec_cfg.policy = row_policy.clone(); // 半失败点重试(修复审查 #2 后续):能走到 synspec 链回退意味着大气已产出 // (无论是否收敛——大气失败会归因 tlusty)。重试只重跑光谱,不再重算大气: // 强制关闭 TLUSTY 阶段(复用既有 .7 产物),与场景 B「仅更新光谱」同一执行路径。 // // 审查修复:保留原始完整 TLUSTY 策略链(而非归一化为 ["cold_run"])—— // pop_stage_strategy_for_fallback 读取「最新已上报行」(ORDER BY created_at DESC), // 若把 retry 行的 tlusty_strategies 改写为 ["cold_run"],后续该 retry 失败且节点 // 把 failed_stage 误报为 "tlusty" 时,弹出 cold_run 后链空 → 误判「不可回退」, // 丢失原本的 TLUSTY 策略链回退能力。executor 已凭 enabled=false 跳过大气计算, // strategies 仅作审计/潜在回退依据,保留原链不改变执行行为。 let mut retry_tlusty_cfg = tlusty_cfg.clone(); retry_tlusty_cfg.enabled = false; let task_spec = TaskSpec { task_id: Uuid::new_v4(), point_name: name.to_string(), params: params.clone(), // synspec-only 任务大气来自既有产物(新节点以 tlusty_config.enabled=false // 跳过大气计算),种子来源为空。 seed_point_name: None, timeout_sec, workflow_name: Some(workflow_name.to_string()), // 回退/重试任务 wave 设 0:出队 ORDER BY wave ASC(sqlite_queue.rs)→ 0 为最高派发优先级, // 重试点优先于正常 pending 重派(重试往往已有部分计算/种子等待)。若设计意图是 // 重试不插队,应改设高 wave 而非 0——此处 0 是"优先",不是"不抢占"。 wave: 0, tlusty_config: retry_tlusty_cfg, synspec_config: next_synspec_cfg, synspec_params, // TLUSTY 已关闭(半失败重试只重跑光谱),不执行 chain/input → None。 tlusty_chain_params: None, seed_chain_params: None, tlusty_input_params: None, atmosphere_ref: Some(name.to_string()), // 不重算大气 → 不做物理正确性校验。 energy_tolerance: None, temp_max_factor: None, temp_floor: None, temp_ceiling: None, emflux_tolerance: None, convergence_min_ratio: None, bfac_max: None, bfac_min: None, linelist: linelist.clone(), }; self.db.insert_task(&task_spec).await?; self.db .update_grid_status(name, common::models::GridPointStatus::Queued, workflow_name) .await?; if let Err(e) = self.queue.push_task(&task_spec).await { let _ = self .db .update_grid_status( name, common::models::GridPointStatus::Pending, workflow_name, ) .await; let _ = self.queue.remove_task(&task_spec.task_id.to_string()).await; let _ = self.db.delete_task(&task_spec.task_id).await; return Err(e); } info!( "策略回退:工作流 {} 网格点 {} SYNSPEC 失败,弹出策略 {},派发下一顺位 {:?}(剩余链 {:?})", workflow_name, name, popped, rest_strategies[0], rest_strategies ); return Ok(true); } // TLUSTY 链回退:找下一顺位可派发策略(seed_step 需注入近邻种子;无种子跳过)。 // 注:纯内存操作(不再调用 DB pop)——调度器拿到 rest 后构造新 TaskSpec, // insert_task 写入新行携带剩余链,旧行保持原状(审计正确)。 // pending_rest:resolve 前的基础剩余链(含全部 seed_step 顺位),H1 空链分支据此记录标记。 let pending_rest = rest_strategies.clone(); let (rest_strategies, seed_point_name) = Self::resolve_dispatchable_chain(&self.db, workflow_name, params, rest_strategies).await?; if rest_strategies.is_empty() { // H1 修复:resolve_dispatchable_chain 仅在「剩余顺位全为无近邻种子的 seed_step」 // 时返回空(调用方已在上方 line 790 排除了真耗尽)。此时**不能**保持 failed 终态: // failed 是吸收态,reclaim/调度都不再触碰,点会永久丢失(即使邻居随后收敛、种子 // 出现也无法救回)。与调度正常派发路径(scheduler.rs 472-486 行)同口径,打回 // pending,待近邻种子出现后由下轮调度自愈派发 seed_step。 tracing::warn!( "策略回退:网格点 {} 剩余策略链全为无近邻种子的 seed_step,打回 pending 待种子出现", name ); self.db .update_grid_status( name, common::models::GridPointStatus::Pending, workflow_name, ) .await?; // H1 活锁修复:记录剩余策略链标记(如 ["seed_step"])。若无此标记,下轮调度会 // 用**完整 YAML 链**重派 → 已失败的 cold_run 被重跑 → 再失败 → 再打回 pending → // 无界失败重试活锁。标记后调度路径用剩余链解析,等种子出现才派 seed_step,不再 // 重跑 cold_run。注意:force_recompute/skip_converged 启动重置不设此标记(走完整链)。 if let Ok(json) = serde_json::to_string(&pending_rest) { self.db .set_pending_strategies(name, workflow_name, &json) .await?; } return Ok(false); } let next_strategy = rest_strategies[0].clone(); let mut next_tlusty_cfg = tlusty_cfg.clone(); next_tlusty_cfg.strategies = rest_strategies.clone(); // 回退任务继承派发时 policy(§4.2 保持用户初始配置,不随 YAML 漂移)。 next_tlusty_cfg.policy = row_policy.clone(); // name 取自权威的 report.point_name(源精度正确,详见旧版同名注释)。 let task_spec = TaskSpec { task_id: Uuid::new_v4(), point_name: name.to_string(), params: params.clone(), seed_point_name, timeout_sec, workflow_name: Some(workflow_name.to_string()), // 回退/重试任务 wave 设 0:出队 ORDER BY wave ASC(sqlite_queue.rs)→ 0 为最高派发优先级, // 重试点优先于正常 pending 重派(重试往往已有部分计算/种子等待)。若设计意图是 // 重试不插队,应改设高 wave 而非 0——此处 0 是"优先",不是"不抢占"。 wave: 0, tlusty_config: next_tlusty_cfg, synspec_config: synspec_cfg.clone(), synspec_params, tlusty_chain_params: tlusty_chain.clone(), seed_chain_params: seed_chain.clone(), tlusty_input_params: tlusty_input.clone(), atmosphere_ref: None, energy_tolerance, temp_max_factor, temp_floor, temp_ceiling, emflux_tolerance, convergence_min_ratio, bfac_max, bfac_min, linelist: linelist.clone(), }; self.db.insert_task(&task_spec).await?; self.db .update_grid_status(name, common::models::GridPointStatus::Queued, workflow_name) .await?; if let Err(e) = self.queue.push_task(&task_spec).await { let _ = self .db .update_grid_status( name, common::models::GridPointStatus::Pending, workflow_name, ) .await; let _ = self.queue.remove_task(&task_spec.task_id.to_string()).await; let _ = self.db.delete_task(&task_spec.task_id).await; return Err(e); } info!( "策略回退:工作流 {} 网格点 {} 弹出失败策略 {},派发下一顺位 {}(剩余链 {:?})", workflow_name, name, popped, next_strategy, rest_strategies ); Ok(true) } } #[cfg(test)] mod tests { use super::*; use common::config::GridAxesConfig; #[tokio::test] async fn test_grid_scheduler_initialization_and_scheduling() { let temp_dir = tempfile::tempdir().unwrap(); let db_path = temp_dir.path().join("sched_db.db"); let queue_db_path = temp_dir.path().join("sched_queue.db"); let db = Database::new(&db_path.to_string_lossy()).await.unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&queue_db_path.to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); #[allow(deprecated)] // results 是死字段,构造时必须填 None let cfg = GridConfig { grid: GridAxesConfig { teff: vec![35000.0.into()], logg: vec![5.5.into()], loghe: vec![(-1.0).into()], logc: vec![(-2.0).into()], logn: vec![(-2.0).into()], logo: vec![(-2.0).into()], }, tlusty_chain: vec![], seed_chain: vec![], tlusty_input: None, synspec_input: None, nworkers: 4, timeout_sec: 3600, resume: true, seed_step_fallback: true, niter: Some(100), template: None, fort55: None, linelist: None, tlusty_stage: None, synspec_stage: None, energy_tolerance: None, temp_max_factor: None, temp_floor: None, temp_ceiling: None, emflux_tolerance: None, convergence_min_ratio: None, bfac_max: None, bfac_min: None, }; scheduler.initialize_grid(&cfg, "test_wf").await.unwrap(); db.upsert_workflow("test_wf", None, "", "running") .await .unwrap(); let pending = db.get_pending_grid_points("test_wf").await.unwrap(); assert_eq!(pending.len(), 1); let dispatched = scheduler.schedule_pending_tasks().await.unwrap(); assert_eq!(dispatched, 1); let popped = queue.pop_task("test-node").await.unwrap(); assert!(popped.is_some()); } /// 多工作流分区调度测试(#3 修复验证): /// 1. wf_a 调度推入队列的任务,在初始化 wf_b 后依然存在(initialize_grid 改用 /// clear_queue_by_workflow,不再全局 clear_queue)。 /// 2. 两个 running 工作流的 pending 点都能被 schedule_pending_tasks 派发。 #[tokio::test] async fn test_multi_workflow_dispatch_isolation() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("mw_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("mw_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); #[allow(deprecated)] // results 是死字段,构造时必须填 None let mk_cfg = |teff: f64| GridConfig { grid: GridAxesConfig { teff: vec![teff.into()], logg: vec![5.5.into()], loghe: vec![(-1.0).into()], logc: vec![(-2.0).into()], logn: vec![(-2.0).into()], logo: vec![(-2.0).into()], }, tlusty_chain: vec![], seed_chain: vec![], tlusty_input: None, synspec_input: None, nworkers: 4, timeout_sec: 3600, resume: true, seed_step_fallback: true, niter: Some(100), template: None, fort55: None, linelist: None, tlusty_stage: None, synspec_stage: None, energy_tolerance: None, temp_max_factor: None, temp_floor: None, temp_ceiling: None, emflux_tolerance: None, convergence_min_ratio: None, bfac_max: None, bfac_min: None, }; // wf_a 初始化并推入队列 scheduler .initialize_grid(&mk_cfg(35000.0), "wf_a") .await .unwrap(); db.upsert_workflow("wf_a", None, "", "running") .await .unwrap(); let d_a = scheduler.schedule_pending_tasks().await.unwrap(); assert_eq!(d_a, 1); // 任务已在队 let popped_first = queue.pop_task("node-a").await.unwrap(); assert!(popped_first.is_some()); // 模拟在途任务结束(2026-08-02 派发去重修复后:claimed 队列行存在即视为在途, // 调度器不再为同点重复派发;须先 remove_task 清凭证,reset 后的重派才会生效)。 queue .remove_task(&popped_first.unwrap().task_id.to_string()) .await .unwrap(); // 重新推一个 wf_a 任务(上一行 pop 掉了),再初始化 wf_b db.update_grid_status( &GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), } .model_name(), common::models::GridPointStatus::Pending, "wf_a", ) .await .unwrap(); let _ = scheduler .schedule_pending_tasks_for_workflow("wf_a", 100) .await .unwrap(); // 此时 wf_a 队列里应有一个任务 let popped_second = queue.pop_task("node-a").await.unwrap(); assert_eq!( popped_second.as_ref().and_then(|t| t.workflow_name.clone()), Some("wf_a".to_string()) ); // 同上:清掉领用凭证,使后续 reset + 重派能产生新的在队任务。 queue .remove_task(&popped_second.unwrap().task_id.to_string()) .await .unwrap(); // 关键断言:把 wf_a 任务重新推回队列后,初始化 wf_b 不应清空它。 db.update_grid_status( &GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), } .model_name(), common::models::GridPointStatus::Pending, "wf_a", ) .await .unwrap(); let _ = scheduler .schedule_pending_tasks_for_workflow("wf_a", 100) .await .unwrap(); // 初始化 wf_b(内部 clear_queue_by_workflow("wf_b"),不该动 wf_a 的任务) scheduler .initialize_grid(&mk_cfg(40000.0), "wf_b") .await .unwrap(); db.upsert_workflow("wf_b", None, "", "running") .await .unwrap(); // wf_a 的任务仍在队:可被 node 弹出,且 workflow_name == wf_a let popped_a = queue.pop_task("node-a").await.unwrap(); assert!(popped_a.is_some(), "初始化 wf_b 不应清空 wf_a 的队列任务"); assert_eq!(popped_a.unwrap().workflow_name, Some("wf_a".to_string())); // wf_b 的点也能被调度(两个 running 工作流并存) let d_b = scheduler.schedule_pending_tasks().await.unwrap(); assert!(d_b >= 1, "wf_b 的 pending 点应被派发"); } /// 测试夹具:单点网格配置。 fn mk_single_point_cfg(teff: f64) -> GridConfig { #[allow(deprecated)] // results 是死字段,构造时必须填 None let cfg = GridConfig { grid: GridAxesConfig { teff: vec![teff.into()], logg: vec![5.5.into()], loghe: vec![(-1.0).into()], logc: vec![(-2.0).into()], logn: vec![(-2.0).into()], logo: vec![(-2.0).into()], }, tlusty_chain: vec![], seed_chain: vec![], tlusty_input: None, synspec_input: None, nworkers: 4, timeout_sec: 3600, resume: true, seed_step_fallback: true, niter: Some(100), template: None, fort55: None, linelist: None, tlusty_stage: None, synspec_stage: None, energy_tolerance: None, temp_max_factor: None, temp_floor: None, temp_ceiling: None, emflux_tolerance: None, convergence_min_ratio: None, bfac_max: None, bfac_min: None, }; cfg } fn mk_cold_spec(point: &str, params: &GridPointParams, wf: &str) -> TaskSpec { TaskSpec { task_id: Uuid::new_v4(), point_name: point.to_string(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some(wf.to_string()), wave: 0, ..Default::default() } } /// D 派发去重(2026-08-02 涡旋事故修复):点在 tasks 表有活任务行(队列行存活) /// 时,调度器跳过重复派发,点保持 queued,在途任务自行结算。 #[tokio::test] async fn test_schedule_skips_point_with_live_task_row() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("skip_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("skip_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_skip").await.unwrap(); db.upsert_workflow("wf_skip", None, "", "running") .await .unwrap(); // 模拟 requeue 残留:tasks 行 pending + 队列行存活(pending 未领用)。 let live = mk_cold_spec(&name, ¶ms, "wf_skip"); db.insert_task(&live).await.unwrap(); queue.push_task(&live).await.unwrap(); let dispatched = scheduler .schedule_pending_tasks_for_workflow("wf_skip", 100) .await .unwrap(); assert_eq!(dispatched, 0, "有活任务行时不应重复派发"); // 点被 claim 置 queued 后保持(在途任务上报自然结算)。 let (status, _) = db .get_grid_point_status(&name, "wf_skip") .await .unwrap() .unwrap(); assert_eq!(status, "queued"); // 队列里仍是原任务(pop 得到原 task_id)。 let popped = queue.pop_task("node-x").await.unwrap().unwrap(); assert_eq!(popped.task_id, live.task_id); } /// D 僵尸自愈(2026-08-02 涡旋事故修复,3b 场景):tasks 行 pending 但无任何 /// 队列行(insert_task 后 push_task 前崩溃遗留)→ 清除僵尸并正常派发新任务。 #[tokio::test] async fn test_schedule_cleans_zombie_and_redispatches() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("zombie_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("zombie_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let params = GridPointParams { teff: 36000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_zom").await.unwrap(); db.upsert_workflow("wf_zom", None, "", "running") .await .unwrap(); // 僵尸行:只 insert_task,从未 push(无队列行)。 let zombie = mk_cold_spec(&name, ¶ms, "wf_zom"); db.insert_task(&zombie).await.unwrap(); let dispatched = scheduler .schedule_pending_tasks_for_workflow("wf_zom", 100) .await .unwrap(); assert_eq!(dispatched, 1, "清除僵尸后应正常派发"); // 僵尸行被删,新任务行 task_id 不同。 let rows = db .has_pending_tasks_for_point(&name, "wf_zom", None) .await .unwrap(); assert_eq!(rows.len(), 1); assert_ne!(rows[0], zombie.task_id.to_string(), "应派发全新任务"); let popped = queue.pop_task("node-x").await.unwrap().unwrap(); assert_eq!(popped.point_name, name); } /// E' 孤儿回收(2026-08-02 涡旋事故重构):以 MQ 活性交叉校验区分真孤儿与在途。 /// 场景:A running+全死行→重置+删行;B running+活行→不重置;C queued+全死行→重置; /// D converged→不触碰;E running+混合行→删死行、不重置。 #[tokio::test] async fn test_reclaim_orphaned_points_mq_liveness() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("reclaim_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("reclaim_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mk = |teff: f64| GridPointParams { teff: teff.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let upsert = |p: GridPointParams, st: common::models::GridPointStatus| { let db = db.clone(); async move { let name = p.model_name(); db.upsert_grid_point(&p, 0, "wf_r").await.unwrap(); db.update_grid_status(&name, st, "wf_r").await.unwrap(); name } }; // A:running + 死行(无队列行)→ 真孤儿。 let pa = mk(30000.0); let na = upsert(pa.clone(), common::models::GridPointStatus::Running).await; let ta = mk_cold_spec(&na, &pa, "wf_r"); db.insert_task(&ta).await.unwrap(); // B:running + 活行(队列行 pending)→ 在途,不触碰。 let pb = mk(31000.0); let nb = upsert(pb.clone(), common::models::GridPointStatus::Running).await; let tb = mk_cold_spec(&nb, &pb, "wf_r"); db.insert_task(&tb).await.unwrap(); queue.push_task(&tb).await.unwrap(); // C:queued + 死行 → 真孤儿(3b 卡死点救援)。 let pc = mk(32000.0); let nc = upsert(pc.clone(), common::models::GridPointStatus::Queued).await; let tc = mk_cold_spec(&nc, &pc, "wf_r"); db.insert_task(&tc).await.unwrap(); // D:converged + 死行 → 终态点不是候选。 let pd = mk(33000.0); let nd = upsert(pd.clone(), common::models::GridPointStatus::Completed).await; let td = mk_cold_spec(&nd, &pd, "wf_r"); db.insert_task(&td).await.unwrap(); // E:running + 混合(一死一活)→ 删死行、不重置。 let pe = mk(34000.0); let ne = upsert(pe.clone(), common::models::GridPointStatus::Running).await; let te_dead = mk_cold_spec(&ne, &pe, "wf_r"); let te_live = mk_cold_spec(&ne, &pe, "wf_r"); db.insert_task(&te_dead).await.unwrap(); db.insert_task(&te_live).await.unwrap(); queue.push_task(&te_live).await.unwrap(); // 回拨全部 tasks 行 created_at 7 小时(datetime('now') 为秒精度,同秒插入的 // 行无法满足 `< now` 判据,故用裸 SQL 构造 stale 条件),使其老于 stale_sec=6h。 { let pool = db.pool.clone(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET created_at = datetime('now', '-7 hours') WHERE workflow_name = 'wf_r'", [], ) .unwrap(); }) .await .unwrap(); } let rescued = scheduler.reclaim_orphaned_points(21600).await.unwrap(); assert_eq!(rescued, 2, "仅 A 与 C 被判定为真孤儿"); assert_eq!( db.get_grid_point_status(&na, "wf_r") .await .unwrap() .unwrap() .0, "pending", "A 应被重置为 pending" ); assert_eq!( db.get_grid_point_status(&nb, "wf_r") .await .unwrap() .unwrap() .0, "running", "B 有活队列行,不得重置" ); assert_eq!( db.get_grid_point_status(&nc, "wf_r") .await .unwrap() .unwrap() .0, "pending", "C 应被重置为 pending" ); assert_eq!( db.get_grid_point_status(&nd, "wf_r") .await .unwrap() .unwrap() .0, "completed", "D 终态点不触碰" ); assert_eq!( db.get_grid_point_status(&ne, "wf_r") .await .unwrap() .unwrap() .0, "running", "E 有活行,不得重置" ); // 死行清除核验:A/C 的行已删;B 的行保留;E 的死行删、活行留。 assert!(db .has_pending_tasks_for_point(&na, "wf_r", None) .await .unwrap() .is_empty()); assert!(db .has_pending_tasks_for_point(&nc, "wf_r", None) .await .unwrap() .is_empty()); assert_eq!( db.has_pending_tasks_for_point(&nb, "wf_r", None) .await .unwrap() .len(), 1 ); let e_rows = db .has_pending_tasks_for_point(&ne, "wf_r", None) .await .unwrap(); assert_eq!(e_rows, vec![te_live.task_id.to_string()], "E 仅活行幸存"); } /// C 回退状态守卫(2026-08-02 涡旋事故修复):仅 failed 点触发回退, /// converged/queued/running/不存在的点一律 Ok(false) 且不产生 seed_step 任务。 #[tokio::test] async fn test_seed_fallback_skips_non_failed_point() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_guard_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_guard_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_fb", None, "", "running") .await .unwrap(); let mk = |teff: f64| GridPointParams { teff: teff.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; // 放入一个可用种子,确保"跳过"是状态守卫的功劳而非无种子可匹配。 let seed_params = mk(35000.0); db.insert_seed_named(&seed_params.model_name(), &seed_params, "/tmp/seed.7") .await .unwrap(); for (teff, st) in [ (36000.0, common::models::GridPointStatus::Completed), (37000.0, common::models::GridPointStatus::Queued), (38000.0, common::models::GridPointStatus::Running), ] { let p = mk(teff); let name = p.model_name(); db.upsert_grid_point(&p, 0, "wf_fb").await.unwrap(); db.update_grid_status(&name, st.clone(), "wf_fb") .await .unwrap(); let fired = scheduler .trigger_tlusty_fallback(&p, &name, "wf_fb") .await .unwrap(); assert!(!fired, "状态 {:?} 不应触发回退", st); assert!( db.has_pending_tasks_for_point(&name, "wf_fb", Some("seed_step")) .await .unwrap() .is_empty(), "状态 {:?} 不应产生 seed_step 任务", st ); } // 不存在的点也不触发。 let ghost = mk(99000.0); let fired = scheduler .trigger_tlusty_fallback(&ghost, &ghost.model_name(), "wf_fb") .await .unwrap(); assert!(!fired); // 对照组:failed 点正常触发回退(有种子匹配)→ 点被复活为 queued。 // 策略链回退需该点有一携带多策略链的 tasks 行(模拟冷启动失败、链含 seed_step)。 // 该 task 行需为已上报态(failed/completed),否则会被僵尸卫生清理。 let pf = mk(35000.0); let nf = pf.model_name(); db.upsert_grid_point(&pf, 0, "wf_fb").await.unwrap(); db.update_grid_status(&nf, common::models::GridPointStatus::Failed, "wf_fb") .await .unwrap(); let cold_failed = TaskSpec { task_id: Uuid::new_v4(), point_name: nf.clone(), params: pf.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_fb".to_string()), wave: 0, tlusty_config: PhaseConfig { strategies: vec!["cold_run".to_string(), "seed_step".to_string()], ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&cold_failed).await.unwrap(); // 模拟该 cold_run 已上报失败(record_task_report 的口径:status→failed)。 { let pool = db.pool.clone(); let tid = cold_failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_tlusty_fallback(&pf, &nf, "wf_fb") .await .unwrap(); assert!(fired, "failed 点应正常触发回退"); assert_eq!( db.get_grid_point_status(&nf, "wf_fb") .await .unwrap() .unwrap() .0, "queued" ); } /// 策略链回退的僵尸卫生与在途去重(策略链版): /// - 场景一:点的最新 task 行携带完整策略链 [cold_run, seed_step] 且无队列凭证 /// (僵尸/已结算)→ 弹出 cold_run 后剩余 [seed_step] 非空、有种子 → 触发回退, /// 派发新 seed_step 任务。 /// - 场景二:点有活队列行(任务真在途)→ 跳过本次回退,不重复派发。 #[tokio::test] async fn test_seed_fallback_zombie_seed_hygiene() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_zom_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_zom_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_fz", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; db.insert_seed_named(¶ms.model_name(), ¶ms, "/tmp/seed.7") .await .unwrap(); // 场景一:failed 点 + 最新 task 行策略链 [cold_run, seed_step] + 无队列行(已上报 failed)。 // 弹出 cold_run → 剩 [seed_step],有种子 → 派发新 seed_step 任务。 // 注:task 行须为已上报态(failed),否则会被僵尸卫生当 pending 清理。 let p1 = params.clone(); let n1 = p1.model_name(); db.upsert_grid_point(&p1, 0, "wf_fz").await.unwrap(); db.update_grid_status(&n1, common::models::GridPointStatus::Failed, "wf_fz") .await .unwrap(); let cold_failed = TaskSpec { task_id: Uuid::new_v4(), point_name: n1.clone(), params: p1.clone(), seed_point_name: Some("old_seed".to_string()), timeout_sec: 3600, workflow_name: Some("wf_fz".to_string()), wave: 0, tlusty_config: PhaseConfig { strategies: vec!["cold_run".to_string(), "seed_step".to_string()], ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&cold_failed).await.unwrap(); { let pool = db.pool.clone(); let tid = cold_failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_tlusty_fallback(&p1, &n1, "wf_fz") .await .unwrap(); assert!( fired, "弹出 cold_run 后剩余 seed_step 非空且有种子 → 应触发回退" ); let seed_rows = db .has_pending_tasks_for_point(&n1, "wf_fz", Some("seed_step")) .await .unwrap(); assert_eq!(seed_rows.len(), 1, "新回退 seed_step 任务已创建"); assert_ne!(seed_rows[0], cold_failed.task_id.to_string()); // 场景二:点有活队列行(任务真在途)→ 跳过本次回退,行保留。 let p2 = GridPointParams { teff: 36000.0.into(), ..params.clone() }; let n2 = p2.model_name(); db.upsert_grid_point(&p2, 0, "wf_fz").await.unwrap(); db.update_grid_status(&n2, common::models::GridPointStatus::Failed, "wf_fz") .await .unwrap(); let live_task = TaskSpec { task_id: Uuid::new_v4(), point_name: n2.clone(), params: p2.clone(), seed_point_name: Some("old_seed".to_string()), timeout_sec: 3600, workflow_name: Some("wf_fz".to_string()), wave: 0, tlusty_config: PhaseConfig { strategies: vec!["cold_run".to_string(), "seed_step".to_string()], ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&live_task).await.unwrap(); queue.push_task(&live_task).await.unwrap(); let fired2 = scheduler .trigger_tlusty_fallback(&p2, &n2, "wf_fz") .await .unwrap(); assert!(!fired2, "活队列行存在时不应重复回退"); let rows2 = db .has_pending_tasks_for_point(&n2, "wf_fz", None) .await .unwrap(); assert_eq!(rows2, vec![live_task.task_id.to_string()], "在途行应保留"); } /// 策略链耗尽守卫:点的最新 task 行策略链仅剩 1 项(已耗尽)→ 保持 failed,不回退。 #[tokio::test] async fn test_strategy_fallback_exhausted_chain_keeps_failed() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_exh_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_exh_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_exh", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; db.insert_seed_named(¶ms.model_name(), ¶ms, "/tmp/seed.7") .await .unwrap(); let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_exh").await.unwrap(); db.update_grid_status(&name, common::models::GridPointStatus::Failed, "wf_exh") .await .unwrap(); // 策略链仅 [seed_step]:弹出后为空 → 耗尽,不回退。 let only_seed = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_exh".to_string()), wave: 0, tlusty_config: PhaseConfig { strategies: vec!["seed_step".to_string()], ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&only_seed).await.unwrap(); { let pool = db.pool.clone(); let tid = only_seed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_strategy_fallback(¶ms, &name, "wf_exh", "tlusty") .await .unwrap(); assert!(!fired, "策略链耗尽(仅剩 1 项弹出后为空)应保持 failed"); assert_eq!( db.get_grid_point_status(&name, "wf_exh") .await .unwrap() .unwrap() .0, "failed", "耗尽后点应保持 failed" ); } /// H1 活锁修复回归:cold_run 失败、剩余链 [seed_step] 无种子时,回退应把点打回 pending /// **并记录 pending_strategies 标记**(["seed_step"])。无此标记,下轮调度会用完整 YAML 链 /// 重派 → 重跑已失败的 cold_run → 再失败 → 再打回 → 无界失败重试活锁。 #[tokio::test] async fn test_h1_fallback_resets_pending_and_sets_marker() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_live_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_live_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_live", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_live").await.unwrap(); db.update_grid_status(&name, common::models::GridPointStatus::Failed, "wf_live") .await .unwrap(); // 最新任务行:cold_run 已执行并失败,链 [cold_run, seed_step](cold_run 弹后剩 seed_step)。 let cold_task = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_live".to_string()), wave: 0, tlusty_config: PhaseConfig { strategies: vec!["cold_run".to_string(), "seed_step".to_string()], ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&cold_task).await.unwrap(); { let pool = db.pool.clone(); let tid = cold_task.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } // 无任何种子 → resolve_dispatchable_chain([seed_step] 弹 cold_run 后) → 空 → H1。 let fired = scheduler .trigger_strategy_fallback(¶ms, &name, "wf_live", "tlusty") .await .unwrap(); assert!(!fired, "H1 路径返回 false(打回 pending 非派发)"); assert_eq!( db.get_grid_point_status(&name, "wf_live") .await .unwrap() .unwrap() .0, "pending", "H1 应把点打回 pending 待种子出现" ); assert_eq!( db.get_pending_strategies(&name, "wf_live").await.unwrap(), Some(r#"["seed_step"]"#.to_string()), "H1 应记录剩余策略链标记,供调度路径跳过已失败的 cold_run" ); } /// F 重启/停止卫生(2026-08-02 涡旋事故修复):initialize_grid 清理 pending 队列行 /// 时同步删除 tasks 表对应行;claimed 队列行与其 tasks 行保留(#6 设计不动)。 #[tokio::test] async fn test_initialize_grid_purges_cleared_tasks_rows() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("init_purge_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new( &temp_dir .path() .join("init_purge_queue.db") .to_string_lossy(), ) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mk = |teff: f64| GridPointParams { teff: teff.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; // 两点网格:X 构造 pending 队列行,Y 构造 claimed 队列行。 let cfg = { let mut c = mk_single_point_cfg(35000.0); c.grid.teff = vec![35000.0.into(), 36000.0.into()]; c }; scheduler.initialize_grid(&cfg, "wf_ip").await.unwrap(); db.upsert_workflow("wf_ip", None, "", "running") .await .unwrap(); let px = mk(35000.0); let py = mk(36000.0); let nx = px.model_name(); let ny = py.model_name(); // Y 先入队并被领用(claimed):pop 按 created_at FIFO,先入者先出。 let ty = mk_cold_spec(&ny, &py, "wf_ip"); db.insert_task(&ty).await.unwrap(); queue.push_task(&ty).await.unwrap(); let claimed = queue.pop_task("node-a").await.unwrap().unwrap(); assert_eq!(claimed.task_id, ty.task_id); // X 后入队,保持 pending(重启时应被成对清除)。 let tx = mk_cold_spec(&nx, &px, "wf_ip"); db.insert_task(&tx).await.unwrap(); queue.push_task(&tx).await.unwrap(); assert_eq!( db.has_pending_tasks_for_point(&nx, "wf_ip", None) .await .unwrap() .len(), 1 ); assert_eq!( db.has_pending_tasks_for_point(&ny, "wf_ip", None) .await .unwrap() .len(), 1 ); // 重新初始化(模拟重启):X 的 pending 队列行与其 tasks 行被同步清除, // Y 的 claimed 行与其 tasks 行保留。 scheduler.initialize_grid(&cfg, "wf_ip").await.unwrap(); assert!( db.has_pending_tasks_for_point(&nx, "wf_ip", None) .await .unwrap() .is_empty(), "X 的僵尸 tasks 行应随队列行同步清除" ); assert_eq!( db.has_pending_tasks_for_point(&ny, "wf_ip", None) .await .unwrap(), vec![ty.task_id.to_string()], "Y 的在途 tasks 行应保留" ); assert!( queue .verify_task_claim(&ty.task_id.to_string(), "node-a") .await .unwrap() .is_some(), "claimed 领用凭证必须保留(节点上报不被 403)" ); } /// 失败阶段归因(docs/task_engine_decoupling_design.md §4.2 注,修复审查 #2): /// SYNSPEC 失败弹 `synspec_strategies` 链重试光谱合成,绝不弹 TLUSTY 链、不做邻居 /// 种子搜索——旧实现把 synspec 失败误归因到 TLUSTY 链,种子门控可能把可重试的 /// 瞬态失败误判为不可回退(synspec-only 工作流因无近邻种子直接判死)。 #[tokio::test] async fn test_synspec_fallback_pops_synspec_chain_only() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_syn_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_syn_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_syn", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_syn").await.unwrap(); db.update_grid_status(&name, common::models::GridPointStatus::Failed, "wf_syn") .await .unwrap(); // 失败任务的 synspec 链含两项(standard, standard → 可重试一次);tlusty 链保持默认。 // 关键:不插入任何近邻种子(synspec 回退不依赖种子)。 let failed = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_syn".to_string()), wave: 0, synspec_config: PhaseConfig { strategies: vec!["standard".to_string(), "standard".to_string()], ..PhaseConfig::default_synspec() }, ..Default::default() }; db.insert_task(&failed).await.unwrap(); { let pool = db.pool.clone(); let tid = failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_strategy_fallback(¶ms, &name, "wf_syn", "synspec") .await .unwrap(); assert!(fired, "synspec 链含多策略时应触发回退"); // 新派发任务:synspec 链剩 [standard],tlusty 链保持不变。 let rows = db .has_pending_tasks_for_point(&name, "wf_syn", None) .await .unwrap(); assert_eq!(rows.len(), 1); assert_ne!(rows[0], failed.task_id.to_string()); // synspec 失败不得弹 TLUSTY 链:原失败任务行的 tlusty 链保持完整(弹栈机制是 // 只读的,回退不修改旧行;新派发的 synspec-only 重试行把 tlusty 链归一化为 // ["cold_run"],但那不是"弹 TLUSTY 链")。故直接查原失败行而非最新行。 let tlusty_of_original: String = { let pool = db.pool.clone(); let tid = failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.query_row( "SELECT tlusty_strategies FROM tasks WHERE task_id = ?1", rusqlite::params![tid], |r| r.get::<_, String>(0), ) .unwrap() }) .await .unwrap() }; assert_eq!( tlusty_of_original, "[\"cold_run\",\"seed_step\"]", "synspec 失败不得改写/弹出 TLUSTY 链(旧行保持派发时完整链)" ); // 审查修复 #9/M3:重试行保留原始完整 TLUSTY 链(而非归一化为 ["cold_run"])。 // 归一化会破坏后续 TLUSTY 链回退——若该重试失败且 failed_stage 被误报为 "tlusty", // pop 读到 ["cold_run"] 弹出后链空 → 误判不可回退,丢失原本的策略链回退能力。 // executor 已凭 enabled=false 跳过大气计算,strategies 仅作审计/潜在回退依据。 assert_eq!( db.get_latest_tlusty_strategies(&name, "wf_syn") .await .unwrap(), vec!["cold_run".to_string(), "seed_step".to_string()], "重试行保留原始完整 TLUSTY 链(enabled=false 已让 executor 跳过大气,不归一化)" ); let popped = queue.pop_task("node-x").await.unwrap().unwrap(); assert_eq!( popped.synspec_config.strategies, vec!["standard".to_string()] ); // 点被重置为 queued(回退重试在途)。 assert_eq!( db.get_grid_point_status(&name, "wf_syn") .await .unwrap() .unwrap() .0, "queued" ); } /// 2026-08-04 语义修正:SkipFailed **不再门控**策略链回退——回退只由启动时的策略链 /// (回退优先级排序)驱动。失败任务行携带可回退链 [cold_run, seed_step] + 近邻种子时, /// 即使 policy 为 skip_failed 也弹出 cold_run 派发 seed_step 重试;SkipFailed 只影响 /// 启动时对失败点的重置(见 initialize_grid 三分支)。 #[tokio::test] async fn test_skip_failed_policy_does_not_gate_fallback() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_skipf_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_skipf_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); db.upsert_workflow("wf_skipf", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_skipf").await.unwrap(); db.update_grid_status(&name, common::models::GridPointStatus::Failed, "wf_skipf") .await .unwrap(); // 近邻种子:保证 seed_step 顺位可派发——回退是否触发由链决定,而非 policy。 db.insert_seed_named(¶ms.model_name(), ¶ms, "/tmp/seed.7") .await .unwrap(); // 失败任务行:policy=skip_failed + 默认可回退链 [cold_run, seed_step]。 let failed = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_skipf".to_string()), wave: 0, tlusty_config: PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_tlusty() }, ..Default::default() }; db.insert_task(&failed).await.unwrap(); { let pool = db.pool.clone(); let tid = failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_strategy_fallback(¶ms, &name, "wf_skipf", "tlusty") .await .unwrap(); assert!( fired, "skip_failed 不再抑制回退:链含下一顺位且有种子 → 应回退派发 seed_step" ); assert_eq!( db.get_grid_point_status(&name, "wf_skipf") .await .unwrap() .unwrap() .0, "queued", "回退任务已派发,点应为 queued" ); let rows = db .has_pending_tasks_for_point(&name, "wf_skipf", Some("seed_step")) .await .unwrap(); assert_eq!(rows.len(), 1, "应派发一个 seed_step 回退任务"); } /// ForceRecompute 策略(设计 §2.1 场景 C,修复审查 #1):start 工作流时把已收敛/已失败 /// 终态点全部打回 pending 强制重算;SkipConverged(默认)保持终态不动。 #[tokio::test] async fn test_force_recompute_policy_resets_terminal_points() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_fr_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_fr_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mut cfg = mk_single_point_cfg(35000.0); cfg.grid.teff = vec![35000.0.into(), 36000.0.into()]; // 初始初始化(默认 skip_converged)。 scheduler.initialize_grid(&cfg, "wf_fr").await.unwrap(); // 手工置终态:nx 收敛、ny 失败。 let px = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let py = GridPointParams { teff: 36000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let nx = px.model_name(); let ny = py.model_name(); db.update_grid_status(&nx, common::models::GridPointStatus::Completed, "wf_fr") .await .unwrap(); db.update_grid_status(&ny, common::models::GridPointStatus::Failed, "wf_fr") .await .unwrap(); // 对照组:SkipFailed(双阶段)重新初始化 → 终态保持(最保守增量:收敛+失败都保留)。 // 注:默认 SkipConverged 现会"重试失败"(把失败点打回 pending),控制组须用 // SkipFailed 才能断言"终态全保留"(2026-08-04 语义修正)。 cfg.tlusty_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_tlusty() }); cfg.synspec_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_synspec() }); scheduler.initialize_grid(&cfg, "wf_fr").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_fr") .await .unwrap() .unwrap() .0, "completed", "skip_failed 不得重置收敛点" ); assert_eq!( db.get_grid_point_status(&ny, "wf_fr") .await .unwrap() .unwrap() .0, "failed", "skip_failed 不得重置失败点" ); // ForceRecompute 重新初始化 → 终态全部打回 pending。 cfg.tlusty_stage = Some(PhaseConfig { policy: ResumePolicy::ForceRecompute, ..PhaseConfig::default_tlusty() }); scheduler.initialize_grid(&cfg, "wf_fr").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_fr") .await .unwrap() .unwrap() .0, "pending", "force_recompute 应重置收敛点" ); assert_eq!( db.get_grid_point_status(&ny, "wf_fr") .await .unwrap() .unwrap() .0, "pending", "force_recompute 应重置失败点" ); } /// 2026-08-04 语义修正:SkipConverged(默认)启动时**跳过收敛、重试失败**——仅把已失败 /// 的点打回 pending 重试,已收敛点保留。此前与 SkipFailed 在启动时行为相同(都不重置), /// 导致"只重试失败点"这一常用增量语义无法表达。 #[tokio::test] async fn test_initialize_grid_skip_converged_retries_failed_only() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_sc_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_sc_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mut cfg = mk_single_point_cfg(35000.0); cfg.grid.teff = vec![35000.0.into(), 36000.0.into()]; scheduler.initialize_grid(&cfg, "wf_sc").await.unwrap(); let px = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let py = GridPointParams { teff: 36000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let nx = px.model_name(); let ny = py.model_name(); db.update_grid_status(&nx, common::models::GridPointStatus::Completed, "wf_sc") .await .unwrap(); db.update_grid_status(&ny, common::models::GridPointStatus::Failed, "wf_sc") .await .unwrap(); // 默认策略 skip_converged 重新初始化:失败点 ny 打回 pending 重试,收敛点 nx 保留。 scheduler.initialize_grid(&cfg, "wf_sc").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_sc") .await .unwrap() .unwrap() .0, "completed", "skip_converged 应保留收敛点" ); assert_eq!( db.get_grid_point_status(&ny, "wf_sc") .await .unwrap() .unwrap() .0, "pending", "skip_converged 应把失败点打回 pending 重试" ); } /// 2026-08-04 语义修正:SkipFailed 启动时**跳过收敛及失败**——收敛点和失败点都保留, /// 只算从未计算过的点(最保守增量)。与 SkipConverged(重试失败)区别开。 #[tokio::test] async fn test_initialize_grid_skip_failed_keeps_all_terminal() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_sf_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_sf_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mut cfg = mk_single_point_cfg(35000.0); cfg.grid.teff = vec![35000.0.into(), 36000.0.into()]; scheduler.initialize_grid(&cfg, "wf_sf").await.unwrap(); let px = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let py = GridPointParams { teff: 36000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let nx = px.model_name(); let ny = py.model_name(); db.update_grid_status(&nx, common::models::GridPointStatus::Completed, "wf_sf") .await .unwrap(); db.update_grid_status(&ny, common::models::GridPointStatus::Failed, "wf_sf") .await .unwrap(); // SkipFailed 重新初始化:收敛 + 失败都保留(不重置任何终态点)。 cfg.tlusty_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_tlusty() }); cfg.synspec_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_synspec() }); scheduler.initialize_grid(&cfg, "wf_sf").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_sf") .await .unwrap() .unwrap() .0, "completed", "skip_failed 应保留收敛点" ); assert_eq!( db.get_grid_point_status(&ny, "wf_sf") .await .unwrap() .unwrap() .0, "failed", "skip_failed 应保留失败点" ); } /// ForceRecompute 对 SYNSPEC 阶段同样生效(设计 §2.2 场景 B「仅更新光谱」: /// TLUSTY 关 + SYNSPEC force_recompute,修复审查 #2 后续):initialize_grid 此前只消费 /// TLUSTY policy,场景 B 的已收敛点永远无法重算光谱(UI 上明明可选却毫无效果)。 /// 现任一步骤为 force_recompute 即重置全部终态点回 pending,由调度器重派光谱任务。 #[tokio::test] async fn test_synspec_force_recompute_resets_terminal_points() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_sfr_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_sfr_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); let mut cfg = mk_single_point_cfg(35000.0); cfg.grid.teff = vec![35000.0.into(), 36000.0.into()]; scheduler.initialize_grid(&cfg, "wf_sfr").await.unwrap(); let px = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let py = GridPointParams { teff: 36000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let nx = px.model_name(); let ny = py.model_name(); db.update_grid_status(&nx, common::models::GridPointStatus::Completed, "wf_sfr") .await .unwrap(); db.update_grid_status(&ny, common::models::GridPointStatus::Failed, "wf_sfr") .await .unwrap(); // 对照组:SkipFailed(双阶段)重新初始化 → 终态保持(收敛+失败都保留)。 // 注:默认 SkipConverged 现会"重试失败"(把失败点打回 pending),控制组须用 // SkipFailed 才能断言"终态全保留"(2026-08-04 语义修正)。 cfg.tlusty_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_tlusty() }); cfg.synspec_stage = Some(PhaseConfig { policy: ResumePolicy::SkipFailed, ..PhaseConfig::default_synspec() }); scheduler.initialize_grid(&cfg, "wf_sfr").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_sfr") .await .unwrap() .unwrap() .0, "completed", "skip_failed 不得重置收敛点" ); assert_eq!( db.get_grid_point_status(&ny, "wf_sfr") .await .unwrap() .unwrap() .0, "failed", "skip_failed 不得重置失败点" ); // 场景 B:TLUSTY 关闭 + SYNSPEC force_recompute → 终态全部打回 pending(重算光谱)。 cfg.tlusty_stage = Some(PhaseConfig { enabled: false, ..PhaseConfig::default_tlusty() }); cfg.synspec_stage = Some(PhaseConfig { policy: ResumePolicy::ForceRecompute, ..PhaseConfig::default_synspec() }); scheduler.initialize_grid(&cfg, "wf_sfr").await.unwrap(); assert_eq!( db.get_grid_point_status(&nx, "wf_sfr") .await .unwrap() .unwrap() .0, "pending", "SYNSPEC force_recompute 应重置收敛点(场景 B 重算光谱)" ); assert_eq!( db.get_grid_point_status(&ny, "wf_sfr") .await .unwrap() .unwrap() .0, "pending", "SYNSPEC force_recompute 应重置失败点" ); } /// 半失败点 synspec 链回退(修复审查 #2 后续):大气已收敛 + 光谱失败时,回退任务 /// 强制关闭 TLUSTY(不重算有效大气)并凭 atmosphere_ref 复用既有 .7——否则重试会 /// 把整个大气链重跑一遍(浪费)甚至因缺种子而误判不可回退。 #[tokio::test] async fn test_synspec_fallback_retry_disables_tlusty() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_sfr2_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_sfr2_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); // 双阶段默认配置(tlusty 启用),synspec 链含两项可回退。 db.upsert_workflow("wf_sfr2", None, "", "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; let name = params.model_name(); db.upsert_grid_point(¶ms, 0, "wf_sfr2").await.unwrap(); db.update_grid_status(&name, common::models::GridPointStatus::Failed, "wf_sfr2") .await .unwrap(); let failed = TaskSpec { task_id: Uuid::new_v4(), point_name: name.clone(), params: params.clone(), seed_point_name: None, timeout_sec: 3600, workflow_name: Some("wf_sfr2".to_string()), wave: 0, synspec_config: PhaseConfig { strategies: vec!["standard".to_string(), "standard".to_string()], ..PhaseConfig::default_synspec() }, ..Default::default() }; db.insert_task(&failed).await.unwrap(); { let pool = db.pool.clone(); let tid = failed.task_id.to_string(); tokio::task::spawn_blocking(move || { let conn = pool.get().unwrap(); conn.execute( "UPDATE tasks SET status = 'failed', completed_at = datetime('now') WHERE task_id = ?1", rusqlite::params![tid], ) .unwrap(); }) .await .unwrap(); } let fired = scheduler .trigger_strategy_fallback(¶ms, &name, "wf_sfr2", "synspec") .await .unwrap(); assert!(fired, "synspec 链含多策略时应触发回退"); let popped = queue.pop_task("node-x").await.unwrap().unwrap(); // 重试任务:tlusty 强制关闭 + atmosphere_ref 绑定本点(大气复用,不重算)。 assert!( !popped.tlusty_config.enabled, "半失败点重试不得重算大气(tlusty 应关闭)" ); assert_eq!(popped.atmosphere_ref.as_deref(), Some(name.as_str())); assert_eq!( popped.synspec_config.strategies, vec!["standard".to_string()] ); assert_eq!( db.get_grid_point_status(&name, "wf_sfr2") .await .unwrap() .unwrap() .0, "queued" ); } /// 初始派发链头为 seed_step 且无近邻种子 → 跳过该顺位派发 cold_run(修复审查 #3): /// 旧实现直接取 strategies[0] 派发 seed_step 任务(seed_point_name=None),节点端 /// 强校验必失败。现 resolve_dispatchable_chain 统一种子门控,链重写为可派发顺位起。 /// 无种子的 seed_step 轮转到链尾保留为回退顺位(修复审查 #6:不丢弃,冷启动失败后 /// 若种子已出现仍可种子步进重试)。 #[tokio::test] async fn test_initial_dispatch_seed_head_falls_back_within_chain() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_sh_db.db").to_string_lossy()) .await .unwrap(); let queue = Arc::new( SqliteTaskQueue::new(&temp_dir.path().join("fb_sh_queue.db").to_string_lossy()) .await .unwrap(), ); let scheduler = GridScheduler::new(db.clone(), queue.clone()); // tlusty 链头为 seed_step(用户把种子步进排最前),无任何种子。 let yaml = "\ grid:\n teff: [35000]\n logg: [5.5]\n loghe: [-1]\n logc: [-2]\n logn: [-2]\n logo: [-2]\n\ tlusty_stage:\n enabled: true\n policy: skip_converged\n strategies: [seed_step, cold_run]\n"; db.upsert_workflow("wf_sh", None, yaml, "running") .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; db.upsert_grid_point(¶ms, 0, "wf_sh").await.unwrap(); let dispatched = scheduler .schedule_pending_tasks_for_workflow("wf_sh", 100) .await .unwrap(); assert_eq!( dispatched, 1, "应跳过无种子的 seed_step 顺位,派发 cold_run" ); let popped = queue.pop_task("node-x").await.unwrap().unwrap(); assert_eq!( popped.tlusty_config.current_strategy("cold_run"), "cold_run", "无种子时派发的首项策略应为 cold_run" ); assert_eq!( popped.tlusty_config.strategies, vec!["cold_run".to_string(), "seed_step".to_string()], "无种子的 seed_step 应轮转到链尾保留为回退顺位,而不是被丢弃" ); assert!(popped.seed_point_name.is_none()); // 有种子时:链头 seed_step 直接派发并注入种子。 // 先结算第一条任务(清队列凭证,模拟任务上报结束),重置点为 pending, // 插入近邻种子后重新派发。 queue .remove_task(&popped.task_id.to_string()) .await .unwrap(); db.update_grid_status( &popped.point_name, common::models::GridPointStatus::Pending, "wf_sh", ) .await .unwrap(); db.insert_seed_named(¶ms.model_name(), ¶ms, "/tmp/seed.7") .await .unwrap(); let dispatched2 = scheduler .schedule_pending_tasks_for_workflow("wf_sh", 100) .await .unwrap(); assert_eq!(dispatched2, 1, "有种子时 seed_step 顺位可派发"); let popped2 = queue.pop_task("node-x").await.unwrap().unwrap(); assert_eq!( popped2.tlusty_config.current_strategy("cold_run"), "seed_step", "有种子时派发的首项策略应为 seed_step" ); assert_eq!( popped2.tlusty_config.strategies, vec!["seed_step".to_string(), "cold_run".to_string()], "有种子时链保持 [seed_step, cold_run] 不变" ); assert_eq!( popped2.seed_point_name.as_deref(), Some(params.model_name().as_str()) ); } /// resolve_dispatchable_chain 终止性回归(修复审查 #6 后续):重复 seed_step 链 /// (如 `[seed_step, seed_step]`)若用「弹出再压回」轮转会因保持原序死循环。 /// 暂存式实现应把无种子 seed_step 全部弹出并判空(无任何可派发顺位)。 #[tokio::test] async fn test_resolve_dispatchable_chain_duplicate_seed_terminates() { let temp_dir = tempfile::tempdir().unwrap(); let db = Database::new(&temp_dir.path().join("fb_ds_db.db").to_string_lossy()) .await .unwrap(); let params = GridPointParams { teff: 35000.0.into(), logg: 5.5.into(), loghe: (-1.0).into(), logc: (-2.0).into(), logn: (-2.0).into(), logo: (-2.0).into(), }; // 无种子 → 重复 seed_step 链应判空(不派发、不死循环)。 let (chain, seed) = GridScheduler::resolve_dispatchable_chain( &db, "wf", ¶ms, vec!["seed_step".to_string(), "seed_step".to_string()], ) .await .unwrap(); assert!(chain.is_empty()); assert!(seed.is_none()); // 无种子 → [seed_step, seed_step, cold_run] 解析为 [cold_run, seed_step, seed_step], // 两个无种子 seed_step 保留为回退顺位且链首可派发。 let (chain2, seed2) = GridScheduler::resolve_dispatchable_chain( &db, "wf", ¶ms, vec![ "seed_step".to_string(), "seed_step".to_string(), "cold_run".to_string(), ], ) .await .unwrap(); assert_eq!(seed2, None); assert_eq!( chain2, vec![ "cold_run".to_string(), "seed_step".to_string(), "seed_step".to_string() ], "无种子 seed_step 应追加到链尾,链首为可派发的 cold_run" ); } }