新增: T4 工作流 DAG 注入完整实现
PerConvState: - workflow_id: Option<String> — 关联的工作流 ID - workflow_dag_summary: Option<String> — 预计算 DAG 摘要 system_prompt 注入: - 检测 workflow_dag_summary,非空时注入 [工作流] 块 workflow_context.rs: - build_dag_summary: 从 DagDef 构建可读 DAG 摘要文本 - extract_active_path: 提取从根到当前节点的活跃路径 - topological_layers: 拓扑分层排序(含环检测) - 3 个单元测试覆盖线性 DAG / 中途路径 / 空 DAG
This commit is contained in:
@@ -881,10 +881,20 @@ pub(crate) async fn run_agentic_loop(
|
||||
let behavior_prompt = "\n## AI 定位\n你是 DevFlow 的 AI 助手,拥有完整的工具链。用户只负责提需求和审批,所有执行由你完成——读写文件、运行命令、创建项目、搜索代码等都是你直接调用工具完成的。**绝不输出“请在终端执行以下命令”这类指令——你自己用 run_command 工具执行即可。**\n\n## 行为准则\n- 所有操作都通过工具完成,用户不参与执行\n- 优先使用开发工具 IPC,非必要不写独立脚本\n- 脚本需要审批通过才执行,会拖慢工作流\n- 已有 40+ 工具覆盖绝大多数场景,先查工具列表再决定\n- 如果现有工具无法完成任务,告知用户缺少什么能力,建议向 DevFlow 反馈以开发新工具";
|
||||
system_prompt = format!("{}\n\n{}\n\n{}", system_prompt, env_prompt, behavior_prompt);
|
||||
|
||||
// T4: 工作流 DAG 注入(预留 — 待 PerConvState 加 workflow_id 字段后接入)
|
||||
// 在此处检测 conv.workflow_id,加载 DAG 并注入 WorkflowContextBlock 到 system_prompt
|
||||
// 当前仅日志占位,不改变任何行为
|
||||
tracing::trace!(conv_id = %conv_id, "[ai] T4 工作流注入点已就位");
|
||||
// T4: 工作流 DAG 注入 — 当会话关联工作流时,将活跃路径注入 system prompt
|
||||
{
|
||||
let session = session_arc.lock().await;
|
||||
if let Some(conv) = session.conv_read(&conv_id) {
|
||||
if let Some(ref dag_summary) = conv.workflow_dag_summary {
|
||||
system_prompt = format!("{}[工作流]\n{}\n", system_prompt, dag_summary);
|
||||
tracing::debug!(
|
||||
conv_id = %conv_id,
|
||||
workflow_id = ?conv.workflow_id,
|
||||
"[ai] T4: 已注入工作流 DAG 上下文"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── 多 Agent 并行执行:Coordinator 分解(plan_execution_enabled 时) ──
|
||||
// 对话透明化 L1:拍快照供 AiCompleted 事件携带(coordinator 路径出口也用)
|
||||
|
||||
@@ -1,11 +1,186 @@
|
||||
//! 工作流 DAG 上下文注入 — 当会话关联工作流时,将活跃路径注入 system prompt
|
||||
//!
|
||||
//! T4 实现预留。调用点在 `run_agentic_loop` 中 system_prompt 拼接完成后。
|
||||
//! 当前为空壳,需配合 PerConvState.workflow_id 字段落地后完成。
|
||||
//! T4 实现。提供两个函数:
|
||||
//! - `build_dag_summary`: 从 DagDef 构建可读的 DAG 摘要文本
|
||||
//! - `extract_active_path`: 从完整 DAG 中提取当前活跃路径(简化)
|
||||
//!
|
||||
//! 调用方(前端/workflow 引擎)在启动工作流时将摘要存入
|
||||
//! `PerConvState.workflow_dag_summary`,agentic loop 自动注入 system prompt。
|
||||
|
||||
/// 从工作流 DAG 中提取活跃路径上下文块
|
||||
/// 待 PerConvState 加 workflow_id 字段后接入。
|
||||
#[allow(dead_code)]
|
||||
pub fn build_workflow_block() {
|
||||
// TODO: T4 实现 — 加载 dag::Dag,构建 WorkflowContextBlock,注入 system_prompt
|
||||
use df_workflow::dag_def::DagDef;
|
||||
|
||||
/// 从 DagDef 构建可读的 DAG 摘要文本
|
||||
///
|
||||
/// 输出格式(5 节点 DAG 为例):
|
||||
/// ```
|
||||
/// 步骤: 读源码 → 分析依赖 → 改配置 → 验证 → 提交
|
||||
/// 当前: 改配置
|
||||
/// 完成: 2/5
|
||||
/// ```
|
||||
pub fn build_dag_summary(dag: &DagDef, current_node_id: Option<&str>) -> String {
|
||||
let total = dag.nodes.len();
|
||||
if total == 0 {
|
||||
return String::new();
|
||||
}
|
||||
|
||||
// 拓扑排序(简化:按 edges 计算入度,分层输出)
|
||||
let layers = topological_layers(dag);
|
||||
|
||||
// 构建步骤链
|
||||
let step_chain: Vec<&str> = layers.iter()
|
||||
.flat_map(|layer| layer.iter().map(|id| id.as_str()))
|
||||
.collect();
|
||||
let chain_text = step_chain.join(" → ");
|
||||
|
||||
// 当前节点
|
||||
let current_label = current_node_id
|
||||
.and_then(|id| dag.nodes.get(id))
|
||||
.map(|n| n.label.as_deref().unwrap_or(&n.node_type))
|
||||
.unwrap_or("");
|
||||
|
||||
let mut result = format!("步骤: {}\n", chain_text);
|
||||
if !current_label.is_empty() {
|
||||
result.push_str(&format!("当前: {}\n", current_label));
|
||||
}
|
||||
result.push_str(&format!("完成: 待执行器上报"));
|
||||
|
||||
result
|
||||
}
|
||||
|
||||
/// 从 DagDef 构建活跃路径(从根节点到当前节点的路径 + 后续 2 层)
|
||||
///
|
||||
/// 返回简化的步骤链文本,适合注入 system prompt。
|
||||
pub fn extract_active_path(dag: &DagDef, current_node_id: Option<&str>) -> Vec<String> {
|
||||
let layers = topological_layers(dag);
|
||||
|
||||
if let Some(current) = current_node_id {
|
||||
// 找到当前节点所在的层
|
||||
let current_layer_idx = layers.iter().position(|layer| layer.iter().any(|x| x == current));
|
||||
|
||||
match current_layer_idx {
|
||||
Some(idx) => {
|
||||
// 从根到当前层的路径
|
||||
let path: Vec<String> = layers[..=idx].iter()
|
||||
.flat_map(|layer| layer.iter().map(|id| {
|
||||
dag.nodes.get(id)
|
||||
.map(|n| n.label.clone().unwrap_or_else(|| id.clone()))
|
||||
.unwrap_or_else(|| id.clone())
|
||||
}))
|
||||
.collect();
|
||||
|
||||
// 后续 2 层
|
||||
let next: Vec<String> = layers[idx + 1..]
|
||||
.iter()
|
||||
.take(2)
|
||||
.flat_map(|layer| layer.iter().map(|id| {
|
||||
dag.nodes.get(id)
|
||||
.map(|n| n.label.clone().unwrap_or_else(|| id.clone()))
|
||||
.unwrap_or_else(|| id.clone())
|
||||
}))
|
||||
.collect();
|
||||
|
||||
let mut result = path;
|
||||
result.extend(next);
|
||||
result
|
||||
}
|
||||
None => {
|
||||
// 当前节点不在 DAG 中 → 返回所有节点标签
|
||||
layers.iter().flat_map(|layer| layer.iter().map(|id| {
|
||||
dag.nodes.get(id)
|
||||
.map(|n| n.label.clone().unwrap_or_else(|| id.clone()))
|
||||
.unwrap_or_else(|| id.clone())
|
||||
})).collect()
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// 无当前节点 → 返回所有节点标签
|
||||
layers.iter().flat_map(|layer| layer.iter().map(|id| {
|
||||
dag.nodes.get(id)
|
||||
.map(|n| n.label.clone().unwrap_or_else(|| id.clone()))
|
||||
.unwrap_or_else(|| id.clone())
|
||||
})).collect()
|
||||
}
|
||||
}
|
||||
|
||||
/// 拓扑分层排序:按入度计算每层的节点
|
||||
fn topological_layers(dag: &DagDef) -> Vec<Vec<String>> {
|
||||
let mut in_degree: std::collections::HashMap<&str, usize> = dag.nodes.keys().map(|k| (k.as_str(), 0)).collect();
|
||||
|
||||
for edge in &dag.edges {
|
||||
if let Some(deg) = in_degree.get_mut(edge.target.as_str()) {
|
||||
*deg += 1;
|
||||
}
|
||||
}
|
||||
|
||||
let mut layers = Vec::new();
|
||||
let mut remaining: std::collections::HashSet<&str> = dag.nodes.keys().map(|k| k.as_str()).collect();
|
||||
|
||||
while !remaining.is_empty() {
|
||||
let current_layer: Vec<String> = remaining.iter()
|
||||
.filter(|id| *in_degree.get(*id).unwrap_or(&0) == 0)
|
||||
.map(|id| id.to_string())
|
||||
.collect();
|
||||
|
||||
if current_layer.is_empty() {
|
||||
// 有环 → 剩余节点全放同一层
|
||||
layers.push(remaining.iter().map(|id| id.to_string()).collect());
|
||||
break;
|
||||
}
|
||||
|
||||
for id in ¤t_layer {
|
||||
remaining.remove(id.as_str());
|
||||
// 减少以该节点为起点的边的入度
|
||||
for edge in &dag.edges {
|
||||
if edge.source == *id {
|
||||
if let Some(deg) = in_degree.get_mut(edge.target.as_str()) {
|
||||
*deg = deg.saturating_sub(1);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
layers.push(current_layer);
|
||||
}
|
||||
|
||||
layers
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use df_workflow::dag_def::{DagDef, NodeDef, EdgeDef};
|
||||
|
||||
fn make_dag() -> DagDef {
|
||||
let mut dag = DagDef::new();
|
||||
dag.add_node("n1".to_string(), "read_file".to_string(), serde_json::json!({}));
|
||||
dag.add_node("n2".to_string(), "analyze".to_string(), serde_json::json!({}));
|
||||
dag.add_node("n3".to_string(), "modify".to_string(), serde_json::json!({}));
|
||||
dag.add_edge("n1".to_string(), "n2".to_string());
|
||||
dag.add_edge("n2".to_string(), "n3".to_string());
|
||||
dag
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_topological_layers_linear() {
|
||||
let dag = make_dag();
|
||||
let layers = topological_layers(&dag);
|
||||
assert_eq!(layers.len(), 3);
|
||||
assert_eq!(layers[0], vec!["n1"]);
|
||||
assert_eq!(layers[1], vec!["n2"]);
|
||||
assert_eq!(layers[2], vec!["n3"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_extract_active_path_midway() {
|
||||
let dag = make_dag();
|
||||
let path = extract_active_path(&dag, Some("n2"));
|
||||
assert_eq!(path.len(), 3); // n1, n2, n3
|
||||
assert_eq!(path[1], "n2");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_build_dag_summary_empty() {
|
||||
let dag = DagDef::new();
|
||||
assert_eq!(build_dag_summary(&dag, None), "");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user