mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
fix(reply): preserve stalled recovery input and final feedback
Keep the claimed follow-up responsible for last-resort feedback and require persisted input before transcript-only continuation. Move the existing continuation block into the execution owner; cover real queue delivery and early preflight cancellation. Co-authored-by: obviyus <22031114+obviyus@users.noreply.github.com>
This commit is contained in:
parent
313f58aeea
commit
e0b05c52af
5 changed files with 186 additions and 62 deletions
|
|
@ -214,7 +214,7 @@ The Control UI **System busyness** overlay and `diagnostics.lanes` report this w
|
|||
- `session.stuck` is reserved for recoverable stale session bookkeeping, including idle queued sessions with stale ownerless model/tool activity.
|
||||
- `session.stuck` always triggers recovery that can release the affected session lane. A `session.stalled` classification past the abort threshold (blocked tool call, stalled model call, or stalled embedded run) can also trigger active-abort recovery, so both classifications can unstick a queue, not only `session.stuck`.
|
||||
- Repeated model requests without semantic progress share one stagnation clock. Fresh transport bytes or another retry cannot renew it indefinitely. Recovery rechecks that evidence before aborting, honors owned tool and provider retry deadlines, and lets the existing run owner settle before the queue drains.
|
||||
- When recovery aborts an interactive turn before it replied, OpenClaw continues the conversation instead of asking the user to retry. When the next queued follow-up is from the same sender and route (for example a `chat.send` with `queueMode: "followup"`), it runs with a note that the previous turn stalled; otherwise one recovery turn starts on the stalled turn's route and answers from the work already in the transcript, without resubmitting the original message or repeating completed actions. Other senders' queued messages and channel messages that arrived during the stalled turn run afterwards as their own turns. The "stopped making progress" notice appears only when that recovery turn stalls too. Heartbeat and cron turns are unaffected.
|
||||
- When recovery aborts an interactive turn before it replied and its request is already saved in the transcript, OpenClaw attempts one continuation instead of immediately asking the user to retry. When the next queued follow-up is from the same sender and route (for example a `chat.send` with `queueMode: "followup"`), it takes over the outstanding request; otherwise a recovery turn starts on the stalled turn's route. Both use the existing transcript and are instructed not to repeat completed actions. Other senders' queued messages and channel messages that arrived during the stalled turn run afterwards as their own turns. The "stopped making progress" notice is the last resort if that continuation also stalls or cannot be scheduled, including when the original request was not yet saved. Heartbeat and cron turns are unaffected.
|
||||
- Repeated `session.stuck` and `session.long_running` warning log lines back off exponentially while the session remains unchanged; recovery attempts still run on every heartbeat tick regardless of that backoff.
|
||||
|
||||
## Related
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
import crypto from "node:crypto";
|
||||
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
|
||||
import { appendCurrentInboundContext } from "../../agents/embedded-agent-runner/run/runtime-context-prompt.js";
|
||||
import type { OpenClawConfig } from "../../config/config.js";
|
||||
import type { SessionEntry } from "../../config/sessions.js";
|
||||
import { withBeforeAgentReplyObserver } from "../../plugins/before-agent-reply.js";
|
||||
|
|
@ -11,6 +12,7 @@ import type { ReplyPayload } from "../types.js";
|
|||
import {
|
||||
resolveReplyRunDeliveryContext,
|
||||
resolveSourceReplyPolicy,
|
||||
scheduleFollowupDrainAfterReplyOperationClear,
|
||||
type RunReplyAgentParams,
|
||||
} from "./agent-runner-core.js";
|
||||
import { executeAgentTurn } from "./agent-runner-execution.js";
|
||||
|
|
@ -26,12 +28,65 @@ import {
|
|||
buildRecoverablePendingFinalDeliveryText,
|
||||
normalizePendingFinalDeliveryPayloads,
|
||||
} from "./pending-final-delivery.js";
|
||||
import { claimNextQueuedFollowupRequestFrom, enqueueFollowupRun } from "./queue.js";
|
||||
import { isReplyOperationSuperseded } from "./reply-operation-abort.js";
|
||||
import { recordReplyOperationAgentTurn } from "./reply-operation-run-state.js";
|
||||
import type { ReplyOperation } from "./reply-run-registry.js";
|
||||
import { replyRunRegistry } from "./reply-run-registry.js";
|
||||
import { createReplyRestartRecoveryClaimController } from "./restart-recovery-claim.js";
|
||||
import { resolveReplySourceTurnId } from "./source-turn-id.js";
|
||||
import { buildStalledTurnRecoveryRun, STALLED_TURN_GUIDANCE } from "./stalled-turn-recovery.js";
|
||||
|
||||
/** Continues a saved stalled request once, retaining its final-feedback obligation. */
|
||||
export function continueStalledReplyTurn({
|
||||
followupRun,
|
||||
queueKey,
|
||||
resolvedQueue,
|
||||
replyOperation,
|
||||
runFollowupTurn,
|
||||
}: Pick<RunReplyAgentParams, "followupRun" | "queueKey" | "resolvedQueue"> & {
|
||||
replyOperation: ReplyOperation;
|
||||
runFollowupTurn: FinalizeReplyAgentRunInput["runFollowupTurn"];
|
||||
}): boolean {
|
||||
// Preflight can stall before admission persists the request. Transcript-only
|
||||
// recovery cannot answer that input; leave its notice with dispatch.
|
||||
if (followupRun.userTurnTranscriptRecorder?.hasPersisted() !== true) {
|
||||
return false;
|
||||
}
|
||||
try {
|
||||
followupRun.operatorAuthority?.assertCurrent();
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
const queuedRequest = claimNextQueuedFollowupRequestFrom(queueKey, followupRun);
|
||||
if (queuedRequest) {
|
||||
// The handoff owns the same last-resort feedback as a dedicated recovery.
|
||||
queuedRequest.stalledTurnRecovery = true;
|
||||
queuedRequest.currentInboundContext = appendCurrentInboundContext(
|
||||
queuedRequest.currentInboundContext,
|
||||
[{ kind: "runtime-instruction", text: STALLED_TURN_GUIDANCE }],
|
||||
);
|
||||
return true;
|
||||
}
|
||||
const enqueued = enqueueFollowupRun(
|
||||
queueKey,
|
||||
buildStalledTurnRecoveryRun(followupRun),
|
||||
resolvedQueue,
|
||||
"none",
|
||||
runFollowupTurn,
|
||||
false,
|
||||
{ position: "front" },
|
||||
);
|
||||
if (enqueued) {
|
||||
scheduleFollowupDrainAfterReplyOperationClear({
|
||||
operation: replyOperation,
|
||||
queueKey,
|
||||
runFollowup: runFollowupTurn,
|
||||
});
|
||||
}
|
||||
return enqueued;
|
||||
}
|
||||
|
||||
type ExecutePreparedReplyAgentRunInput = Omit<
|
||||
FinalizeReplyAgentRunInput,
|
||||
"activeSessionEntry" | "preflightCompactionApplied" | "execution" | "runId" | "runStartedAt"
|
||||
|
|
|
|||
|
|
@ -1,6 +1,5 @@
|
|||
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
|
||||
import { resolveDefaultAgentId } from "../../agents/agent-scope-config.js";
|
||||
import { appendCurrentInboundContext } from "../../agents/embedded-agent-runner/run/runtime-context-prompt.js";
|
||||
import { resolveReplyCompletion } from "../../agents/reply-completion.js";
|
||||
import { readChannelContextGatewayContextResolver } from "../../channels/message-access/admission-evidence.js";
|
||||
import { settleProgressVisibilityCallbackResult } from "../../channels/progress-visibility.js";
|
||||
|
|
@ -31,6 +30,7 @@ import {
|
|||
scheduleFollowupDrainAfterReplyOperationClear,
|
||||
} from "./agent-runner-core.js";
|
||||
import {
|
||||
continueStalledReplyTurn,
|
||||
createReplyAgentRestartRecoveryController,
|
||||
executePreparedReplyAgentRun,
|
||||
} from "./agent-runner-execute.js";
|
||||
|
|
@ -55,11 +55,7 @@ import { createFollowupRunner } from "./followup-runner.js";
|
|||
import { REPLY_RUN_STILL_SHUTTING_DOWN_TEXT } from "./get-reply-run-queue.js";
|
||||
import { resolveOriginMessageProvider } from "./origin-routing.js";
|
||||
import { resolveActiveRunQueueAction } from "./queue-policy.js";
|
||||
import {
|
||||
claimNextQueuedFollowupRequestFrom,
|
||||
enqueueFollowupRun,
|
||||
scheduleFollowupDrain,
|
||||
} from "./queue.js";
|
||||
import { enqueueFollowupRun, scheduleFollowupDrain } from "./queue.js";
|
||||
import { resolveFollowupAbortSignal } from "./queue/types.js";
|
||||
import { REPLY_ADMISSION_TICKET } from "./reply-admission-ticket.js";
|
||||
import { createReplyMediaContext } from "./reply-media-paths.js";
|
||||
|
|
@ -76,7 +72,6 @@ import {
|
|||
import { resolveRoutedDeliveryThreadId } from "./routed-delivery-thread.js";
|
||||
import { resolveSourceReplyExpectation } from "./source-reply-delivery-mode.js";
|
||||
import { readChannelSourceTurnId } from "./source-turn-id.js";
|
||||
import { buildStalledTurnRecoveryRun, STALLED_TURN_GUIDANCE } from "./stalled-turn-recovery.js";
|
||||
import { createTypingSignaler } from "./typing-mode.js";
|
||||
export async function runReplyAgent(
|
||||
input: RunReplyAgentParams,
|
||||
|
|
@ -602,38 +597,14 @@ export async function runReplyAgent(
|
|||
// Dispatch owns the stall notice; this owner holds the queue facts needed to answer
|
||||
// instead. The same sender's next queued request inherits the guidance; otherwise one
|
||||
// recovery run bound to this turn's route and authority is queued.
|
||||
replyOperationRunState.continueStalledTurn = () => {
|
||||
try {
|
||||
followupRun.operatorAuthority?.assertCurrent();
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
const queuedRequest = claimNextQueuedFollowupRequestFrom(queueKey, followupRun);
|
||||
if (queuedRequest) {
|
||||
queuedRequest.currentInboundContext = appendCurrentInboundContext(
|
||||
queuedRequest.currentInboundContext,
|
||||
[{ kind: "runtime-instruction", text: STALLED_TURN_GUIDANCE }],
|
||||
);
|
||||
return true;
|
||||
}
|
||||
const enqueued = enqueueFollowupRun(
|
||||
replyOperationRunState.continueStalledTurn = () =>
|
||||
continueStalledReplyTurn({
|
||||
followupRun,
|
||||
queueKey,
|
||||
buildStalledTurnRecoveryRun(followupRun),
|
||||
resolvedQueue,
|
||||
"none",
|
||||
replyOperation,
|
||||
runFollowupTurn,
|
||||
false,
|
||||
{ position: "front" },
|
||||
);
|
||||
if (enqueued) {
|
||||
scheduleFollowupDrainAfterReplyOperationClear({
|
||||
operation: replyOperation,
|
||||
queueKey,
|
||||
runFollowup: runFollowupTurn,
|
||||
});
|
||||
}
|
||||
return enqueued;
|
||||
};
|
||||
});
|
||||
}
|
||||
const {
|
||||
admitUserTurn,
|
||||
|
|
|
|||
|
|
@ -1,9 +1,15 @@
|
|||
// Tests how an admitted interactive run continues a turn the stale watchdog dropped.
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { createDeferred } from "../../../test/helpers/promise.js";
|
||||
import {
|
||||
createAdmittedRunOperatorAuthority,
|
||||
type AdmittedRunOperatorAuthority,
|
||||
} from "../../agents/admitted-run-context.js";
|
||||
import { createUserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.js";
|
||||
import {
|
||||
createSqliteTranscriptTarget,
|
||||
readTranscriptMessages,
|
||||
} from "../../sessions/user-turn-transcript.test-support.js";
|
||||
import type { TemplateContext } from "../templating.js";
|
||||
import type * as AgentRunnerExecution from "./agent-runner-execution.js";
|
||||
import { runReplyAgent } from "./agent-runner.js";
|
||||
|
|
@ -26,13 +32,18 @@ import { createMockTypingController } from "./test-helpers.js";
|
|||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
executeAgentTurn: vi.fn(),
|
||||
preflight: vi.fn(async (_params: { abortSignal: AbortSignal }) => undefined),
|
||||
drainedRuns: vi.fn(async (_run: FollowupRun) => {}),
|
||||
executeFollowups: false,
|
||||
routeReply: vi.fn(async (..._args: unknown[]) => ({ ok: true, delivered: true })),
|
||||
followupSettled: () => {},
|
||||
}));
|
||||
const executeAgentTurnMock = mocks.executeAgentTurn;
|
||||
const drainedRuns = mocks.drainedRuns;
|
||||
let executionStarted = createDeferred();
|
||||
|
||||
vi.mock("./agent-runner-memory.js", () => ({
|
||||
runSessionCompactionIfNeeded: async () => undefined,
|
||||
runSessionCompactionIfNeeded: (params: { abortSignal: AbortSignal }) => mocks.preflight(params),
|
||||
runMemoryFlushIfNeeded: async () => ({ sessionEntry: undefined, outcome: "skipped" }),
|
||||
}));
|
||||
|
||||
|
|
@ -41,14 +52,36 @@ vi.mock("./agent-runner-execution.js", async () => ({
|
|||
executeAgentTurn: (...args: unknown[]) => mocks.executeAgentTurn(...args),
|
||||
}));
|
||||
|
||||
vi.mock("./followup-runner.js", () => ({
|
||||
createFollowupRunner: () => mocks.drainedRuns,
|
||||
vi.mock("./followup-runner.js", async (importOriginal) => {
|
||||
const { createFollowupRunner } = await importOriginal<typeof import("./followup-runner.js")>();
|
||||
return {
|
||||
createFollowupRunner: (...args: Parameters<typeof createFollowupRunner>) => {
|
||||
const runFollowup = createFollowupRunner(...args);
|
||||
return async (queued: FollowupRun) => {
|
||||
await mocks.drainedRuns(queued);
|
||||
if (mocks.executeFollowups) {
|
||||
try {
|
||||
await runFollowup(queued);
|
||||
} finally {
|
||||
mocks.followupSettled();
|
||||
}
|
||||
}
|
||||
};
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
vi.mock("./route-reply.js", () => ({
|
||||
isRoutableChannel: (channel: string | undefined) => channel === "telegram",
|
||||
routeReply: (...args: unknown[]) => mocks.routeReply(...args),
|
||||
}));
|
||||
|
||||
type StalledRun = {
|
||||
operation: ReplyOperation;
|
||||
run: Promise<unknown>;
|
||||
runState: ReplyOperationRunState;
|
||||
recorder: ReturnType<typeof createUserTurnTranscriptRecorder>;
|
||||
transcriptTarget: ReturnType<typeof createSqliteTranscriptTarget>;
|
||||
};
|
||||
|
||||
const queueKey = "agent:main:telegram:direct:stalled";
|
||||
|
|
@ -68,6 +101,16 @@ function createStalledRun(
|
|||
followupRun.operatorAuthority = options.operatorAuthority;
|
||||
followupRun.images = [{ type: "image", data: "aW1n", mimeType: "image/png" }];
|
||||
followupRun.transcriptPrompt = "what is good at the hotel restaurant?";
|
||||
const transcriptTarget = createSqliteTranscriptTarget({
|
||||
dir: followupRun.run.workspaceDir,
|
||||
sessionId: "stalled-session",
|
||||
sessionKey: queueKey,
|
||||
});
|
||||
const recorder = createUserTurnTranscriptRecorder({
|
||||
input: { text: followupRun.transcriptPrompt },
|
||||
target: transcriptTarget,
|
||||
});
|
||||
followupRun.userTurnTranscriptRecorder = recorder;
|
||||
const operation = createReplyOperation({
|
||||
sessionKey: queueKey,
|
||||
sessionId: "stalled-session",
|
||||
|
|
@ -104,11 +147,12 @@ function createStalledRun(
|
|||
shouldInjectGroupIntro: false,
|
||||
typingMode: "instant",
|
||||
});
|
||||
return { operation, run, runState };
|
||||
return { operation, run, runState, recorder, transcriptTarget };
|
||||
}
|
||||
|
||||
async function stallBeforeOutput(stalled: StalledRun) {
|
||||
await vi.waitFor(() => expect(executeAgentTurnMock).toHaveBeenCalledOnce());
|
||||
await executionStarted.promise;
|
||||
expect(executeAgentTurnMock).toHaveBeenCalledOnce();
|
||||
expect(expireStaleReplyOperation(stalled.operation, "stuck_recovery")).toBe(false);
|
||||
}
|
||||
|
||||
|
|
@ -137,9 +181,15 @@ describe("runReplyAgent stalled turn continuation", () => {
|
|||
replyRunTesting.resetReplyRunRegistry();
|
||||
clearSessionQueues([queueKey]);
|
||||
drainedRuns.mockClear();
|
||||
mocks.executeFollowups = false;
|
||||
executionStarted = createDeferred();
|
||||
mocks.preflight.mockReset().mockResolvedValue(undefined);
|
||||
mocks.routeReply.mockClear();
|
||||
mocks.followupSettled = () => {};
|
||||
executeAgentTurnMock
|
||||
.mockReset()
|
||||
.mockImplementation(async (params: { replyOperation: { abortSignal: AbortSignal } }) => {
|
||||
executionStarted.resolve();
|
||||
await new Promise<void>((resolve) => {
|
||||
params.replyOperation.abortSignal.addEventListener("abort", () => resolve(), {
|
||||
once: true,
|
||||
|
|
@ -174,30 +224,53 @@ describe("runReplyAgent stalled turn continuation", () => {
|
|||
expect(recovery?.images).toBeUndefined();
|
||||
expect(recovery?.abortSignal).toBeUndefined();
|
||||
expect(recovery?.run.sessionKey).toBe(queueKey);
|
||||
expect(stalled.recorder.hasPersisted()).toBe(true);
|
||||
expect(await readTranscriptMessages(stalled.transcriptTarget)).toEqual([
|
||||
expect.objectContaining({ content: "what is good at the hotel restaurant?" }),
|
||||
]);
|
||||
expect(getFollowupQueueDepth(queueKey)).toBe(0);
|
||||
});
|
||||
|
||||
it("gives the same sender's already-queued request the interruption guidance instead", async () => {
|
||||
const stalled = createStalledRun();
|
||||
const queued = createQueuedRequest({ senderId: "traveler", to: "12345" });
|
||||
expect(enqueueFollowupRun(queueKey, queued, settings, "message-id", drainedRuns, false)).toBe(
|
||||
true,
|
||||
);
|
||||
await stallBeforeOutput(stalled);
|
||||
it.each(["followup", "collect"] as const)(
|
||||
"sends one last-resort notice when the claimed %s request also stalls",
|
||||
async (mode) => {
|
||||
mocks.executeFollowups = true;
|
||||
const settled = createDeferred();
|
||||
mocks.followupSettled = settled.resolve;
|
||||
const stalled = createStalledRun();
|
||||
const queued = createQueuedRequest({ senderId: "traveler", to: "12345" });
|
||||
queued.messageId = "msg-double-stall-" + mode;
|
||||
expect(
|
||||
enqueueFollowupRun(queueKey, queued, { ...settings, mode }, "message-id", undefined, false),
|
||||
).toBe(true);
|
||||
await stallBeforeOutput(stalled);
|
||||
executeAgentTurnMock.mockImplementationOnce(async ({ replyOperation }) => {
|
||||
expireStaleReplyOperation(replyOperation, "stuck_recovery");
|
||||
return { runId: "recovery-run", outcome: { kind: "aborted", reason: "user" } };
|
||||
});
|
||||
|
||||
expect(stalled.runState.continueStalledTurn?.()).toBe(true);
|
||||
expect(getFollowupQueueDepth(queueKey)).toBe(1);
|
||||
expect(queued.prompt).toBe("answer already");
|
||||
expect(queued.currentInboundContext?.fragments).toContainEqual({
|
||||
kind: "runtime-instruction",
|
||||
text: expect.stringContaining("previous turn stopped making progress"),
|
||||
});
|
||||
expect(stalled.runState.continueStalledTurn?.()).toBe(true);
|
||||
expect(queued.currentInboundContext?.fragments).toContainEqual({
|
||||
kind: "runtime-instruction",
|
||||
text: expect.stringContaining("previous turn stopped making progress"),
|
||||
});
|
||||
await settleStalledOwner(stalled);
|
||||
await settled.promise;
|
||||
expect(drainedRuns).toHaveBeenCalledExactlyOnceWith(queued);
|
||||
|
||||
await settleStalledOwner(stalled);
|
||||
await vi.waitFor(() => expect(drainedRuns).toHaveBeenCalledOnce());
|
||||
expect(drainedRuns.mock.calls[0]?.[0]).toBe(queued);
|
||||
expect(drainedRuns.mock.calls[0]?.[0].stalledTurnRecovery).toBeUndefined();
|
||||
});
|
||||
expect(mocks.routeReply).toHaveBeenCalledExactlyOnceWith(
|
||||
expect.objectContaining({
|
||||
channel: "telegram",
|
||||
to: "12345",
|
||||
payload: expect.objectContaining({
|
||||
text: "⚠️ This turn was interrupted because it stopped making progress. Please try again.",
|
||||
isError: true,
|
||||
}),
|
||||
}),
|
||||
);
|
||||
expect(executeAgentTurnMock).toHaveBeenCalledTimes(2);
|
||||
},
|
||||
);
|
||||
|
||||
// Default DM scope shares one main session across senders.
|
||||
it("never hands the stalled request past another sender's earlier queued request", async () => {
|
||||
|
|
@ -228,6 +301,31 @@ describe("runReplyAgent stalled turn continuation", () => {
|
|||
expect(last).toBe(sameSenderLater);
|
||||
});
|
||||
|
||||
it("leaves dispatch responsible for a stall before the request reaches the transcript", async () => {
|
||||
const preflightStarted = createDeferred();
|
||||
mocks.preflight.mockImplementationOnce(async ({ abortSignal }) => {
|
||||
preflightStarted.resolve();
|
||||
await new Promise<void>((resolve) => {
|
||||
abortSignal.addEventListener("abort", () => resolve(), { once: true });
|
||||
});
|
||||
abortSignal.throwIfAborted();
|
||||
return undefined;
|
||||
});
|
||||
const stalled = createStalledRun();
|
||||
const settled = stalled.run.catch(() => undefined);
|
||||
await preflightStarted.promise;
|
||||
expireStaleReplyOperation(stalled.operation, "stuck_recovery");
|
||||
|
||||
const continued = stalled.runState.continueStalledTurn?.();
|
||||
await settled;
|
||||
stalled.operation.complete();
|
||||
expect(stalled.recorder.hasPersisted()).toBe(false);
|
||||
expect(continued).toBe(false);
|
||||
expect(getFollowupQueueDepth(queueKey)).toBe(0);
|
||||
expect(executeAgentTurnMock).not.toHaveBeenCalled();
|
||||
expect(await readTranscriptMessages(stalled.transcriptTarget)).toEqual([]);
|
||||
});
|
||||
|
||||
it("falls back to the notice once the stalled turn's authority is revoked", async () => {
|
||||
let revoked = false;
|
||||
const operatorAuthority = createAdmittedRunOperatorAuthority({
|
||||
|
|
|
|||
|
|
@ -166,7 +166,7 @@ export type FollowupRun = {
|
|||
};
|
||||
/** Internal marker for the one-shot stranded final recovery retry. */
|
||||
strandedReplyRetry?: boolean;
|
||||
/** Internal marker for the one-shot recovery run after a stalled interactive turn. */
|
||||
/** This continuation owes last-resort feedback if it also stalls, including claimed input. */
|
||||
stalledTurnRecovery?: boolean;
|
||||
/** Preserve priority runs when old-item queue overflow eviction runs before drain. */
|
||||
protectFromQueueOverflow?: boolean;
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue