From a67761e8439bb11766089ba9abe9430e53b417eb Mon Sep 17 00:00:00 2001 From: temp Date: Mon, 24 Aug 2026 14:43:48 +0800 Subject: [PATCH] fix: surface original Eino retry errors --- internal/multiagent/eino_run_error_handler.go | 89 ++++++++++++++++++- .../multiagent/eino_run_error_handler_test.go | 78 ++++++++++++++++ 2 files changed, 166 insertions(+), 1 deletion(-) diff --git a/internal/multiagent/eino_run_error_handler.go b/internal/multiagent/eino_run_error_handler.go index 0e816c8b..606196b9 100644 --- a/internal/multiagent/eino_run_error_handler.go +++ b/internal/multiagent/eino_run_error_handler.go @@ -3,6 +3,8 @@ package multiagent import ( "context" "errors" + "fmt" + "strings" "github.com/cloudwego/eino/adk" ) @@ -82,12 +84,97 @@ func (h *einoRunErrorHandler) emitError(err error, kind string) { if h == nil || h.progress == nil || err == nil { return } + userErr := einoUserFacingRunError(err) data := map[string]interface{}{ "conversationId": h.conversationID, "source": "eino", + "error": err.Error(), } if kind != "" { data["errorKind"] = kind + } else if userErr.kind != "" { + data["errorKind"] = userErr.kind } - h.progress("error", err.Error(), data) + if userErr.summary != "" { + data["errorSummary"] = userErr.summary + } + if userErr.retryExhausted { + data["retryExhausted"] = true + if userErr.totalRetries > 0 { + data["totalRetries"] = userErr.totalRetries + } + } + if userErr.rawLastError != "" { + data["lastError"] = userErr.rawLastError + } + message := err.Error() + if userErr.message != "" { + message = userErr.message + } + h.progress("error", message, data) +} + +type einoRunUserError struct { + message string + kind string + summary string + rawLastError string + retryExhausted bool + totalRetries int +} + +func einoUserFacingRunError(err error) einoRunUserError { + var out einoRunUserError + if err == nil { + return out + } + var retryErr *adk.RetryExhaustedError + if !errors.As(err, &retryErr) { + return out + } + out.retryExhausted = true + out.totalRetries = retryErr.TotalRetries + lastErr := retryErr.LastErr + if lastErr == nil { + out.kind = "model_retry_exhausted" + out.summary = "模型调用多次重试后仍未成功。" + out.message = out.summary + return out + } + out.rawLastError = strings.TrimSpace(lastErr.Error()) + if isEinoShouldRetryOutputRejected(lastErr) { + out.kind = "empty_model_output" + out.summary = "模型输出被重试策略拒绝;常见原因是空内容或缺少有效输出。" + out.message = formatEinoRetryExhaustedMessage(out.rawLastError, retryErr.TotalRetries) + return out + } + kind, summary := einoTransientRunErrorUserDetail(lastErr) + if strings.TrimSpace(summary) == "" { + summary = einoTrimRetryErrorSummary(lastErr.Error()) + } + if kind == "" { + kind = "model_retry_exhausted" + } + out.kind = kind + out.summary = summary + out.message = formatEinoRetryExhaustedMessage(summary, retryErr.TotalRetries) + return out +} + +func isEinoShouldRetryOutputRejected(err error) bool { + if err == nil { + return false + } + return strings.Contains(strings.ToLower(err.Error()), "model output rejected by shouldretry") +} + +func formatEinoRetryExhaustedMessage(summary string, totalRetries int) string { + summary = strings.TrimSpace(summary) + if summary == "" { + summary = "模型调用多次重试后仍未成功。" + } + if totalRetries > 0 { + return fmt.Sprintf("模型调用重试已耗尽(已重试 %d 次):%s", totalRetries, summary) + } + return "模型调用重试已耗尽:" + summary } diff --git a/internal/multiagent/eino_run_error_handler_test.go b/internal/multiagent/eino_run_error_handler_test.go index c310af64..b1d374c1 100644 --- a/internal/multiagent/eino_run_error_handler_test.go +++ b/internal/multiagent/eino_run_error_handler_test.go @@ -3,6 +3,7 @@ package multiagent import ( "context" "errors" + "strings" "testing" "github.com/cloudwego/eino/adk" @@ -61,6 +62,83 @@ func TestEinoRunErrorHandlerTimeoutAndGeneralErrorProgress(t *testing.T) { } } +func TestEinoRunErrorHandlerRetryExhaustedEmptyOutputProgress(t *testing.T) { + err := &adk.RetryExhaustedError{ + LastErr: errors.New("model output rejected by ShouldRetry at attempt 5"), + TotalRetries: 4, + } + var message string + var data map[string]interface{} + + got := newEinoRunErrorHandler(einoRunErrorHandlerConfig{ + ConversationID: "conv-1", + Progress: func(eventType, msg string, raw interface{}) { + if eventType == "error" { + message = msg + data, _ = raw.(map[string]interface{}) + } + }, + }).Handle(err) + + if !errors.Is(got, err) { + t.Fatalf("err = %v", got) + } + if !strings.Contains(message, "模型调用重试已耗尽") || + !strings.Contains(message, "model output rejected by ShouldRetry at attempt 5") { + t.Fatalf("message = %q", message) + } + if data["errorKind"] != "empty_model_output" { + t.Fatalf("errorKind = %#v", data["errorKind"]) + } + if data["errorSummary"] != "模型输出被重试策略拒绝;常见原因是空内容或缺少有效输出。" { + t.Fatalf("errorSummary = %#v", data["errorSummary"]) + } + if data["retryExhausted"] != true || data["totalRetries"] != 4 { + t.Fatalf("retry metadata = %#v", data) + } + if data["lastError"] != "model output rejected by ShouldRetry at attempt 5" { + t.Fatalf("lastError = %#v", data["lastError"]) + } + if data["error"] != err.Error() { + t.Fatalf("raw error = %#v, want %#v", data["error"], err.Error()) + } +} + +func TestEinoRunErrorHandlerRetryExhaustedOriginalErrorProgress(t *testing.T) { + err := &adk.RetryExhaustedError{ + LastErr: errors.New("HTTP 429 Too Many Requests"), + TotalRetries: 3, + } + var message string + var data map[string]interface{} + + got := newEinoRunErrorHandler(einoRunErrorHandlerConfig{ + ConversationID: "conv-1", + Progress: func(eventType, msg string, raw interface{}) { + if eventType == "error" { + message = msg + data, _ = raw.(map[string]interface{}) + } + }, + }).Handle(err) + + if !errors.Is(got, err) { + t.Fatalf("err = %v", got) + } + if !strings.Contains(message, "HTTP 429 Too Many Requests") { + t.Fatalf("message = %q", message) + } + if data["errorKind"] != "rate_limit" { + t.Fatalf("errorKind = %#v", data["errorKind"]) + } + if data["errorSummary"] != "HTTP 429 Too Many Requests" { + t.Fatalf("errorSummary = %#v", data["errorSummary"]) + } + if data["lastError"] != "HTTP 429 Too Many Requests" { + t.Fatalf("lastError = %#v", data["lastError"]) + } +} + func TestEinoRunErrorHandlerIterationLimitProgress(t *testing.T) { var events []string var errorKind interface{}