From e9d19c8617f6ccbd9204ee118dd8d45a96dd71fb Mon Sep 17 00:00:00 2001 From: lyon Date: Sun, 21 Jun 2026 05:08:08 +0800 Subject: [PATCH] fix: retry transient workbench facts reads --- internal/cloud/server-workbench-http.ts | 17 ++- internal/cloud/workbench-facts-store.ts | 146 +++++++++++++++++++++--- 2 files changed, 145 insertions(+), 18 deletions(-) diff --git a/internal/cloud/server-workbench-http.ts b/internal/cloud/server-workbench-http.ts index 1dc0fe2e..f072d69e 100644 --- a/internal/cloud/server-workbench-http.ts +++ b/internal/cloud/server-workbench-http.ts @@ -828,16 +828,31 @@ function workbenchProjectionStoreError(error) { const data = objectValue(error?.data) ?? {}; const retryAfterNumber = Number(data.retryAfterMs); const retryAfterMs = Number.isFinite(retryAfterNumber) && retryAfterNumber >= 0 ? Math.trunc(retryAfterNumber) : null; + const retryAttempt = nonNegativeIntegerOrNull(data.retryAttempt ?? data.retryAttempts); + const retryMax = nonNegativeIntegerOrNull(data.retryMax); + const retryDelaysMs = Array.isArray(data.retryDelaysMs) + ? data.retryDelaysMs.map(nonNegativeIntegerOrNull).filter((value) => value !== null) + : []; return workbenchError("projection_store_unavailable", "Workbench durable projection store is unavailable.", { causeCode: error?.code ?? null, queryResult: textValue(data.queryResult) || null, blockedLayer: textValue(data.blockedLayer) || null, retryable: data.retryable !== false, transient: data.transient === true || data.retryable === true, - retryAfterMs: retryAfterMs ?? null + retryAfterMs: retryAfterMs ?? null, + retryAttempt, + retryMax, + retryAttempts: retryAttempt, + retryExhausted: data.retryExhausted === true, + retryDelaysMs: retryDelaysMs.length > 0 ? retryDelaysMs : null }); } +function nonNegativeIntegerOrNull(value) { + const parsed = Number(value); + return Number.isFinite(parsed) && parsed >= 0 ? Math.trunc(parsed) : null; +} + function factSessionId(session) { return safeSessionId(session?.sessionId ?? session?.id) ?? null; } diff --git a/internal/cloud/workbench-facts-store.ts b/internal/cloud/workbench-facts-store.ts index 75ad4e9f..ba6bbde8 100644 --- a/internal/cloud/workbench-facts-store.ts +++ b/internal/cloud/workbench-facts-store.ts @@ -7,6 +7,9 @@ import { safeConversationId, safeSessionId, safeTraceId } from "./server-http-ut import { durableTraceStatus, normalizeWorkbenchStatus, TERMINAL_STATUSES, traceTerminalEvidence } from "./workbench-turn-projection.ts"; const TERMINAL_DURABLE_TRACE_CACHE_MAX = 256; +const DEFAULT_WORKBENCH_FACTS_QUERY_MAX_ATTEMPTS = 5; +const DEFAULT_WORKBENCH_FACTS_QUERY_BACKOFF_BASE_MS = 80; +const DEFAULT_WORKBENCH_FACTS_QUERY_BACKOFF_MAX_MS = 1000; const terminalDurableTraceCache = new Map(); export function createWorkbenchFactsStore(options = {}, actor = null) { @@ -18,23 +21,44 @@ export function createWorkbenchFactsStore(options = {}, actor = null) { async function queryFacts(params = {}) { if (typeof runtimeStore?.queryWorkbenchFacts === "function") { - try { - const result = await runtimeStore.queryWorkbenchFacts(params); - return { - facts: normalizeWorkbenchFactsResult(result?.facts), - count: Number.isFinite(Number(result?.count)) ? Number(result.count) : 0, - durable: true, - persistence: result?.persistence ?? null - }; - } catch (error) { - logWorkbenchFactsQueryFailure(options.logger ?? console, params, error); - return { - facts: emptyWorkbenchFacts(), - count: 0, - durable: true, - error, - projection: projectionStoreUnavailableTrace(safeTraceId(params.traceId) ?? "trc_unassigned", error) - }; + const retryPolicy = workbenchFactsQueryRetryPolicy(options.env ?? process.env); + const retryDelaysMs = []; + for (let attempt = 1; attempt <= retryPolicy.maxAttempts; attempt += 1) { + try { + const result = await runtimeStore.queryWorkbenchFacts(params); + return { + facts: normalizeWorkbenchFactsResult(result?.facts), + count: Number.isFinite(Number(result?.count)) ? Number(result.count) : 0, + durable: true, + persistence: result?.persistence ?? null, + retry: attempt > 1 ? { attempts: attempt, maxAttempts: retryPolicy.maxAttempts, delaysMs: retryDelaysMs.slice(), valuesRedacted: true } : null + }; + } catch (error) { + const retryable = isRetryableWorkbenchFactsQueryError(error); + const finalAttempt = !retryable || attempt >= retryPolicy.maxAttempts; + const retryAfterMs = finalAttempt ? null : workbenchFactsRetryDelayMs(retryPolicy, attempt); + const annotatedError = annotateWorkbenchFactsQueryError(error, { + retryable, + retryAttempt: attempt, + retryMax: retryPolicy.maxAttempts, + retryAfterMs, + retryDelaysMs, + retryExhausted: finalAttempt && retryable + }); + if (finalAttempt) { + logWorkbenchFactsQueryFailure(options.logger ?? console, params, annotatedError); + return { + facts: emptyWorkbenchFacts(), + count: 0, + durable: true, + error: annotatedError, + projection: projectionStoreUnavailableTrace(safeTraceId(params.traceId) ?? "trc_unassigned", annotatedError) + }; + } + logWorkbenchFactsQueryRetry(options.logger ?? console, params, annotatedError, { attempt, maxAttempts: retryPolicy.maxAttempts, retryAfterMs }); + retryDelaysMs.push(retryAfterMs); + await sleepMs(retryAfterMs); + } } } return { @@ -240,6 +264,7 @@ function rememberTerminalDurableTrace(snapshot) { function projectionStoreUnavailableTrace(traceId, error) { const message = "运行记录投影存储暂不可读,Trace 更新暂不可见。"; + const retry = workbenchFactsRetryMetadata(error); const projection = { projectionStatus: "blocked", projectionHealth: "unavailable", @@ -250,6 +275,10 @@ function projectionStoreUnavailableTrace(traceId, error) { message, causeCode: error?.code ?? null, retryable: true, + retryAttempt: retry.retryAttempt, + retryMax: retry.retryMax, + retryExhausted: retry.retryExhausted, + retryAfterMs: retry.retryAfterMs, valuesPrinted: false }, valuesRedacted: true @@ -268,6 +297,89 @@ function projectionStoreUnavailableTrace(traceId, error) { }; } +function workbenchFactsQueryRetryPolicy(env = {}) { + return { + maxAttempts: boundedInteger(env.HWLAB_WORKBENCH_FACTS_QUERY_MAX_ATTEMPTS, DEFAULT_WORKBENCH_FACTS_QUERY_MAX_ATTEMPTS, 1, 10), + baseDelayMs: boundedInteger(env.HWLAB_WORKBENCH_FACTS_QUERY_BACKOFF_BASE_MS, DEFAULT_WORKBENCH_FACTS_QUERY_BACKOFF_BASE_MS, 0, 5000), + maxDelayMs: boundedInteger(env.HWLAB_WORKBENCH_FACTS_QUERY_BACKOFF_MAX_MS, DEFAULT_WORKBENCH_FACTS_QUERY_BACKOFF_MAX_MS, 1, 10000) + }; +} + +function workbenchFactsRetryDelayMs(policy, failedAttempt) { + const delay = policy.baseDelayMs * Math.pow(2, Math.max(0, failedAttempt - 1)); + return Math.max(0, Math.min(policy.maxDelayMs, Math.trunc(delay))); +} + +function isRetryableWorkbenchFactsQueryError(error) { + const message = String(error?.message ?? error ?? ""); + const code = String(error?.code ?? "").toLowerCase(); + return /timeout|connection terminated|terminating connection|connection closed|connection ended|econnreset|econnrefused|etimedout/iu.test(message) + || /timeout|econnreset|econnrefused|etimedout|57p01|08006|08003|53300/u.test(code); +} + +function annotateWorkbenchFactsQueryError(error, metadata) { + const normalized = error instanceof Error ? error : new Error(String(error ?? "Workbench facts query failed.")); + const data = normalized.data && typeof normalized.data === "object" && !Array.isArray(normalized.data) ? normalized.data : {}; + normalized.data = { + ...data, + retryable: metadata.retryable, + transient: metadata.retryable === true, + retryAttempt: metadata.retryAttempt, + retryAttempts: metadata.retryAttempt, + retryMax: metadata.retryMax, + retryAfterMs: metadata.retryAfterMs, + retryDelaysMs: metadata.retryDelaysMs.slice(), + retryExhausted: metadata.retryExhausted === true, + valuesRedacted: true + }; + return normalized; +} + +function workbenchFactsRetryMetadata(error) { + const data = error?.data && typeof error.data === "object" && !Array.isArray(error.data) ? error.data : {}; + return { + retryAttempt: finiteNonNegativeNumber(data.retryAttempt), + retryMax: finiteNonNegativeNumber(data.retryMax), + retryAfterMs: finiteNonNegativeNumber(data.retryAfterMs), + retryExhausted: data.retryExhausted === true + }; +} + +function logWorkbenchFactsQueryRetry(logger, params, error, retry) { + try { + const line = JSON.stringify({ + event: "workbench_facts_query_retry", + ok: false, + attempt: retry.attempt, + maxAttempts: retry.maxAttempts, + retryAfterMs: retry.retryAfterMs, + traceId: safeTraceId(params.traceId) ?? null, + sessionId: safeSessionId(params.sessionId) ?? null, + errorCode: error?.code ?? null, + message: error instanceof Error ? error.message : String(error ?? "unknown"), + valuesRedacted: true + }); + if (typeof logger?.warn === "function") logger.warn(line); + } catch { + // Retry logging must not affect read-model responses. + } +} + +function sleepMs(ms) { + return new Promise((resolve) => setTimeout(resolve, Math.max(0, Math.trunc(Number(ms) || 0)))); +} + +function boundedInteger(value, fallback, min, max) { + const parsed = Number(value); + if (!Number.isFinite(parsed)) return fallback; + return Math.max(min, Math.min(max, Math.trunc(parsed))); +} + +function finiteNonNegativeNumber(value) { + const parsed = Number(value); + return Number.isFinite(parsed) && parsed >= 0 ? Math.trunc(parsed) : null; +} + function isProjectionDiagnosticTrace(snapshot) { return Boolean(snapshot?.projection?.blocker || snapshot?.projectionHealth || snapshot?.blocker); }