Skip to content

Commit 7838811

Browse files
Scope workflow execution to the active trigger graph (#60)
* WIP: add browserless FMHY site adapters [gstack-context] * spec: constrain trigger-scoped workflow execution * fix: scope workflow execution to trigger graph --------- Co-authored-by: lunnt <2276214182@qq.com>
1 parent cfe74aa commit 7838811

12 files changed

Lines changed: 1677 additions & 17 deletions

File tree

backend/api/v1/studio_lifecycle.py

Lines changed: 90 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
from backend.schemas import workflow as workflow_schemas
2929
from backend.schemas.common import ApiResponse
3030
from backend.workflow.compiler import compile_workflow_project
31+
from backend.workflow.trigger_scope import scoped_project, select_active_union
3132

3233
router = APIRouter()
3334

@@ -57,6 +58,60 @@ def _isolated_source_errors(
5758
]
5859

5960

61+
def _parked_diagnostics(
62+
project: workflow_schemas.WorkflowProject,
63+
parked_ids: list[str],
64+
) -> list[workflow_schemas.WorkflowCompileError]:
65+
"""Emit node-anchored diagnostics for every parked canvas node.
66+
67+
Membership comes first (one ``parked_node`` row per parked id, in authored
68+
order). Any original configuration diagnostic for the same node follows in
69+
its existing order so the UI can render each failure cause individually.
70+
"""
71+
72+
parked_set = set(parked_ids)
73+
diagnostics: list[workflow_schemas.WorkflowCompileError] = []
74+
for node_id in parked_ids:
75+
diagnostics.append(
76+
workflow_schemas.WorkflowCompileError(
77+
code="parked_node",
78+
message=f'Workflow node "{node_id}" is not connected to a supported trigger.',
79+
node_id=node_id,
80+
path=["nodes", node_id],
81+
)
82+
)
83+
84+
# Compile parked nodes in isolation to surface configuration diagnostics
85+
# (unknown bindings, missing params, etc.) as warnings without edges that
86+
# would produce irrelevant port-mismatch noise.
87+
if not parked_set:
88+
return diagnostics
89+
parked_nodes = [n for n in project.nodes if n.id in parked_set]
90+
parked_project = workflow_schemas.WorkflowProject(
91+
id=project.id,
92+
name=project.name,
93+
profile=project.profile,
94+
version=project.version,
95+
nodes=parked_nodes,
96+
edges=[],
97+
settings=project.settings,
98+
adapters=list(project.adapters),
99+
agentPermissions=project.agentPermissions,
100+
)
101+
parked_result = compile_workflow_project(parked_project)
102+
for error in parked_result.errors:
103+
if error.node_id and error.node_id in parked_set:
104+
diagnostics.append(
105+
workflow_schemas.WorkflowCompileError(
106+
code=error.code,
107+
message=error.message,
108+
node_id=error.node_id,
109+
path=error.path,
110+
)
111+
)
112+
return diagnostics
113+
114+
60115
def _image_generation_nodes(
61116
nodes: object,
62117
*,
@@ -185,6 +240,9 @@ async def validate_draft(
185240
project_id=project_id,
186241
workflow_id=workflow_id,
187242
)
243+
warnings: list[workflow_schemas.WorkflowCompileError] = []
244+
valid = False
245+
stored_graph: dict[str, Any] | None = None
188246
try:
189247
project = workflow_schemas.WorkflowProject.model_validate(resolved_graph)
190248
except ValidationError as exc:
@@ -196,25 +254,47 @@ async def validate_draft(
196254
)
197255
for error in exc.errors()
198256
)
199-
valid = False
200257
else:
201-
errors.extend(_isolated_source_errors(project))
202-
if errors:
203-
valid = False
258+
active_union = select_active_union(project)
259+
if not active_union.has_supported_trigger:
260+
# Legacy / media-canvas / non-trigger workflows preserve the
261+
# existing full-graph validation path unchanged.
262+
errors.extend(_isolated_source_errors(project))
263+
if not errors:
264+
result = compile_workflow_project(project)
265+
errors = list(result.errors)
266+
if result.valid and result.plan is not None:
267+
valid = True
268+
stored_graph = resolved_graph
204269
else:
205-
result = compile_workflow_project(project)
206-
errors = result.errors
207-
valid = result.valid
208-
270+
scoped = scoped_project(
271+
project=project,
272+
active_ids=active_union.active_node_ids,
273+
external_ids={
274+
node.id
275+
for node in project.nodes
276+
if isinstance(node.params.get("externalWorkflow"), dict)
277+
},
278+
)
279+
errors.extend(_isolated_source_errors(scoped))
280+
if not errors:
281+
scoped_result = compile_workflow_project(scoped)
282+
errors = list(scoped_result.errors)
283+
if scoped_result.valid and scoped_result.plan is not None:
284+
valid = True
285+
stored_graph = scoped.model_dump(mode="json")
286+
warnings.extend(
287+
_parked_diagnostics(project, active_union.parked_node_ids)
288+
)
209289
row = StudioWorkflowValidationRun(
210290
workflow_id=workflow_id,
211291
draft_revision=draft.revision,
212292
status="completed" if valid else "failed",
213293
valid=valid,
214294
errors=[error.model_dump(mode="json") for error in errors],
215-
warnings=[],
295+
warnings=[warning.model_dump(mode="json") for warning in warnings],
216296
compile_version=workflow_schemas.WORKFLOW_COMPILE_VERSION,
217-
resolved_graph=resolved_graph if valid else None,
297+
resolved_graph=stored_graph,
218298
)
219299
db.add(row)
220300
await db.flush()

backend/workflow/opencli_hda_tracer.py

Lines changed: 56 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -305,14 +305,63 @@ async def start_workflow_run(
305305
trace_id = body.traceId or str(uuid.uuid4())
306306
started_at = _utcnow()
307307
prior_events = list(existing_events or [])
308+
# Source-level trigger scope selection runs before authoritative compilation
309+
# so a disconnected, incomplete canvas node cannot block a valid
310+
# trigger-reachable component. The compiled-runtime selector remains as a
311+
# defensive assertion against post-compile drift (e.g. template expansion
312+
# producing a second matching trigger entry).
313+
from backend.workflow.trigger_scope import has_supported_triggers, select_trigger_scope
314+
315+
has_triggers = has_supported_triggers(body.project)
316+
if has_triggers:
317+
scope_result = select_trigger_scope(
318+
body.project,
319+
trigger_kind=body.trigger.kind,
320+
trigger_node_id=body.trigger.triggerNodeId,
321+
)
322+
if scope_result.selection_error is not None:
323+
scope_project = body.project
324+
runtime_nodes: list[CompiledWorkflowNode] = []
325+
errors = [scope_result.selection_error]
326+
events = _compile_failure_events(
327+
workflow_id=body.project.id,
328+
run_id=run_id,
329+
trace_id=trace_id,
330+
errors=errors,
331+
)
332+
stored_events = [*prior_events, *events]
333+
projection = _build_projection(
334+
workflow_id=body.project.id,
335+
run_id=run_id,
336+
trace_id=trace_id,
337+
package_node_id=body.packageNodeId,
338+
started_at=started_at,
339+
valid=False,
340+
errors=errors,
341+
runtime_nodes=[],
342+
events=stored_events,
343+
)
344+
await _store_workflow_run(
345+
run_id,
346+
request=body,
347+
projection=projection,
348+
events=stored_events,
349+
session=session,
350+
workflow_version_id=workflow_version_id,
351+
studio_workflow_version_id=studio_workflow_version_id,
352+
)
353+
return projection
354+
scope_project = scope_result.project
355+
else:
356+
scope_project = body.project
308357
compile_result = (
309358
await compile_managed_dify_workflow_project(
310-
body.project,
359+
scope_project,
311360
graphon_client=graphon_client,
312361
session=session,
313362
)
314363
if graphon_client is not None
315-
else compile_workflow_project(body.project)
364+
else compile_workflow_project(scope_project)
316365
)
317366

318367
if not compile_result.valid or compile_result.plan is None:
@@ -428,7 +477,7 @@ async def start_workflow_run(
428477
) and (body.packageNodeId is not None or _select_package_id(runtime_nodes, None) is not None)
429478
trace = (
430479
build_opencli_hda_trace(
431-
body.project,
480+
scope_project,
432481
package_node_id=body.packageNodeId,
433482
run_id=run_id,
434483
trace_id=trace_id,
@@ -1565,8 +1614,8 @@ async def start_workflow_run(
15651614
trace_id=trace_id,
15661615
package_node_id=(trace.packageNodeId if trace else None) or body.packageNodeId,
15671616
started_at=started_at,
1568-
valid=trace.valid if trace else True,
1569-
errors=trace.errors if trace else [],
1617+
valid=compile_result.valid,
1618+
errors=list(compile_result.errors) if trace is None else compile_result.errors,
15701619
runtime_nodes=runtime_nodes,
15711620
events=events,
15721621
)
@@ -1589,8 +1638,8 @@ async def start_workflow_run(
15891638
package_node_id=(trace.packageNodeId if trace else None)
15901639
or body.packageNodeId,
15911640
started_at=started_at,
1592-
valid=trace.valid if trace else True,
1593-
errors=trace.errors if trace else [],
1641+
valid=compile_result.valid,
1642+
errors=list(compile_result.errors) if trace is None else compile_result.errors,
15941643
runtime_nodes=runtime_nodes,
15951644
events=stored.events,
15961645
)

0 commit comments

Comments
 (0)