feat(orchestrator): op-orchestrator crate — S3a design orchestrator

New crate restoring the design orchestrator (TS orchestrator.ts
phases 1-4, single-screen sequential) that the Rust migration
dropped. Pure crate — no winit/casement/agent dependency.

- types.rs: DocSink + LlmClient traits, CallRequest/LlmChunk/
  Progress/RunSummary and friends
- intent.rs: classify_intent (design vs chat)
- plan.rs: OrchestratorPlan parse + heuristic fallback
- parse.rs: LLM output -> canonical PenNode trees
- plan_normalize.rs: single-screen plan normalization
- variables.rs: plan-derived $color-* seed/snapshot/rollback
- prompt.rs: planning + sub-agent prompt assembly via op-ai-skills
- scaffold.rs: single-screen canvas scaffold (InsertSubtree)
- subagent.rs: sequential sub-agent execution
- cleanup.rs: descendant_count + run_cleanup_passes (reusable for S3b)
- run.rs: Orchestrator::run() — 4-phase spine, spec §6 error/abort/
  zero-content semantics
- test_support.rs: VecDocSink + ScriptedLlm test stubs

NOT verified locally: op-orchestrator depends on op-editor-core,
which does not currently build — vendor/jian is pinned to unpushed
commit 80121906 (see prior commit). Unit tests written throughout;
they run once the jian build is restored.

S3a Plan B + Plan C (orchestrator modules).
This commit is contained in:
Fini 2026-05-22 22:56:11 +08:00
parent b1c3136deb
commit 5d015ee570
14 changed files with 1763 additions and 0 deletions

View file

@ -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"

View file

@ -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::<usize>()
}
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());
}
}

View file

@ -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}");
}
}
}

View file

@ -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::*;

View file

@ -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<Vec<PenNode>, ParseError> {
let json = extract_json_array(text)
.ok_or_else(|| ParseError("no JSON array found".into()))?;
let nodes: Vec<PenNode> = 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());
}
}

View file

@ -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<String>,
#[serde(default)]
pub gap: Option<f64>,
#[serde(default)]
pub padding: Option<f64>,
#[serde(default)]
pub fill: Option<String>,
}
/// 一个生成子任务 —— 对应设计里的一个区块。
#[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<String>,
}
/// 设计系统提示 —— 调色板等。
#[derive(Debug, Clone, PartialEq, Default, Deserialize)]
pub struct StyleGuide {
/// `名称 -> hex`,阶段 2 据此 seed `$color-*` 文档变量。
#[serde(default)]
pub palette: BTreeMap<String, String>,
}
/// 规划阶段的完整产物。
#[derive(Debug, Clone, PartialEq, Deserialize)]
pub struct OrchestratorPlan {
pub root_frame: RootFrameSpec,
pub subtasks: Vec<Subtask>,
#[serde(default)]
pub style_guide: Option<StyleGuide>,
}
/// 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<OrchestratorPlan, PlanParseError> {
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<Subtask> = (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());
}
}

View file

@ -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<Subtask>) -> 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);
}
}

View file

@ -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": "<name>", "width": <px>, "height": <px>,
"layout": "vertical", "gap": <px>, "padding": <px>, "fill": "#RRGGBB" },
"subtasks": [
{ "id": "<kebab-id>", "label": "<human label>",
"region": { "width": <px>, "height": <px> } }
],
"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": "<prefix>-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::<Vec<_>>()
.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::<Vec<_>>()
.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"));
}
}

View file

@ -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<RunSummary, OrchestratorError> {
// -- 阶段 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<SubtaskOutcome> = 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<LlmChunk, crate::types::LlmError>>,
) -> Result<String, bool> {
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<Progress> = 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);
}
}

View file

@ -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<Vec<EditorCommand>, 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:?}"),
}
}
}

View file

@ -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"));
}
}

View file

@ -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<EditorCommand>,
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<VecDeque<ScriptResponse>>,
}
impl ScriptedLlm {
pub(crate) fn new(responses: Vec<ScriptResponse>) -> Self {
Self {
responses: Mutex::new(responses.into()),
}
}
}
impl LlmClient for ScriptedLlm {
fn call(&self, _req: CallRequest) -> BoxStream<'static, Result<LlmChunk, LlmError>> {
let next = self.responses.lock().unwrap().pop_front();
let items: Vec<Result<LlmChunk, LlmError>> = 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))
}
}

View file

@ -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<LlmChunk, LlmError>>;
}
/// 一次 LLM 调用的完整输入。字段一次定全,S3b 并发 / 流式不必
/// 改 trait 签名。
#[derive(Debug, Clone)]
pub struct CallRequest {
pub system_prompt: String,
pub user_prompt: String,
pub model: Option<String>,
pub provider: Option<String>,
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<AtomicBool>` 语义)。
#[derive(Debug, Clone, Default)]
pub struct AbortFlag(Arc<AtomicBool>);
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<String>,
}
/// `run()` 成功返回的汇总。
#[derive(Debug, Clone)]
pub struct RunSummary {
pub root_frame_id: String,
pub subtasks: Vec<SubtaskOutcome>,
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<String>,
pub provider: Option<String>,
}
#[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());
}
}

View file

@ -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<String>,
}
/// 在 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<EditorCommand> {
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<Item = (&String, &String)> {
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());
}
}