mirror of
https://github.com/Ed1s0nZ/CyberStrikeAI.git
synced 2026-08-16 08:00:29 +02:00
Add files via upload
This commit is contained in:
@@ -725,7 +725,9 @@ func (h *AgentHandler) ProcessMessageForRobot(ctx context.Context, conversationI
|
|||||||
"deep",
|
"deep",
|
||||||
)
|
)
|
||||||
if errMA != nil {
|
if errMA != nil {
|
||||||
h.persistEinoAgentTraceForResume(conversationID, resultMA)
|
if shouldPersistEinoAgentTraceAfterRunError(ctx) {
|
||||||
|
h.persistEinoAgentTraceForResume(conversationID, resultMA)
|
||||||
|
}
|
||||||
errMsg := "执行失败: " + errMA.Error()
|
errMsg := "执行失败: " + errMA.Error()
|
||||||
if assistantMessageID != "" {
|
if assistantMessageID != "" {
|
||||||
_, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", errMsg, time.Now(), assistantMessageID)
|
_, _ = h.db.Exec("UPDATE messages SET content = ?, updated_at = ? WHERE id = ?", errMsg, time.Now(), assistantMessageID)
|
||||||
@@ -2576,7 +2578,7 @@ func (h *AgentHandler) executeBatchQueue(queueID string) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if runErr != nil {
|
if runErr != nil {
|
||||||
if useRunResult {
|
if useRunResult && shouldPersistEinoAgentTraceAfterRunError(baseCtx) {
|
||||||
h.persistEinoAgentTraceForResume(conversationID, resultMA)
|
h.persistEinoAgentTraceForResume(conversationID, resultMA)
|
||||||
}
|
}
|
||||||
// 检查是否是取消错误
|
// 检查是否是取消错误
|
||||||
@@ -2632,16 +2634,6 @@ func (h *AgentHandler) executeBatchQueue(queueID string) {
|
|||||||
h.logger.Warn("保存取消消息失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.Error(errMsg))
|
h.logger.Warn("保存取消消息失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.Error(errMsg))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// 保存代理轨迹(如果存在)
|
|
||||||
if result != nil && (result.LastAgentTraceInput != "" || result.LastAgentTraceOutput != "") {
|
|
||||||
if err := h.db.SaveAgentTrace(conversationID, result.LastAgentTraceInput, result.LastAgentTraceOutput); err != nil {
|
|
||||||
h.logger.Warn("保存取消任务的代理轨迹失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.Error(err))
|
|
||||||
}
|
|
||||||
} else if useRunResult && resultMA != nil && (resultMA.LastAgentTraceInput != "" || resultMA.LastAgentTraceOutput != "") {
|
|
||||||
if err := h.db.SaveAgentTrace(conversationID, resultMA.LastAgentTraceInput, resultMA.LastAgentTraceOutput); err != nil {
|
|
||||||
h.logger.Warn("保存取消任务的代理轨迹失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.Error(err))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
h.batchTaskManager.UpdateTaskStatusWithConversationID(queueID, task.ID, "cancelled", cancelMsg, "", conversationID)
|
h.batchTaskManager.UpdateTaskStatusWithConversationID(queueID, task.ID, "cancelled", cancelMsg, "", conversationID)
|
||||||
} else {
|
} else {
|
||||||
h.logger.Error("批量任务执行失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.String("conversationId", conversationID), zap.Error(runErr))
|
h.logger.Error("批量任务执行失败", zap.String("queueId", queueID), zap.String("taskId", task.ID), zap.String("conversationId", conversationID), zap.Error(runErr))
|
||||||
|
|||||||
@@ -123,7 +123,14 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) {
|
|||||||
roleTools := prep.RoleTools
|
roleTools := prep.RoleTools
|
||||||
|
|
||||||
taskStatus := "completed"
|
taskStatus := "completed"
|
||||||
defer h.tasks.FinishTask(conversationID, taskStatus)
|
// 仅在成功 StartTask 后再 FinishTask。若 StartTask 因 ErrTaskAlreadyRunning 失败仍 defer FinishTask,
|
||||||
|
// 会误删其他连接上正在运行的同会话任务,导致「第一次拦截、第二次却放行」。
|
||||||
|
taskOwned := false
|
||||||
|
defer func() {
|
||||||
|
if taskOwned {
|
||||||
|
h.tasks.FinishTask(conversationID, taskStatus)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
sendEvent("progress", "正在启动 Eino ADK 单代理(ChatModelAgent)...", map[string]interface{}{
|
sendEvent("progress", "正在启动 Eino ADK 单代理(ChatModelAgent)...", map[string]interface{}{
|
||||||
"conversationId": conversationID,
|
"conversationId": conversationID,
|
||||||
@@ -168,6 +175,7 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) {
|
|||||||
timeoutCancel()
|
timeoutCancel()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
taskOwned = true
|
||||||
firstRun = false
|
firstRun = false
|
||||||
} else {
|
} else {
|
||||||
if err := h.tasks.ResetTaskCancelForContinue(conversationID, cancelWithCause); err != nil {
|
if err := h.tasks.ResetTaskCancelForContinue(conversationID, cancelWithCause); err != nil {
|
||||||
@@ -203,8 +211,10 @@ func (h *AgentHandler) EinoSingleAgentLoopStream(c *gin.Context) {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
h.persistEinoAgentTraceForResume(conversationID, result)
|
|
||||||
cause := context.Cause(baseCtx)
|
cause := context.Cause(baseCtx)
|
||||||
|
if shouldPersistEinoAgentTraceAfterRunError(baseCtx) {
|
||||||
|
h.persistEinoAgentTraceForResume(conversationID, result)
|
||||||
|
}
|
||||||
if errors.Is(cause, ErrUserInterruptContinue) {
|
if errors.Is(cause, ErrUserInterruptContinue) {
|
||||||
reason := h.tasks.TakeInterruptContinueReason(conversationID)
|
reason := h.tasks.TakeInterruptContinueReason(conversationID)
|
||||||
prepNext, perr := h.prepareSessionAfterUserInterrupt(conversationID, assistantMessageID, reason, roleTools)
|
prepNext, perr := h.prepareSessionAfterUserInterrupt(conversationID, assistantMessageID, reason, roleTools)
|
||||||
@@ -374,7 +384,9 @@ func (h *AgentHandler) EinoSingleAgentLoop(c *gin.Context) {
|
|||||||
progressCallback,
|
progressCallback,
|
||||||
)
|
)
|
||||||
if runErr != nil {
|
if runErr != nil {
|
||||||
h.persistEinoAgentTraceForResume(prep.ConversationID, result)
|
if shouldPersistEinoAgentTraceAfterRunError(baseCtx) {
|
||||||
|
h.persistEinoAgentTraceForResume(prep.ConversationID, result)
|
||||||
|
}
|
||||||
c.JSON(http.StatusInternalServerError, gin.H{"error": runErr.Error()})
|
c.JSON(http.StatusInternalServerError, gin.H{"error": runErr.Error()})
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -141,7 +141,13 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) {
|
|||||||
orch := strings.TrimSpace(req.Orchestration)
|
orch := strings.TrimSpace(req.Orchestration)
|
||||||
|
|
||||||
taskStatus := "completed"
|
taskStatus := "completed"
|
||||||
defer h.tasks.FinishTask(conversationID, taskStatus)
|
// 仅在成功 StartTask 后再 FinishTask;避免「任务已存在」分支 return 时误删正在运行的同会话任务。
|
||||||
|
taskOwned := false
|
||||||
|
defer func() {
|
||||||
|
if taskOwned {
|
||||||
|
h.tasks.FinishTask(conversationID, taskStatus)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
sendEvent("progress", "正在启动 Eino 多代理...", map[string]interface{}{
|
sendEvent("progress", "正在启动 Eino 多代理...", map[string]interface{}{
|
||||||
"conversationId": conversationID,
|
"conversationId": conversationID,
|
||||||
@@ -178,6 +184,7 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) {
|
|||||||
timeoutCancel()
|
timeoutCancel()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
taskOwned = true
|
||||||
firstRun = false
|
firstRun = false
|
||||||
} else {
|
} else {
|
||||||
if err := h.tasks.ResetTaskCancelForContinue(conversationID, cancelWithCause); err != nil {
|
if err := h.tasks.ResetTaskCancelForContinue(conversationID, cancelWithCause); err != nil {
|
||||||
@@ -215,8 +222,10 @@ func (h *AgentHandler) MultiAgentLoopStream(c *gin.Context) {
|
|||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
h.persistEinoAgentTraceForResume(conversationID, result)
|
|
||||||
cause := context.Cause(baseCtx)
|
cause := context.Cause(baseCtx)
|
||||||
|
if shouldPersistEinoAgentTraceAfterRunError(baseCtx) {
|
||||||
|
h.persistEinoAgentTraceForResume(conversationID, result)
|
||||||
|
}
|
||||||
if errors.Is(cause, ErrUserInterruptContinue) {
|
if errors.Is(cause, ErrUserInterruptContinue) {
|
||||||
reason := h.tasks.TakeInterruptContinueReason(conversationID)
|
reason := h.tasks.TakeInterruptContinueReason(conversationID)
|
||||||
prepNext, perr := h.prepareSessionAfterUserInterrupt(conversationID, assistantMessageID, reason, roleTools)
|
prepNext, perr := h.prepareSessionAfterUserInterrupt(conversationID, assistantMessageID, reason, roleTools)
|
||||||
@@ -388,7 +397,9 @@ func (h *AgentHandler) MultiAgentLoop(c *gin.Context) {
|
|||||||
strings.TrimSpace(req.Orchestration),
|
strings.TrimSpace(req.Orchestration),
|
||||||
)
|
)
|
||||||
if runErr != nil {
|
if runErr != nil {
|
||||||
h.persistEinoAgentTraceForResume(prep.ConversationID, result)
|
if shouldPersistEinoAgentTraceAfterRunError(baseCtx) {
|
||||||
|
h.persistEinoAgentTraceForResume(prep.ConversationID, result)
|
||||||
|
}
|
||||||
h.logger.Error("Eino DeepAgent 执行失败", zap.Error(runErr))
|
h.logger.Error("Eino DeepAgent 执行失败", zap.Error(runErr))
|
||||||
errMsg := "执行失败: " + runErr.Error()
|
errMsg := "执行失败: " + runErr.Error()
|
||||||
if prep.AssistantMessageID != "" {
|
if prep.AssistantMessageID != "" {
|
||||||
|
|||||||
@@ -17,6 +17,15 @@ var ErrUserInterruptContinue = errors.New("user interrupt with continue")
|
|||||||
// ErrTaskAlreadyRunning 会话已有任务正在执行
|
// ErrTaskAlreadyRunning 会话已有任务正在执行
|
||||||
var ErrTaskAlreadyRunning = errors.New("agent task already running for conversation")
|
var ErrTaskAlreadyRunning = errors.New("agent task already running for conversation")
|
||||||
|
|
||||||
|
// shouldPersistEinoAgentTraceAfterRunError:Eino 相关 Run 非成功返回时,是否仍写入 last_react_* 供下轮 loadHistoryFromAgentTrace。
|
||||||
|
// 用户主动停止(WithCancelCause(ErrTaskCancelled))时不写入半截轨迹,避免下一轮 Plan-Execute 等因损坏/不完整 tool 上下文出现 no tool call 等异常。
|
||||||
|
func shouldPersistEinoAgentTraceAfterRunError(baseCtx context.Context) bool {
|
||||||
|
if baseCtx == nil {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
return !errors.Is(context.Cause(baseCtx), ErrTaskCancelled)
|
||||||
|
}
|
||||||
|
|
||||||
// AgentTask 描述正在运行的Agent任务
|
// AgentTask 描述正在运行的Agent任务
|
||||||
type AgentTask struct {
|
type AgentTask struct {
|
||||||
ConversationID string `json:"conversationId"`
|
ConversationID string `json:"conversationId"`
|
||||||
|
|||||||
Reference in New Issue
Block a user