From 3b5868689ce61c947f92604b694f7ebdd33935d1 Mon Sep 17 00:00:00 2001 From: mlogclub Date: Mon, 17 Aug 2026 23:15:32 +0800 Subject: [PATCH] fix(workflow): harden publishing and editor testing Require high-risk actions to reference human confirmation outputs, filter deleted workflows from dashboard lists, and keep template metadata synchronized. Add safe browser-side debug execution for business nodes without triggering real side effects. --- .../browser-client/business-debug-runtime.ts | 214 ++++++++++++++++++ .../client/browser-client/index.ts | 28 ++- internal/ai/workflow/validator/validator.go | 4 +- .../ai/workflow/validator/validator_test.go | 25 ++ .../handlers/dashboard/ai_workflow_handler.go | 3 +- internal/services/ai_workflow_service_test.go | 27 +++ .../_components/official-workflow-editor.tsx | 11 + .../_components/workflow-workbench.tsx | 4 +- 8 files changed, 311 insertions(+), 5 deletions(-) create mode 100644 flowgram-editor/src/plugins/runtime-plugin/client/browser-client/business-debug-runtime.ts diff --git a/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/business-debug-runtime.ts b/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/business-debug-runtime.ts new file mode 100644 index 0000000..ed5ad10 --- /dev/null +++ b/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/business-debug-runtime.ts @@ -0,0 +1,214 @@ +import { + WorkflowStatus, + type IReport, + type TaskRunInput, + type WorkflowInputs, + type WorkflowOutputs, +} from '@flowgram.ai/runtime-interface'; + +type DebugNode = { + id: string; + type: string; + data?: { + config?: Record; + inputsValues?: Record; + }; +}; + +type DebugEdge = { + sourceNodeID: string; + targetNodeID: string; +}; + +type DebugSchema = { + nodes?: DebugNode[]; + edges?: DebugEdge[]; +}; + +type DebugValue = { + type?: string; + content?: unknown; +}; + +const BUSINESS_NODE_TYPES = new Set([ + 'conversation_understanding', + 'reply_policy', + 'knowledge_retrieve', + 'answerability_gate', + 'llm_reply', + 'analyze_conversation', + 'prepare_ticket_draft', + 'human_confirm', + 'create_ticket', + 'handoff_to_human', + 'send_reply', +]); + +export function isBusinessDebugSchema(input: TaskRunInput): boolean { + const schema = parseSchema(input.schema); + return (schema.nodes ?? []).some((node) => BUSINESS_NODE_TYPES.has(node.type)); +} + +export function buildBusinessDebugReport(input: TaskRunInput, taskID: string): IReport { + const schema = parseSchema(input.schema); + const nodes = schema.nodes ?? []; + const nodesByID = new Map(nodes.map((node) => [node.id, node])); + const outgoing = new Map(); + for (const edge of schema.edges ?? []) { + outgoing.set(edge.sourceNodeID, [...(outgoing.get(edge.sourceNodeID) ?? []), edge]); + } + + const values = new Map(); + const reports: IReport['reports'] = {}; + const startedAt = Date.now(); + let current = nodes.find((node) => node.type === 'start'); + let workflowOutputs: WorkflowOutputs = { status: 'completed', debug: true }; + + for (let step = 0; current && step < 128; step += 1) { + const nodeInputs = resolveNodeInputs(current, values); + const nodeOutputs = simulateNode(current, nodeInputs, input.inputs); + values.set(current.id, nodeOutputs); + const now = Date.now(); + reports[current.id] = { + id: current.id, + status: WorkflowStatus.Succeeded, + terminated: true, + startTime: now, + endTime: now, + timeCost: 0, + snapshots: [ + { + id: `${taskID}-${current.id}`, + nodeID: current.id, + inputs: nodeInputs, + outputs: nodeOutputs, + data: { debug: true }, + }, + ], + }; + if (current.type === 'end') { + workflowOutputs = { ...nodeOutputs, status: 'completed', debug: true }; + break; + } + const nextEdge = (outgoing.get(current.id) ?? [])[0]; + current = nextEdge ? nodesByID.get(nextEdge.targetNodeID) : undefined; + } + + const endedAt = Date.now(); + return { + id: taskID, + inputs: input.inputs, + outputs: workflowOutputs, + workflowStatus: { + status: WorkflowStatus.Succeeded, + terminated: true, + startTime: startedAt, + endTime: endedAt, + timeCost: endedAt - startedAt, + }, + reports, + messages: { + log: [], + info: [], + debug: [], + error: [], + warning: [], + }, + }; +} + +function parseSchema(raw: string): DebugSchema { + try { + return JSON.parse(raw) as DebugSchema; + } catch { + return {}; + } +} + +function resolveNodeInputs( + node: DebugNode, + values: Map +): WorkflowInputs { + return Object.fromEntries( + Object.entries(node.data?.inputsValues ?? {}).map(([name, value]) => [ + name, + resolveValue(value, values), + ]) + ); +} + +function resolveValue(value: DebugValue, values: Map): unknown { + if (value?.type === 'ref' && Array.isArray(value.content)) { + const [nodeID, field] = value.content; + if (typeof nodeID === 'string' && typeof field === 'string') { + return values.get(nodeID)?.[field]; + } + } + return value?.content; +} + +function simulateNode( + node: DebugNode, + inputs: WorkflowInputs, + workflowInputs: WorkflowInputs +): WorkflowOutputs { + const userMessage = String(workflowInputs.userMessage ?? workflowInputs.query ?? '测试消息'); + switch (node.type) { + case 'start': + return { ...workflowInputs, userMessage, query: userMessage }; + case 'conversation_understanding': + return { + normalizedMessage: userMessage, + messageIntent: 'ticket_request', + answerScope: 'needs_ticket', + confidence: 1, + riskSignals: [], + reason: '调试模拟结果', + }; + case 'reply_policy': + return { action: 'prepare_ticket', requiresFlow: true, targetFlow: 'prepare_ticket' }; + case 'knowledge_retrieve': + return { documents: [], count: 0, query: String(inputs.query ?? userMessage) }; + case 'answerability_gate': + return { answerability: 'answerable', reason: '调试模拟结果' }; + case 'analyze_conversation': + return { intent: 'ticket_request', riskLevel: 'low', needTicket: true }; + case 'prepare_ticket_draft': { + const issue = String(inputs.issue ?? userMessage).trim() || '测试问题'; + const title = issue.length > 30 ? `${issue.slice(0, 30)}…` : issue; + const ticketDraft = { + ready: true, + title, + description: issue, + missingFields: [], + followUpQuestions: [], + conversationFacts: [issue], + }; + return { ticketDraft, ...ticketDraft }; + } + case 'human_confirm': + return { confirmed: true, responseText: '调试运行自动确认' }; + case 'create_ticket': + return { + ticketId: 0, + ticketNo: '', + created: false, + skipped: true, + message: '调试运行不会创建工单。', + }; + case 'handoff_to_human': + return { handoffId: 0, skipped: true, message: '调试运行不会转人工。' }; + case 'llm_reply': { + const staticReply = String(node.data?.config?.staticReply ?? '').trim(); + return { replyText: staticReply || '调试回复' }; + } + case 'send_reply': + return { sent: Boolean(inputs.replyText), replyMessageId: 0 }; + case 'condition': + return { matched: true }; + case 'end': + return { ...inputs, status: 'completed' }; + default: + return { ...inputs, debug: true }; + } +} diff --git a/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/index.ts b/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/index.ts index de08189..29d3dad 100644 --- a/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/index.ts +++ b/flowgram-editor/src/plugins/runtime-plugin/client/browser-client/index.ts @@ -4,14 +4,27 @@ */ /* eslint-disable no-console */ -import { FlowGramAPIName, IRuntimeClient } from '@flowgram.ai/runtime-interface'; +import { + FlowGramAPIName, + IRuntimeClient, + type TaskReportOutput, +} from '@flowgram.ai/runtime-interface'; import { injectable } from '@flowgram.ai/free-layout-editor'; +import { buildBusinessDebugReport, isBusinessDebugSchema } from './business-debug-runtime'; + @injectable() export class WorkflowRuntimeBrowserClient implements IRuntimeClient { + private businessDebugTasks = new Map(); + constructor() {} public [FlowGramAPIName.TaskRun]: IRuntimeClient[FlowGramAPIName.TaskRun] = async (input) => { + if (isBusinessDebugSchema(input)) { + const taskID = `business-debug-${Date.now()}`; + this.businessDebugTasks.set(taskID, buildBusinessDebugReport(input, taskID)); + return { taskID }; + } const { TaskRunAPI } = await import('@flowgram.ai/runtime-js'); // Load on demand - 按需加载 return TaskRunAPI(input); }; @@ -19,6 +32,9 @@ export class WorkflowRuntimeBrowserClient implements IRuntimeClient { public [FlowGramAPIName.TaskReport]: IRuntimeClient[FlowGramAPIName.TaskReport] = async ( input ) => { + if (this.businessDebugTasks.has(input.taskID)) { + return this.businessDebugTasks.get(input.taskID); + } const { TaskReportAPI } = await import('@flowgram.ai/runtime-js'); // Load on demand - 按需加载 return TaskReportAPI(input); }; @@ -26,6 +42,10 @@ export class WorkflowRuntimeBrowserClient implements IRuntimeClient { public [FlowGramAPIName.TaskResult]: IRuntimeClient[FlowGramAPIName.TaskResult] = async ( input ) => { + const report = this.businessDebugTasks.get(input.taskID); + if (report) { + return report.outputs; + } const { TaskResultAPI } = await import('@flowgram.ai/runtime-js'); // Load on demand - 按需加载 return TaskResultAPI(input); }; @@ -33,6 +53,9 @@ export class WorkflowRuntimeBrowserClient implements IRuntimeClient { public [FlowGramAPIName.TaskCancel]: IRuntimeClient[FlowGramAPIName.TaskCancel] = async ( input ) => { + if (this.businessDebugTasks.delete(input.taskID)) { + return { success: true }; + } const { TaskCancelAPI } = await import('@flowgram.ai/runtime-js'); // Load on demand - 按需加载 return TaskCancelAPI(input); }; @@ -40,6 +63,9 @@ export class WorkflowRuntimeBrowserClient implements IRuntimeClient { public [FlowGramAPIName.TaskValidate]: IRuntimeClient[FlowGramAPIName.TaskValidate] = async ( input ) => { + if (isBusinessDebugSchema(input)) { + return { valid: true, errors: [] }; + } const { TaskValidateAPI } = await import('@flowgram.ai/runtime-js'); // Load on demand - 按需加载 return TaskValidateAPI(input); }; diff --git a/internal/ai/workflow/validator/validator.go b/internal/ai/workflow/validator/validator.go index dc3f9e8..598111b 100644 --- a/internal/ai/workflow/validator/validator.go +++ b/internal/ai/workflow/validator/validator.go @@ -169,7 +169,9 @@ func (v *definitionValidator) validateConfirmationGuards() { func (v *definitionValidator) validateConfirmedInput(nodeID string, node dsl.Node) { value, ok := node.Data.InputsValues["confirmed"] sourceNodeID, sourceField, refOK := value.Ref() + field := "nodes." + nodeID + ".data.inputsValues.confirmed" if !ok || !refOK || strings.TrimSpace(sourceNodeID) == "" || strings.TrimSpace(sourceField) == "" { + v.addError(field, "confirmed input must come from human_confirm.confirmed") return } sourceNode, ok := v.nodesByID[sourceNodeID] @@ -177,7 +179,7 @@ func (v *definitionValidator) validateConfirmedInput(nodeID string, node dsl.Nod return } if sourceNode.Type != registry.NodeTypeHumanConfirm || strings.TrimSpace(sourceField) != "confirmed" { - v.addError("nodes."+nodeID+".data.inputsValues.confirmed", "confirmed input must come from human_confirm.confirmed") + v.addError(field, "confirmed input must come from human_confirm.confirmed") } } diff --git a/internal/ai/workflow/validator/validator_test.go b/internal/ai/workflow/validator/validator_test.go index edba6a7..f8b22e4 100644 --- a/internal/ai/workflow/validator/validator_test.go +++ b/internal/ai/workflow/validator/validator_test.go @@ -106,6 +106,31 @@ func TestValidateDefinitionRejectsMissingRequiredInputValue(t *testing.T) { } } +func TestValidateDefinitionRejectsConstantConfirmationForHighRiskNode(t *testing.T) { + def := dsl.Definition{ + Nodes: []dsl.Node{ + node("start_1", "start", nil, nil), + node("confirm_1", "human_confirm", inputs("prompt", dsl.ConstantValue("请确认")), nil), + node("create_1", "create_ticket", map[string]dsl.Value{ + "ticketDraft": dsl.ConstantValue(map[string]any{"title": "测试", "description": "测试描述"}), + "confirmed": dsl.ConstantValue(false), + }, nil), + node("end_1", "end", nil, nil), + }, + Edges: []dsl.Edge{ + edge("start_1", "confirm_1"), + edge("confirm_1", "create_1"), + edge("create_1", "end_1"), + }, + } + + result := validator.ValidateDefinition(def, registry.DefaultRegistry()) + + if result.Valid || !hasValidationMessage(result, "confirmed input must come from human_confirm.confirmed") { + t.Fatalf("expected confirmation-source error, got %#v", result.Errors) + } +} + func TestValidateDefinitionRejectsUnknownInputSourceNode(t *testing.T) { def := minimalDefinition() def.Nodes[1].Data.InputsValues["replyText"] = dsl.RefValue("missing_1", "replyText") diff --git a/internal/handlers/dashboard/ai_workflow_handler.go b/internal/handlers/dashboard/ai_workflow_handler.go index b4c6327..72fbd48 100644 --- a/internal/handlers/dashboard/ai_workflow_handler.go +++ b/internal/handlers/dashboard/ai_workflow_handler.go @@ -6,6 +6,7 @@ import ( "agent-desk/internal/pkg/constants" "agent-desk/internal/pkg/dto/request" "agent-desk/internal/pkg/dto/response" + "agent-desk/internal/pkg/enums" "agent-desk/internal/pkg/httpx" "agent-desk/internal/pkg/httpx/params" "agent-desk/internal/services" @@ -22,7 +23,7 @@ func AIWorkflowAnyList(ctx *gin.Context) { cnd := params.NewPagedSqlCnd(ctx, params.QueryFilter{ParamName: "status"}, params.QueryFilter{ParamName: "name", Op: params.Like}, - ).Desc("id") + ).NotEq("status", enums.StatusDeleted).Desc("id") list, paging := services.AIWorkflowService.FindPageByCnd(cnd) httpx.WriteJSON(ctx, &web.PageResult{Results: builders.BuildAIWorkflowList(list), Page: paging}) } diff --git a/internal/services/ai_workflow_service_test.go b/internal/services/ai_workflow_service_test.go index d595138..3e8b648 100644 --- a/internal/services/ai_workflow_service_test.go +++ b/internal/services/ai_workflow_service_test.go @@ -188,6 +188,33 @@ func TestAIWorkflowServicePublishRejectsInvalidDSL(t *testing.T) { } } +func TestAIWorkflowServiceListExcludesDeletedWorkflows(t *testing.T) { + setupAIWorkflowTestDB(t) + operator := aiWorkflowTestOperator() + active, err := AIWorkflowService.CreateWorkflow(request.CreateAIWorkflowRequest{ + Name: "active workflow", + Definition: validAIWorkflowDefinition(), + }, operator) + if err != nil { + t.Fatalf("CreateWorkflow(active) error = %v", err) + } + deleted, err := AIWorkflowService.CreateWorkflow(request.CreateAIWorkflowRequest{ + Name: "deleted workflow", + Definition: validAIWorkflowDefinition(), + }, operator) + if err != nil { + t.Fatalf("CreateWorkflow(deleted) error = %v", err) + } + if err := AIWorkflowService.DeleteWorkflow(deleted.ID, operator); err != nil { + t.Fatalf("DeleteWorkflow() error = %v", err) + } + + list, paging := AIWorkflowService.FindPageByCnd(sqls.NewCnd().NotEq("status", enums.StatusDeleted).Desc("id").Page(1, 20)) + if paging.Total != 1 || len(list) != 1 || list[0].ID != active.ID { + t.Fatalf("deleted workflow must be excluded: total=%d list=%#v", paging.Total, list) + } +} + func TestAIWorkflowServiceRunListAndDetail(t *testing.T) { setupAIWorkflowTestDB(t) now := time.Now() diff --git a/web/app/dashboard/ai-workflows/_components/official-workflow-editor.tsx b/web/app/dashboard/ai-workflows/_components/official-workflow-editor.tsx index 14c15ae..bd85aaf 100644 --- a/web/app/dashboard/ai-workflows/_components/official-workflow-editor.tsx +++ b/web/app/dashboard/ai-workflows/_components/official-workflow-editor.tsx @@ -7,6 +7,13 @@ import type { AIWorkflowDefinition, AIWorkflowNodeSpec } from "@/lib/api/admin" const MESSAGE_SOURCE = "agent-desk" const EMPTY_NODE_SPECS: AIWorkflowNodeSpec[] = [] +function definitionsEqual( + left: AIWorkflowDefinition, + right: AIWorkflowDefinition +) { + return JSON.stringify(left) === JSON.stringify(right) +} + type EditorMessage = | { source: typeof MESSAGE_SOURCE @@ -64,6 +71,10 @@ export function OfficialWorkflowEditor({ if (event.data.type === "workflow:ready") { loadDocument() } else if (event.data.type === "workflow:change") { + if (definitionsEqual(event.data.document, definitionRef.current)) { + return + } + definitionRef.current = event.data.document onDefinitionChange(event.data.document) } } diff --git a/web/app/dashboard/ai-workflows/_components/workflow-workbench.tsx b/web/app/dashboard/ai-workflows/_components/workflow-workbench.tsx index 5df711e..ac4fc33 100644 --- a/web/app/dashboard/ai-workflows/_components/workflow-workbench.tsx +++ b/web/app/dashboard/ai-workflows/_components/workflow-workbench.tsx @@ -387,8 +387,8 @@ export function WorkflowWorkbench({ if (!template) return setSelectedTemplate(code) setDefinition(structuredClone(template.definition)) - if (!name.trim()) setName(template.name) - if (!description.trim()) setDescription(template.description) + setName(template.name) + setDescription(template.description) setDirty(true) }} />