diff --git a/crates/op-orchestrator/Cargo.toml b/crates/op-orchestrator/Cargo.toml new file mode 100644 index 000000000..90c29ed6e --- /dev/null +++ b/crates/op-orchestrator/Cargo.toml @@ -0,0 +1,16 @@ +[package] +name = "op-orchestrator" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +authors.workspace = true +repository.workspace = true + +[dependencies] +op-editor-core = { path = "../op-editor-core" } +op-ai-skills = { path = "../op-ai-skills" } +jian-ops-schema = { path = "../../vendor/jian/crates/jian-ops-schema" } +serde = { workspace = true } +serde_json = { workspace = true } +futures = "0.3" diff --git a/crates/op-orchestrator/src/cleanup.rs b/crates/op-orchestrator/src/cleanup.rs new file mode 100644 index 000000000..130764dd2 --- /dev/null +++ b/crates/op-orchestrator/src/cleanup.rs @@ -0,0 +1,111 @@ +//! 阶段 4 —— 清理 pass。 +//! +//! [`run_cleanup_passes`] 在所有 subtask 插入完成后运行,是独立 +//! 函数 —— S3a 顺序路径与 S3b 并发路径都复用它(spec §9)。 +//! +//! [`descendant_count`] 给 `run()` 的"零内容"判定提供基线: +//! scaffold 之后数一次,subtask 全跑完再数一次,没涨即零内容。 + +use crate::plan::OrchestratorPlan; +use crate::types::DocSink; +use jian_ops_schema::node::PenNode; +use op_editor_core::{EditorState, PenNodeExt}; + +/// 递归统计 `node` 下的后代数(不含自身)。 +fn count_descendants(node: &PenNode) -> usize { + match node.children() { + Some(children) => { + children.len() + children.iter().map(count_descendants).sum::() + } + None => 0, + } +} + +/// 统计活动页里 id 为 `root_id` 的节点的后代总数。节点不存在 +/// 时返回 0。 +pub fn descendant_count(state: &EditorState, root_id: &str) -> usize { + state + .active_children() + .iter() + .find(|n| n.id_str() == root_id) + .map(count_descendants) + .unwrap_or(0) +} + +/// 阶段 4 清理 pass —— 在全部 subtask 插入完成后运行。 +/// +/// S3a 骨架版本目前不改文档:三个 TS 清理 pass(移动端重复状态栏 +/// 去重 `removeDuplicateStatusBars`、单组件 section root unwrap +/// `unwrapSingleComponentSectionRoot`、高度自适应 +/// `adjustRootFrameHeightToContent`)需对照 TS 源码逐条移植并配合 +/// 画布验证 —— 列为 S3a 的后续细化项。函数签名与调用点在此固定, +/// 使 S3b 并发路径可直接复用(spec §9)。 +pub fn run_cleanup_passes(_sink: &mut dyn DocSink, _plan: &OrchestratorPlan) { + // 后续细化:① removeDuplicateStatusBars ② unwrapSingleComponentSectionRoot + // ③ adjustRootFrameHeightToContent —— 逐条对照 TS orchestrator.ts。 +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{OrchestratorPlan, RootFrameSpec}; + use crate::test_support::VecDocSink; + use op_editor_core::{EditorCommand, NodeId}; + use serde_json::json; + + fn frame_json(id: &str, children: serde_json::Value) -> PenNode { + serde_json::from_value(json!({ + "type": "frame", "id": id, "name": id, + "x": 0, "y": 0, "width": 100, "height": 100, + "children": children, + })) + .expect("frame json") + } + + fn plan() -> OrchestratorPlan { + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "P".into(), + width: 1200.0, + height: 800.0, + layout: None, + gap: None, + padding: None, + fill: None, + }, + subtasks: vec![], + style_guide: None, + } + } + + #[test] + fn descendant_count_counts_nested() { + let mut sink = VecDocSink::new(); + // root 套 child 套 grandchild + let tree = frame_json("root", json!([frame_json_value("c", json!([frame_json_value("gc", json!([]))]))])); + sink.state.apply(EditorCommand::InsertSubtree { + nodes: vec![tree], + parent_id: NodeId::NONE, + }); + assert_eq!(descendant_count(&sink.state, "root"), 2); + assert_eq!(descendant_count(&sink.state, "missing"), 0); + } + + /// 同 `frame_json` 但返回 `serde_json::Value`(供嵌套构造)。 + fn frame_json_value(id: &str, children: serde_json::Value) -> serde_json::Value { + json!({ + "type": "frame", "id": id, "name": id, + "x": 0, "y": 0, "width": 100, "height": 100, + "children": children, + }) + } + + #[test] + fn run_cleanup_passes_is_callable() { + let mut sink = VecDocSink::new(); + run_cleanup_passes(&mut sink, &plan()); + // 骨架版不改文档。 + assert!(sink.applied.is_empty()); + } +} diff --git a/crates/op-orchestrator/src/intent.rs b/crates/op-orchestrator/src/intent.rs new file mode 100644 index 000000000..f6ffe1c0e --- /dev/null +++ b/crates/op-orchestrator/src/intent.rs @@ -0,0 +1,64 @@ +//! 意图分类 —— design vs chat 路由。 +//! +//! 行为忠实 TS `ai-service.ts` 的路由:有明确设计意图 → Design, +//! 否则 Chat。S3a 用轻量关键词启发式(可单测);LLM 分类留后续。 + +use crate::types::Intent; + +/// 命中任一即判为设计意图。中英双语,覆盖动词 + 设计名词。 +const DESIGN_KEYWORDS: &[&str] = &[ + // 英文动词 + "design", "create", "build", "make a", "make me", "generate", "draw", + "mock up", "mockup", "wireframe", "prototype", "lay out", + // 英文设计名词 + "landing page", "dashboard", "pricing card", "login page", "sign up", + "ui for", "screen for", "app screen", "web page", + // 中文动词 + "设计", "生成", "做一个", "做个", "创建", "画一个", "画个", "搞一个", + "帮我做", "来个", + // 中文设计名词 + "页面", "界面", "仪表盘", "落地页", "原型", "登录页", "注册页", +]; + +/// 把用户消息分类为 [`Intent::Design`] 或 [`Intent::Chat`]。 +/// +/// 命中任一设计关键词即 Design;否则 Chat。大小写不敏感。 +pub fn classify_intent(message: &str) -> Intent { + let lower = message.to_lowercase(); + if DESIGN_KEYWORDS.iter().any(|k| lower.contains(k)) { + Intent::Design + } else { + Intent::Chat + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn design_verbs_classify_as_design() { + for p in [ + "design a login page", + "生成一个仪表盘", + "做一个落地页", + "create a pricing card", + "Generate a dashboard for sales", + "帮我做个注册页", + ] { + assert_eq!(classify_intent(p), Intent::Design, "{p}"); + } + } + + #[test] + fn questions_classify_as_chat() { + for p in [ + "what is a frame?", + "这个怎么用", + "解释一下布局", + "thanks, that looks good", + ] { + assert_eq!(classify_intent(p), Intent::Chat, "{p}"); + } + } +} diff --git a/crates/op-orchestrator/src/lib.rs b/crates/op-orchestrator/src/lib.rs new file mode 100644 index 000000000..9cf48f278 --- /dev/null +++ b/crates/op-orchestrator/src/lib.rs @@ -0,0 +1,28 @@ +//! `op-orchestrator` — S3a 设计编排器(单屏顺序骨架)。 +//! +//! 把 TS `apps/web/src/services/ai/orchestrator.ts` 的阶段 1-4 +//! 单屏路径补回 Rust。副作用只走 [`DocSink`] / [`LlmClient`] 两个 +//! trait,核心逻辑全是纯函数,不依赖 winit/casement/agent。 +//! +//! Plan B 提供"零件"(类型 + intent/plan/parse/normalize/variables); +//! Plan C 在 `run` 模块接出四阶段主轴。 + +pub mod intent; +pub mod parse; +pub mod plan; +pub mod plan_normalize; +pub mod types; +pub mod variables; + +pub mod cleanup; +pub mod prompt; +pub mod run; +pub mod scaffold; +pub mod subagent; + +#[cfg(test)] +mod test_support; + +pub use intent::classify_intent; +pub use run::Orchestrator; +pub use types::*; diff --git a/crates/op-orchestrator/src/parse.rs b/crates/op-orchestrator/src/parse.rs new file mode 100644 index 000000000..bc9d19f88 --- /dev/null +++ b/crates/op-orchestrator/src/parse.rs @@ -0,0 +1,96 @@ +//! 节点解析 —— sub-agent 的 LLM 文本输出 → canonical `PenNode` 树。 +//! +//! sub-agent prompt(见 `prompt` 模块)要求模型输出 canonical +//! PenNode JSON 数组。本模块抽出 JSON 数组并 serde 反序列化。 + +use jian_ops_schema::node::PenNode; + +/// 节点解析错误。 +#[derive(Debug, Clone)] +pub struct ParseError(pub String); + +impl std::fmt::Display for ParseError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "node parse error: {}", self.0) + } +} + +/// 从 LLM 文本里抽 JSON 数组并反序列化为 canonical `PenNode` 树。 +/// +/// 容忍 ```json 围栏 + 前后散文 —— 抽第一个平衡的 `[...]`。 +/// 空数组视为失败(零节点)。 +pub fn parse_nodes(text: &str) -> Result, ParseError> { + let json = extract_json_array(text) + .ok_or_else(|| ParseError("no JSON array found".into()))?; + let nodes: Vec = serde_json::from_str(json) + .map_err(|e| ParseError(format!("deserialize: {e}")))?; + if nodes.is_empty() { + return Err(ParseError("empty node array".into())); + } + Ok(nodes) +} + +/// 抽第一个平衡的 `[...]`(忽略字符串内的方括号)。 +fn extract_json_array(text: &str) -> Option<&str> { + let start = text.find('[')?; + let bytes = text.as_bytes(); + let mut depth = 0i32; + let mut in_str = false; + let mut esc = false; + for i in start..bytes.len() { + let c = bytes[i]; + if in_str { + if esc { + esc = false; + } else if c == b'\\' { + esc = true; + } else if c == b'"' { + in_str = false; + } + continue; + } + match c { + b'"' => in_str = true, + b'[' => depth += 1, + b']' => { + depth -= 1; + if depth == 0 { + return Some(&text[start..=i]); + } + } + _ => {} + } + } + None +} + +#[cfg(test)] +mod tests { + use super::*; + use op_editor_core::PenNodeExt as _; + + #[test] + fn parse_nodes_reads_fenced_pennode_array() { + // 一个最小的 frame 节点(canonical schema:`type` 标签 + + // base 字段)。 + let text = r#"Sure, here are the nodes: +```json +[ + { "type": "frame", "id": "hero", "name": "Hero", "x": 0, "y": 0, "width": 1200, "height": 400, "children": [] } +] +```"#; + let nodes = parse_nodes(text).expect("parse"); + assert_eq!(nodes.len(), 1); + assert_eq!(nodes[0].id_str(), "hero"); + } + + #[test] + fn parse_nodes_rejects_empty_array() { + assert!(parse_nodes("[]").is_err()); + } + + #[test] + fn parse_nodes_rejects_no_array() { + assert!(parse_nodes("the model wrote prose only").is_err()); + } +} diff --git a/crates/op-orchestrator/src/plan.rs b/crates/op-orchestrator/src/plan.rs new file mode 100644 index 000000000..88343f664 --- /dev/null +++ b/crates/op-orchestrator/src/plan.rs @@ -0,0 +1,222 @@ +//! `OrchestratorPlan` —— 规划阶段的产物。 +//! +//! 规划 LLM 调用产出一段 JSON,`parse_plan` 解析它;调用失败 / 解析 +//! 失败时 `build_fallback_plan` 给一个启发式的可跑 plan(对齐 TS +//! `buildFallbackPlanFromPrompt`)。 + +use crate::types::DesignRequest; +use serde::Deserialize; +use std::collections::BTreeMap; + +/// 一个 subtask 区域的尺寸。 +#[derive(Debug, Clone, PartialEq, Deserialize)] +pub struct Region { + pub width: f64, + pub height: f64, +} + +/// 根 frame 规格。 +#[derive(Debug, Clone, PartialEq, Deserialize)] +pub struct RootFrameSpec { + pub id: String, + pub name: String, + pub width: f64, + pub height: f64, + #[serde(default)] + pub layout: Option, + #[serde(default)] + pub gap: Option, + #[serde(default)] + pub padding: Option, + #[serde(default)] + pub fill: Option, +} + +/// 一个生成子任务 —— 对应设计里的一个区块。 +#[derive(Debug, Clone, PartialEq, Deserialize)] +pub struct Subtask { + pub id: String, + pub label: String, + pub region: Region, + /// 规范化后赋值 = `id`。 + #[serde(default)] + pub id_prefix: String, + /// 规范化后赋值 = 根 frame id。 + #[serde(default)] + pub parent_frame_id: Option, +} + +/// 设计系统提示 —— 调色板等。 +#[derive(Debug, Clone, PartialEq, Default, Deserialize)] +pub struct StyleGuide { + /// `名称 -> hex`,阶段 2 据此 seed `$color-*` 文档变量。 + #[serde(default)] + pub palette: BTreeMap, +} + +/// 规划阶段的完整产物。 +#[derive(Debug, Clone, PartialEq, Deserialize)] +pub struct OrchestratorPlan { + pub root_frame: RootFrameSpec, + pub subtasks: Vec, + #[serde(default)] + pub style_guide: Option, +} + +/// plan 解析错误。 +#[derive(Debug, Clone)] +pub struct PlanParseError(pub String); + +impl std::fmt::Display for PlanParseError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "plan parse error: {}", self.0) + } +} + +/// 从规划 LLM 的文本输出里抽出 plan JSON 并反序列化。 +/// +/// 容忍 ```json 围栏 + 前后散文 —— 抽第一个平衡的 `{...}`。 +/// `subtasks` 为空视为解析失败。 +pub fn parse_plan(text: &str) -> Result { + let json = extract_json_object(text) + .ok_or_else(|| PlanParseError("no JSON object found".into()))?; + let plan: OrchestratorPlan = serde_json::from_str(json) + .map_err(|e| PlanParseError(format!("deserialize: {e}")))?; + if plan.subtasks.is_empty() { + return Err(PlanParseError("plan has no subtasks".into())); + } + Ok(plan) +} + +/// 抽第一个平衡的 `{...}`(忽略字符串内的花括号)。 +fn extract_json_object(text: &str) -> Option<&str> { + let start = text.find('{')?; + let bytes = text.as_bytes(); + let mut depth = 0i32; + let mut in_str = false; + let mut esc = false; + for i in start..bytes.len() { + let c = bytes[i]; + if in_str { + if esc { + esc = false; + } else if c == b'\\' { + esc = true; + } else if c == b'"' { + in_str = false; + } + continue; + } + match c { + b'"' => in_str = true, + b'{' => depth += 1, + b'}' => { + depth -= 1; + if depth == 0 { + return Some(&text[start..=i]); + } + } + _ => {} + } + } + None +} + +/// 启发式兜底 plan —— 规划 LLM 不可用 / 解析失败时用,保证编排器 +/// 仍能跑出点东西。对齐 TS `buildFallbackPlanFromPrompt`:固定宽度 +/// 根 frame + 按 prompt 规模切 1-3 个等高区块。 +pub fn build_fallback_plan(req: &DesignRequest) -> OrchestratorPlan { + const WIDTH: f64 = 1200.0; + const SECTION_HEIGHT: f64 = 360.0; + + // prompt 越长 → 区块越多,封顶 3(单屏骨架不做更复杂的切分)。 + let section_count = match req.prompt.chars().count() { + 0..=80 => 1, + 81..=200 => 2, + _ => 3, + }; + + let subtasks: Vec = (0..section_count) + .map(|i| { + let id = format!("section-{}", i + 1); + Subtask { + id: id.clone(), + label: format!("Section {}", i + 1), + region: Region { + width: WIDTH, + height: SECTION_HEIGHT, + }, + id_prefix: id, + parent_frame_id: None, + } + }) + .collect(); + + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "Design".into(), + width: WIDTH, + height: SECTION_HEIGHT * section_count as f64, + layout: Some("vertical".into()), + gap: Some(0.0), + padding: Some(0.0), + fill: Some("#FFFFFF".into()), + }, + subtasks, + style_guide: None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn req(prompt: &str) -> DesignRequest { + DesignRequest { + prompt: prompt.into(), + model: None, + provider: None, + } + } + + #[test] + fn parse_plan_reads_fenced_json() { + let text = r#"Here is the plan: +```json +{ + "root_frame": { "id": "root", "name": "Page", "width": 1200, "height": 800 }, + "subtasks": [ + { "id": "hero", "label": "Hero", "region": { "width": 1200, "height": 400 } } + ] +} +``` +done."#; + let plan = parse_plan(text).expect("parse"); + assert_eq!(plan.root_frame.id, "root"); + assert_eq!(plan.subtasks.len(), 1); + assert_eq!(plan.subtasks[0].id, "hero"); + } + + #[test] + fn parse_plan_rejects_no_subtasks() { + let text = r#"{ "root_frame": { "id": "r", "name": "P", "width": 1, "height": 1 }, "subtasks": [] }"#; + assert!(parse_plan(text).is_err()); + } + + #[test] + fn parse_plan_rejects_garbage() { + assert!(parse_plan("the model refused to answer").is_err()); + } + + #[test] + fn fallback_plan_is_usable() { + let short = build_fallback_plan(&req("a button")); + assert_eq!(short.subtasks.len(), 1); + assert!(short.root_frame.width > 0.0); + + let long = build_fallback_plan(&req(&"x".repeat(300))); + assert_eq!(long.subtasks.len(), 3); + assert!(!long.subtasks.is_empty()); + } +} diff --git a/crates/op-orchestrator/src/plan_normalize.rs b/crates/op-orchestrator/src/plan_normalize.rs new file mode 100644 index 000000000..9c1fc44cd --- /dev/null +++ b/crates/op-orchestrator/src/plan_normalize.rs @@ -0,0 +1,130 @@ +//! plan 规范化 —— 单屏路径。 +//! +//! 行为忠实 TS,但**一次性算分类、干净派生**,不做 TS 那种 +//! in-place strip-then-reclassify(`orchestrator.ts:838-845` 自标 +//! fragile)。 + +use crate::plan::OrchestratorPlan; +use crate::types::DesignRequest; + +/// 规范化产出的派生信息。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct NormInfo { + /// 根 frame 窄到移动端宽度 —— scaffold 阶段据此注入固定状态栏。 + pub is_mobile: bool, +} + +/// 移动端宽度上限(含)—— ≤ 此值视为移动端单屏。 +const MOBILE_MAX_WIDTH: f64 = 480.0; + +/// subtask 的 id / label 命中即视为"状态栏"区块 —— 移动端由 +/// scaffold 注入固定状态栏,plan 里若带状态栏 subtask 则剔除。 +fn is_status_bar_subtask(id: &str, label: &str) -> bool { + let hay = format!("{} {}", id.to_lowercase(), label.to_lowercase()); + hay.contains("status bar") || hay.contains("status-bar") || hay.contains("statusbar") +} + +/// 就地规范化 `plan`: +/// - 一次性判定 `is_mobile`(根 frame 宽度); +/// - 移动端剔除 plan 自带的状态栏 subtask(状态栏改由 scaffold 注入); +/// - 给每个 subtask 赋 `id_prefix = id`、`parent_frame_id = 根 id`。 +pub fn normalize(plan: &mut OrchestratorPlan, _req: &DesignRequest) -> NormInfo { + let is_mobile = plan.root_frame.width <= MOBILE_MAX_WIDTH; + + if is_mobile { + plan.subtasks + .retain(|st| !is_status_bar_subtask(&st.id, &st.label)); + } + + let root_id = plan.root_frame.id.clone(); + for st in &mut plan.subtasks { + st.id_prefix = st.id.clone(); + st.parent_frame_id = Some(root_id.clone()); + } + + NormInfo { is_mobile } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{OrchestratorPlan, Region, RootFrameSpec, Subtask}; + + fn req() -> DesignRequest { + DesignRequest { + prompt: "x".into(), + model: None, + provider: None, + } + } + + fn subtask(id: &str, label: &str) -> Subtask { + Subtask { + id: id.into(), + label: label.into(), + region: Region { + width: 100.0, + height: 100.0, + }, + id_prefix: String::new(), + parent_frame_id: None, + } + } + + fn plan(width: f64, subtasks: Vec) -> OrchestratorPlan { + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "P".into(), + width, + height: 800.0, + layout: None, + gap: None, + padding: None, + fill: None, + }, + subtasks, + style_guide: None, + } + } + + #[test] + fn normalize_assigns_id_prefix_and_parent() { + let mut p = plan(1200.0, vec![subtask("hero", "Hero"), subtask("feat", "Features")]); + let info = normalize(&mut p, &req()); + assert!(!info.is_mobile); + for st in &p.subtasks { + assert_eq!(st.id_prefix, st.id); + assert_eq!(st.parent_frame_id.as_deref(), Some("root")); + } + } + + #[test] + fn normalize_flags_mobile_by_width() { + let mut p = plan(390.0, vec![subtask("hero", "Hero")]); + let info = normalize(&mut p, &req()); + assert!(info.is_mobile); + } + + #[test] + fn normalize_strips_status_bar_subtask_on_mobile() { + let mut p = plan( + 390.0, + vec![subtask("status-bar", "Status Bar"), subtask("hero", "Hero")], + ); + normalize(&mut p, &req()); + assert_eq!(p.subtasks.len(), 1); + assert_eq!(p.subtasks[0].id, "hero"); + } + + #[test] + fn normalize_keeps_status_bar_subtask_on_desktop() { + // 桌面端不剔除(只有移动端 scaffold 注入固定状态栏)。 + let mut p = plan( + 1200.0, + vec![subtask("status-bar", "Status Bar"), subtask("hero", "Hero")], + ); + normalize(&mut p, &req()); + assert_eq!(p.subtasks.len(), 2); + } +} diff --git a/crates/op-orchestrator/src/prompt.rs b/crates/op-orchestrator/src/prompt.rs new file mode 100644 index 000000000..0705161f8 --- /dev/null +++ b/crates/op-orchestrator/src/prompt.rs @@ -0,0 +1,163 @@ +//! Prompt 构造 —— 规划 prompt 与 sub-agent prompt。 +//! +//! skill 正文经 `op_ai_skills::resolve_skills` 解析:规划用 +//! `Phase::Planning`,sub-agent 用 `Phase::Generation`。两段格式 +//! 指令(产 OrchestratorPlan JSON / 产 PenNode JSON 数组)由本 +//! 模块附加在 skill 正文之后。 +//! +//! 注:格式指令文本是功能性的最小版;逐条对齐 TS +//! `orchestrator-prompt-optimizer.ts` / `orchestrator-sub-agent.ts` +//! 的细则(style-guide 注入、移动端禁 phone-wrapper 等)是后续 +//! 细化项。 + +use crate::plan::{OrchestratorPlan, Subtask}; +use crate::types::{AbortFlag, CallRequest, DesignRequest}; +use std::time::Duration; + +const PLANNING_TIMEOUT: Duration = Duration::from_secs(300); +const SUBAGENT_TIMEOUT: Duration = Duration::from_secs(420); + +/// 规划阶段要求模型产出的 JSON 形状说明。 +const PLAN_FORMAT: &str = r#" +Respond with a single JSON object describing the design plan: +{ + "root_frame": { "id": "root", "name": "", "width": , "height": , + "layout": "vertical", "gap": , "padding": , "fill": "#RRGGBB" }, + "subtasks": [ + { "id": "", "label": "", + "region": { "width": , "height": } } + ], + "style_guide": { "palette": { "color-1": "#RRGGBB", "color-2": "#RRGGBB" } } +} +Each subtask is one visual section. Use 1-6 subtasks. Output ONLY the JSON object."#; + +/// sub-agent 阶段要求模型产出的 JSON 形状说明。 +const NODE_FORMAT: &str = r#" +Respond with a JSON array of canonical PenNode objects for THIS section only. +Each node is tagged by "type" (frame/group/rectangle/ellipse/line/polygon/path/ +text/text_input/image/icon_font). Frames/groups nest children via "children". +Example: [ { "type": "frame", "id": "-1", "name": "Card", "x": 0, "y": 0, +"width": 1200, "height": 200, "children": [] } ] +Output ONLY the JSON array."#; + +/// 把解析出的 skill 正文 join 成一段。 +fn skill_preamble(phase: op_ai_skills::Phase, message: &str) -> String { + let ctx = op_ai_skills::resolve_skills(phase, message, &op_ai_skills::ResolveOptions::default()); + ctx.skills + .iter() + .map(|s| s.content.as_str()) + .collect::>() + .join("\n\n") +} + +/// 规划阶段的 LLM 调用输入。 +pub fn build_orchestrator_prompt(req: &DesignRequest, abort: AbortFlag) -> CallRequest { + let mut system_prompt = skill_preamble(op_ai_skills::Phase::Planning, &req.prompt); + system_prompt.push_str("\n\n"); + system_prompt.push_str(PLAN_FORMAT); + CallRequest { + system_prompt, + user_prompt: req.prompt.clone(), + model: req.model.clone(), + provider: req.provider.clone(), + timeout: PLANNING_TIMEOUT, + abort, + } +} + +/// 单个 sub-agent 的 LLM 调用输入。 +pub fn build_subagent_prompt( + subtask: &Subtask, + plan: &OrchestratorPlan, + req: &DesignRequest, + abort: AbortFlag, +) -> CallRequest { + let mut system_prompt = skill_preamble(op_ai_skills::Phase::Generation, &subtask.label); + system_prompt.push_str("\n\n"); + system_prompt.push_str(NODE_FORMAT); + + let palette = plan + .style_guide + .as_ref() + .map(|sg| { + sg.palette + .iter() + .map(|(k, v)| format!("{k}={v}")) + .collect::>() + .join(", ") + }) + .unwrap_or_default(); + + let user_prompt = format!( + "Overall design: {}\n\nGenerate the section \"{}\" \ + (区块 id 前缀 `{}-`). Target region: {:.0}x{:.0} px.\nPalette: {}", + req.prompt, subtask.label, subtask.id_prefix, subtask.region.width, subtask.region.height, + if palette.is_empty() { "(default)" } else { &palette }, + ); + + CallRequest { + system_prompt, + user_prompt, + model: req.model.clone(), + provider: req.provider.clone(), + timeout: SUBAGENT_TIMEOUT, + abort, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{Region, RootFrameSpec}; + + fn req() -> DesignRequest { + DesignRequest { + prompt: "a pricing page".into(), + model: Some("claude".into()), + provider: None, + } + } + + fn plan() -> OrchestratorPlan { + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "P".into(), + width: 1200.0, + height: 800.0, + layout: None, + gap: None, + padding: None, + fill: None, + }, + subtasks: vec![], + style_guide: None, + } + } + + #[test] + fn orchestrator_prompt_carries_request_and_format() { + let cr = build_orchestrator_prompt(&req(), AbortFlag::new()); + assert!(cr.user_prompt.contains("a pricing page")); + assert!(cr.system_prompt.contains("subtasks")); + assert_eq!(cr.model.as_deref(), Some("claude")); + } + + #[test] + fn subagent_prompt_carries_subtask_and_node_format() { + let st = Subtask { + id: "hero".into(), + label: "Hero".into(), + region: Region { + width: 1200.0, + height: 400.0, + }, + id_prefix: "hero".into(), + parent_frame_id: Some("root".into()), + }; + let cr = build_subagent_prompt(&st, &plan(), &req(), AbortFlag::new()); + assert!(cr.user_prompt.contains("Hero")); + assert!(cr.user_prompt.contains("hero-")); + assert!(cr.system_prompt.contains("PenNode")); + } +} diff --git a/crates/op-orchestrator/src/run.rs b/crates/op-orchestrator/src/run.rs new file mode 100644 index 000000000..4d0e0cd70 --- /dev/null +++ b/crates/op-orchestrator/src/run.rs @@ -0,0 +1,268 @@ +//! `Orchestrator::run()` —— 四阶段编排主轴(spec §4)。 +//! +//! 规划 → 画布搭建 → 顺序子 agent → 清理。副作用全经 +//! [`DocSink`] / [`LlmClient`]。错误 / abort / 零内容语义见 +//! spec §6。 + +use crate::cleanup::{descendant_count, run_cleanup_passes}; +use crate::plan::{build_fallback_plan, parse_plan}; +use crate::plan_normalize::normalize; +use crate::prompt::build_orchestrator_prompt; +use crate::scaffold::build_scaffold; +use crate::subagent::run_subtask; +use crate::types::{ + AbortFlag, DesignRequest, DocSink, LlmChunk, LlmClient, OrchestratorError, Progress, RunSummary, + SubtaskOutcome, +}; +use crate::variables::{rollback, seed_commands, snapshot_plan_vars}; +use futures::StreamExt; +use op_editor_core::{EditorCommand, NodeId}; + +/// 设计编排器。S3a 阶段无构造期配置 —— 保留 struct 以便 S3b/S3c +/// 挂选项。 +#[derive(Debug, Default, Clone, Copy)] +pub struct Orchestrator; + +impl Orchestrator { + pub fn new() -> Self { + Orchestrator + } + + /// 跑一次完整编排。见 spec §4 数据流。 + pub async fn run( + &self, + request: DesignRequest, + sink: &mut dyn DocSink, + llm: &dyn LlmClient, + on_progress: &mut dyn FnMut(Progress), + abort: &AbortFlag, + ) -> Result { + // -- 阶段 1:规划 -- + on_progress(Progress::Planning); + let plan_call = build_orchestrator_prompt(&request, abort.clone()); + let mut plan = match collect_text(llm.call(plan_call)).await { + Ok(text) => parse_plan(&text).unwrap_or_else(|_| build_fallback_plan(&request)), + Err(aborted) => { + if aborted { + // abort 在规划阶段:尚未动文档,直接返回。 + return Err(OrchestratorError::Aborted); + } + // 非 abort 失败 → 启发式兜底 plan。 + build_fallback_plan(&request) + } + }; + + // -- 阶段 1.5:规范化 -- + let norm = normalize(&mut plan, &request); + let root_id = plan.root_frame.id.clone(); + + // -- 进入"已动文档"区,全程 undo batch 包裹 -- + sink.begin_undo_batch(); + let var_snapshot = snapshot_plan_vars(sink, &plan); + + // -- 阶段 2:画布搭建 -- + for cmd in seed_commands(&plan) { + sink.apply(cmd); + } + match build_scaffold(&plan, norm.is_mobile) { + Ok(cmds) => { + for cmd in cmds { + sink.apply(cmd); + } + } + Err(e) => { + // scaffold 模板 bug —— 收尾后报内部错误。 + rollback(sink, &var_snapshot); + sink.end_undo_batch(); + return Err(OrchestratorError::Internal(e)); + } + } + let scaffold_baseline = descendant_count(sink.state(), &root_id); + on_progress(Progress::ScaffoldDone); + + // -- 阶段 3:顺序子 agent -- + let mut outcomes: Vec = Vec::new(); + let mut aborted_mid = false; + let mut zero_node_failure = false; + for subtask in &plan.subtasks { + if abort.is_set() { + aborted_mid = true; + break; + } + on_progress(Progress::SubtaskStarted { + id: subtask.id.clone(), + label: subtask.label.clone(), + }); + let outcome = run_subtask(subtask, &plan, &request, llm, sink, abort).await; + if outcome.node_count == 0 { + // 零节点失败 —— 停止后续 subtask(spec §6.2)。 + on_progress(Progress::SubtaskFailed { + id: subtask.id.clone(), + error: outcome.error.clone().unwrap_or_default(), + }); + outcomes.push(outcome); + zero_node_failure = true; + break; + } + on_progress(Progress::SubtaskDone { + id: subtask.id.clone(), + node_count: outcome.node_count, + }); + outcomes.push(outcome); + } + + // -- 阶段 4:清理 -- + run_cleanup_passes(sink, &plan); + on_progress(Progress::CleanupDone); + + // -- 阶段 4.5:收尾(spec §6.3 三路径)-- + let zero_content = descendant_count(sink.state(), &root_id) <= scaffold_baseline; + if zero_content { + // 错误路径才移除空 scaffold root;abort / 正常零内容只回滚变量。 + if zero_node_failure { + sink.apply(EditorCommand::DeleteNode { + node_id: NodeId::new(root_id.clone()), + }); + } + rollback(sink, &var_snapshot); + } + sink.end_undo_batch(); + + // -- 返回 -- + if zero_content { + return Err(if aborted_mid { + OrchestratorError::Aborted + } else { + OrchestratorError::NoContent + }); + } + let total_nodes = outcomes.iter().map(|o| o.node_count).sum(); + Ok(RunSummary { + root_frame_id: root_id, + subtasks: outcomes, + total_nodes, + }) + } +} + +/// 消费一次 LLM 调用的流 —— 拼接所有 `Text` chunk,丢弃 `Thinking`。 +/// `Err(true)` 表示中止,`Err(false)` 表示真实错误。 +async fn collect_text( + mut stream: futures::stream::BoxStream<'static, Result>, +) -> Result { + let mut text = String::new(); + while let Some(item) = stream.next().await { + match item { + Ok(LlmChunk::Text(t)) => text.push_str(&t), + Ok(LlmChunk::Thinking(_)) => {} + Err(e) => return Err(e.aborted), + } + } + Ok(text) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::test_support::{ScriptResponse, ScriptedLlm, VecDocSink}; + + fn req() -> DesignRequest { + DesignRequest { + prompt: "a landing page".into(), + model: None, + provider: None, + } + } + + const PLAN_JSON: &str = r#"{ + "root_frame": { "id": "root", "name": "Page", "width": 1200, "height": 800, + "layout": "vertical", "gap": 0, "fill": "#FFFFFF" }, + "subtasks": [ + { "id": "hero", "label": "Hero", "region": { "width": 1200, "height": 400 } }, + { "id": "feat", "label": "Features", "region": { "width": 1200, "height": 400 } } + ] + }"#; + + fn node_json(prefix: &str) -> String { + format!( + r#"[{{"type":"frame","id":"{prefix}-1","name":"Sec","x":0,"y":0,"width":1200,"height":300,"children":[]}}]"# + ) + } + + #[test] + fn run_happy_path_applies_scaffold_and_subtasks() { + let llm = ScriptedLlm::new(vec![ + ScriptResponse::Text(PLAN_JSON.into()), + ScriptResponse::Text(node_json("hero")), + ScriptResponse::Text(node_json("feat")), + ]); + let mut sink = VecDocSink::new(); + let mut events: Vec = Vec::new(); + let mut on_progress = |p: Progress| events.push(p); + + let summary = futures::executor::block_on(Orchestrator::new().run( + req(), + &mut sink, + &llm, + &mut on_progress, + &AbortFlag::new(), + )) + .expect("run ok"); + + assert_eq!(summary.root_frame_id, "root"); + assert_eq!(summary.subtasks.len(), 2); + assert!(summary.total_nodes >= 2); + // undo batch 配对。 + assert_eq!(sink.batch_depth, 0); + // 至少有 scaffold + 两个 subtask 的 InsertSubtree。 + let inserts = sink + .applied + .iter() + .filter(|c| matches!(c, EditorCommand::InsertSubtree { .. })) + .count(); + assert!(inserts >= 3, "expected >=3 InsertSubtree, got {inserts}"); + assert!(matches!(events.first(), Some(Progress::Planning))); + assert!(matches!(events.last(), Some(Progress::CleanupDone))); + } + + #[test] + fn run_zero_node_subtask_stops_and_errors() { + // 规划 OK,但第一个 subtask 吐垃圾 → 零节点失败 → 停止。 + let llm = ScriptedLlm::new(vec![ + ScriptResponse::Text(PLAN_JSON.into()), + ScriptResponse::Text("the model refused".into()), + ]); + let mut sink = VecDocSink::new(); + let mut on_progress = |_p: Progress| {}; + let result = futures::executor::block_on(Orchestrator::new().run( + req(), + &mut sink, + &llm, + &mut on_progress, + &AbortFlag::new(), + )); + assert!(matches!(result, Err(OrchestratorError::NoContent))); + // undo batch 仍配对。 + assert_eq!(sink.batch_depth, 0); + } + + #[test] + fn run_planning_failure_uses_fallback_plan() { + // 规划吐垃圾 → fallback plan;subtask 正常 → 成功。 + let llm = ScriptedLlm::new(vec![ + ScriptResponse::Text("no json here".into()), + ScriptResponse::Text(node_json("section-1")), + ]); + let mut sink = VecDocSink::new(); + let mut on_progress = |_p: Progress| {}; + let summary = futures::executor::block_on(Orchestrator::new().run( + req(), + &mut sink, + &llm, + &mut on_progress, + &AbortFlag::new(), + )) + .expect("fallback run ok"); + assert!(summary.total_nodes >= 1); + } +} diff --git a/crates/op-orchestrator/src/scaffold.rs b/crates/op-orchestrator/src/scaffold.rs new file mode 100644 index 000000000..f2afff670 --- /dev/null +++ b/crates/op-orchestrator/src/scaffold.rs @@ -0,0 +1,115 @@ +//! 阶段 2 —— 单屏画布搭建。 +//! +//! 产出一条 `InsertSubtree`:把根 frame(移动端再带一个固定状态 +//! 栏 child)插到活动页根。根 frame 用 JSON 构建后反序列化为 +//! canonical `PenNode` —— 避免在 Rust 侧硬写富 schema 的每个字段, +//! 且与 `parse` 模块的解析路径一致。 + +use crate::plan::OrchestratorPlan; +use jian_ops_schema::node::PenNode; +use op_editor_core::{EditorCommand, NodeId}; +use serde_json::json; + +/// 移动端固定状态栏的高度(iOS 风格 chrome,通用 mockup 约定)。 +const STATUS_BAR_HEIGHT: f64 = 44.0; + +/// 构建阶段 2 的画布命令。`is_mobile` 为真时根 frame 带一个固定 +/// 状态栏 child。返回 `Err` 表示根 frame JSON 模板有问题(实现 +/// bug,非用户输入问题)。 +pub fn build_scaffold( + plan: &OrchestratorPlan, + is_mobile: bool, +) -> Result, String> { + let rf = &plan.root_frame; + let layout = rf.layout.as_deref().unwrap_or("vertical"); + let fill_hex = rf.fill.as_deref().unwrap_or("#FFFFFF"); + + let children = if is_mobile { + json!([{ + "type": "frame", + "id": format!("{}-status-bar", rf.id), + "name": "Status Bar", + "x": 0, + "y": 0, + "width": rf.width, + "height": STATUS_BAR_HEIGHT, + "fill": [{ "type": "solid", "color": fill_hex }], + "children": [], + }]) + } else { + json!([]) + }; + + let frame = json!({ + "type": "frame", + "id": rf.id, + "name": rf.name, + "x": 0, + "y": 0, + "width": rf.width, + "height": rf.height, + "layout": layout, + "gap": rf.gap.unwrap_or(0.0), + "fill": [{ "type": "solid", "color": fill_hex }], + "children": children, + }); + + let node: PenNode = + serde_json::from_value(frame).map_err(|e| format!("scaffold root frame: {e}"))?; + + Ok(vec![EditorCommand::InsertSubtree { + nodes: vec![node], + parent_id: NodeId::NONE, + }]) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{OrchestratorPlan, RootFrameSpec}; + use op_editor_core::PenNodeExt; + + fn plan() -> OrchestratorPlan { + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "Design".into(), + width: 1200.0, + height: 800.0, + layout: Some("vertical".into()), + gap: Some(0.0), + padding: Some(0.0), + fill: Some("#FFFFFF".into()), + }, + subtasks: vec![], + style_guide: None, + } + } + + #[test] + fn build_scaffold_desktop_one_root_no_children() { + let cmds = build_scaffold(&plan(), false).expect("scaffold"); + assert_eq!(cmds.len(), 1); + match &cmds[0] { + EditorCommand::InsertSubtree { nodes, parent_id } => { + assert_eq!(nodes.len(), 1); + assert!(!parent_id.is_real()); // NONE → page root + assert_eq!(nodes[0].id_str(), "root"); + assert!(nodes[0].children().map(|c| c.is_empty()).unwrap_or(true)); + } + other => panic!("expected InsertSubtree, got {other:?}"), + } + } + + #[test] + fn build_scaffold_mobile_injects_status_bar() { + let cmds = build_scaffold(&plan(), true).expect("scaffold"); + match &cmds[0] { + EditorCommand::InsertSubtree { nodes, .. } => { + let children = nodes[0].children().expect("frame children"); + assert_eq!(children.len(), 1); + } + other => panic!("expected InsertSubtree, got {other:?}"), + } + } +} diff --git a/crates/op-orchestrator/src/subagent.rs b/crates/op-orchestrator/src/subagent.rs new file mode 100644 index 000000000..1a156a695 --- /dev/null +++ b/crates/op-orchestrator/src/subagent.rs @@ -0,0 +1,177 @@ +//! 阶段 3 —— 单个 sub-agent 的顺序执行。 +//! +//! 一个 subtask:构 prompt → 调 `LlmClient` → 收集文本 → 解析成 +//! `PenNode` 树 → 经 `DocSink` 发一条 `InsertSubtree`。 +//! +//! 返回的 [`SubtaskOutcome`] 用 `node_count` 区分(见 spec §6.2): +//! - `node_count == 0` —— 零节点失败,调用方应停止后续 subtask; +//! - `node_count > 0`(`error` 可带软错误)—— 部分产出,继续后续。 + +use crate::parse::parse_nodes; +use crate::plan::{OrchestratorPlan, Subtask}; +use crate::prompt::build_subagent_prompt; +use crate::types::{AbortFlag, DesignRequest, DocSink, LlmChunk, LlmClient, SubtaskOutcome}; +use futures::StreamExt; +use op_editor_core::{EditorCommand, NodeId}; + +/// 执行一个 subtask。总是返回 [`SubtaskOutcome`];调用方据 +/// `node_count` 决定继续/停止。 +pub async fn run_subtask( + subtask: &Subtask, + plan: &OrchestratorPlan, + req: &DesignRequest, + llm: &dyn LlmClient, + sink: &mut dyn DocSink, + abort: &AbortFlag, +) -> SubtaskOutcome { + let fail = |msg: String| SubtaskOutcome { + id: subtask.id.clone(), + node_count: 0, + error: Some(msg), + }; + + // 收集 LLM 文本输出。 + let call_req = build_subagent_prompt(subtask, plan, req, abort.clone()); + let mut stream = llm.call(call_req); + let mut text = String::new(); + while let Some(item) = stream.next().await { + match item { + Ok(LlmChunk::Text(t)) => text.push_str(&t), + Ok(LlmChunk::Thinking(_)) => {} + Err(e) => { + return fail(if e.aborted { + "aborted".into() + } else { + e.message + }); + } + } + } + + // 解析成 PenNode 树。 + let nodes = match parse_nodes(&text) { + Ok(n) => n, + Err(e) => return fail(e.to_string()), + }; + let node_count = nodes.len(); + + // 经 DocSink 发 InsertSubtree。 + let parent_id = match &subtask.parent_frame_id { + Some(id) => NodeId::new(id.clone()), + None => NodeId::NONE, + }; + let applied = sink.apply(EditorCommand::InsertSubtree { nodes, parent_id }); + if !applied { + return fail("InsertSubtree rejected by document".into()); + } + + SubtaskOutcome { + id: subtask.id.clone(), + node_count, + error: None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{OrchestratorPlan, Region, RootFrameSpec}; + use crate::test_support::{ScriptResponse, ScriptedLlm, VecDocSink}; + use crate::types::LlmError; + use futures::executor::block_on; + + fn req() -> DesignRequest { + DesignRequest { + prompt: "a page".into(), + model: None, + provider: None, + } + } + + fn plan() -> OrchestratorPlan { + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "P".into(), + width: 1200.0, + height: 800.0, + layout: None, + gap: None, + padding: None, + fill: None, + }, + subtasks: vec![], + style_guide: None, + } + } + + fn subtask() -> Subtask { + Subtask { + id: "hero".into(), + label: "Hero".into(), + region: Region { + width: 1200.0, + height: 400.0, + }, + id_prefix: "hero".into(), + parent_frame_id: None, + } + } + + const NODE_JSON: &str = r#"[{"type":"frame","id":"hero-1","name":"Card","x":0,"y":0,"width":1200,"height":200,"children":[]}]"#; + + #[test] + fn run_subtask_ok_applies_insert_subtree() { + let llm = ScriptedLlm::new(vec![ScriptResponse::Text(NODE_JSON.into())]); + let mut sink = VecDocSink::new(); + let outcome = block_on(run_subtask( + &subtask(), + &plan(), + &req(), + &llm, + &mut sink, + &AbortFlag::new(), + )); + assert_eq!(outcome.node_count, 1); + assert!(outcome.error.is_none()); + assert!(matches!( + sink.applied.last(), + Some(EditorCommand::InsertSubtree { .. }) + )); + } + + #[test] + fn run_subtask_zero_node_on_garbage() { + let llm = ScriptedLlm::new(vec![ScriptResponse::Text("the model refused".into())]); + let mut sink = VecDocSink::new(); + let outcome = block_on(run_subtask( + &subtask(), + &plan(), + &req(), + &llm, + &mut sink, + &AbortFlag::new(), + )); + assert_eq!(outcome.node_count, 0); + assert!(outcome.error.is_some()); + } + + #[test] + fn run_subtask_zero_node_on_llm_error() { + let llm = ScriptedLlm::new(vec![ScriptResponse::Fail(LlmError { + message: "rate limited".into(), + aborted: false, + })]); + let mut sink = VecDocSink::new(); + let outcome = block_on(run_subtask( + &subtask(), + &plan(), + &req(), + &llm, + &mut sink, + &AbortFlag::new(), + )); + assert_eq!(outcome.node_count, 0); + assert_eq!(outcome.error.as_deref(), Some("rate limited")); + } +} diff --git a/crates/op-orchestrator/src/test_support.rs b/crates/op-orchestrator/src/test_support.rs new file mode 100644 index 000000000..d89e182d9 --- /dev/null +++ b/crates/op-orchestrator/src/test_support.rs @@ -0,0 +1,76 @@ +//! 测试桩 —— crate 内各测试模块共用。仅在 `cfg(test)` 下编译。 + +use crate::types::{CallRequest, DocSink, LlmChunk, LlmClient, LlmError}; +use futures::stream::BoxStream; +use op_editor_core::{EditorCommand, EditorState}; +use std::collections::VecDeque; +use std::sync::Mutex; + +/// 录命令 + 持内存 `EditorState` 的 `DocSink` 实现。 +pub(crate) struct VecDocSink { + pub state: EditorState, + pub applied: Vec, + pub batch_depth: i32, +} + +impl VecDocSink { + pub(crate) fn new() -> Self { + Self { + state: EditorState::new(), + applied: Vec::new(), + batch_depth: 0, + } + } +} + +impl DocSink for VecDocSink { + fn state(&self) -> &EditorState { + &self.state + } + fn apply(&mut self, cmd: EditorCommand) -> bool { + self.applied.push(cmd.clone()); + self.state.apply(cmd) + } + fn begin_undo_batch(&mut self) { + self.batch_depth += 1; + } + fn end_undo_batch(&mut self) { + self.batch_depth -= 1; + } +} + +/// 一次脚本化的 LLM 响应。 +pub(crate) enum ScriptResponse { + /// 一段成功文本(作为单个 `Text` chunk 返回)。 + Text(String), + /// 一次失败。 + Fail(LlmError), +} + +/// 按脚本逐次返回响应的 `LlmClient` 桩 —— 每次 `call` 弹一条。 +pub(crate) struct ScriptedLlm { + responses: Mutex>, +} + +impl ScriptedLlm { + pub(crate) fn new(responses: Vec) -> Self { + Self { + responses: Mutex::new(responses.into()), + } + } +} + +impl LlmClient for ScriptedLlm { + fn call(&self, _req: CallRequest) -> BoxStream<'static, Result> { + let next = self.responses.lock().unwrap().pop_front(); + let items: Vec> = match next { + Some(ScriptResponse::Text(t)) => vec![Ok(LlmChunk::Text(t))], + Some(ScriptResponse::Fail(e)) => vec![Err(e)], + None => vec![Err(LlmError { + message: "scripted LLM exhausted".into(), + aborted: false, + })], + }; + Box::pin(futures::stream::iter(items)) + } +} diff --git a/crates/op-orchestrator/src/types.rs b/crates/op-orchestrator/src/types.rs new file mode 100644 index 000000000..a0cf11abe --- /dev/null +++ b/crates/op-orchestrator/src/types.rs @@ -0,0 +1,168 @@ +//! 公共 trait 与类型 —— 编排器与外界的全部接口。 +//! +//! 副作用只从 [`DocSink`](文档变更)与 [`LlmClient`](LLM 调用) +//! 两个 trait 进出;其余全是数据类型。 + +use futures::stream::BoxStream; +use op_editor_core::{EditorCommand, EditorState}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +/// 文档出口。写经 [`apply`](DocSink::apply);读经 +/// [`state`](DocSink::state)。host 实现把 apply 与存盘 / 重绘 / +/// undo 批界绑在一处(对齐 `op-mcp` 的 applier 模式)。 +pub trait DocSink: Send { + /// 只读访问当前文档 —— cleanup 判定空 scaffold、未来 S3c + /// 校验都要读。 + fn state(&self) -> &EditorState; + /// 应用一条编辑命令;返回 `false` 表示命令被拒(文档未变)。 + fn apply(&mut self, cmd: EditorCommand) -> bool; + /// 开启一个 undo 批 —— 批内的所有 apply 合并为一次 undo。 + fn begin_undo_batch(&mut self); + /// 关闭当前 undo 批。 + fn end_undo_batch(&mut self); +} + +/// LLM 调用出口。每次 [`call`](LlmClient::call) 是一次独立、无累积 +/// 上下文的 LLM turn —— host 实现应为每次调用新建引擎,使规划与 +/// 各 sub-agent 拿到隔离上下文。 +pub trait LlmClient: Send + Sync { + fn call(&self, req: CallRequest) -> BoxStream<'static, Result>; +} + +/// 一次 LLM 调用的完整输入。字段一次定全,S3b 并发 / 流式不必 +/// 改 trait 签名。 +#[derive(Debug, Clone)] +pub struct CallRequest { + pub system_prompt: String, + pub user_prompt: String, + pub model: Option, + pub provider: Option, + pub timeout: Duration, + pub abort: AbortFlag, +} + +/// 流元素 —— 区分文本与思考,与 TS 的 text/thinking/error 三分 +/// 一致。错误走 `Result` 的 `Err(LlmError)`,不混入 chunk。 +#[derive(Debug, Clone)] +pub enum LlmChunk { + Text(String), + Thinking(String), +} + +/// 一次 LLM 调用的失败。`aborted` 区分用户中止与真实错误。 +#[derive(Debug, Clone)] +pub struct LlmError { + pub message: String, + pub aborted: bool, +} + +/// 廉价可克隆的中止句柄(`Arc` 语义)。 +#[derive(Debug, Clone, Default)] +pub struct AbortFlag(Arc); + +impl AbortFlag { + pub fn new() -> Self { + Self::default() + } + /// 置位 —— 之后所有 `is_set` 返回 `true`。 + pub fn set(&self) { + self.0.store(true, Ordering::SeqCst); + } + /// 是否已被中止。 + pub fn is_set(&self) -> bool { + self.0.load(Ordering::SeqCst) + } +} + +/// 用户消息的意图分类 —— 决定走编排器还是普通聊天。 +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Intent { + Design, + Chat, +} + +/// `run()` 的进度回调载荷。字段从 S3a 定全,S3b/S3c 只填新分支。 +#[derive(Debug, Clone)] +pub enum Progress { + Planning, + ScaffoldDone, + SubtaskStarted { id: String, label: String }, + SubtaskDone { id: String, node_count: usize }, + SubtaskFailed { id: String, error: String }, + CleanupDone, +} + +/// 单个 subtask 的执行结果。`error` 带值但 `node_count > 0` 表示 +/// "部分产出"(软错误);`node_count == 0` 表示零节点失败。 +#[derive(Debug, Clone)] +pub struct SubtaskOutcome { + pub id: String, + pub node_count: usize, + pub error: Option, +} + +/// `run()` 成功返回的汇总。 +#[derive(Debug, Clone)] +pub struct RunSummary { + pub root_frame_id: String, + pub subtasks: Vec, + pub total_nodes: usize, +} + +/// `run()` 的失败。 +#[derive(Debug, Clone)] +pub enum OrchestratorError { + /// 用户中途取消。 + Aborted, + /// 跑完但未产出任何真实内容。 + NoContent, + /// 内部错误(意外情况)。 + Internal(String), +} + +impl std::fmt::Display for OrchestratorError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + OrchestratorError::Aborted => write!(f, "orchestration aborted by user"), + OrchestratorError::NoContent => write!(f, "orchestration produced no content"), + OrchestratorError::Internal(m) => write!(f, "orchestration internal error: {m}"), + } + } +} + +impl std::error::Error for OrchestratorError {} + +/// 编排器输入 —— 一次设计请求。 +#[derive(Debug, Clone)] +pub struct DesignRequest { + pub prompt: String, + pub model: Option, + pub provider: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::test_support::VecDocSink; + + #[test] + fn vec_doc_sink_implements_docsink() { + let mut sink = VecDocSink::new(); + sink.begin_undo_batch(); + sink.end_undo_batch(); + assert_eq!(sink.batch_depth, 0); + assert!(sink.applied.is_empty()); + let _ = sink.state(); + } + + #[test] + fn abort_flag_sets_and_reads() { + let flag = AbortFlag::new(); + assert!(!flag.is_set()); + let clone = flag.clone(); + flag.set(); + assert!(clone.is_set()); + } +} diff --git a/crates/op-orchestrator/src/variables.rs b/crates/op-orchestrator/src/variables.rs new file mode 100644 index 000000000..4281577ec --- /dev/null +++ b/crates/op-orchestrator/src/variables.rs @@ -0,0 +1,129 @@ +//! plan 派生变量 —— seed / 快照 / 回滚。 +//! +//! 规划产出的 `style_guide.palette` 在阶段 2 被 seed 成文档里的 +//! 颜色变量;sub-agent 生成的节点会带 `$color-*` 引用。若本轮零 +//! 内容,这些 seed 出来的变量必须回滚,否则留下 dangling 引用。 +//! +//! 设计:`seed` 只用 `CreateVariable`(只创建不存在的变量,不 +//! 覆盖用户既有变量);快照记下"哪些 plan 变量名在 seed 前不 +//! 存在",回滚就 `DeleteVariable` 这些 —— 正好删掉 seed 真正新建 +//! 的那批,既有变量不受影响。 + +use crate::plan::OrchestratorPlan; +use crate::types::DocSink; +use op_editor_core::EditorCommand; + +/// 回滚快照 —— seed 前"不存在"的 plan 变量名集合(即 seed 会 +/// 真正新建的那批)。 +#[derive(Debug, Clone, Default)] +pub struct VarSnapshot { + /// seed 前不存在的变量名 —— 回滚时删除这些。 + pub created: Vec, +} + +/// 在 seed *之前* 调用 —— 记下 plan 调色板里哪些变量名当前 +/// 不存在(seed 会新建它们)。 +pub fn snapshot_plan_vars(sink: &dyn DocSink, plan: &OrchestratorPlan) -> VarSnapshot { + let state = sink.state(); + let created = palette_iter(plan) + .filter(|(name, _)| state.find_variable(name).is_none()) + .map(|(name, _)| name.clone()) + .collect(); + VarSnapshot { created } +} + +/// plan 调色板 → `CreateVariable` 命令(每个颜色一条)。 +/// `CreateVariable` 对已存在的同名变量会被 applier 拒(返回 +/// `false`),故不会覆盖用户既有变量。 +pub fn seed_commands(plan: &OrchestratorPlan) -> Vec { + palette_iter(plan) + .map(|(name, hex)| EditorCommand::CreateVariable { + name: name.clone(), + kind: "color".into(), + default_value: hex.clone(), + }) + .collect() +} + +/// 回滚 —— 删除 seed 真正新建的那批变量(快照里记下的)。 +/// 既有变量不在快照里,不受影响。 +pub fn rollback(sink: &mut dyn DocSink, snap: &VarSnapshot) { + for name in &snap.created { + sink.apply(EditorCommand::DeleteVariable { name: name.clone() }); + } +} + +/// plan 调色板的 `(名称, hex)` 迭代器 —— `style_guide` 缺省时为空。 +fn palette_iter(plan: &OrchestratorPlan) -> impl Iterator { + plan.style_guide + .as_ref() + .into_iter() + .flat_map(|sg| sg.palette.iter()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::plan::{OrchestratorPlan, RootFrameSpec, StyleGuide}; + use crate::test_support::VecDocSink; + use std::collections::BTreeMap; + + fn plan_with_palette(pairs: &[(&str, &str)]) -> OrchestratorPlan { + let mut palette = BTreeMap::new(); + for (k, v) in pairs { + palette.insert((*k).to_string(), (*v).to_string()); + } + OrchestratorPlan { + root_frame: RootFrameSpec { + id: "root".into(), + name: "P".into(), + width: 1200.0, + height: 800.0, + layout: None, + gap: None, + padding: None, + fill: None, + }, + subtasks: vec![], + style_guide: Some(StyleGuide { palette }), + } + } + + #[test] + fn seed_commands_one_per_palette_color() { + let plan = plan_with_palette(&[("color-1", "#2563EB"), ("color-2", "#0EA5E9")]); + let cmds = seed_commands(&plan); + assert_eq!(cmds.len(), 2); + assert!(matches!( + &cmds[0], + EditorCommand::CreateVariable { kind, .. } if kind == "color" + )); + } + + #[test] + fn seed_then_rollback_round_trips() { + let plan = plan_with_palette(&[("color-1", "#2563EB")]); + let mut sink = VecDocSink::new(); + + // seed 前快照:color-1 不存在 → 进 created 集 + let snap = snapshot_plan_vars(&sink, &plan); + assert_eq!(snap.created, vec!["color-1".to_string()]); + + // seed + for cmd in seed_commands(&plan) { + sink.apply(cmd); + } + assert!(sink.state().find_variable("color-1").is_some()); + + // 回滚 → 变量被删 + rollback(&mut sink, &snap); + assert!(sink.state().find_variable("color-1").is_none()); + } + + #[test] + fn no_style_guide_yields_no_commands() { + let mut plan = plan_with_palette(&[]); + plan.style_guide = None; + assert!(seed_commands(&plan).is_empty()); + } +}