From c50de7770b604bed8f12546de4c97596a89517ac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=85=AC=E6=98=8E?= <83812544+Ed1s0nZ@users.noreply.github.com> Date: Mon, 20 Jul 2026 17:38:38 +0800 Subject: [PATCH] Add files via upload --- internal/multiagent/eino_adk_run_loop.go | 10 ++- internal/multiagent/eino_transient_retry.go | 78 +++++++++++++++++++ .../multiagent/eino_transient_retry_test.go | 40 ++++++++++ 3 files changed, 127 insertions(+), 1 deletion(-) diff --git a/internal/multiagent/eino_adk_run_loop.go b/internal/multiagent/eino_adk_run_loop.go index e29c9d85..44d15b10 100644 --- a/internal/multiagent/eino_adk_run_loop.go +++ b/internal/multiagent/eino_adk_run_loop.go @@ -588,11 +588,14 @@ func runEinoADKAgentLoop(ctx context.Context, args *einoADKRunLoopArgs, baseMsgs zap.Duration("backoff", backoff)) } if progress != nil { - progress("eino_run_retry", fmt.Sprintf("遇到临时错误(限流或网络波动),%d 秒后第 %d/%d 次重试…", int(backoff.Seconds()), attemptNo, maxAttempts), map[string]interface{}{ + errorKind, errorSummary := einoTransientRunErrorUserDetail(runErr) + progress("eino_run_retry", fmt.Sprintf("遇到临时错误,%d 秒后第 %d/%d 次重试。原因:%s", int(backoff.Seconds()), attemptNo, maxAttempts, errorSummary), map[string]interface{}{ "conversationId": conversationID, "source": "eino", "orchestration": orchMode, "error": runErr.Error(), + "errorKind": errorKind, + "errorSummary": errorSummary, "attempt": attemptNo, "maxAttempts": maxAttempts, "backoffSec": int(backoff.Seconds()), @@ -601,7 +604,12 @@ func runEinoADKAgentLoop(ctx context.Context, args *einoADKRunLoopArgs, baseMsgs "conversationId": conversationID, "source": "eino", "orchestration": orchMode, + "error": runErr.Error(), + "errorKind": errorKind, + "errorSummary": errorSummary, "attempt": attemptNo, + "maxAttempts": maxAttempts, + "backoffSec": int(backoff.Seconds()), "contextSource": string(ctxSource), }) } diff --git a/internal/multiagent/eino_transient_retry.go b/internal/multiagent/eino_transient_retry.go index fc033e47..751a90f1 100644 --- a/internal/multiagent/eino_transient_retry.go +++ b/internal/multiagent/eino_transient_retry.go @@ -94,6 +94,84 @@ func isRetryableHTTPStatus(status int) bool { } } +func einoTransientRunErrorUserDetail(err error) (kind, summary string) { + if err == nil { + return "", "" + } + msg := strings.TrimSpace(err.Error()) + lower := strings.ToLower(msg) + if status := httpStatusFromErrorText(lower); status > 0 { + switch { + case status == 429: + kind = "rate_limit" + case status == 408 || status == 409 || status == 425: + kind = "retryable_http" + case status >= 500 && status <= 599: + kind = "upstream_server" + default: + kind = "http_error" + } + } else { + var apiErr *einoopenai.APIError + if errors.As(err, &apiErr) && apiErr.HTTPStatusCode > 0 { + switch { + case apiErr.HTTPStatusCode == 429: + kind = "rate_limit" + case apiErr.HTTPStatusCode == 408 || apiErr.HTTPStatusCode == 409 || apiErr.HTTPStatusCode == 425: + kind = "retryable_http" + case apiErr.HTTPStatusCode >= 500 && apiErr.HTTPStatusCode <= 599: + kind = "upstream_server" + default: + kind = "http_error" + } + } + } + if kind == "" { + switch { + case strings.Contains(lower, "too many requests") || + strings.Contains(lower, "rate limit") || + strings.Contains(lower, "rate_limit") || + strings.Contains(lower, "ratelimit"): + kind = "rate_limit" + case strings.Contains(lower, "overloaded") || + strings.Contains(lower, "capacity") || + strings.Contains(lower, "temporarily unavailable") || + strings.Contains(lower, "service unavailable"): + kind = "upstream_busy" + case strings.Contains(lower, "connection reset") || + strings.Contains(lower, "connection refused") || + strings.Contains(lower, "connection closed") || + strings.Contains(lower, "i/o timeout") || + strings.Contains(lower, "no such host") || + strings.Contains(lower, "network is unreachable") || + strings.Contains(lower, "broken pipe") || + strings.Contains(lower, "read tcp") || + strings.Contains(lower, "write tcp") || + strings.Contains(lower, "dial tcp") || + strings.Contains(lower, "tls handshake timeout") || + strings.Contains(lower, "goaway") || + strings.Contains(lower, "unexpected eof"): + kind = "network" + case strings.Contains(lower, "stream error") || + strings.Contains(lower, "unexpected end of json"): + kind = "stream" + default: + kind = "transient" + } + } + return kind, einoTrimRetryErrorSummary(msg) +} + +func einoTrimRetryErrorSummary(msg string) string { + msg = strings.Join(strings.Fields(strings.TrimSpace(msg)), " ") + const maxRunes = 500 + runes := []rune(msg) + if len(runes) <= maxRunes { + return msg + } + return string(runes[:maxRunes]) + "..." +} + func httpStatusFromErrorText(msg string) int { match := httpStatusInErrorPattern.FindStringSubmatch(msg) if len(match) != 2 { diff --git a/internal/multiagent/eino_transient_retry_test.go b/internal/multiagent/eino_transient_retry_test.go index 5f8f04bc..696a4b7b 100644 --- a/internal/multiagent/eino_transient_retry_test.go +++ b/internal/multiagent/eino_transient_retry_test.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "io" + "strings" "testing" "time" @@ -62,6 +63,45 @@ func TestEinoTransientRetryBackoff(t *testing.T) { } } +func TestEinoTransientRunErrorUserDetail(t *testing.T) { + t.Parallel() + cases := []struct { + name string + err error + wantKind string + }{ + {"rate limit", errors.New("HTTP 429 Too Many Requests"), "rate_limit"}, + {"upstream", errors.New("upstream returned 503"), "upstream_server"}, + {"network", errors.New("read tcp: connection reset by peer"), "network"}, + {"stream", errors.New("unexpected end of JSON"), "stream"}, + } + for _, tc := range cases { + tc := tc + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + kind, summary := einoTransientRunErrorUserDetail(tc.err) + if kind != tc.wantKind { + t.Fatalf("kind=%q, want %q", kind, tc.wantKind) + } + if summary == "" { + t.Fatal("summary should not be empty") + } + }) + } +} + +func TestEinoTrimRetryErrorSummary(t *testing.T) { + t.Parallel() + raw := strings.Repeat("报错 ", 260) + got := einoTrimRetryErrorSummary(raw) + if len([]rune(got)) > 503 { + t.Fatalf("summary too long: %d runes", len([]rune(got))) + } + if !strings.HasSuffix(got, "...") { + t.Fatal("trimmed summary should end with ellipsis") + } +} + func TestEinoMessagesForRunRestart(t *testing.T) { t.Parallel() base := []adk.Message{schema.UserMessage("hi")}