From 89326b30ae1ead26123205863d5f8d6ca365addf Mon Sep 17 00:00:00 2001 From: qhd <1766646056@qq.com> Date: Fri, 14 Aug 2026 16:35:44 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E4=BF=AE=E5=A4=8D=E7=BD=91=E5=85=B3?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E5=B9=B6=E5=8F=91=E7=AD=89=E5=BE=85=E7=AB=9E?= =?UTF-8?q?=E6=80=81=E5=AF=BC=E8=87=B4=E6=B0=B8=E4=B9=85=E9=98=BB=E5=A1=9E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- workflow/service/flow/lambda_node_util.go | 22 ++++++++++------------ 1 file changed, 10 insertions(+), 12 deletions(-) diff --git a/workflow/service/flow/lambda_node_util.go b/workflow/service/flow/lambda_node_util.go index 1670ec9..7d2c9f7 100644 --- a/workflow/service/flow/lambda_node_util.go +++ b/workflow/service/flow/lambda_node_util.go @@ -455,27 +455,25 @@ func GetModelResult(ctx context.Context, sessionId string, nodeInput *flowDto.No //updateTokenCount(ctx, nodeInput.NodeExecutionId, modelInfo.Model.ResponseTokenField, taskResult) } } else { - taskIdList := make([]string, len(composeResult.Messages.Rounds)) - - for idx, item := range composeResult.Messages.Rounds { - taskId, err := createGatewayTaskOnly(ctx, composeResult.EpicycleId, nodeInput.Config.ModelConfig.ModelName, item) - if err != nil { - return nil, err - } - taskIdList[idx] = taskId - } - // 全局共享子上下文,实现一处报错全部终止 subCtx, globalCancel := context.WithCancel(ctx) defer globalCancel() // 函数退出兜底释放 var wg sync.WaitGroup - errChan := make(chan error, len(taskIdList)) + errChan := make(chan error, len(composeResult.Messages.Rounds)) // 加互斥锁保护结果map var mu sync.Mutex - for idx, taskId := range taskIdList { + // 每个任务创建后立即启动等待协程:把「回调 vs Wait 注册」的竞争窗口从整个创建循环 + // 压缩到微秒级,避免创建期间完成的回调被 Notify 静默丢弃导致 Wait 永久阻塞 + for idx, item := range composeResult.Messages.Rounds { + taskId, err := createGatewayTaskOnly(ctx, composeResult.EpicycleId, nodeInput.Config.ModelConfig.ModelName, item) + if err != nil { + globalCancel() // 取消已启动的等待协程,避免泄漏 + return nil, err + } + wg.Add(1) go func(idx int, taskId string) {