fix: recover repeated model requests without progress

This commit is contained in:
joshavant
2026-08-04 16:46:44 -05:00
committed by Josh Avant
parent a8f314d587
commit 517d0b1c2c
16 changed files with 698 additions and 72 deletions
@@ -525,6 +525,11 @@ describe("wrapStreamFnWithDiagnosticModelCallEvents", () => {
expect(snapshot.lastProgressReason).toBe("model_call:stream_progress");
expect(snapshot.lastProgressAgeMs).toBe(0);
expect(runProgressEvents).toHaveLength(2);
expect(
runProgressEvents.every(
(event) => event.type === "run.progress" && event.progressKind === "liveness",
),
).toBe(true);
} finally {
await iterator.return?.();
await waitForDiagnosticEventsDrained();
@@ -323,6 +323,7 @@ function maybeEmitModelCallStreamProgress(
...(eventBase.sessionKey ? { sessionKey: eventBase.sessionKey } : {}),
...(eventBase.sessionId ? { sessionId: eventBase.sessionId } : {}),
reason: MODEL_CALL_STREAM_PROGRESS_REASON,
progressKind: "liveness" as const,
};
markDiagnosticRunProgress(progressFields);
if (
@@ -357,6 +357,7 @@ export async function runEmbeddedFallbackCandidate(params: {
messageToolDeliveryState: params.messageToolDeliveryState,
provider: params.provider,
model: params.model,
runId: params.runId,
effectiveSessionId: params.effectiveRun.sessionId,
notifyUserAboutCompaction: params.notifyUserAboutCompaction,
onCompactionCompleted: () => {
@@ -3,6 +3,7 @@ import { isMessagingToolSendAction } from "../../agents/embedded-agent-messaging
import type { RunEmbeddedAgentParams } from "../../agents/embedded-agent-runner/run/params.js";
import { normalizeAgentPlanSteps } from "../../channels/streaming.js";
import { logVerbose } from "../../globals.js";
import { markDiagnosticRunProgress } from "../../logging/diagnostic-run-activity.js";
import { createSubsystemLogger } from "../../logging/subsystem.js";
import type { ReplyPayload } from "../types.js";
import type { AgentLifecycleTerminalBackstop } from "./agent-lifecycle-terminal.js";
@@ -35,6 +36,7 @@ export function createAgentRunEventHandler(params: {
sourceRepliesAreToolOnly: boolean;
provider: string;
model: string;
runId: string;
effectiveSessionId?: string;
notifyUserAboutCompaction: boolean;
onCompactionCompleted: () => number;
@@ -83,6 +85,17 @@ export function createAgentRunEventHandler(params: {
return async (evt) => {
params.turn.replyOperation?.recordActivity();
// Agent outputs are portable forward-progress facts. Usage and lifecycle
// bookkeeping stay mechanical so repeated model attempts cannot self-refresh.
if (evt.stream !== "usage" && evt.stream !== "lifecycle") {
markDiagnosticRunProgress({
runId: params.runId,
sessionId: params.effectiveSessionId,
sessionKey: params.turn.sessionKey,
reason: `agent_event:${evt.stream}`,
progressKind: "semantic",
});
}
params.lifecycleBackstop.note(evt);
const hasLifecyclePhase = evt.stream === "lifecycle" && typeof evt.data.phase === "string";
if (evt.stream !== "lifecycle" || hasLifecyclePhase) {
@@ -1,5 +1,11 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { createDraftStreamLoop } from "../../channels/draft-stream-loop.js";
import {
getDiagnosticSessionActivitySnapshot,
markDiagnosticEmbeddedRunStarted,
resetDiagnosticRunActivityForTest,
} from "../../logging/diagnostic-run-activity.js";
import { markDiagnosticModelStartedForTest } from "../../logging/diagnostic-run-activity.test-support.js";
import type { PartialReplyPayload } from "../get-reply-options.types.js";
import type { GetReplyOptions } from "../types.js";
import {
@@ -39,6 +45,7 @@ const state = setupAgentRunnerExecutionTestState();
beforeEach(() => {
sanitizerState.sanitizeUserFacingText.mockClear();
resetDiagnosticRunActivityForTest();
});
async function executeTestTurn(
@@ -50,6 +57,40 @@ async function executeTestTurn(
}
describe("executeAgentTurn: lifecycle progress", () => {
it("records assistant events as semantic run progress", async () => {
state.runEmbeddedAgentMock.mockImplementationOnce(async (params: EmbeddedAgentParams) => {
const sessionId = params.sessionId ?? "session";
const sessionKey = params.sessionKey ?? "main";
markDiagnosticEmbeddedRunStarted({ sessionId, sessionKey, runId: params.runId });
for (let attempt = 0; attempt < 2; attempt += 1) {
markDiagnosticModelStartedForTest({
sessionId,
sessionKey,
runId: params.runId,
provider: "mock",
model: "request-model",
observationUnit: "request",
});
}
expect(
getDiagnosticSessionActivitySnapshot({ sessionId, sessionKey })
.repeatedRequestNoProgressAgeMs,
).toBe(0);
await params.onAgentEvent?.({
stream: "assistant",
data: { phase: "commentary", text: "Working" },
});
expect(
getDiagnosticSessionActivitySnapshot({ sessionId, sessionKey })
.repeatedRequestNoProgressAgeMs,
).toBeUndefined();
return { payloads: [{ text: "final" }], meta: {} };
});
await executeTestTurn();
});
it("forwards item lifecycle events to reply options", async () => {
const onItemEvent = vi.fn();
state.runEmbeddedAgentMock.mockImplementationOnce(async (params: EmbeddedAgentParams) => {
@@ -361,6 +361,9 @@ export type FallbackRunnerParams = {
};
export type EmbeddedAgentParams = {
runId: string;
sessionId?: string;
sessionKey?: string;
prompt?: string;
transcriptPrompt?: string;
lifecycleGeneration?: string;
+3
View File
@@ -284,6 +284,7 @@ type DiagnosticSessionAttentionBaseEvent = DiagnosticBaseEvent & {
activeToolName?: string;
activeToolCallId?: string;
activeToolAgeMs?: number;
repeatedRequestNoProgressAgeMs?: number;
terminalProgressStale?: boolean;
};
@@ -371,6 +372,8 @@ export type DiagnosticRunProgressEvent = DiagnosticBaseEvent & {
sessionId?: string;
runId?: string;
reason: string;
/** Semantic progress resets no-forward-progress evidence; liveness only keeps work alive. */
progressKind?: "semantic" | "liveness";
};
/**
@@ -0,0 +1,109 @@
// Mechanical request retries stay continuous for one run owner. Only typed
// semantic progress or owner teardown can clear the clock that recovery reads.
type RepeatedRequestOwner = { runId: string; sequence: number };
export type DiagnosticRepeatedRequestActivity = {
repeatedRequestOwnerRunId?: string;
repeatedRequestFirstStartedAt?: number;
repeatedRequestCount?: number;
repeatedRequestMutationSequence?: number;
};
let mutationSequence = 0;
function nextMutationSequence(): number {
mutationSequence += 1;
return mutationSequence;
}
function currentOwner(owners: Iterable<RepeatedRequestOwner>): RepeatedRequestOwner | undefined {
let current: RepeatedRequestOwner | undefined;
for (const owner of owners) {
if (!current || owner.sequence > current.sequence) {
current = owner;
}
}
return current;
}
export function recordRepeatedRequestObservation(
activity: DiagnosticRepeatedRequestActivity,
owners: Iterable<RepeatedRequestOwner>,
params: {
runId?: string;
observationUnit?: "request" | "turn";
now?: number;
},
): void {
if (params.observationUnit === "turn") {
return;
}
const owner = currentOwner(owners);
const runId = params.runId?.trim();
if (!owner || !runId || owner.runId !== runId) {
return;
}
if (activity.repeatedRequestOwnerRunId !== runId) {
activity.repeatedRequestOwnerRunId = runId;
activity.repeatedRequestFirstStartedAt = params.now ?? Date.now();
activity.repeatedRequestCount = 1;
} else {
activity.repeatedRequestCount = (activity.repeatedRequestCount ?? 0) + 1;
}
activity.repeatedRequestMutationSequence = nextMutationSequence();
}
export function clearRepeatedRequestActivity(
activity: DiagnosticRepeatedRequestActivity,
params: { runId?: string } = {},
): boolean {
if (
params.runId !== undefined &&
activity.repeatedRequestOwnerRunId !== undefined &&
activity.repeatedRequestOwnerRunId !== params.runId
) {
return false;
}
const cleared = activity.repeatedRequestCount !== undefined;
if (!cleared && params.runId !== undefined) {
return false;
}
activity.repeatedRequestOwnerRunId = undefined;
activity.repeatedRequestFirstStartedAt = undefined;
activity.repeatedRequestCount = undefined;
activity.repeatedRequestMutationSequence = nextMutationSequence();
return cleared;
}
export function mergeRepeatedRequestActivity(
target: DiagnosticRepeatedRequestActivity,
source: DiagnosticRepeatedRequestActivity,
): void {
if (
source.repeatedRequestMutationSequence === undefined ||
(target.repeatedRequestMutationSequence ?? 0) >= source.repeatedRequestMutationSequence
) {
return;
}
target.repeatedRequestOwnerRunId = source.repeatedRequestOwnerRunId;
target.repeatedRequestFirstStartedAt = source.repeatedRequestFirstStartedAt;
target.repeatedRequestCount = source.repeatedRequestCount;
target.repeatedRequestMutationSequence = source.repeatedRequestMutationSequence;
}
export function resolveRepeatedRequestNoProgressAgeMs(
activity: DiagnosticRepeatedRequestActivity,
owners: Iterable<RepeatedRequestOwner>,
now: number,
): number | undefined {
const owner = currentOwner(owners);
if (
!owner ||
owner.runId !== activity.repeatedRequestOwnerRunId ||
(activity.repeatedRequestCount ?? 0) < 2 ||
activity.repeatedRequestFirstStartedAt === undefined
) {
return undefined;
}
return Math.max(0, now - activity.repeatedRequestFirstStartedAt);
}
@@ -0,0 +1,69 @@
import type { DiagnosticSessionActiveWorkKind } from "../infra/diagnostic-events.js";
import {
type DiagnosticArgumentChurnActivity,
resolveArgumentChurnProgress,
} from "./diagnostic-argument-churn-activity.js";
import {
type DiagnosticRepeatedRequestActivity,
resolveRepeatedRequestNoProgressAgeMs,
} from "./diagnostic-repeated-request-activity.js";
export type DiagnosticSessionActivitySnapshot = {
activeWorkKind?: DiagnosticSessionActiveWorkKind;
hasActiveEmbeddedRun?: boolean;
activeToolName?: string;
activeToolCallId?: string;
activeToolAgeMs?: number;
lastProgressAgeMs?: number;
lastProgressReason?: string;
repeatedRequestNoProgressAgeMs?: number;
};
type SnapshotTool = { toolName: string; toolCallId?: string; startedAt: number };
type SnapshotActivity = DiagnosticArgumentChurnActivity &
DiagnosticRepeatedRequestActivity & {
activeEmbeddedRuns: ReadonlyMap<string, { runId: string; sequence: number }>;
activeModelCalls: ReadonlyMap<string, unknown>;
activeTools: ReadonlyMap<string, SnapshotTool>;
lastProgressAt: number;
lastProgressReason?: string;
};
export function buildDiagnosticSessionActivitySnapshot(
activity: SnapshotActivity,
now: number,
): DiagnosticSessionActivitySnapshot {
const activeWorkKind: DiagnosticSessionActiveWorkKind | undefined =
activity.activeTools.size > 0
? "tool_call"
: activity.activeModelCalls.size > 0
? "model_call"
: activity.activeEmbeddedRuns.size > 0
? "embedded_run"
: undefined;
let activeTool: SnapshotTool | undefined;
for (const tool of activity.activeTools.values()) {
if (!activeTool || tool.startedAt < activeTool.startedAt) {
activeTool = tool;
}
}
const churnProgress = resolveArgumentChurnProgress(
activity,
activity.activeEmbeddedRuns.values(),
now,
);
return {
activeWorkKind,
...(activity.activeEmbeddedRuns.size > 0 ? { hasActiveEmbeddedRun: true } : {}),
activeToolName: activeTool?.toolName,
activeToolCallId: activeTool?.toolCallId,
activeToolAgeMs: activeTool ? Math.max(0, now - activeTool.startedAt) : undefined,
lastProgressAgeMs: Math.max(0, now - churnProgress.lastProgressAt),
lastProgressReason: churnProgress.lastProgressReason,
repeatedRequestNoProgressAgeMs: resolveRepeatedRequestNoProgressAgeMs(
activity,
activity.activeEmbeddedRuns.values(),
now,
),
};
}
@@ -3,12 +3,12 @@ import "./diagnostic-run-activity.js";
type DiagnosticModelStartedActivityEvent = Pick<
Extract<DiagnosticEventPayload, { type: "model.call.started" }>,
"runId" | "sessionId" | "sessionKey" | "provider" | "model"
"runId" | "sessionId" | "sessionKey" | "provider" | "model" | "observationUnit"
> & { seq?: number };
type DiagnosticRunProgressActivityEvent = Pick<
Extract<DiagnosticEventPayload, { type: "run.progress" }>,
"runId" | "sessionId" | "sessionKey" | "reason"
"runId" | "sessionId" | "sessionKey" | "reason" | "progressKind"
>;
type DiagnosticRunActivityTestApi = {
+200 -1
View File
@@ -15,15 +15,17 @@ import {
markDiagnosticEmbeddedRunEnded,
markDiagnosticEmbeddedRunStarted,
markDiagnosticRunProgress,
resetDiagnosticRunActivityForTest,
resolveRunStaleThresholdMs,
RUN_STALE_TAKEOVER_MS,
startDiagnosticRunActivityTracking,
stopDiagnosticRunActivityTracking,
} from "./diagnostic-run-activity.js";
import { markDiagnosticModelStartedForTest } from "./diagnostic-run-activity.test-support.js";
afterEach(() => {
vi.useRealTimers();
stopDiagnosticRunActivityTracking();
resetDiagnosticRunActivityForTest();
resetDiagnosticEventsForTest();
});
@@ -369,6 +371,203 @@ describe("argument-churn liveness", () => {
});
});
describe("repeated request liveness", () => {
it("ages repeated requests across mechanical progress until semantic progress arrives", () => {
vi.useFakeTimers();
const startedAt = Date.parse("2026-08-04T00:00:00Z");
vi.setSystemTime(startedAt);
const ref = { sessionId: "retry-session", sessionKey: "agent:main:retry" };
const runId = "retry-run";
markDiagnosticEmbeddedRunStarted({ ...ref, runId });
markDiagnosticModelStartedForTest({
...ref,
runId,
provider: "mock",
model: "retrying-model",
observationUnit: "request",
});
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
for (let attempt = 2; attempt <= 11; attempt += 1) {
vi.setSystemTime(startedAt + (attempt - 1) * 30_000);
markDiagnosticModelStartedForTest({
...ref,
runId,
provider: "mock",
model: "retrying-model",
observationUnit: "request",
});
markDiagnosticRunProgress({
...ref,
runId,
reason: "model_call:stream_progress",
progressKind: "liveness",
});
}
expect(getDiagnosticSessionActivitySnapshot(ref)).toMatchObject({
activeWorkKind: "model_call",
lastProgressAgeMs: 0,
repeatedRequestNoProgressAgeMs: 5 * 60_000,
});
markDiagnosticRunProgress({
...ref,
runId,
reason: "assistant:progress",
progressKind: "semantic",
});
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
});
it("ignores turn observations and clears request evidence across owner lifecycle", () => {
vi.useFakeTimers();
const startedAt = Date.parse("2026-08-04T01:00:00Z");
vi.setSystemTime(startedAt);
const ref = { sessionId: "owner-session", sessionKey: "agent:main:owner" };
markDiagnosticEmbeddedRunStarted({ ...ref, runId: "first-owner" });
for (let attempt = 0; attempt < 2; attempt += 1) {
markDiagnosticModelStartedForTest({
...ref,
runId: "first-owner",
provider: "cli",
model: "turn-model",
observationUnit: "turn",
});
}
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
for (let attempt = 0; attempt < 2; attempt += 1) {
markDiagnosticModelStartedForTest({
...ref,
runId: "first-owner",
provider: "mock",
model: "request-model",
observationUnit: "request",
});
}
vi.setSystemTime(startedAt + 6 * 60_000);
expect(getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs).toBe(
6 * 60_000,
);
markDiagnosticEmbeddedRunStarted({ ...ref, runId: "replacement-owner" });
markDiagnosticModelStartedForTest({
...ref,
runId: "first-owner",
provider: "mock",
model: "delayed-request",
observationUnit: "request",
});
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
markDiagnosticModelStartedForTest({
...ref,
runId: "replacement-owner",
provider: "mock",
model: "request-model",
observationUnit: "request",
});
markDiagnosticModelStartedForTest({
...ref,
runId: "replacement-owner",
provider: "mock",
model: "request-model",
observationUnit: "request",
});
vi.setSystemTime(startedAt + 7 * 60_000);
markDiagnosticRunProgress({
...ref,
runId: "first-owner",
reason: "delayed-old-owner-output",
progressKind: "semantic",
});
expect(getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs).toBe(60_000);
expect(
clearDiagnosticEmbeddedRunActivityForSession({
...ref,
activeSessionId: "replacement-owner",
}).cleared,
).toBe(true);
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
});
it("orders semantic progress across merged session aliases", () => {
vi.useFakeTimers();
const startedAt = Date.parse("2026-08-04T02:00:00Z");
vi.setSystemTime(startedAt);
const sessionId = "merge-session";
const sessionKey = "agent:main:merge";
const runId = "merge-run";
markDiagnosticEmbeddedRunStarted({ sessionId, runId });
markDiagnosticModelStartedForTest({
sessionId,
runId,
provider: "mock",
model: "request-model",
observationUnit: "request",
});
markDiagnosticRunProgress({
sessionKey,
reason: "reply:delivered",
progressKind: "semantic",
});
vi.setSystemTime(startedAt + 6 * 60_000);
expect(
getDiagnosticSessionActivitySnapshot({ sessionId, sessionKey })
.repeatedRequestNoProgressAgeMs,
).toBeUndefined();
});
it("clears repeated request evidence on run completion and listener restart", async () => {
const ref = { sessionId: "completed-session", sessionKey: "agent:main:completed" };
const runId = "completed-run";
startDiagnosticRunActivityTracking();
markDiagnosticEmbeddedRunStarted({ ...ref, runId });
for (let attempt = 0; attempt < 2; attempt += 1) {
markDiagnosticModelStartedForTest({
...ref,
runId,
provider: "mock",
model: "request-model",
observationUnit: "request",
});
}
expect(getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs).toBe(0);
emitTrustedDiagnosticEvent({
type: "run.completed",
...ref,
runId,
durationMs: 1,
outcome: "completed",
});
await waitForDiagnosticEventsDrained();
expect(
getDiagnosticSessionActivitySnapshot(ref).repeatedRequestNoProgressAgeMs,
).toBeUndefined();
stopDiagnosticRunActivityTracking();
startDiagnosticRunActivityTracking();
expect(getDiagnosticSessionActivitySnapshot(ref)).toEqual({});
});
});
describe("resolveRunStaleThresholdMs", () => {
it.each([
{
+67 -67
View File
@@ -3,7 +3,6 @@ import {
getInternalDiagnosticEventSequence,
onInternalDiagnosticEvent,
type DiagnosticEventPayload,
type DiagnosticSessionActiveWorkKind,
} from "../infra/diagnostic-events.js";
import {
applyArgumentChurnObservation,
@@ -13,20 +12,32 @@ import {
type DiagnosticArgumentChurnObservationParams,
mergeArgumentChurnActivity,
recordDiagnosticActivityProgress,
resolveArgumentChurnProgress,
} from "./diagnostic-argument-churn-activity.js";
import { createDiagnosticEmbeddedRunIndex } from "./diagnostic-embedded-run-index.js";
import {
clearRepeatedRequestActivity,
type DiagnosticRepeatedRequestActivity,
mergeRepeatedRequestActivity,
recordRepeatedRequestObservation,
} from "./diagnostic-repeated-request-activity.js";
import {
buildDiagnosticSessionActivitySnapshot,
type DiagnosticSessionActivitySnapshot,
} from "./diagnostic-run-activity-snapshot.js";
type SessionActivity = DiagnosticArgumentChurnActivity & {
sessionId?: string;
sessionKey?: string;
activeEmbeddedRuns: Map<string, ActiveEmbeddedRun>;
activeTools: Map<string, ActiveTool>;
activeModelCalls: Map<string, ActiveModelCall>;
recoveredOwnerStartEventCutoffs: Map<string, number>;
lastProgressAt: number;
lastProgressReason?: string;
};
export type { DiagnosticSessionActivitySnapshot } from "./diagnostic-run-activity-snapshot.js";
type SessionActivity = DiagnosticArgumentChurnActivity &
DiagnosticRepeatedRequestActivity & {
sessionId?: string;
sessionKey?: string;
activeEmbeddedRuns: Map<string, ActiveEmbeddedRun>;
activeTools: Map<string, ActiveTool>;
activeModelCalls: Map<string, ActiveModelCall>;
recoveredOwnerStartEventCutoffs: Map<string, number>;
lastProgressAt: number;
lastProgressReason?: string;
};
type ActiveEmbeddedRun = {
runId: string;
@@ -60,12 +71,12 @@ type DiagnosticToolStartedActivityEvent = Pick<
type DiagnosticModelStartedActivityEvent = Pick<
Extract<DiagnosticEventPayload, { type: "model.call.started" }>,
"runId" | "sessionId" | "sessionKey" | "provider" | "model"
"runId" | "sessionId" | "sessionKey" | "provider" | "model" | "observationUnit"
> & { seq?: number };
type DiagnosticRunProgressActivityEvent = Pick<
Extract<DiagnosticEventPayload, { type: "run.progress" }>,
"runId" | "sessionId" | "sessionKey" | "reason"
"runId" | "sessionId" | "sessionKey" | "reason" | "progressKind"
>;
// Quiet-but-alive tools are normal agent behavior; the CLI byte watchdog kills
@@ -77,16 +88,6 @@ export const BLOCKED_TOOL_CALL_ABORT_FLOOR_MS = 15 * 60_000;
// Default quiet-run reclaim window for steer/takeover. Evidence clocks stay local.
export const RUN_STALE_TAKEOVER_MS = 10 * 60_000;
export type DiagnosticSessionActivitySnapshot = {
activeWorkKind?: DiagnosticSessionActiveWorkKind;
hasActiveEmbeddedRun?: boolean;
activeToolName?: string;
activeToolCallId?: string;
activeToolAgeMs?: number;
lastProgressAgeMs?: number;
lastProgressReason?: string;
};
// Quiet-but-alive tool phases get the blocked-tool floor so a human message
// cannot reclaim a healthy long tool that stuck recovery would not touch yet.
export function resolveRunStaleThresholdMs(
@@ -175,6 +176,7 @@ function mergeSessionActivity(target: SessionActivity, source: SessionActivity):
target.lastProgressSequence = source.lastProgressSequence;
}
mergeArgumentChurnActivity(target, source);
mergeRepeatedRequestActivity(target, source);
replaceSessionActivityReferences(source, target);
}
@@ -232,6 +234,15 @@ function touchSessionActivity(activity: SessionActivity, reason: string, now = D
recordDiagnosticActivityProgress(activity);
}
function touchSemanticSessionActivity(
activity: SessionActivity,
reason: string,
params: { runId?: string; now?: number } = {},
): void {
clearRepeatedRequestActivity(activity, { runId: params.runId });
touchSessionActivity(activity, reason, params.now);
}
function toolKey(event: {
runId?: string;
sessionId?: string;
@@ -267,7 +278,10 @@ function recordToolStarted(event: DiagnosticToolStartedActivityEvent): void {
startedAt: now,
lastProgressAt: now,
});
touchSessionActivity(activity, `tool:${event.toolName}:started`, now);
touchSemanticSessionActivity(activity, `tool:${event.toolName}:started`, {
runId: event.runId,
now,
});
}
function recordToolEnded(
@@ -281,7 +295,7 @@ function recordToolEnded(
return;
}
activity.activeTools.delete(toolKey(event));
touchSessionActivity(activity, `tool:${event.toolName}:ended`);
touchSemanticSessionActivity(activity, `tool:${event.toolName}:ended`, { runId: event.runId });
}
function recordModelStarted(event: DiagnosticModelStartedActivityEvent): void {
@@ -292,6 +306,7 @@ function recordModelStarted(event: DiagnosticModelStartedActivityEvent): void {
if (shouldIgnoreRecoveredOwnerStartEvent(activity, event)) {
return;
}
recordRepeatedRequestObservation(activity, activity.activeEmbeddedRuns.values(), event);
activity.activeModelCalls.set(modelCallKey(event), {
runId: event.runId,
sessionId: event.sessionId,
@@ -330,7 +345,11 @@ export function markDiagnosticRunProgress(params: DiagnosticRunProgressActivityE
if (!activity) {
return;
}
touchSessionActivity(activity, params.reason);
if (params.progressKind === "liveness") {
touchSessionActivity(activity, params.reason);
return;
}
touchSemanticSessionActivity(activity, params.reason, { runId: params.runId });
}
function recordRunCompleted(
@@ -346,7 +365,7 @@ function recordRunCompleted(
embeddedRunIndex.clear(activity);
clearArgumentChurnActivity(activity, { runId: event.runId });
clearArgumentChurnPolicyWaits(activity, { runId: event.runId });
touchSessionActivity(activity, "run:completed");
touchSemanticSessionActivity(activity, "run:completed", { runId: event.runId });
}
export function markDiagnosticEmbeddedRunStarted(params: {
@@ -362,6 +381,7 @@ export function markDiagnosticEmbeddedRunStarted(params: {
}
// Registration is the ownership boundary. A replacement or re-armed run
// must never inherit the prior owner's semantic-stall clock.
clearRepeatedRequestActivity(activity);
if (activity.argumentChurnStartedAt !== undefined) {
clearArgumentChurnActivity(activity, { runId: ownerRunId });
}
@@ -398,8 +418,9 @@ export function markDiagnosticEmbeddedRunEnded(params: {
if (activity.activeEmbeddedRuns.size === 0) {
clearArgumentChurnActivity(activity);
clearArgumentChurnPolicyWaits(activity);
clearRepeatedRequestActivity(activity);
}
touchSessionActivity(activity, "embedded_run:ended");
touchSemanticSessionActivity(activity, "embedded_run:ended");
}
function resolveEmbeddedRunWorkKey(params: { sessionId: string; workKey?: string }): string {
@@ -619,8 +640,9 @@ export function clearDiagnosticEmbeddedRunActivityForSession(params: {
const clearedPolicyWait = clearArgumentChurnPolicyWaits(activity, {
runId: params.activeSessionId,
});
const clearedRepeatedRequests = clearRepeatedRequestActivity(activity);
return {
cleared: clearedChurn || clearedPolicyWait,
cleared: clearedChurn || clearedPolicyWait || clearedRepeatedRequests,
blockedByActiveEmbeddedRun: false,
};
}
@@ -650,7 +672,8 @@ export function clearDiagnosticEmbeddedRunActivityForSession(params: {
activity.activeModelCalls.clear();
clearArgumentChurnActivity(activity, { runId: params.activeSessionId });
clearArgumentChurnPolicyWaits(activity, { runId: params.activeSessionId });
touchSessionActivity(activity, "embedded_run:ended");
clearRepeatedRequestActivity(activity);
touchSemanticSessionActivity(activity, "embedded_run:ended");
return { cleared: true, blockedByActiveEmbeddedRun: false };
}
@@ -663,35 +686,7 @@ export function getDiagnosticSessionActivitySnapshot(
return {};
}
let activeWorkKind: DiagnosticSessionActiveWorkKind | undefined;
if (activity.activeTools.size > 0) {
activeWorkKind = "tool_call";
} else if (activity.activeModelCalls.size > 0) {
activeWorkKind = "model_call";
} else if (activity.activeEmbeddedRuns.size > 0) {
activeWorkKind = "embedded_run";
}
let activeTool: ActiveTool | undefined;
for (const tool of activity.activeTools.values()) {
if (!activeTool || tool.startedAt < activeTool.startedAt) {
activeTool = tool;
}
}
const churnProgress = resolveArgumentChurnProgress(
activity,
activity.activeEmbeddedRuns.values(),
now,
);
return {
activeWorkKind,
...(activity.activeEmbeddedRuns.size > 0 ? { hasActiveEmbeddedRun: true } : {}),
activeToolName: activeTool?.toolName,
activeToolCallId: activeTool?.toolCallId,
activeToolAgeMs: activeTool ? Math.max(0, now - activeTool.startedAt) : undefined,
lastProgressAgeMs: Math.max(0, now - churnProgress.lastProgressAt),
lastProgressReason: churnProgress.lastProgressReason,
};
return buildDiagnosticSessionActivitySnapshot(activity, now);
}
export function getDiagnosticEmbeddedRunActivitySequence(): number {
@@ -718,6 +713,17 @@ function markDiagnosticModelStartedForTest(params: DiagnosticModelStartedActivit
export function resetDiagnosticRunActivityForTest(): void {
stopDiagnosticRunActivityTracking();
installDiagnosticRunActivityTestApi();
}
function installDiagnosticRunActivityTestApi(): void {
(globalThis as Record<PropertyKey, unknown>)[
Symbol.for("openclaw.diagnosticRunActivityTestApi")
] = {
markDiagnosticModelStartedForTest,
markDiagnosticRunProgressForTest,
markDiagnosticToolStartedForTest,
};
}
let unregisterDiagnosticRunActivityListener: (() => void) | undefined;
@@ -769,11 +775,5 @@ export function stopDiagnosticRunActivityTracking(): void {
}
if (process.env.VITEST || process.env.NODE_ENV === "test") {
(globalThis as Record<PropertyKey, unknown>)[
Symbol.for("openclaw.diagnosticRunActivityTestApi")
] = {
markDiagnosticModelStartedForTest,
markDiagnosticRunProgressForTest,
markDiagnosticToolStartedForTest,
};
installDiagnosticRunActivityTestApi();
}
@@ -34,6 +34,19 @@ export function classifySessionAttention(params: {
}): SessionAttentionClassification {
if (params.activity.activeWorkKind) {
const lastProgressAgeMs = params.activity.lastProgressAgeMs ?? 0;
if (
params.activity.hasActiveEmbeddedRun === true &&
typeof params.stuckSessionAbortMs === "number" &&
(params.activity.repeatedRequestNoProgressAgeMs ?? 0) >= params.stuckSessionAbortMs
) {
return {
eventType: "session.stalled",
reason: "repeated_model_requests_without_progress",
classification: "stalled_agent_run",
activeWorkKind: params.activity.activeWorkKind,
recoveryEligible: false,
};
}
// Idle session with queued work and stale orphaned activity (no active
// embedded owner) should be classified as recoverable stuck state, not as
@@ -8,6 +8,11 @@ import {
import { testing as embeddedRunTesting } from "../agents/embedded-agent-runner/runs.test-support.js";
import { createReplyOperation } from "../auto-reply/reply/reply-run-registry.js";
import { testing as replyRunTesting } from "../auto-reply/reply/reply-run-registry.test-support.js";
import {
onDiagnosticEvent,
resetDiagnosticEventsForTest,
type DiagnosticEventPayload,
} from "../infra/diagnostic-events.js";
import { enqueueCommandInLane, getQueueSize, resetCommandLane } from "../process/command-queue.js";
import { resetCommandQueueStateForTest } from "../process/command-queue.test-support.js";
import {
@@ -15,12 +20,17 @@ import {
markDiagnosticArgumentChurnObservation,
markDiagnosticEmbeddedRunStarted,
markDiagnosticRunProgress,
resetDiagnosticRunActivityForTest,
} from "./diagnostic-run-activity.js";
import { markDiagnosticModelStartedForTest } from "./diagnostic-run-activity.test-support.js";
import {
testing as recoveryTesting,
recoverStuckDiagnosticSession,
} from "./diagnostic-stuck-session-recovery.runtime.js";
import {
logSessionStateChange,
resetDiagnosticStateForTest,
startDiagnosticHeartbeat,
} from "./diagnostic.js";
async function expectPendingAfterEventLoopTurn(promise: Promise<unknown>): Promise<void> {
let settled = false;
@@ -44,7 +54,93 @@ describe("stuck session recovery integration", () => {
embeddedRunTesting.resetActiveEmbeddedRuns();
replyRunTesting.resetReplyRunRegistry();
resetCommandQueueStateForTest();
resetDiagnosticRunActivityForTest();
resetDiagnosticStateForTest();
resetDiagnosticEventsForTest();
});
it("recovers repeated paid-call-shaped activity once without duplicate queued delivery", async () => {
vi.useFakeTimers();
vi.setSystemTime(Date.parse("2026-08-04T03:00:00Z"));
const sessionKey = "agent:main:repeated-requests";
const sessionId = "repeated-requests-session";
const lane = resolveEmbeddedSessionLane(sessionKey);
const operation = createReplyOperation({ sessionKey, sessionId, resetTriggered: false });
operation.setPhase("running");
let markActiveStarted!: () => void;
const activeStarted = new Promise<void>((resolve) => {
markActiveStarted = resolve;
});
const active = enqueueCommandInLane(
lane,
() =>
new Promise<"aborted">((resolve) => {
markActiveStarted();
operation.abortSignal.addEventListener(
"abort",
() => {
operation.complete();
resolve("aborted");
},
{ once: true },
);
}),
{ warnAfterMs: Number.MAX_SAFE_INTEGER },
);
let deliveries = 0;
const queued = enqueueCommandInLane(
lane,
async () => {
deliveries += 1;
return "delivered";
},
{ warnAfterMs: Number.MAX_SAFE_INTEGER },
);
await activeStarted;
const events: DiagnosticEventPayload[] = [];
const unsubscribe = onDiagnosticEvent((event) => events.push(event));
startDiagnosticHeartbeat(
{ diagnostics: { enabled: true } },
{
recoverStuckSession: recoverStuckDiagnosticSession,
testTimings: { stuckSessionWarnMs: 30_000, stuckSessionAbortMs: 90_000 },
},
);
logSessionStateChange({ sessionId, sessionKey, state: "processing" });
markDiagnosticEmbeddedRunStarted({ sessionId, sessionKey, runId: sessionId });
markDiagnosticModelStartedForTest({
sessionId,
sessionKey,
runId: sessionId,
provider: "mock",
model: "repeated-request-model",
observationUnit: "request",
});
for (let attempt = 2; attempt <= 3; attempt += 1) {
await vi.advanceTimersByTimeAsync(30_000);
markDiagnosticModelStartedForTest({
sessionId,
sessionKey,
runId: sessionId,
provider: "mock",
model: "repeated-request-model",
observationUnit: "request",
});
}
await vi.advanceTimersByTimeAsync(30_000);
await Promise.resolve();
await expect(active).resolves.toBe("aborted");
await expect(queued).resolves.toBe("delivered");
await vi.advanceTimersByTimeAsync(1);
expect(deliveries).toBe(1);
expect(getQueueSize(lane)).toBe(0);
expect(events.filter((event) => event.type === "session.recovery.requested")).toHaveLength(1);
expect(events.find((event) => event.type === "session.recovery.completed")).toMatchObject({
status: "aborted",
action: "abort_embedded_run",
});
unsubscribe();
});
it("does not reset a blocked lane while a reply operation is still active", async () => {
+61
View File
@@ -1046,6 +1046,67 @@ describe("stuck session diagnostics threshold", () => {
);
});
it("recovers repeated request attempts despite fresh mechanical activity", async () => {
const events: DiagnosticEventPayload[] = [];
const recoverStuckSession = vi.fn(() => new Promise<never>(() => {}));
const stuckSessionWarnMs = 30_000;
const stuckSessionAbortMs = 90_000;
const unsubscribe = onDiagnosticEvent((event) => {
events.push(event);
});
try {
startDiagnosticHeartbeat(
{ diagnostics: { enabled: true } },
{
recoverStuckSession,
testTimings: { stuckSessionWarnMs, stuckSessionAbortMs },
},
);
logSessionStateChange({ sessionId: "s1", sessionKey: "main", state: "processing" });
markDiagnosticEmbeddedRunStarted({ sessionId: "s1", sessionKey: "main", runId: "run-1" });
markDiagnosticModelStartedForTest({
sessionId: "s1",
sessionKey: "main",
runId: "run-1",
provider: "mock",
model: "retrying-model",
observationUnit: "request",
});
for (let attempt = 2; attempt <= 6; attempt += 1) {
vi.advanceTimersByTime(30_000);
markDiagnosticModelStartedForTest({
sessionId: "s1",
sessionKey: "main",
runId: "run-1",
provider: "mock",
model: "retrying-model",
observationUnit: "request",
});
}
} finally {
unsubscribe();
}
expectRecordFields(
requireRecord(
events.find((event) => event.type === "session.stalled"),
"stalled event",
),
{
classification: "stalled_agent_run",
reason: "repeated_model_requests_without_progress",
repeatedRequestNoProgressAgeMs: stuckSessionAbortMs,
},
);
expect(recoverStuckSession).toHaveBeenCalledTimes(1);
expectRecoveryCall(
recoverStuckSession,
{ sessionId: "s1", sessionKey: "main", queueDepth: 0, allowActiveAbort: true },
["ageMs", "stateGeneration"],
);
});
it("reports silent model calls as long-running before the abort threshold", async () => {
const events: DiagnosticEventPayload[] = [];
const recoverStuckSession = vi.fn();
+12
View File
@@ -573,6 +573,10 @@ function isActiveAbortRecoveryEligible(params: {
stuckSessionAbortMs: number;
}): boolean {
return (
(params.classification?.eventType === "session.stalled" &&
params.classification.classification === "stalled_agent_run" &&
params.activity?.hasActiveEmbeddedRun === true &&
(params.activity.repeatedRequestNoProgressAgeMs ?? 0) >= params.stuckSessionAbortMs) ||
isStalledEmbeddedRunRecoveryEligible(params) ||
isBlockedToolCallRecoveryEligible(params) ||
isStalledModelCallRecoveryEligible(params)
@@ -958,6 +962,9 @@ function sessionAttentionFields(params: {
...(params.activity.activeToolAgeMs !== undefined
? { activeToolAgeMs: params.activity.activeToolAgeMs }
: {}),
...(params.activity.repeatedRequestNoProgressAgeMs !== undefined
? { repeatedRequestNoProgressAgeMs: params.activity.repeatedRequestNoProgressAgeMs }
: {}),
...(terminalProgressStale ? { terminalProgressStale: true } : {}),
};
}
@@ -979,6 +986,11 @@ function formatSessionActivityLogFields(activity: DiagnosticSessionActivitySnaps
if (activity.activeToolAgeMs !== undefined) {
fields.push(`activeToolAge=${Math.round(activity.activeToolAgeMs / 1000)}s`);
}
if (activity.repeatedRequestNoProgressAgeMs !== undefined) {
fields.push(
`repeatedRequestNoProgressAge=${Math.round(activity.repeatedRequestNoProgressAgeMs / 1000)}s`,
);
}
if (isTerminalDiagnosticProgressReason(activity.lastProgressReason)) {
fields.push("terminalProgressStale=true");
}