Merge pull request #2694 from pikasTech/fix/workbench-long-retry-visibility
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
Pipelines as Code CI / hwlab-nc01-v03-ci-poll- Success
修复 Workbench 长程重试无可见进度
This commit is contained in:
@@ -216,6 +216,64 @@ test("projects AgentRun semantic infrastructure retry without reclassification",
|
||||
assert.equal(projected.event.runnerJobId, "rjob_retry");
|
||||
});
|
||||
|
||||
test("normalizes sparse AgentRun provider retry metadata for Workbench visibility", () => {
|
||||
const projected = projectAgentRunKafkaEventToHwlabEvent({
|
||||
schema: "agentrun.event.v1",
|
||||
eventType: "agentrun.run.event",
|
||||
producedAt: "2026-07-20T19:32:07.000Z",
|
||||
run: { runId: "run_provider_retry", sessionId: "ses_provider_retry", status: "running" },
|
||||
command: { commandId: "cmd_provider_retry", state: "running" },
|
||||
event: {
|
||||
id: "evt_provider_retry",
|
||||
seq: 52,
|
||||
type: "error",
|
||||
createdAt: "2026-07-20T19:32:07.000Z",
|
||||
payload: {
|
||||
failureKind: "provider-stream-disconnected",
|
||||
message: "Provider stream disconnected",
|
||||
willRetry: true,
|
||||
retryAttempt: 1,
|
||||
retryMax: 5,
|
||||
retryDelayMs: 240_000
|
||||
}
|
||||
}
|
||||
}, { source: "hwlab-test" });
|
||||
|
||||
assert.equal(projected.event.status, "retrying");
|
||||
assert.equal(projected.event.failureDomain, "upstream");
|
||||
assert.equal(projected.event.retryPhase, "retryScheduled");
|
||||
assert.equal(projected.event.retryAttempt, 1);
|
||||
assert.equal(projected.event.retryMaxAttempts, 5);
|
||||
assert.equal(projected.event.retryBackoffMs, 240_000);
|
||||
assert.equal(projected.event.nextRetryAt, "2026-07-20T19:36:07.000Z");
|
||||
assert.equal(projected.event.terminal, false);
|
||||
});
|
||||
|
||||
test("keeps sparse non-retryable provider errors failed", () => {
|
||||
const projected = projectAgentRunKafkaEventToHwlabEvent({
|
||||
schema: "agentrun.event.v1",
|
||||
eventType: "agentrun.run.event",
|
||||
producedAt: "2026-07-20T19:32:07.000Z",
|
||||
run: { runId: "run_provider_failed", sessionId: "ses_provider_failed", status: "failed" },
|
||||
command: { commandId: "cmd_provider_failed", state: "failed" },
|
||||
event: {
|
||||
id: "evt_provider_failed",
|
||||
seq: 53,
|
||||
type: "error",
|
||||
createdAt: "2026-07-20T19:32:07.000Z",
|
||||
payload: {
|
||||
failureKind: "provider-unavailable",
|
||||
message: "Provider unavailable",
|
||||
willRetry: false
|
||||
}
|
||||
}
|
||||
}, { source: "hwlab-test" });
|
||||
|
||||
assert.equal(projected.event.status, "failed");
|
||||
assert.equal(projected.event.failureDomain, "upstream");
|
||||
assert.equal(projected.event.retryPhase, null);
|
||||
});
|
||||
|
||||
test("projects AgentRun Git mirror fetch as running preparation", () => {
|
||||
const projected = projectAgentRunKafkaEventToHwlabEvent({
|
||||
schema: "agentrun.event.v1",
|
||||
|
||||
@@ -1610,7 +1610,7 @@ function mapAgentRunSourceEventToHwlabEvent({ input, sourceEvent, payload, run,
|
||||
}
|
||||
if (type === "backend_status") {
|
||||
const phase = firstText(payload.phase) || "status";
|
||||
const semanticRetry = semanticFailureRetryProjection(payload);
|
||||
const semanticRetry = semanticFailureRetryProjection(payload, sourceEvent.createdAt);
|
||||
return {
|
||||
...base,
|
||||
type: "backend",
|
||||
@@ -1678,26 +1678,31 @@ function mapAgentRunSourceEventToHwlabEvent({ input, sourceEvent, payload, run,
|
||||
}
|
||||
if (type === "error") {
|
||||
const failureKind = firstText(payload.failureKind, payload.errorCode, sourceEvent.failureKind, sourceEvent.errorCode) || "backend";
|
||||
const semanticRetry = semanticFailureRetryProjection(payload);
|
||||
const semanticRetry = semanticFailureRetryProjection(payload, sourceEvent.createdAt);
|
||||
return { ...base, type: "error", eventType: "error", status: semanticRetry.status ?? "failed", label: `agentrun:error:${failureKind}`, errorCode: firstText(payload.code, failureKind), failureKind, message: textPayload(payload, "AgentRun error"), ...semanticRetry.fields, terminal: false };
|
||||
}
|
||||
return { ...base, type: "backend", eventType: "backend", status: "running", label: `agentrun:event:${type}`, message: textPayload(payload, firstText(input.eventType, type) || "AgentRun event") };
|
||||
}
|
||||
|
||||
function semanticFailureRetryProjection(payload) {
|
||||
const retryPhase = firstText(payload.retryPhase);
|
||||
const failureDomain = firstText(payload.failureDomain);
|
||||
function semanticFailureRetryProjection(payload, eventCreatedAt = null) {
|
||||
const failureKind = firstText(payload.failureKind, payload.errorCode);
|
||||
const willRetry = payload.willRetry === true;
|
||||
const retryPhase = firstText(payload.retryPhase)
|
||||
|| (payload.retryExhausted === true ? "retryExhausted" : willRetry ? "retryScheduled" : null);
|
||||
const failureDomain = firstText(payload.failureDomain) || inferredFailureDomain(failureKind);
|
||||
const component = firstText(payload.component);
|
||||
const code = firstText(payload.code);
|
||||
const code = firstText(payload.code, failureKind);
|
||||
if (!retryPhase && !failureDomain && !component && !code) return { status: null, fields: {} };
|
||||
const normalized = String(retryPhase ?? "").toLowerCase();
|
||||
const status = normalized.includes("exhausted")
|
||||
? "failed"
|
||||
: normalized.includes("recovered")
|
||||
? "running"
|
||||
: normalized.includes("scheduled") || (normalized.includes("retry") && normalized.includes("started"))
|
||||
? "retrying"
|
||||
: "running";
|
||||
const status = !normalized
|
||||
? null
|
||||
: normalized.includes("exhausted")
|
||||
? "failed"
|
||||
: normalized.includes("recovered")
|
||||
? "running"
|
||||
: normalized.includes("scheduled") || (normalized.includes("retry") && normalized.includes("started"))
|
||||
? "retrying"
|
||||
: "running";
|
||||
return {
|
||||
status,
|
||||
fields: {
|
||||
@@ -1708,11 +1713,11 @@ function semanticFailureRetryProjection(payload) {
|
||||
failureCode: code,
|
||||
code,
|
||||
summary: firstText(payload.summary, payload.message),
|
||||
retryable: explicitBoolean(payload.retryable),
|
||||
retryable: explicitBoolean(payload.retryable) ?? (willRetry ? true : null),
|
||||
retryAttempt: integerValue(payload.attempt ?? payload.retryAttempt),
|
||||
retryMaxAttempts: integerValue(payload.maxAttempts ?? payload.retryMaxAttempts),
|
||||
retryMaxAttempts: integerValue(payload.maxAttempts ?? payload.retryMaxAttempts ?? payload.retryMax),
|
||||
retryBackoffMs: integerValue(payload.backoffMs ?? payload.retryDelayMs),
|
||||
nextRetryAt: timestampValue(payload.nextRetryAt),
|
||||
nextRetryAt: retryTimestamp(payload.nextRetryAt, eventCreatedAt, payload.backoffMs ?? payload.retryDelayMs),
|
||||
firstObservedAt: timestampValue(payload.firstObservedAt),
|
||||
observedAt: timestampValue(payload.observedAt),
|
||||
runnerJobId: firstText(payload.runnerJobId),
|
||||
@@ -1721,6 +1726,22 @@ function semanticFailureRetryProjection(payload) {
|
||||
};
|
||||
}
|
||||
|
||||
function inferredFailureDomain(failureKind) {
|
||||
const normalized = String(failureKind ?? "").trim().toLowerCase().replace(/_/gu, "-");
|
||||
if (!normalized) return null;
|
||||
if (normalized.startsWith("provider-") || normalized.includes("upstream")) return "upstream";
|
||||
return "infrastructure";
|
||||
}
|
||||
|
||||
function retryTimestamp(explicitValue, eventCreatedAt, retryDelayMs) {
|
||||
const explicit = timestampValue(explicitValue);
|
||||
if (explicit) return explicit;
|
||||
const createdAt = timestampValue(eventCreatedAt);
|
||||
const delayMs = integerValue(retryDelayMs);
|
||||
if (!createdAt || delayMs === null || delayMs < 0) return null;
|
||||
return new Date(Date.parse(createdAt) + delayMs).toISOString();
|
||||
}
|
||||
|
||||
function eventMatchesFilters(value, filters) {
|
||||
return Object.values(eventFilterResults(value, filters)).every(Boolean);
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import { workbenchTraceTimelinePolicy } from "@/config/runtime";
|
||||
import { messageDiagnosticView } from "@/utils/workbench-error-runtime";
|
||||
import MessageMarkdown from "./MessageMarkdown.vue";
|
||||
import { traceLifecycleExpanded } from "./trace-lifecycle";
|
||||
import { semanticRuntimeStatusView } from "./workbench-semantic-runtime-status";
|
||||
|
||||
const props = withDefaults(defineProps<{
|
||||
message: ChatMessage;
|
||||
@@ -123,58 +124,6 @@ function traceStorageKey(message: ChatMessage): string {
|
||||
return `${props.storageKeyPrefix}.${message.traceId ?? message.runnerTrace?.traceId ?? message.id}`;
|
||||
}
|
||||
|
||||
function semanticRuntimeStatusView(message: ChatMessage, currentTimeMs: number): { title: string; detail: string; tone: string; emphasis: "quiet" | "alert"; code: string | null } | null {
|
||||
const events = Array.isArray(message.runnerTrace?.events) ? message.runnerTrace.events : [];
|
||||
const event = [...events].reverse().find((item) => typeof item.failureDomain === "string" && typeof item.retryPhase === "string");
|
||||
if (!event) return null;
|
||||
const phase = String(event.retryPhase).toLowerCase().replace(/[^a-z]/gu, "");
|
||||
const domain = event.failureDomain === "upstream" ? "上游故障" : "基础设施故障";
|
||||
const summary = firstText(event.summary, event.message, event.failureCode, event.code, "执行依赖暂时不可用");
|
||||
const attempt = finiteInteger(event.retryAttempt ?? event.attempt);
|
||||
const maxAttempts = finiteInteger(event.retryMaxAttempts ?? event.maxAttempts);
|
||||
const progress = attempt !== null && maxAttempts !== null ? `第 ${attempt}/${maxAttempts} 次` : "本次";
|
||||
const code = firstText(event.failureCode, event.code);
|
||||
if (code === "git-mirror-fetch-in-progress") {
|
||||
return { title: `运行环境准备:${summary}`, detail: `${progress}源码获取进行中`, tone: "info", emphasis: "quiet", code };
|
||||
}
|
||||
if (code === "git-mirror-fetch-completed") {
|
||||
return { title: `运行环境准备:${summary}`, detail: "源码获取完成,正在启动 Runner", tone: "recovered", emphasis: "quiet", code };
|
||||
}
|
||||
let detail = "正在判定有限重试";
|
||||
let tone = "warning";
|
||||
if (phase.includes("scheduled")) {
|
||||
const nextRetryMs = timestampMs(event.nextRetryAt);
|
||||
const remainingSeconds = nextRetryMs === null ? null : Math.max(0, Math.ceil((nextRetryMs - currentTimeMs) / 1_000));
|
||||
detail = `${progress}重试已安排${remainingSeconds === null ? "" : `,预计 ${remainingSeconds} 秒后开始`}`;
|
||||
} else if (phase.includes("started")) {
|
||||
detail = `${progress}重试已开始`;
|
||||
} else if (phase.includes("recovered")) {
|
||||
detail = "故障已恢复,继续执行";
|
||||
tone = "recovered";
|
||||
} else if (phase.includes("exhausted")) {
|
||||
detail = "有限重试已耗尽,执行已失败";
|
||||
tone = "failed";
|
||||
}
|
||||
return {
|
||||
title: `${domain}:${summary}`,
|
||||
detail,
|
||||
tone,
|
||||
emphasis: "alert",
|
||||
code
|
||||
};
|
||||
}
|
||||
|
||||
function firstText(...values: unknown[]): string {
|
||||
for (const value of values) {
|
||||
if (typeof value === "string" && value.trim()) return value.trim();
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
function finiteInteger(value: unknown): number | null {
|
||||
const number = Number(value);
|
||||
return Number.isSafeInteger(number) && number >= 0 ? number : null;
|
||||
}
|
||||
</script>
|
||||
|
||||
<template>
|
||||
|
||||
+43
@@ -0,0 +1,43 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import type { ChatMessage } from "@/types";
|
||||
import { semanticRuntimeStatusView } from "./workbench-semantic-runtime-status";
|
||||
|
||||
test("long provider retry remains semantically visible while Kafka terminal authority stays open", () => {
|
||||
const message = {
|
||||
id: "msg_long_retry",
|
||||
role: "agent",
|
||||
title: "Code Agent",
|
||||
text: "",
|
||||
status: "running",
|
||||
createdAt: "2026-07-20T19:26:00.000Z",
|
||||
timing: {
|
||||
startedAt: "2026-07-20T19:26:00.000Z",
|
||||
lastEventAt: "2026-07-20T19:37:49.000Z",
|
||||
finishedAt: null,
|
||||
durationMs: null
|
||||
},
|
||||
runnerTrace: {
|
||||
events: [{
|
||||
type: "error",
|
||||
failureDomain: "upstream",
|
||||
failureCode: "provider-stream-disconnected",
|
||||
summary: "Provider stream disconnected",
|
||||
retryPhase: "retryScheduled",
|
||||
retryAttempt: 2,
|
||||
retryMaxAttempts: 5,
|
||||
nextRetryAt: "2026-07-20T19:41:49.000Z",
|
||||
createdAt: "2026-07-20T19:36:16.000Z"
|
||||
}]
|
||||
}
|
||||
} as ChatMessage;
|
||||
|
||||
const view = semanticRuntimeStatusView(message, Date.parse("2026-07-20T19:43:13.000Z"));
|
||||
|
||||
assert.equal(view?.title, "上游故障:Provider stream disconnected");
|
||||
assert.equal(view?.detail, "第 2/5 次重试等待已持续 5 分 24 秒,仍未收到新事件");
|
||||
assert.equal(view?.emphasis, "alert");
|
||||
assert.equal(view?.code, "provider-stream-disconnected");
|
||||
assert.equal(message.status, "running", "view-only visibility must not synthesize a terminal state");
|
||||
});
|
||||
@@ -0,0 +1,89 @@
|
||||
import type { ChatMessage, TraceEvent } from "@/types";
|
||||
|
||||
export interface WorkbenchSemanticRuntimeStatus {
|
||||
title: string;
|
||||
detail: string;
|
||||
tone: string;
|
||||
emphasis: "quiet" | "alert";
|
||||
code: string | null;
|
||||
}
|
||||
|
||||
export function semanticRuntimeStatusView(message: ChatMessage, currentTimeMs: number): WorkbenchSemanticRuntimeStatus | null {
|
||||
const events = Array.isArray(message.runnerTrace?.events) ? message.runnerTrace.events : [];
|
||||
const event = [...events].reverse().find(isSemanticRuntimeEvent);
|
||||
if (!event) return null;
|
||||
const phase = String(event.retryPhase).toLowerCase().replace(/[^a-z]/gu, "");
|
||||
const domain = event.failureDomain === "upstream" ? "上游故障" : "基础设施故障";
|
||||
const summary = firstText(event.summary, event.message, event.failureCode, event.code, "执行依赖暂时不可用");
|
||||
const attempt = finiteInteger(event.retryAttempt ?? event.attempt);
|
||||
const maxAttempts = finiteInteger(event.retryMaxAttempts ?? event.maxAttempts);
|
||||
const progress = attempt !== null && maxAttempts !== null ? `第 ${attempt}/${maxAttempts} 次` : "本次";
|
||||
const code = firstText(event.failureCode, event.code) || null;
|
||||
if (code === "git-mirror-fetch-in-progress") {
|
||||
return { title: `运行环境准备:${summary}`, detail: `${progress}源码获取进行中`, tone: "info", emphasis: "quiet", code };
|
||||
}
|
||||
if (code === "git-mirror-fetch-completed") {
|
||||
return { title: `运行环境准备:${summary}`, detail: "源码获取完成,正在启动 Runner", tone: "recovered", emphasis: "quiet", code };
|
||||
}
|
||||
let detail = "正在判定有限重试";
|
||||
let tone = "warning";
|
||||
if (phase.includes("scheduled")) {
|
||||
detail = scheduledRetryDetail(message, event, currentTimeMs, progress);
|
||||
} else if (phase.includes("started")) {
|
||||
detail = `${progress}重试已开始`;
|
||||
} else if (phase.includes("recovered")) {
|
||||
detail = "故障已恢复,继续执行";
|
||||
tone = "recovered";
|
||||
} else if (phase.includes("exhausted")) {
|
||||
detail = "有限重试已耗尽,执行已失败";
|
||||
tone = "failed";
|
||||
}
|
||||
return { title: `${domain}:${summary}`, detail, tone, emphasis: "alert", code };
|
||||
}
|
||||
|
||||
function scheduledRetryDetail(message: ChatMessage, event: TraceEvent, currentTimeMs: number, progress: string): string {
|
||||
const nextRetryMs = timestampMs(event.nextRetryAt);
|
||||
if (nextRetryMs !== null && nextRetryMs > currentTimeMs) {
|
||||
return `${progress}重试已安排,预计 ${Math.max(0, Math.ceil((nextRetryMs - currentTimeMs) / 1_000))} 秒后开始`;
|
||||
}
|
||||
if (nextRetryMs !== null) {
|
||||
const lastEventMs = timestampMs(message.timing?.lastEventAt ?? message.lastEventAt ?? event.createdAt);
|
||||
if (lastEventMs !== null) {
|
||||
return `${progress}重试等待已持续 ${formatDuration(Math.max(0, currentTimeMs - lastEventMs))},仍未收到新事件`;
|
||||
}
|
||||
}
|
||||
return `${progress}重试已安排,等待上游恢复`;
|
||||
}
|
||||
|
||||
function isSemanticRuntimeEvent(event: TraceEvent): boolean {
|
||||
return typeof event.failureDomain === "string" && typeof event.retryPhase === "string";
|
||||
}
|
||||
|
||||
function timestampMs(value: unknown): number | null {
|
||||
if (typeof value !== "string" || !value.trim()) return null;
|
||||
const parsed = Date.parse(value);
|
||||
return Number.isFinite(parsed) ? parsed : null;
|
||||
}
|
||||
|
||||
function finiteInteger(value: unknown): number | null {
|
||||
const number = Number(value);
|
||||
return Number.isSafeInteger(number) && number >= 0 ? number : null;
|
||||
}
|
||||
|
||||
function firstText(...values: unknown[]): string {
|
||||
for (const value of values) {
|
||||
if (typeof value === "string" && value.trim()) return value.trim();
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
function formatDuration(ms: number): string {
|
||||
const seconds = ms > 0 && ms < 1000 ? 1 : Math.max(0, Math.floor(ms / 1000));
|
||||
if (seconds < 60) return `${seconds} 秒`;
|
||||
const minutes = Math.floor(seconds / 60);
|
||||
if (minutes < 60) return `${minutes} 分 ${seconds % 60} 秒`;
|
||||
const hours = Math.floor(minutes / 60);
|
||||
if (hours < 24) return `${hours} 小时 ${minutes % 60} 分`;
|
||||
const days = Math.floor(hours / 24);
|
||||
return `${days} 天 ${hours % 24} 小时`;
|
||||
}
|
||||
Reference in New Issue
Block a user