diff --git a/internal/handler/agent.go b/internal/handler/agent.go index 6d16ef2b..1a1067b5 100644 --- a/internal/handler/agent.go +++ b/internal/handler/agent.go @@ -333,17 +333,24 @@ type ChatReasoningRequest struct { Effort string `json:"effort,omitempty"` } +// ChatFinalizationRequest is a caller-provided delivery policy. The server does +// not infer execution intent from natural-language user text. +type ChatFinalizationRequest struct { + RequireExecutionEvidence *bool `json:"requireExecutionEvidence,omitempty"` +} + // ChatRequest 聊天请求 type ChatRequest struct { - Message string `json:"message" binding:"required"` - ConversationID string `json:"conversationId,omitempty"` - ProjectID string `json:"projectId,omitempty"` // 新对话绑定的项目(可选;未指定时可用 config.project.default_project_id) - Role string `json:"role,omitempty"` // 角色名称 - Attachments []ChatAttachment `json:"attachments,omitempty"` - WebShellConnectionID string `json:"webshellConnectionId,omitempty"` // WebShell 管理 - AI 助手:当前选中的连接 ID,仅使用 webshell_* 工具 - AIChannelID string `json:"aiChannelId,omitempty"` // 会话级 AI 通道;空则使用 ai.default_channel - Hitl *HITLRequest `json:"hitl,omitempty"` - Reasoning *ChatReasoningRequest `json:"reasoning,omitempty"` + Message string `json:"message" binding:"required"` + ConversationID string `json:"conversationId,omitempty"` + ProjectID string `json:"projectId,omitempty"` // 新对话绑定的项目(可选;未指定时可用 config.project.default_project_id) + Role string `json:"role,omitempty"` // 角色名称 + Attachments []ChatAttachment `json:"attachments,omitempty"` + WebShellConnectionID string `json:"webshellConnectionId,omitempty"` // WebShell 管理 - AI 助手:当前选中的连接 ID,仅使用 webshell_* 工具 + AIChannelID string `json:"aiChannelId,omitempty"` // 会话级 AI 通道;空则使用 ai.default_channel + Hitl *HITLRequest `json:"hitl,omitempty"` + Reasoning *ChatReasoningRequest `json:"reasoning,omitempty"` + Finalization ChatFinalizationRequest `json:"finalization,omitempty"` // Orchestration 仅对 /api/multi-agent、/api/multi-agent/stream:deep | plan_execute | supervisor;空则等同 deep。机器人/批量等无请求体时由服务端默认 deep。/api/eino-agent* 不使用此字段。 Orchestration string `json:"orchestration,omitempty"` } @@ -668,10 +675,18 @@ func (h *AgentHandler) mergeAssistantMessagePartialOnCancel(messageID, partial s // ChatResponse 聊天响应 type ChatResponse struct { - Response string `json:"response"` - MCPExecutionIDs []string `json:"mcpExecutionIds,omitempty"` // 本次对话中执行的MCP调用ID列表 - ConversationID string `json:"conversationId"` // 对话ID - Time time.Time `json:"time"` + Response string `json:"response"` + MCPExecutionIDs []string `json:"mcpExecutionIds,omitempty"` // 本次对话中执行的MCP调用ID列表 + ConversationID string `json:"conversationId"` // 对话ID + Time time.Time `json:"time"` + Finalizable bool `json:"finalizable"` + Finalized bool `json:"finalized"` + Status string `json:"status,omitempty"` + CompletionReason string `json:"completionReason,omitempty"` + EvidenceVerified bool `json:"evidenceVerified"` + EvidenceRefs []string `json:"evidenceRefs,omitempty"` + PendingExecutionIDs []string `json:"pendingExecutionIds,omitempty"` + MissingChecks []string `json:"missingChecks,omitempty"` } func (h *AgentHandler) finalizeRobotAgentError(ctx context.Context, assistantMessageID, conversationID string, resultMA *multiagent.RunResult, errMA error) (string, string, error) { @@ -687,19 +702,20 @@ func (h *AgentHandler) finalizeRobotAgentError(ctx context.Context, assistantMes } func (h *AgentHandler) finalizeRobotAgentSuccess(assistantMessageID, conversationID string, resultMA *multiagent.RunResult) (string, string, error) { - if assistantMessageID != "" { - if errU := h.db.UpdateAssistantMessageFinalize(assistantMessageID, resultMA.Response, resultMA.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(resultMA.LastAgentTraceInput)); errU != nil { - h.logger.Warn("机器人:更新助手消息失败", zap.Error(errU)) - } - } else { - if _, err := h.db.AddMessage(conversationID, "assistant", resultMA.Response, resultMA.MCPExecutionIDs); err != nil { + decision := h.finalizeAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, "robot", resultMA, resultMA.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(resultMA.LastAgentTraceInput), true) + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + } + if assistantMessageID == "" { + if _, err := h.db.AddMessage(conversationID, "assistant", responseText, resultMA.MCPExecutionIDs); err != nil { h.logger.Warn("机器人:保存助手消息失败", zap.Error(err)) } } if resultMA.LastAgentTraceInput != "" || resultMA.LastAgentTraceOutput != "" { _ = h.db.SaveAgentTrace(conversationID, resultMA.LastAgentTraceInput, resultMA.LastAgentTraceOutput) } - return resultMA.Response, conversationID, nil + return responseText, conversationID, nil } func (h *AgentHandler) runRobotEinoSingleWithRetry( @@ -830,6 +846,9 @@ func (h *AgentHandler) ProcessMessageForRobot(ctx context.Context, platform stri progressCallback := h.createProgressCallback(taskCtx, cancelWithCause, conversationID, assistantMessageID, nil) robotMode := config.NormalizeAgentMode(agentMode) + if err := h.db.SetConversationAgentMode(conversationID, robotMode); err != nil { + h.logger.Warn("机器人:更新对话模式失败", zap.String("conversationId", conversationID), zap.String("agentMode", robotMode), zap.Error(err)) + } switch robotMode { case "eino_single": return h.runRobotEinoSingleWithRetry(taskCtx, conversationID, finalMessage, agentHistoryMessages, roleTools, progressCallback, assistantMessageID, &taskStatus) diff --git a/internal/handler/batch_queue_executor.go b/internal/handler/batch_queue_executor.go index 7f770e31..4daaa0a2 100644 --- a/internal/handler/batch_queue_executor.go +++ b/internal/handler/batch_queue_executor.go @@ -238,6 +238,11 @@ func (h *AgentHandler) executeOneBatchSubTask(queueID string, queue *BatchTaskQu useBatchMulti = true batchOrch = "deep" } + if useBatchMulti { + _ = h.db.SetConversationAgentMode(conversationID, batchOrch) + } else { + _ = h.db.SetConversationAgentMode(conversationID, "eino_single") + } var resultMA *multiagent.RunResult var runErr error @@ -268,19 +273,38 @@ func (h *AgentHandler) executeOneBatchSubTask(queueID string, queue *BatchTaskQu h.logger.Info("批量任务执行成功", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.String("conversationId", conversationID)) - resText := resultMA.Response mcpIDs := resultMA.MCPExecutionIDs lastIn := resultMA.LastAgentTraceInput lastOut := resultMA.LastAgentTraceOutput + reasoningContent := multiagent.AggregatedReasoningFromTraceJSON(lastIn) + agentMode := "batch_eino_single" + if useBatchMulti { + agentMode = "batch_eino_" + batchOrch + } + decision := h.finalizeAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, agentMode, resultMA, mcpIDs, reasoningContent, true) + resText := decision.FinalText + if !decision.Finalizable { + resText = finalizationBlockedMessage(decision) + finishStatus = decision.Status + sendEvent("finalization_check", resText, decision) + } + sendEvent("response", resText, finalizationResponsePayload(decision, map[string]interface{}{ + "conversationId": conversationID, + "messageId": assistantMessageID, + "agentMode": agentMode, + "mcpExecutionIds": mcpIDs, + "batchQueueId": queueID, + "batchTaskId": task.ID, + "batchTaskStatus": map[bool]string{true: string(BatchTaskStatusCompleted), false: string(BatchTaskStatusFailed)}[decision.Finalizable], + "candidatePreview": safeTruncateString(resultMA.Response, 500), + })) - if assistantMessageID != "" { - if updateErr := h.db.UpdateAssistantMessageFinalize(assistantMessageID, resText, mcpIDs, multiagent.AggregatedReasoningFromTraceJSON(lastIn)); updateErr != nil { - h.logger.Warn("更新助手消息失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.Error(updateErr)) - if _, err = h.db.AddMessage(conversationID, "assistant", resText, mcpIDs); err != nil { - h.logger.Error("保存助手消息失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.String("conversationId", conversationID), zap.Error(err)) - } - } - } else if _, err = h.db.AddMessage(conversationID, "assistant", resText, mcpIDs); err != nil { + if assistantMessageID == "" { + _, err = h.db.AddMessage(conversationID, "assistant", resText, mcpIDs) + } else if !decision.Finalizable { + err = nil + } + if err != nil { h.logger.Error("保存助手消息失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.String("conversationId", conversationID), zap.Error(err)) } @@ -290,6 +314,10 @@ func (h *AgentHandler) executeOneBatchSubTask(queueID string, queue *BatchTaskQu } } + if !decision.Finalizable { + h.batchTaskManager.UpdateTaskStatusWithConversationID(queueID, task.ID, BatchTaskStatusFailed, resText, finalizationCheckMessage(decision), conversationID) + return + } h.batchTaskManager.UpdateTaskStatusWithConversationID(queueID, task.ID, BatchTaskStatusCompleted, resText, "", conversationID) } diff --git a/internal/handler/eino_empty_response_continue.go b/internal/handler/eino_empty_response_continue.go index 98139854..2c57df86 100644 --- a/internal/handler/eino_empty_response_continue.go +++ b/internal/handler/eino_empty_response_continue.go @@ -68,15 +68,15 @@ func (h *AgentHandler) tryContinueOnEinoEmptyResponse( case <-time.After(backoff): } - inject := multiagent.FormatEmptyResponseContinueUserMessage() - h.applyEinoTraceResumeSegment(conversationID, result, curHistory, curFinalMessage, inject) + h.applyEinoTraceResumeSegment(conversationID, result, curHistory, curFinalMessage, "") if progressCallback != nil { progressCallback("eino_empty_response_continue", "已恢复上下文,正在续跑…", map[string]interface{}{ - "conversationId": conversationID, - "source": "eino", - "attempt": *attempt, - "maxAttempts": maxAttempts, - "contextSource": "empty_response_continue", + "conversationId": conversationID, + "source": "eino", + "attempt": *attempt, + "maxAttempts": maxAttempts, + "contextSource": "empty_response_continue", + "contextInjection": false, }) } return true diff --git a/internal/handler/eino_single_agent.go b/internal/handler/eino_single_agent.go index 762861dd..6d965813 100644 --- a/internal/handler/eino_single_agent.go +++ b/internal/handler/eino_single_agent.go @@ -10,6 +10,7 @@ import ( "sync" "time" + "cyberstrike-ai/internal/agentfinalizer" "cyberstrike-ai/internal/mcp" "cyberstrike-ai/internal/multiagent" @@ -189,6 +190,8 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) { // 同一请求内分段续跑时,主代理 iteration 事件按偏移累计,避免 UI 出现「第3轮 → 第1轮」回跳。 var mainIterationOffset int var emptyResponseContinueAttempt int + var finalizationAutoContinueAttempt int + var decision agentfinalizer.Decision for { segmentMainIterationMax := 0 @@ -258,6 +261,13 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) { baseCtx, cancelWithCause, taskCtx, timeoutCancel = h.rebindEinoRunningTask(taskCtx, conversationID, timeoutCancel) continue } + decision = h.decideAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, "eino_single", result, cumulativeMCPExecutionIDs, requestRequiresExecutionEvidence(&req)) + if h.tryAutoContinueAfterFinalization(taskCtx, conversationID, result, decision, &finalizationAutoContinueAttempt, &curHistory, &curFinalMessage, progressCallback) { + mainIterationOffset += segmentMainIterationMax + timeoutCancel() + baseCtx, cancelWithCause, taskCtx, timeoutCancel = h.rebindEinoRunningTask(taskCtx, conversationID, timeoutCancel) + continue + } timeoutCancel() break } @@ -358,9 +368,10 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) { timeoutCancel() - if assistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(assistantMessageID, result.Response, cumulativeMCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput)) + if decision.CompletionReason == "" { + decision = h.decideAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, "eino_single", result, cumulativeMCPExecutionIDs, requestRequiresExecutionEvidence(&req)) } + h.persistFinalizationDecision(conversationID, assistantMessageID, "eino_single", cumulativeMCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput), decision) if result.LastAgentTraceInput != "" || result.LastAgentTraceOutput != "" { if err := h.db.SaveAgentTrace(conversationID, result.LastAgentTraceInput, result.LastAgentTraceOutput); err != nil { @@ -368,12 +379,19 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) { } } - sendEvent("response", result.Response, map[string]interface{}{ + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + sendEvent("finalization_check", responseText, decision) + taskStatus = decision.Status + h.tasks.UpdateTaskStatus(conversationID, taskStatus) + } + sendEvent("response", responseText, finalizationResponsePayload(decision, map[string]interface{}{ "mcpExecutionIds": cumulativeMCPExecutionIDs, "conversationId": conversationID, "messageId": assistantMessageID, "agentMode": "eino_single", - }) + })) sendEvent("done", "", map[string]interface{}{"conversationId": conversationID}) } @@ -429,6 +447,9 @@ func (h *AgentHandler) EinoSingleAgentLoop(c *gin.Context) { curMsg := prep.FinalMessage var result *multiagent.RunResult var runErr error + var emptyResponseContinueAttempt int + var finalizationAutoContinueAttempt int + var decision agentfinalizer.Decision for { result, runErr = multiagent.RunEinoSingleChatModelAgent( taskCtx, @@ -446,28 +467,46 @@ func (h *AgentHandler) EinoSingleAgentLoop(c *gin.Context) { chatReasoningToClientIntent(req.Reasoning), h.agentSessionContextBlock(prep.ConversationID), ) - if runErr == nil { - break + if runErr != nil { + if shouldPersistEinoAgentTraceAfterRunError(baseCtx) { + h.persistEinoAgentTraceForResume(prep.ConversationID, result) + } + c.JSON(http.StatusInternalServerError, gin.H{"error": runErr.Error()}) + return } - if shouldPersistEinoAgentTraceAfterRunError(baseCtx) { - h.persistEinoAgentTraceForResume(prep.ConversationID, result) + mw := &h.config.MultiAgent.EinoMiddleware + if h.tryContinueOnEinoEmptyResponse(taskCtx, mw, prep.ConversationID, result, &emptyResponseContinueAttempt, &curHist, &curMsg, progressCallback) { + continue } - c.JSON(http.StatusInternalServerError, gin.H{"error": runErr.Error()}) - return + decision = h.decideAgentRunForDeliveryWithPolicy(prep.ConversationID, prep.AssistantMessageID, "eino_single", result, result.MCPExecutionIDs, requestRequiresExecutionEvidence(&req)) + if h.tryAutoContinueAfterFinalization(taskCtx, prep.ConversationID, result, decision, &finalizationAutoContinueAttempt, &curHist, &curMsg, progressCallback) { + continue + } + break } - if prep.AssistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(prep.AssistantMessageID, result.Response, result.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput)) - } + h.persistFinalizationDecision(prep.ConversationID, prep.AssistantMessageID, "eino_single", result.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput), decision) if result.LastAgentTraceInput != "" || result.LastAgentTraceOutput != "" { _ = h.db.SaveAgentTrace(prep.ConversationID, result.LastAgentTraceInput, result.LastAgentTraceOutput) } + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + } c.JSON(http.StatusOK, gin.H{ - "response": result.Response, - "conversationId": prep.ConversationID, - "mcpExecutionIds": result.MCPExecutionIDs, - "assistantMessageId": prep.AssistantMessageID, - "agentMode": "eino_single", + "response": responseText, + "conversationId": prep.ConversationID, + "mcpExecutionIds": result.MCPExecutionIDs, + "assistantMessageId": prep.AssistantMessageID, + "agentMode": "eino_single", + "finalized": decision.Finalized, + "finalizable": decision.Finalizable, + "status": decision.Status, + "completionReason": decision.CompletionReason, + "evidenceVerified": decision.EvidenceVerified, + "evidenceRefs": decision.EvidenceRefs, + "pendingExecutionIds": decision.PendingExecutionIDs, + "missingChecks": decision.MissingChecks, }) } diff --git a/internal/handler/finalization_auto_continue.go b/internal/handler/finalization_auto_continue.go new file mode 100644 index 00000000..2aca373c --- /dev/null +++ b/internal/handler/finalization_auto_continue.go @@ -0,0 +1,77 @@ +package handler + +import ( + "context" + "time" + + "cyberstrike-ai/internal/agent" + "cyberstrike-ai/internal/agentfinalizer" + "cyberstrike-ai/internal/multiagent" + + "go.uber.org/zap" +) + +const finalizationAutoContinueMaxAttempts = 2 + +func shouldAutoContinueAfterFinalization(d agentfinalizer.Decision, attempt int) bool { + if d.Finalizable || d.Finalized { + return false + } + if attempt >= finalizationAutoContinueMaxAttempts { + return false + } + return d.CompletionReason == agentfinalizer.ReasonMissingEvidence +} + +func (h *AgentHandler) tryAutoContinueAfterFinalization( + taskCtx context.Context, + conversationID string, + result *multiagent.RunResult, + decision agentfinalizer.Decision, + attempt *int, + curHistory *[]agent.ChatMessage, + curFinalMessage *string, + progressCallback func(eventType, message string, data interface{}), +) bool { + if !shouldAutoContinueAfterFinalization(decision, *attempt) || result == nil || !multiagent.HasEinoResumeTrace(result) { + return false + } + *attempt++ + h.persistEinoAgentTraceForResume(conversationID, result) + if hist, err := h.loadHistoryFromAgentTrace(conversationID); err == nil && len(hist) > 0 { + *curHistory = hist + } else if h.logger != nil { + h.logger.Warn("finalization auto-continue could not restore trace", + zap.String("conversationId", conversationID), + zap.Error(err)) + return false + } + // Agent 无感续跑:不追加新的 user/system 文案,只使用上一段模型可见轨迹继续 Runner。 + *curFinalMessage = "" + if progressCallback != nil { + progressCallback("finalization_auto_continue", "最终回复检查尚未收敛,正在基于已有轨迹继续执行…", map[string]interface{}{ + "conversationId": conversationID, + "source": "finalizer", + "attempt": *attempt, + "maxAttempts": finalizationAutoContinueMaxAttempts, + "status": decision.Status, + "completionReason": decision.CompletionReason, + "missingChecks": decision.MissingChecks, + "pendingExecutionIds": decision.PendingExecutionIDs, + "contextInjection": false, + }) + } + select { + case <-taskCtx.Done(): + return false + case <-time.After(finalizationAutoContinueBackoff(*attempt)): + return true + } +} + +func finalizationAutoContinueBackoff(attempt int) time.Duration { + if attempt <= 1 { + return 500 * time.Millisecond + } + return time.Duration(attempt) * time.Second +} diff --git a/internal/handler/finalization_auto_continue_test.go b/internal/handler/finalization_auto_continue_test.go new file mode 100644 index 00000000..1ce68d4c --- /dev/null +++ b/internal/handler/finalization_auto_continue_test.go @@ -0,0 +1,59 @@ +package handler + +import ( + "testing" + + "cyberstrike-ai/internal/agentfinalizer" +) + +func TestShouldAutoContinueAfterFinalization(t *testing.T) { + missingEvidence := agentfinalizer.Decision{ + Status: agentfinalizer.StatusBlocked, + CompletionReason: agentfinalizer.ReasonMissingEvidence, + } + if !shouldAutoContinueAfterFinalization(missingEvidence, 0) { + t.Fatal("missing execution evidence should trigger auto-continue") + } + if shouldAutoContinueAfterFinalization(missingEvidence, finalizationAutoContinueMaxAttempts) { + t.Fatal("auto-continue should stop at max attempts") + } + + finalized := agentfinalizer.Decision{ + Status: agentfinalizer.StatusCompleted, + CompletionReason: agentfinalizer.ReasonVerified, + Finalizable: true, + Finalized: true, + } + if shouldAutoContinueAfterFinalization(finalized, 0) { + t.Fatal("finalized decision should not auto-continue") + } + + awaitingHITL := agentfinalizer.Decision{ + Status: agentfinalizer.StatusAwaitingHITL, + CompletionReason: agentfinalizer.ReasonAwaitingHITL, + } + if shouldAutoContinueAfterFinalization(awaitingHITL, 0) { + t.Fatal("awaiting HITL should not auto-continue without approval") + } +} + +func TestRequestRequiresExecutionEvidenceUsesExplicitPolicyOnly(t *testing.T) { + if requestRequiresExecutionEvidence(nil) { + t.Fatal("nil request should not require execution evidence") + } + if requestRequiresExecutionEvidence(&ChatRequest{}) { + t.Fatal("missing finalization policy should not require execution evidence") + } + require := true + if !requestRequiresExecutionEvidence(&ChatRequest{ + Finalization: ChatFinalizationRequest{RequireExecutionEvidence: &require}, + }) { + t.Fatal("explicit true policy should require execution evidence") + } + require = false + if requestRequiresExecutionEvidence(&ChatRequest{ + Finalization: ChatFinalizationRequest{RequireExecutionEvidence: &require}, + }) { + t.Fatal("explicit false policy should not require execution evidence") + } +} diff --git a/internal/handler/finalization_helpers.go b/internal/handler/finalization_helpers.go new file mode 100644 index 00000000..f9fdb8ba --- /dev/null +++ b/internal/handler/finalization_helpers.go @@ -0,0 +1,171 @@ +package handler + +import ( + "fmt" + "strings" + "time" + + "cyberstrike-ai/internal/agentfinalizer" + "cyberstrike-ai/internal/multiagent" + + "go.uber.org/zap" +) + +func (h *AgentHandler) finalizeAgentRunForDelivery( + conversationID string, + assistantMessageID string, + agentMode string, + result *multiagent.RunResult, + mcpExecutionIDs []string, + reasoningContent string, +) agentfinalizer.Decision { + return h.finalizeAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, agentMode, result, mcpExecutionIDs, reasoningContent, false) +} + +func (h *AgentHandler) finalizeAgentRunForDeliveryWithPolicy( + conversationID string, + assistantMessageID string, + agentMode string, + result *multiagent.RunResult, + mcpExecutionIDs []string, + reasoningContent string, + requireExecutionEvidence bool, +) agentfinalizer.Decision { + decision := agentfinalizer.FromRunResult(h.db, result, agentfinalizer.Input{ + ConversationID: conversationID, + AssistantMessageID: assistantMessageID, + AgentMode: agentMode, + MCPExecutionIDs: mcpExecutionIDs, + RequireExecutionEvidence: requireExecutionEvidence, + }) + h.persistFinalizationDecision(conversationID, assistantMessageID, agentMode, mcpExecutionIDs, reasoningContent, decision) + return decision +} + +func (h *AgentHandler) decideAgentRunForDeliveryWithPolicy( + conversationID string, + assistantMessageID string, + agentMode string, + result *multiagent.RunResult, + mcpExecutionIDs []string, + requireExecutionEvidence bool, +) agentfinalizer.Decision { + return agentfinalizer.FromRunResult(h.db, result, agentfinalizer.Input{ + ConversationID: conversationID, + AssistantMessageID: assistantMessageID, + AgentMode: agentMode, + MCPExecutionIDs: mcpExecutionIDs, + RequireExecutionEvidence: requireExecutionEvidence, + }) +} + +func (h *AgentHandler) decideAgentRunForDelivery( + conversationID string, + assistantMessageID string, + agentMode string, + result *multiagent.RunResult, + mcpExecutionIDs []string, +) agentfinalizer.Decision { + return agentfinalizer.FromRunResult(h.db, result, agentfinalizer.Input{ + ConversationID: conversationID, + AssistantMessageID: assistantMessageID, + AgentMode: agentMode, + MCPExecutionIDs: mcpExecutionIDs, + RequireExecutionEvidence: false, + }) +} + +func (h *AgentHandler) persistFinalizationDecision( + conversationID string, + assistantMessageID string, + agentMode string, + mcpExecutionIDs []string, + reasoningContent string, + decision agentfinalizer.Decision, +) { + if assistantMessageID == "" || h.db == nil { + return + } + _ = h.db.AddProcessDetail(assistantMessageID, conversationID, "finalization_check", finalizationCheckMessage(decision), decision) + if decision.Finalizable { + if err := h.db.UpdateAssistantMessageFinalize(assistantMessageID, decision.FinalText, mcpExecutionIDs, reasoningContent); err != nil && h.logger != nil { + h.logger.Warn("更新最终助手消息失败", zap.Error(err), zap.String("conversationId", conversationID), zap.String("agentMode", agentMode)) + } + return + } + _, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", finalizationBlockedMessage(decision), time.Now(), assistantMessageID) +} + +func (h *AgentHandler) finalizeCandidateForDelivery( + conversationID string, + assistantMessageID string, + agentMode string, + response string, + mcpExecutionIDs []string, + awaitingHITL bool, + reasoningContent string, +) agentfinalizer.Decision { + return h.finalizeCandidateForDeliveryWithPolicy(conversationID, assistantMessageID, agentMode, response, mcpExecutionIDs, awaitingHITL, reasoningContent, false) +} + +func (h *AgentHandler) finalizeCandidateForDeliveryWithPolicy( + conversationID string, + assistantMessageID string, + agentMode string, + response string, + mcpExecutionIDs []string, + awaitingHITL bool, + reasoningContent string, + requireExecutionEvidence bool, +) agentfinalizer.Decision { + decision := agentfinalizer.Decide(h.db, agentfinalizer.Input{ + Response: response, + ConversationID: conversationID, + AssistantMessageID: assistantMessageID, + AgentMode: agentMode, + MCPExecutionIDs: mcpExecutionIDs, + AwaitingHITL: awaitingHITL, + RequireExecutionEvidence: requireExecutionEvidence, + }) + if assistantMessageID == "" || h.db == nil { + return decision + } + _ = h.db.AddProcessDetail(assistantMessageID, conversationID, "finalization_check", finalizationCheckMessage(decision), decision) + if decision.Finalizable { + if err := h.db.UpdateAssistantMessageFinalize(assistantMessageID, decision.FinalText, mcpExecutionIDs, reasoningContent); err != nil && h.logger != nil { + h.logger.Warn("更新最终助手消息失败", zap.Error(err), zap.String("conversationId", conversationID), zap.String("agentMode", agentMode)) + } + return decision + } + _, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", finalizationBlockedMessage(decision), time.Now(), assistantMessageID) + return decision +} + +func finalizationCheckMessage(d agentfinalizer.Decision) string { + if d.Finalizable { + return "最终回复检查通过。" + } + return finalizationBlockedMessage(d) +} + +func finalizationBlockedMessage(d agentfinalizer.Decision) string { + parts := []string{"任务尚未达到最终回复条件,暂不生成成功结论。"} + if d.CompletionReason != "" { + parts = append(parts, "原因: "+d.CompletionReason) + } + if len(d.PendingExecutionIDs) > 0 { + parts = append(parts, fmt.Sprintf("仍有 %d 个工具执行未结束: %s", len(d.PendingExecutionIDs), strings.Join(d.PendingExecutionIDs, ", "))) + } + if len(d.MissingChecks) > 0 { + parts = append(parts, "缺失检查: "+strings.Join(d.MissingChecks, "; ")) + } + return strings.Join(parts, "\n") +} + +func finalizationResponsePayload(d agentfinalizer.Decision, extra map[string]interface{}) map[string]interface{} { + return agentfinalizer.ResponsePayload(d, extra) +} + +func requestRequiresExecutionEvidence(req *ChatRequest) bool { + return req != nil && req.Finalization.RequireExecutionEvidence != nil && *req.Finalization.RequireExecutionEvidence +} diff --git a/internal/handler/multi_agent.go b/internal/handler/multi_agent.go index f2a3a22d..e1d2c7bb 100644 --- a/internal/handler/multi_agent.go +++ b/internal/handler/multi_agent.go @@ -10,6 +10,7 @@ import ( "sync" "time" + "cyberstrike-ai/internal/agentfinalizer" "cyberstrike-ai/internal/config" "cyberstrike-ai/internal/mcp" "cyberstrike-ai/internal/multiagent" @@ -197,6 +198,13 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) { // 同一请求内分段续跑时,主代理 iteration 事件按偏移累计,避免 UI 出现「第3轮 → 第1轮」回跳。 var mainIterationOffset int var emptyResponseContinueAttempt int + var finalizationAutoContinueAttempt int + effectiveOrch := config.NormalizeMultiAgentOrchestration(h.config.MultiAgent.Orchestration) + if o := strings.TrimSpace(req.Orchestration); o != "" { + effectiveOrch = config.NormalizeMultiAgentOrchestration(o) + } + agentMode := "eino_" + effectiveOrch + var decision agentfinalizer.Decision for { segmentMainIterationMax := 0 @@ -267,6 +275,13 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) { baseCtx, cancelWithCause, taskCtx, timeoutCancel = h.rebindEinoRunningTask(taskCtx, conversationID, timeoutCancel) continue } + decision = h.decideAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, agentMode, result, cumulativeMCPExecutionIDs, requestRequiresExecutionEvidence(&req)) + if h.tryAutoContinueAfterFinalization(taskCtx, conversationID, result, decision, &finalizationAutoContinueAttempt, &curHistory, &curFinalMessage, progressCallback) { + mainIterationOffset += segmentMainIterationMax + timeoutCancel() + baseCtx, cancelWithCause, taskCtx, timeoutCancel = h.rebindEinoRunningTask(taskCtx, conversationID, timeoutCancel) + continue + } timeoutCancel() break } @@ -367,9 +382,10 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) { timeoutCancel() - if assistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(assistantMessageID, result.Response, cumulativeMCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput)) + if decision.CompletionReason == "" { + decision = h.decideAgentRunForDeliveryWithPolicy(conversationID, assistantMessageID, agentMode, result, cumulativeMCPExecutionIDs, requestRequiresExecutionEvidence(&req)) } + h.persistFinalizationDecision(conversationID, assistantMessageID, agentMode, cumulativeMCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput), decision) if result.LastAgentTraceInput != "" || result.LastAgentTraceOutput != "" { if err := h.db.SaveAgentTrace(conversationID, result.LastAgentTraceInput, result.LastAgentTraceOutput); err != nil { @@ -377,16 +393,19 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) { } } - effectiveOrch := config.NormalizeMultiAgentOrchestration(h.config.MultiAgent.Orchestration) - if o := strings.TrimSpace(req.Orchestration); o != "" { - effectiveOrch = config.NormalizeMultiAgentOrchestration(o) + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + sendEvent("finalization_check", responseText, decision) + taskStatus = decision.Status + h.tasks.UpdateTaskStatus(conversationID, taskStatus) } - sendEvent("response", result.Response, map[string]interface{}{ + sendEvent("response", responseText, finalizationResponsePayload(decision, map[string]interface{}{ "mcpExecutionIds": cumulativeMCPExecutionIDs, "conversationId": conversationID, "messageId": assistantMessageID, - "agentMode": "eino_" + effectiveOrch, - }) + "agentMode": agentMode, + })) sendEvent("done", "", map[string]interface{}{"conversationId": conversationID}) } @@ -437,6 +456,14 @@ func (h *AgentHandler) MultiAgentLoop(c *gin.Context) { curMsg := prep.FinalMessage var result *multiagent.RunResult var runErr error + var emptyResponseContinueAttempt int + var finalizationAutoContinueAttempt int + effectiveOrch := config.NormalizeMultiAgentOrchestration(h.config.MultiAgent.Orchestration) + if o := strings.TrimSpace(req.Orchestration); o != "" { + effectiveOrch = config.NormalizeMultiAgentOrchestration(o) + } + agentMode := "eino_" + effectiveOrch + var decision agentfinalizer.Decision for { result, runErr = multiagent.RunDeepAgent( taskCtx, @@ -456,24 +483,30 @@ func (h *AgentHandler) MultiAgentLoop(c *gin.Context) { chatReasoningToClientIntent(req.Reasoning), h.agentSessionContextBlock(prep.ConversationID), ) - if runErr == nil { - break + if runErr != nil { + if shouldPersistEinoAgentTraceAfterRunError(baseCtx) { + h.persistEinoAgentTraceForResume(prep.ConversationID, result) + } + h.logger.Error("Eino DeepAgent 执行失败", zap.Error(runErr)) + errMsg := "执行失败: " + runErr.Error() + if prep.AssistantMessageID != "" { + _, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", errMsg, time.Now(), prep.AssistantMessageID) + } + c.JSON(http.StatusInternalServerError, gin.H{"error": errMsg}) + return } - if shouldPersistEinoAgentTraceAfterRunError(baseCtx) { - h.persistEinoAgentTraceForResume(prep.ConversationID, result) + mw := &h.config.MultiAgent.EinoMiddleware + if h.tryContinueOnEinoEmptyResponse(taskCtx, mw, prep.ConversationID, result, &emptyResponseContinueAttempt, &curHist, &curMsg, progressCallback) { + continue } - h.logger.Error("Eino DeepAgent 执行失败", zap.Error(runErr)) - errMsg := "执行失败: " + runErr.Error() - if prep.AssistantMessageID != "" { - _, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", errMsg, time.Now(), prep.AssistantMessageID) + decision = h.decideAgentRunForDeliveryWithPolicy(prep.ConversationID, prep.AssistantMessageID, agentMode, result, result.MCPExecutionIDs, requestRequiresExecutionEvidence(&req)) + if h.tryAutoContinueAfterFinalization(taskCtx, prep.ConversationID, result, decision, &finalizationAutoContinueAttempt, &curHist, &curMsg, progressCallback) { + continue } - c.JSON(http.StatusInternalServerError, gin.H{"error": errMsg}) - return + break } - if prep.AssistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(prep.AssistantMessageID, result.Response, result.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput)) - } + h.persistFinalizationDecision(prep.ConversationID, prep.AssistantMessageID, agentMode, result.MCPExecutionIDs, multiagent.AggregatedReasoningFromTraceJSON(result.LastAgentTraceInput), decision) if result.LastAgentTraceInput != "" || result.LastAgentTraceOutput != "" { if err := h.db.SaveAgentTrace(prep.ConversationID, result.LastAgentTraceInput, result.LastAgentTraceOutput); err != nil { @@ -481,11 +514,23 @@ func (h *AgentHandler) MultiAgentLoop(c *gin.Context) { } } + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + } c.JSON(http.StatusOK, ChatResponse{ - Response: result.Response, - MCPExecutionIDs: result.MCPExecutionIDs, - ConversationID: prep.ConversationID, - Time: time.Now(), + Response: responseText, + MCPExecutionIDs: result.MCPExecutionIDs, + ConversationID: prep.ConversationID, + Time: time.Now(), + Finalizable: decision.Finalizable, + Finalized: decision.Finalized, + Status: decision.Status, + CompletionReason: decision.CompletionReason, + EvidenceVerified: decision.EvidenceVerified, + EvidenceRefs: decision.EvidenceRefs, + PendingExecutionIDs: decision.PendingExecutionIDs, + MissingChecks: decision.MissingChecks, }) } diff --git a/internal/handler/multi_agent_prepare.go b/internal/handler/multi_agent_prepare.go index 660d5f04..08564beb 100644 --- a/internal/handler/multi_agent_prepare.go +++ b/internal/handler/multi_agent_prepare.go @@ -6,6 +6,7 @@ import ( "cyberstrike-ai/internal/agent" "cyberstrike-ai/internal/audit" + "cyberstrike-ai/internal/config" "cyberstrike-ai/internal/database" "cyberstrike-ai/internal/mcp/builtin" "cyberstrike-ai/internal/security" @@ -25,6 +26,13 @@ type multiAgentPrepared struct { UserMessageID string } +func chatRequestAgentMode(req *ChatRequest, source string) string { + if strings.HasPrefix(strings.TrimSpace(source), "multi_agent") { + return config.NormalizeMultiAgentOrchestration(req.Orchestration) + } + return "eino_single" +} + func (h *AgentHandler) prepareMultiAgentSession(req *ChatRequest, c *gin.Context, source string) (*multiAgentPrepared, error) { if len(req.Attachments) > maxAttachments { return nil, fmt.Errorf("附件最多 %d 个", maxAttachments) @@ -57,6 +65,7 @@ func (h *AgentHandler) prepareMultiAgentSession(req *ChatRequest, c *gin.Context meta := audit.ConversationCreateMetaFromGin(c, source) meta.ProjectID = projectID meta.RoleName = req.Role + meta.AgentMode = chatRequestAgentMode(req, source) if webshellID != "" { meta.Source = source + "_webshell" meta.WebShellConnectionID = webshellID @@ -84,6 +93,9 @@ func (h *AgentHandler) prepareMultiAgentSession(req *ChatRequest, c *gin.Context if err := h.db.SetConversationRoleName(conversationID, req.Role); err != nil { h.logger.Warn("更新对话角色失败", zap.String("conversationId", conversationID), zap.String("role", req.Role), zap.Error(err)) } + if err := h.db.SetConversationAgentMode(conversationID, chatRequestAgentMode(req, source)); err != nil { + h.logger.Warn("更新对话模式失败", zap.String("conversationId", conversationID), zap.String("source", source), zap.String("orchestration", req.Orchestration), zap.Error(err)) + } agentHistoryMessages, err := h.loadHistoryFromAgentTrace(conversationID) if err != nil { diff --git a/internal/handler/openapi.go b/internal/handler/openapi.go index 2e317393..fad523d0 100644 --- a/internal/handler/openapi.go +++ b/internal/handler/openapi.go @@ -35,6 +35,17 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { scheme = "https" } + finalizationRequestSchema := map[string]interface{}{ + "type": "object", + "description": "最终回复交付策略。后端不会从自然语言内容推断执行意图;执行入口应显式声明是否要求 completed 工具证据。", + "properties": map[string]interface{}{ + "requireExecutionEvidence": map[string]interface{}{ + "type": "boolean", + "description": "为 true 时,缺少 completed 工具执行记录会触发无注入续跑或最终阻断;普通聊天可省略或设为 false。", + }, + }, + } + spec := map[string]interface{}{ "openapi": "3.0.0", "info": map[string]interface{}{ @@ -85,6 +96,70 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { }, "required": []string{"projectId"}, }, + "AgentChatResponse": map[string]interface{}{ + "type": "object", + "description": "Agent 非流式响应。response 只是交付文本;是否为成功最终回复必须以 finalized/finalizable/status 为准。", + "properties": map[string]interface{}{ + "response": map[string]interface{}{ + "type": "string", + "description": "交付给用户的文本。finalized=false 时为阻断/未完成说明,不是成功结论。", + }, + "conversationId": map[string]interface{}{ + "type": "string", + "description": "对话 ID", + }, + "assistantMessageId": map[string]interface{}{ + "type": "string", + "description": "助手消息 ID(部分接口返回)", + }, + "mcpExecutionIds": map[string]interface{}{ + "type": "array", + "description": "本轮关联的 MCP 工具执行 ID", + "items": map[string]interface{}{"type": "string"}, + }, + "agentMode": map[string]interface{}{ + "type": "string", + "description": "agent 模式,例如 eino_single、eino_deep、workflow", + }, + "finalized": map[string]interface{}{ + "type": "boolean", + "description": "是否已经通过最终回复检查。只有 true 才能当成功最终回复。", + }, + "finalizable": map[string]interface{}{ + "type": "boolean", + "description": "候选输出是否可提升为最终回复。", + }, + "status": map[string]interface{}{ + "type": "string", + "description": "最终化状态", + "enum": []string{"completed", "in_progress", "blocked", "failed", "cancelled", "awaiting_hitl"}, + }, + "completionReason": map[string]interface{}{ + "type": "string", + "description": "最终化或阻断原因,例如 verified、pending_tool_executions、missing_execution_evidence", + }, + "evidenceVerified": map[string]interface{}{ + "type": "boolean", + "description": "证据是否满足最终化要求", + }, + "evidenceRefs": map[string]interface{}{ + "type": "array", + "description": "证据引用,例如 mcp_execution:", + "items": map[string]interface{}{"type": "string"}, + }, + "pendingExecutionIds": map[string]interface{}{ + "type": "array", + "description": "仍处于 queued/running 的工具执行 ID", + "items": map[string]interface{}{"type": "string"}, + }, + "missingChecks": map[string]interface{}{ + "type": "array", + "description": "未通过最终化检查的原因列表", + "items": map[string]interface{}{"type": "string"}, + }, + }, + "required": []string{"response", "conversationId", "finalized", "finalizable", "status", "evidenceVerified"}, + }, "Conversation": map[string]interface{}{ "type": "object", "properties": map[string]interface{}{ @@ -1581,6 +1656,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "conversationId": map[string]interface{}{"type": "string"}, "role": map[string]interface{}{"type": "string"}, "webshellConnectionId": map[string]interface{}{"type": "string"}, + "finalization": finalizationRequestSchema, }, "required": []string{"message"}, }, @@ -1588,7 +1664,14 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { }, }, "responses": map[string]interface{}{ - "200": map[string]interface{}{"description": "成功,响应格式同 /api/eino-agent"}, + "200": map[string]interface{}{ + "description": "成功。只有 finalized=true 表示成功最终回复;finalized=false 时 response 为未完成/阻断说明。", + "content": map[string]interface{}{ + "application/json": map[string]interface{}{ + "schema": map[string]interface{}{"$ref": "#/components/schemas/AgentChatResponse"}, + }, + }, + }, "400": map[string]interface{}{"description": "参数错误"}, "401": map[string]interface{}{"description": "未授权"}, "500": map[string]interface{}{"description": "执行失败"}, @@ -1599,7 +1682,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "post": map[string]interface{}{ "tags": []string{"对话交互"}, "summary": "发送消息并获取 AI 回复(Eino ADK 单代理,SSE)", - "description": "向 AI 发送消息并获取流式回复(SSE)。由 Eino **单代理** ADK 执行;事件类型与多代理流式一致(含 `tool_call` / `response_delta` / `thinking` 等)。**不依赖** `multi_agent.enabled`。", + "description": "向 AI 发送消息并获取流式回复(SSE)。由 Eino **单代理** ADK 执行;事件类型与多代理流式一致(含 `tool_call` / `response_delta` / `thinking` 等)。`response_start` / `response_delta` 仅为候选/过程输出;只有 `type: response` 且 `data.finalized=true` 才表示成功最终回复。缺 completed 执行证据时可能先发送 `finalization_auto_continue`,表示服务端基于已有 trace 无注入续跑。`data.finalized=false` 时 message 为未完成/阻断说明。**不依赖** `multi_agent.enabled`。", "operationId": "sendMessageEinoSingleAgentStream", "requestBody": map[string]interface{}{ "required": true, @@ -1612,6 +1695,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "conversationId": map[string]interface{}{"type": "string"}, "role": map[string]interface{}{"type": "string"}, "webshellConnectionId": map[string]interface{}{"type": "string"}, + "finalization": finalizationRequestSchema, }, "required": []string{"message"}, }, @@ -1625,7 +1709,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "text/event-stream": map[string]interface{}{ "schema": map[string]interface{}{ "type": "string", - "description": "SSE 流", + "description": "SSE 流。终态 response 事件 data 包含 finalized、finalizable、status、completionReason、evidenceVerified、evidenceRefs、pendingExecutionIds、missingChecks;过程事件可能包含 finalization_auto_continue。", }, }, }, @@ -1663,6 +1747,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "type": "string", "description": "WebShell 连接 ID(可选,与 Eino 单/多代理流式行为一致)", }, + "finalization": finalizationRequestSchema, "orchestration": map[string]interface{}{ "type": "string", "description": "Eino 预置编排:deep | plan_execute | supervisor;缺省 deep", @@ -1676,7 +1761,12 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { }, "responses": map[string]interface{}{ "200": map[string]interface{}{ - "description": "成功,响应格式同 /api/eino-agent", + "description": "成功。只有 finalized=true 表示成功最终回复;finalized=false 时 response 为未完成/阻断说明。", + "content": map[string]interface{}{ + "application/json": map[string]interface{}{ + "schema": map[string]interface{}{"$ref": "#/components/schemas/AgentChatResponse"}, + }, + }, }, "400": map[string]interface{}{"description": "参数错误"}, "401": map[string]interface{}{"description": "未授权"}, @@ -1689,7 +1779,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "post": map[string]interface{}{ "tags": []string{"对话交互"}, "summary": "发送消息并获取 AI 回复(Eino 多代理,SSE)", - "description": "与 `POST /api/eino-agent/stream` 类似;由 Eino 多代理执行。`orchestration` 指定 deep / plan_execute / supervisor,缺省 deep。**前提**:`multi_agent.enabled: true`;未启用时 SSE 内首条为 `type: error` 后接 `done`。支持 `webshellConnectionId`。", + "description": "与 `POST /api/eino-agent/stream` 类似;由 Eino 多代理执行。`orchestration` 指定 deep / plan_execute / supervisor,缺省 deep。`response_start` / `response_delta` 仅为候选/过程输出;只有 `type: response` 且 `data.finalized=true` 才表示成功最终回复。缺 completed 执行证据时可能先发送 `finalization_auto_continue`,表示服务端基于已有 trace 无注入续跑。**前提**:`multi_agent.enabled: true`;未启用时 SSE 内首条为 `type: error` 后接 `done`。支持 `webshellConnectionId`。", "operationId": "sendMessageMultiAgentStream", "requestBody": map[string]interface{}{ "required": true, @@ -1702,6 +1792,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "conversationId": map[string]interface{}{"type": "string"}, "role": map[string]interface{}{"type": "string"}, "webshellConnectionId": map[string]interface{}{"type": "string"}, + "finalization": finalizationRequestSchema, "orchestration": map[string]interface{}{ "type": "string", "description": "deep | plan_execute | supervisor;缺省 deep", @@ -1720,7 +1811,7 @@ func (h *OpenAPIHandler) GetOpenAPISpec(c *gin.Context) { "text/event-stream": map[string]interface{}{ "schema": map[string]interface{}{ "type": "string", - "description": "SSE 流", + "description": "SSE 流。终态 response 事件 data 包含 finalized、finalizable、status、completionReason、evidenceVerified、evidenceRefs、pendingExecutionIds、missingChecks;过程事件可能包含 finalization_auto_continue。", }, }, }, diff --git a/internal/handler/workflow_integration.go b/internal/handler/workflow_integration.go index 3acf5d73..bdbdc894 100644 --- a/internal/handler/workflow_integration.go +++ b/internal/handler/workflow_integration.go @@ -152,20 +152,37 @@ func (h *AgentHandler) runRoleWorkflowStreamIfBound( sendEvent("done", "", map[string]interface{}{"conversationId": conversationID}) return true } - if prep.AssistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(prep.AssistantMessageID, result.Response, nil, "") + decision := h.finalizeCandidateForDeliveryWithPolicy( + prep.ConversationID, + prep.AssistantMessageID, + "workflow", + result.Response, + nil, + result.AwaitingHITL, + "", + true, + ) + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + taskStatus = decision.Status + h.tasks.UpdateTaskStatus(conversationID, taskStatus) + sendEvent("finalization_check", responseText, decision) } - payload := map[string]interface{}{ + payload := finalizationResponsePayload(decision, map[string]interface{}{ "conversationId": prep.ConversationID, "messageId": prep.AssistantMessageID, "agentMode": "workflow", "workflowRunId": result.RunID, - } + }) if result.AwaitingHITL { payload["workflowStatus"] = "awaiting_hitl" payload["awaitingHitl"] = true + } else { + payload["workflowStatus"] = result.Status + payload["awaitingHitl"] = false } - sendEvent("response", result.Response, payload) + sendEvent("response", responseText, payload) sendEvent("done", "", map[string]interface{}{"conversationId": prep.ConversationID}) return true } @@ -251,17 +268,37 @@ func (h *AgentHandler) runRoleWorkflowJSONIfBound(c *gin.Context, req *ChatReque c.JSON(http.StatusInternalServerError, gin.H{"error": errMsg, "conversationId": conversationID}) return true } - if prep.AssistantMessageID != "" { - _ = h.db.UpdateAssistantMessageFinalize(prep.AssistantMessageID, result.Response, nil, "") + decision := h.finalizeCandidateForDeliveryWithPolicy( + prep.ConversationID, + prep.AssistantMessageID, + "workflow", + result.Response, + nil, + result.AwaitingHITL, + "", + true, + ) + responseText := decision.FinalText + if !decision.Finalizable { + responseText = finalizationBlockedMessage(decision) + taskStatus = decision.Status } c.JSON(http.StatusOK, gin.H{ - "response": result.Response, - "conversationId": prep.ConversationID, - "assistantMessageId": prep.AssistantMessageID, - "agentMode": "workflow", - "workflowRunId": result.RunID, - "workflowStatus": result.Status, - "awaitingHitl": result.AwaitingHITL, + "response": responseText, + "conversationId": prep.ConversationID, + "assistantMessageId": prep.AssistantMessageID, + "agentMode": "workflow", + "workflowRunId": result.RunID, + "workflowStatus": result.Status, + "awaitingHitl": result.AwaitingHITL, + "finalized": decision.Finalized, + "finalizable": decision.Finalizable, + "status": decision.Status, + "completionReason": decision.CompletionReason, + "evidenceVerified": decision.EvidenceVerified, + "evidenceRefs": decision.EvidenceRefs, + "pendingExecutionIds": decision.PendingExecutionIDs, + "missingChecks": decision.MissingChecks, }) return true }