fix: retry transient workbench facts reads
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user