mirror of
https://github.com/langgenius/dify.git
synced 2026-09-24 23:22:26 +08:00
feat: harden /create and /refine workflow generation for edge cases (#37336)
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
f9911ab3ef
commit
09bb87d089
@@ -526,13 +526,77 @@ Now emit the complete workflow graph JSON.
|
||||
"""
|
||||
|
||||
|
||||
# Node wrapper fields that carry no meaning the builder needs: pure canvas /
|
||||
# selection state, plus geometry the runner's postprocess recomputes anyway.
|
||||
# Stripping them out of the refine prompt cuts its size roughly in half on
|
||||
# hand-edited graphs — fewer tokens in, and (because the builder echoes
|
||||
# untouched nodes verbatim) far fewer tokens out, which is where the latency
|
||||
# lives.
|
||||
_PRUNED_NODE_KEYS = frozenset(
|
||||
{
|
||||
"positionAbsolute",
|
||||
"sourcePosition",
|
||||
"targetPosition",
|
||||
"selected",
|
||||
"dragging",
|
||||
"measured",
|
||||
}
|
||||
)
|
||||
|
||||
# Additionally pruned from TOP-LEVEL nodes only: the layered auto-layout
|
||||
# recomputes their position and size defaults, so the builder never needs to
|
||||
# reproduce them. Container children keep ``position`` (relative to the
|
||||
# parent, which we cannot recompute) and containers keep ``width`` /
|
||||
# ``height`` (their canvas size is real config, not a default).
|
||||
_PRUNED_TOP_LEVEL_NODE_KEYS = _PRUNED_NODE_KEYS | {"position", "width", "height"}
|
||||
|
||||
_CONTAINER_DATA_TYPES = frozenset({"iteration", "loop"})
|
||||
|
||||
# Edge fields the builder must echo; everything else (ids, zIndex,
|
||||
# sourceType / targetType, isInIteration / isInLoop markers) is recomputed
|
||||
# by the runner's postprocess from the node topology.
|
||||
_KEPT_EDGE_KEYS = ("source", "target", "sourceHandle", "targetHandle")
|
||||
|
||||
|
||||
def compact_graph_for_builder(current_graph: dict) -> dict:
|
||||
"""
|
||||
Strip canvas noise out of a draft graph before prompt injection.
|
||||
|
||||
Keeps everything semantically meaningful — ids, wrapper ``type``,
|
||||
``parentId``, the full ``data`` config, child positions, container
|
||||
sizes — and drops geometry / selection state the postprocess pass
|
||||
recomputes. The builder echoes untouched nodes verbatim, so every byte
|
||||
removed here is removed twice (prompt AND completion).
|
||||
"""
|
||||
nodes_out: list[dict] = []
|
||||
for node in current_graph.get("nodes") or []:
|
||||
if not isinstance(node, dict):
|
||||
continue
|
||||
is_child = bool(node.get("parentId"))
|
||||
is_container = isinstance(node.get("data"), dict) and node["data"].get("type") in _CONTAINER_DATA_TYPES
|
||||
pruned = _PRUNED_NODE_KEYS if (is_child or is_container) else _PRUNED_TOP_LEVEL_NODE_KEYS
|
||||
compact = {k: v for k, v in node.items() if k not in pruned}
|
||||
if is_container:
|
||||
# Container position is still recomputed by the layout pass.
|
||||
compact.pop("position", None)
|
||||
nodes_out.append(compact)
|
||||
edges_out = [
|
||||
{k: edge[k] for k in _KEPT_EDGE_KEYS if k in edge}
|
||||
for edge in (current_graph.get("edges") or [])
|
||||
if isinstance(edge, dict)
|
||||
]
|
||||
return {"nodes": nodes_out, "edges": edges_out}
|
||||
|
||||
|
||||
def format_builder_existing_graph_section(current_graph: dict | None) -> str:
|
||||
"""
|
||||
Refine mode: give the builder the FULL existing graph JSON so it can keep
|
||||
Refine mode: give the builder the existing graph JSON so it can keep
|
||||
every node and edge the user's change does not touch byte-for-byte — same
|
||||
ids, same config, same prompt templates. Without the full config the
|
||||
builder would regenerate untouched nodes from scratch and silently drop
|
||||
the user's hand-tuned settings.
|
||||
the user's hand-tuned settings. Canvas-only fields are stripped first
|
||||
(see ``compact_graph_for_builder``) — they're recomputed in postprocess,
|
||||
so carrying them only slows the call down.
|
||||
|
||||
Returns an empty string in create mode (no ``current_graph``); the builder
|
||||
then behaves exactly as before, constructing the graph purely from the
|
||||
@@ -540,7 +604,7 @@ def format_builder_existing_graph_section(current_graph: dict | None) -> str:
|
||||
"""
|
||||
if not current_graph:
|
||||
return ""
|
||||
graph_json = json.dumps(current_graph, ensure_ascii=False, separators=(",", ":"))
|
||||
graph_json = json.dumps(compact_graph_for_builder(current_graph), ensure_ascii=False, separators=(",", ":"))
|
||||
return (
|
||||
"# Existing graph to refine (JSON)\n\n"
|
||||
"You are REFINING this existing graph, NOT building from scratch. Apply "
|
||||
|
||||
@@ -70,6 +70,10 @@ logger = logging.getLogger(__name__)
|
||||
_NODE_X_OFFSET = 80
|
||||
_NODE_X_STEP = 320
|
||||
_NODE_Y = 280
|
||||
# Vertical gap between lanes when two branches share the same topological
|
||||
# depth (e.g. the two arms of an if-else). Default node height is 100, so
|
||||
# 160 leaves clear air between stacked nodes.
|
||||
_NODE_Y_STEP = 160
|
||||
_DEFAULT_VIEWPORT: GraphViewportDict = {"x": 0.0, "y": 0.0, "zoom": 0.7}
|
||||
_DEFAULT_NODE_WIDTH = 244
|
||||
_DEFAULT_NODE_HEIGHT = 100
|
||||
@@ -575,34 +579,32 @@ class WorkflowGenerator:
|
||||
|
||||
# Defensive ID remap: Dify's run-time placeholder regex only accepts
|
||||
# ``[a-zA-Z0-9_]`` in the node-id slot, so anything the LLM emits with
|
||||
# hyphens (``node-1``, ``node-Kstart``, etc.) would break every
|
||||
# placeholder pointing at it. Strip hyphens out of every id + every
|
||||
# hyphens, dots, or spaces (``node-1``, ``node.2``, etc.) would break
|
||||
# every placeholder pointing at it. Sanitize every id + every
|
||||
# cross-reference (edges' ``source`` / ``target``, ``parentId``,
|
||||
# ``start_node_id`` / ``iteration_id`` / ``loop_id`` on data, and the
|
||||
# ``{{#…#}}`` and ``["node-id", "var"]`` references) BEFORE the rest
|
||||
# of the postprocess pass touches them.
|
||||
cls._strip_hyphens_from_node_ids(nodes=nodes, edges=edges)
|
||||
cls._sanitize_node_ids(nodes=nodes, edges=edges)
|
||||
|
||||
# Container-child nodes carry their own relative positions inside the
|
||||
# parent and have a special ``type`` (custom-iteration-start /
|
||||
# custom-loop-start). We must not override their positions or wrapper
|
||||
# ``type``; only top-level (parentId-less) nodes get the left-to-right
|
||||
# auto layout.
|
||||
top_level_index = 0
|
||||
# ``type``; only top-level (parentId-less) nodes get the layered
|
||||
# auto layout (x = topological depth, y = lane within the layer).
|
||||
cls._layout_top_level_nodes(nodes=nodes, edges=edges)
|
||||
for node in nodes:
|
||||
cls._fill_node_defaults(node)
|
||||
if node.get("parentId"):
|
||||
# Inner node — keep whatever the LLM emitted; only fill the
|
||||
# absolutely-required defaults so the canvas can render it.
|
||||
node.setdefault("position", {"x": 0.0, "y": 0.0})
|
||||
node.setdefault("zIndex", 1002)
|
||||
node.setdefault("extent", "parent")
|
||||
else:
|
||||
node["position"] = {
|
||||
"x": float(_NODE_X_OFFSET + _NODE_X_STEP * top_level_index),
|
||||
"y": float(_NODE_Y),
|
||||
}
|
||||
top_level_index += 1
|
||||
# Inner nodes keep their LLM-emitted relative position; top-level
|
||||
# nodes were positioned by the layered layout. The setdefault only
|
||||
# fires for inner nodes without a position and for a (broken)
|
||||
# id-less node the layout pass couldn't see.
|
||||
node.setdefault("position", {"x": 0.0, "y": 0.0})
|
||||
node.setdefault("positionAbsolute", dict(node["position"]))
|
||||
node.setdefault("width", _DEFAULT_NODE_WIDTH)
|
||||
node.setdefault("height", _DEFAULT_NODE_HEIGHT)
|
||||
@@ -620,6 +622,12 @@ class WorkflowGenerator:
|
||||
if n.get("id") in inner_node_to_parent.values():
|
||||
parent_type[n["id"]] = n.get("data", {}).get("type", "")
|
||||
|
||||
# Branch nodes (if-else / question-classifier) emit one handle per
|
||||
# case; an edge leaving them on the default "source" handle dangles
|
||||
# off a handle that doesn't exist on the canvas. Repair the
|
||||
# unambiguous cases before edge ids are computed from the handles.
|
||||
cls._repair_branch_edge_handles(nodes=nodes, edges=edges)
|
||||
|
||||
# Dedupe edges (LLMs sometimes emit the same edge twice).
|
||||
seen: set[tuple[str, str, str, str]] = set()
|
||||
deduped_edges = []
|
||||
@@ -712,14 +720,19 @@ class WorkflowGenerator:
|
||||
r"\{\{#([a-zA-Z0-9_]{1,50})\.([a-zA-Z_][a-zA-Z0-9_]{0,29}(?:\.[a-zA-Z_][a-zA-Z0-9_]{0,29}){0,9})#\}\}"
|
||||
)
|
||||
|
||||
# Lenient sibling used only by the defensive hyphen-strip pass — it
|
||||
# allows hyphens in the node-id slot so we can rewrite the LLM's
|
||||
# ``{{#node-1.var#}}`` outputs BEFORE the strict walker sees them.
|
||||
# Lenient sibling used only by the defensive id-sanitize pass — it
|
||||
# accepts ANY character in the node-id slot (except the ``.`` separator
|
||||
# and ``#`` terminator) so we can rewrite the LLM's ``{{#node-1.var#}}``
|
||||
# / ``{{#node 2.var#}}`` outputs BEFORE the strict walker sees them.
|
||||
# Never use this for validation, only for rewriting.
|
||||
_LENIENT_VAR_REF_RE: ClassVar = re.compile(r"\{\{#([A-Za-z0-9_-]+)\.([^#]+)#\}\}")
|
||||
_LENIENT_VAR_REF_RE: ClassVar = re.compile(r"\{\{#([^#.{}]+)\.([^#]+)#\}\}")
|
||||
|
||||
# Characters the run-time placeholder regex rejects in the node-id slot.
|
||||
# Anything matching this in a node id must be sanitized away.
|
||||
_INVALID_ID_CHARS_RE: ClassVar = re.compile(r"[^a-zA-Z0-9_]")
|
||||
|
||||
# Strings inside ``data`` that look like node-id slugs and need
|
||||
# remapping when we defensively strip hyphens out of LLM-emitted ids.
|
||||
# remapping when we defensively sanitize LLM-emitted ids.
|
||||
_ID_FIELDS: ClassVar = frozenset({"start_node_id", "iteration_id", "loop_id", "parentId"})
|
||||
|
||||
# ``data`` keys whose value is a plain string list, never a
|
||||
@@ -860,34 +873,51 @@ class WorkflowGenerator:
|
||||
return False
|
||||
|
||||
@classmethod
|
||||
def _strip_hyphens_from_node_ids(cls, *, nodes: list[dict[str, Any]], edges: list[dict[str, Any]]) -> None:
|
||||
def _sanitize_node_ids(cls, *, nodes: list[dict[str, Any]], edges: list[dict[str, Any]]) -> None:
|
||||
"""
|
||||
Strip ``-`` out of every node id and rewrite every cross-reference.
|
||||
Rewrite every node id to ``[a-zA-Z0-9_]`` and fix every cross-reference.
|
||||
|
||||
Dify's run-time ``VARIABLE_PATTERN`` accepts only ``[a-zA-Z0-9_]`` in
|
||||
the node-id slot of ``{{#…#}}`` placeholders. The builder LLM often
|
||||
emits ``node-1`` style ids; left unfixed those make every placeholder
|
||||
silently fail at run time, the literal ``{{#node-1.var#}}`` survives
|
||||
into the prompt, and the LLM at run time echoes it back as the user's
|
||||
output — the bug we are here to kill.
|
||||
emits ``node-1`` style ids (and occasionally dots or spaces); left
|
||||
unfixed those make every placeholder silently fail at run time, the
|
||||
literal ``{{#node-1.var#}}`` survives into the prompt, and the LLM at
|
||||
run time echoes it back as the user's output — the bug we are here
|
||||
to kill.
|
||||
|
||||
Approach: build a one-to-one ``old → new`` map by removing hyphens,
|
||||
then rewrite (a) every node ``id``, (b) every edge ``source`` /
|
||||
``target``, (c) every ``parentId`` / ``start_node_id`` /
|
||||
``iteration_id`` / ``loop_id`` inside ``data``, (d) every
|
||||
``{{#…#}}`` reference in any string, (e) every ``["node-id", "var"]``
|
||||
value-selector list. We do NOT rename variable names — only ids.
|
||||
Approach: build a one-to-one ``old → new`` map by dropping the invalid
|
||||
characters — collision-safe: when the sanitized id is already taken
|
||||
(e.g. the builder emitted BOTH ``node-1`` and ``node1``) a numeric
|
||||
suffix keeps the two distinct instead of silently merging every
|
||||
reference onto one node. Then rewrite (a) every node ``id``, (b) every
|
||||
edge ``source`` / ``target``, (c) every ``parentId`` /
|
||||
``start_node_id`` / ``iteration_id`` / ``loop_id`` inside ``data``,
|
||||
(d) every ``{{#…#}}`` reference in any string, (e) every
|
||||
``["node-id", "var"]`` value-selector list. We do NOT rename variable
|
||||
names — only ids.
|
||||
"""
|
||||
# Build id rewrite map. Collision-safe because we just strip a single
|
||||
# character class — two different hyphenated ids ``node-1`` and
|
||||
# ``node1`` would collide, but the builder LLM has been instructed
|
||||
# to pick one style so in practice it's one or the other.
|
||||
id_map: dict[str, str] = {}
|
||||
# Ids that are already valid are reserved up front so a sanitized id
|
||||
# can never collide with an untouched sibling.
|
||||
used: set[str] = {
|
||||
n["id"] for n in nodes if isinstance(n.get("id"), str) and not cls._INVALID_ID_CHARS_RE.search(n["id"])
|
||||
}
|
||||
fallback_seq = 0
|
||||
for node in nodes:
|
||||
old = node.get("id")
|
||||
if not isinstance(old, str) or "-" not in old:
|
||||
if not isinstance(old, str) or not cls._INVALID_ID_CHARS_RE.search(old):
|
||||
continue
|
||||
new = old.replace("-", "")
|
||||
base = cls._INVALID_ID_CHARS_RE.sub("", old)
|
||||
if not base:
|
||||
# Id was nothing but invalid characters (e.g. "节点", "--").
|
||||
fallback_seq += 1
|
||||
base = f"node_{fallback_seq}"
|
||||
new = base
|
||||
suffix = 2
|
||||
while new in used:
|
||||
new = f"{base}_{suffix}"
|
||||
suffix += 1
|
||||
used.add(new)
|
||||
id_map[old] = new
|
||||
node["id"] = new
|
||||
if not id_map:
|
||||
@@ -901,10 +931,12 @@ class WorkflowGenerator:
|
||||
edge[key] = id_map[v]
|
||||
# Also rewrite the edge id if the builder emitted one referencing
|
||||
# the old ids; the dedupe pass later recomputes it anyway, but
|
||||
# rewriting here keeps logs sane.
|
||||
# rewriting here keeps logs sane. Longest-first so an id that is
|
||||
# a substring of another (``node-1`` in ``node-12``) can't corrupt
|
||||
# the longer match.
|
||||
eid = edge.get("id")
|
||||
if isinstance(eid, str):
|
||||
for old, new in id_map.items():
|
||||
for old, new in sorted(id_map.items(), key=lambda kv: -len(kv[0])):
|
||||
eid = eid.replace(old, new)
|
||||
edge["id"] = eid
|
||||
|
||||
@@ -955,6 +987,116 @@ class WorkflowGenerator:
|
||||
new_id = id_map.get(node_id, node_id)
|
||||
return f"{{{{#{new_id}.{rest}#}}}}"
|
||||
|
||||
@classmethod
|
||||
def _repair_branch_edge_handles(cls, *, nodes: list[dict[str, Any]], edges: list[dict[str, Any]]) -> None:
|
||||
"""
|
||||
Re-home edges that leave a branch node on the default "source" handle.
|
||||
|
||||
if-else exposes one source handle per ``case_id`` plus the implicit
|
||||
"false" (ELSE) handle; question-classifier exposes one per class id.
|
||||
The builder prompt documents this, but LLMs still emit the default
|
||||
handle, which renders as an edge hanging off a handle that doesn't
|
||||
exist and the branch silently never runs.
|
||||
|
||||
Repair only when unambiguous: default-handle edges are assigned to the
|
||||
node's UNUSED branch handles in declaration order, and only when there
|
||||
are at least as many unused handles as edges to fix. Anything
|
||||
ambiguous is left alone — a wrong guess that swaps the IF and ELSE
|
||||
arms is worse than a visible dangling edge.
|
||||
"""
|
||||
for node in nodes:
|
||||
data = node.get("data") or {}
|
||||
node_type = data.get("type")
|
||||
if node_type == BuiltinNodeTypes.IF_ELSE:
|
||||
branch_handles = [
|
||||
str(case["case_id"])
|
||||
for case in (data.get("cases") or [])
|
||||
if isinstance(case, dict) and case.get("case_id")
|
||||
]
|
||||
# ELSE is implicit — it has a handle even though no case
|
||||
# declares it.
|
||||
branch_handles.append("false")
|
||||
elif node_type == BuiltinNodeTypes.QUESTION_CLASSIFIER:
|
||||
branch_handles = [
|
||||
str(klass["id"])
|
||||
for klass in (data.get("classes") or [])
|
||||
if isinstance(klass, dict) and klass.get("id")
|
||||
]
|
||||
else:
|
||||
continue
|
||||
|
||||
node_id = node.get("id")
|
||||
outgoing = [e for e in edges if e.get("source") == node_id]
|
||||
taken = {e.get("sourceHandle") for e in outgoing if e.get("sourceHandle") in branch_handles}
|
||||
unused = [h for h in branch_handles if h not in taken]
|
||||
defaulted = [e for e in outgoing if e.get("sourceHandle") in (None, "", "source")]
|
||||
if not defaulted or len(defaulted) > len(unused):
|
||||
continue
|
||||
for edge, handle in zip(defaulted, unused):
|
||||
edge["sourceHandle"] = handle
|
||||
logger.info(
|
||||
"Workflow generator: re-homed default-handle edge %s -> %s onto branch handle %r",
|
||||
node_id,
|
||||
edge.get("target"),
|
||||
handle,
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def _layout_top_level_nodes(cls, *, nodes: list[dict[str, Any]], edges: list[dict[str, Any]]) -> None:
|
||||
"""
|
||||
Lay out top-level nodes by graph topology instead of array order.
|
||||
|
||||
x = longest-path depth from the entry layer, y = lane within the
|
||||
layer — so an if-else's two arms render as two parallel rows instead
|
||||
of overlapping on one line, and a builder that emits nodes out of
|
||||
execution order still gets a left-to-right canvas. Longest-path (not
|
||||
BFS) layering keeps a join node (variable-aggregator, end) to the
|
||||
right of its deepest branch.
|
||||
|
||||
Cycle-safe: Kahn's algorithm simply never reaches nodes on a cycle,
|
||||
and those get parked one layer past the deepest laid-out node in
|
||||
declaration order — the cycle itself is flagged by the structural
|
||||
validator afterwards.
|
||||
"""
|
||||
top_level = [n for n in nodes if not n.get("parentId") and isinstance(n.get("id"), str) and n.get("id")]
|
||||
id_set = {n["id"] for n in top_level}
|
||||
|
||||
succs: dict[str, list[str]] = {node_id: [] for node_id in id_set}
|
||||
indegree: dict[str, int] = dict.fromkeys(id_set, 0)
|
||||
seen_pairs: set[tuple[str, str]] = set()
|
||||
for edge in edges:
|
||||
src, tgt = edge.get("source"), edge.get("target")
|
||||
if not isinstance(src, str) or not isinstance(tgt, str):
|
||||
continue
|
||||
if src not in id_set or tgt not in id_set or src == tgt or (src, tgt) in seen_pairs:
|
||||
continue
|
||||
seen_pairs.add((src, tgt))
|
||||
succs[src].append(tgt)
|
||||
indegree[tgt] += 1
|
||||
|
||||
depth: dict[str, int] = {}
|
||||
queue = [n["id"] for n in top_level if indegree[n["id"]] == 0]
|
||||
for node_id in queue:
|
||||
depth[node_id] = 0
|
||||
while queue:
|
||||
cur = queue.pop(0)
|
||||
for nxt in succs[cur]:
|
||||
depth[nxt] = max(depth.get(nxt, 0), depth[cur] + 1)
|
||||
indegree[nxt] -= 1
|
||||
if indegree[nxt] == 0:
|
||||
queue.append(nxt)
|
||||
|
||||
overflow_depth = (max(depth.values()) + 1) if depth else 0
|
||||
lanes: dict[int, int] = {}
|
||||
for node in top_level:
|
||||
d = depth.get(node["id"], overflow_depth)
|
||||
lane = lanes.get(d, 0)
|
||||
lanes[d] = lane + 1
|
||||
node["position"] = {
|
||||
"x": float(_NODE_X_OFFSET + _NODE_X_STEP * d),
|
||||
"y": float(_NODE_Y + _NODE_Y_STEP * lane),
|
||||
}
|
||||
|
||||
@classmethod
|
||||
def _inject_start_variable(cls, start_node: dict[str, Any], var: str) -> None:
|
||||
"""Add a default ``paragraph`` input so ``{{#start.<var>#}}`` resolves."""
|
||||
@@ -1122,6 +1264,24 @@ class WorkflowGenerator:
|
||||
errors.append(_err(WorkflowGenerateErrorCode.INVALID_SCHEMA, "Generated graph has no nodes"))
|
||||
return errors
|
||||
|
||||
# Duplicate ids make every cross-reference ambiguous (edges, variable
|
||||
# placeholders, parentId all resolve to "whichever node wins"), so a
|
||||
# graph with them is unusable no matter how the canvas renders it.
|
||||
id_counts: dict[str, int] = {}
|
||||
for node in nodes:
|
||||
node_id = node.get("id", "")
|
||||
if node_id:
|
||||
id_counts[node_id] = id_counts.get(node_id, 0) + 1
|
||||
for node_id, count in id_counts.items():
|
||||
if count > 1:
|
||||
errors.append(
|
||||
_err(
|
||||
WorkflowGenerateErrorCode.DUPLICATE_NODE_ID,
|
||||
f"Duplicate node id {node_id!r} ({count} nodes share it)",
|
||||
node_id=node_id,
|
||||
)
|
||||
)
|
||||
|
||||
types = [node.get("data", {}).get("type", "") for node in nodes]
|
||||
starts = [t for t in types if t == BuiltinNodeTypes.START]
|
||||
if len(starts) != 1:
|
||||
@@ -1156,6 +1316,11 @@ class WorkflowGenerator:
|
||||
if tgt not in known_ids:
|
||||
errors.append(_err(WorkflowGenerateErrorCode.DANGLING_EDGE, f"Edge target node not found: {tgt!r}"))
|
||||
|
||||
# Workflow graphs must be DAGs — a directed cycle hangs or errors the
|
||||
# run, and nothing downstream of the cycle ever executes. (A "loop"
|
||||
# container is the sanctioned way to iterate; its edges are internal.)
|
||||
errors.extend(cls._collect_edge_cycle_errors(graph=graph, known_ids=known_ids))
|
||||
|
||||
# Dangling node-id references in node ``data`` (parentId, start_node_id, iteration_id, loop_id).
|
||||
errors.extend(cls._collect_dangling_id_refs(nodes=nodes, known_ids=known_ids))
|
||||
|
||||
@@ -1175,6 +1340,54 @@ class WorkflowGenerator:
|
||||
|
||||
return errors
|
||||
|
||||
@classmethod
|
||||
def _collect_edge_cycle_errors(cls, *, graph: GraphDict, known_ids: set[str]) -> list[WorkflowGenerateErrorDict]:
|
||||
"""
|
||||
Flag directed cycles among the graph's edges (Kahn's algorithm).
|
||||
|
||||
Self-loops are reported per node; a longer cycle is reported once,
|
||||
naming every node Kahn's peeling never reaches (cycle members plus
|
||||
anything downstream of them). Edges into unknown ids are ignored
|
||||
here — the dangling-edge check already covers those.
|
||||
"""
|
||||
out: list[WorkflowGenerateErrorDict] = []
|
||||
succs: dict[str, list[str]] = {node_id: [] for node_id in known_ids}
|
||||
indegree: dict[str, int] = dict.fromkeys(known_ids, 0)
|
||||
for edge in graph.get("edges", []):
|
||||
src, tgt = edge.get("source"), edge.get("target")
|
||||
if src not in known_ids or tgt not in known_ids:
|
||||
continue
|
||||
if src == tgt:
|
||||
out.append(
|
||||
_err(
|
||||
WorkflowGenerateErrorCode.GRAPH_CYCLE,
|
||||
f"Node {src!r} has an edge pointing at itself",
|
||||
node_id=src,
|
||||
)
|
||||
)
|
||||
continue
|
||||
succs[src].append(tgt)
|
||||
indegree[tgt] += 1
|
||||
|
||||
queue = [node_id for node_id, deg in indegree.items() if deg == 0]
|
||||
visited = 0
|
||||
while queue:
|
||||
cur = queue.pop()
|
||||
visited += 1
|
||||
for nxt in succs[cur]:
|
||||
indegree[nxt] -= 1
|
||||
if indegree[nxt] == 0:
|
||||
queue.append(nxt)
|
||||
if visited < len(known_ids):
|
||||
trapped = sorted(node_id for node_id, deg in indegree.items() if deg > 0)
|
||||
out.append(
|
||||
_err(
|
||||
WorkflowGenerateErrorCode.GRAPH_CYCLE,
|
||||
f"Workflow graph contains a cycle; affected nodes: {', '.join(trapped)}",
|
||||
)
|
||||
)
|
||||
return out
|
||||
|
||||
@classmethod
|
||||
def _collect_dangling_id_refs(
|
||||
cls, *, nodes: list[dict[str, Any]], known_ids: set[str]
|
||||
|
||||
@@ -19,6 +19,9 @@ class WorkflowGenerateErrorCode:
|
||||
INVALID_JSON: Final = "INVALID_JSON"
|
||||
INVALID_SCHEMA: Final = "INVALID_SCHEMA"
|
||||
EMPTY_INSTRUCTION: Final = "EMPTY_INSTRUCTION"
|
||||
INSTRUCTION_TOO_LONG: Final = "INSTRUCTION_TOO_LONG"
|
||||
DUPLICATE_NODE_ID: Final = "DUPLICATE_NODE_ID"
|
||||
GRAPH_CYCLE: Final = "GRAPH_CYCLE"
|
||||
EMPTY_PLAN: Final = "EMPTY_PLAN"
|
||||
UNKNOWN_NODE_REFERENCE: Final = "UNKNOWN_NODE_REFERENCE"
|
||||
INVALID_CONTAINER: Final = "INVALID_CONTAINER"
|
||||
|
||||
Reference in New Issue
Block a user