refactor(workflow): simplify edge configuration

This commit is contained in:
Xin Wang
2026-08-03 15:17:01 +08:00
parent dd40bfb90a
commit 8064bb7152
10 changed files with 53 additions and 205 deletions

View File

@@ -785,22 +785,9 @@ class WorkflowBrain(BaseBrain):
triggering_user_message: dict[str, Any] | None = None,
) -> NodeConfig:
await self._begin_edge_transition(edge)
context_messages = list(leading_messages or [])
speech = self._engine.edge_transition_speech(edge)
if speech:
content = self._store.render(speech).strip()
if content:
await self._queue_visible_speech(
content,
source="workflow-edge-transition",
node_id=str(edge.get("target") or "") or None,
)
context_message = fixed_speech_context_message(content)
if context_message is not None:
context_messages.append(context_message)
return await self._resolve_path(
str(edge.get("target") or ""),
leading_messages=context_messages,
leading_messages=list(leading_messages or []),
triggering_user_text=triggering_user_text,
triggering_user_message=triggering_user_message,
)
@@ -863,19 +850,6 @@ class WorkflowBrain(BaseBrain):
if not edge:
return self._passive_node_config(node_id, context_messages)
await self._begin_edge_transition(edge)
speech = self._engine.edge_transition_speech(edge)
if speech:
content = self._store.render(speech).strip()
if content:
target_id = str(edge.get("target") or "")
await self._queue_visible_speech(
content,
source="workflow-edge-transition",
node_id=target_id or None,
)
context_message = fixed_speech_context_message(content)
if context_message is not None:
context_messages.append(context_message)
node_id = str(edge.get("target") or "")
raise RuntimeError("工作流连续自动跳转超过安全上限")

View File

@@ -138,17 +138,17 @@ def node_types_response() -> dict[str, Any]:
def _edge_data_v3(edge: dict) -> dict:
data = deepcopy(edge.get("data") or {})
data.pop("transitionSpeech", None)
data.pop("transition_speech", None)
if data.get("mode") in EDGE_MODES:
data.setdefault("priority", 10)
return data
condition = str(data.pop("condition", "") or "").strip()
transition = data.pop("transition_speech", data.get("transitionSpeech", ""))
data.update(
{
"mode": "llm" if condition else "always",
"priority": 10,
"condition": condition,
"transitionSpeech": transition,
}
)
return data
@@ -234,6 +234,8 @@ def normalize_graph(graph: dict[str, Any] | None) -> dict[str, Any]:
_normalize_settings(settings)
source.setdefault("nodes", [])
source.setdefault("edges", [])
for edge in source["edges"]:
edge["data"] = _edge_data_v3(edge)
for node in source["nodes"]:
data = node.setdefault("data", {})
if node.get("type") == "start":
@@ -327,7 +329,7 @@ def normalize_graph(graph: dict[str, Any] | None) -> dict[str, Any]:
"id": f"e-{start_id}-{synthetic_id}",
"source": start_id,
"target": synthetic_id,
"data": {"mode": "always", "priority": 0, "transitionSpeech": ""},
"data": {"mode": "always", "priority": 0},
}
)

View File

@@ -113,14 +113,6 @@ class WorkflowEngine:
return f"当满足以下条件时转到「{target}」:{condition}"
return f"当当前阶段任务完成时转到「{target}」。"
def edge_transition_speech(self, edge: dict | None) -> str:
if not edge:
return ""
data = edge.get("data") or {}
return str(
data.get("transitionSpeech") or data.get("transition_speech") or ""
)
def global_prompt(self) -> str:
return str(self.settings.get("globalPrompt") or "").strip()

View File

@@ -2393,7 +2393,6 @@ class WorkflowBrainTests(unittest.IsolatedAsyncioTestCase):
"mode": "llm",
"priority": 10,
"condition": "需求已收集",
"transitionSpeech": "正在为你结束流程",
},
}
],
@@ -2569,29 +2568,6 @@ class WorkflowBrainTests(unittest.IsolatedAsyncioTestCase):
self.assertTrue(call_end.ending)
self.assertTrue(call_end.armed)
self.assertTrue(any(getattr(frame, "text", "") == "感谢来电" for frame in queued))
transition_context_frames = [
frame
for frame in worker.frames
if isinstance(frame, LLMMessagesAppendFrame)
and frame.messages
== [
{
"role": "system",
"content": (
f"{FIXED_SPEECH_CONTEXT_MARKER}\n正在为你结束流程"
),
}
]
]
self.assertTrue(transition_context_frames)
transition_events = [
frame.message
for frame in queued
if isinstance(frame, OutputTransportMessageUrgentFrame)
and frame.message.get("source") == "workflow-edge-transition"
]
self.assertEqual(transition_events[0]["content"], "正在为你结束流程")
self.assertEqual(transition_events[0]["nodeId"], "end")
assistant_transcripts = [
frame.message.get("content")
for frame in queued
@@ -2601,11 +2577,7 @@ class WorkflowBrainTests(unittest.IsolatedAsyncioTestCase):
]
self.assertEqual(
assistant_transcripts,
["正在为你结束流程", "感谢来电"],
)
self.assertIn(
"正在为你结束流程",
brain._store.values["system__conversation_history"],
["感谢来电"],
)
self.assertIn(
"感谢来电",

View File

@@ -64,7 +64,7 @@ class _Brain:
async def on_client_ready(self):
for content, timestamp in (
("Start Edge 过渡语", "2026-07-14T10:00:00.200+00:00"),
("Message 节点播报", "2026-07-14T10:00:00.200+00:00"),
("Agent 固定进入语", "2026-07-14T10:00:00.300+00:00"),
):
await self.worker.queue_frame(
@@ -118,7 +118,7 @@ class PipelineEventTest(unittest.IsolatedAsyncioTestCase):
ordered = sorted(transcripts, key=lambda message: message["timestamp"])
self.assertEqual(
[message["content"] for message in ordered],
["助手开场白", "Start Edge 过渡语", "Agent 固定进入语"],
["助手开场白", "Message 节点播报", "Agent 固定进入语"],
)
self.assertEqual(transcripts[0]["timestamp"], greeting_time)
self.assertEqual(brain.prepared_greeting, "助手开场白")

View File

@@ -132,6 +132,17 @@ class WorkflowGraphTests(unittest.TestCase):
self.assertEqual(cleaned_agent["data"]["entryMode"], "wait_user")
self.assertNotIn("entrySpeech", cleaned_agent["data"])
def test_edge_speech_fields_are_removed(self):
graph = valid_graph()
graph["edges"][0]["data"]["transitionSpeech"] = "不再播放"
graph["edges"][1]["data"]["transition_speech"] = "旧字段也不再播放"
normalized = normalize_graph(graph)
for edge in normalized["edges"]:
self.assertNotIn("transitionSpeech", edge["data"])
self.assertNotIn("transition_speech", edge["data"])
def test_action_defaults_preserve_legacy_result_assignment_behavior(self):
graph = valid_graph()
graph["nodes"].extend(