mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-07 18:42:25 +00:00
fix(copilot): retain timed-out sessions until idle
This commit is contained in:
@@ -468,6 +468,7 @@ describe("runCopilotAttempt", () => {
|
||||
|
||||
releaseBeforeCompaction.resolve();
|
||||
activeSession?.emit("session.compaction_complete", { success: true });
|
||||
activeSession?.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(activeSession?.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
@@ -1814,6 +1815,10 @@ describe("runCopilotAttempt", () => {
|
||||
}),
|
||||
expect.anything(),
|
||||
);
|
||||
sdk.sessions[0]?.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(sdk.sessions[0]?.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
it("marks a timeout during active SDK compaction", async () => {
|
||||
@@ -1837,6 +1842,7 @@ describe("runCopilotAttempt", () => {
|
||||
expect(sdk.sessions[0]?.disconnect).not.toHaveBeenCalled();
|
||||
|
||||
sdk.sessions[0]?.emit("session.compaction_complete", { messagesRemoved: 3, success: true });
|
||||
sdk.sessions[0]?.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(sdk.sessions[0]?.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
@@ -1848,12 +1854,54 @@ describe("runCopilotAttempt", () => {
|
||||
);
|
||||
});
|
||||
|
||||
it("retains a timed-out session until later compaction reaches session.idle", async () => {
|
||||
const afterCompaction = vi.fn();
|
||||
const onDeferredCompaction = vi.fn();
|
||||
initializeGlobalHookRunner(
|
||||
createMockPluginRegistry([{ hookName: "after_compaction", handler: afterCompaction }]),
|
||||
);
|
||||
let activeSession: FakeSession | undefined;
|
||||
const sdk = makeFakeSdk({
|
||||
onCreateSession: (session) => {
|
||||
activeSession = session;
|
||||
session.sendAndWait.mockRejectedValueOnce(
|
||||
new Error("Timeout after 60000ms waiting for session.idle"),
|
||||
);
|
||||
},
|
||||
});
|
||||
|
||||
const result = await runCopilotAttempt(makeParams(), {
|
||||
onDeferredCompaction,
|
||||
pool: makeFakePool(sdk),
|
||||
});
|
||||
|
||||
expect(result.timedOut).toBe(true);
|
||||
expect(result.timedOutDuringCompaction).toBe(false);
|
||||
expect(onDeferredCompaction).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ sdkSessionId: "sess-1" }),
|
||||
);
|
||||
expect(activeSession?.disconnect).not.toHaveBeenCalled();
|
||||
|
||||
activeSession?.emit("session.compaction_start", {});
|
||||
activeSession?.emit("session.compaction_complete", { messagesRemoved: 3, success: true });
|
||||
await vi.waitFor(() => {
|
||||
expect(afterCompaction).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
expect(activeSession?.disconnect).not.toHaveBeenCalled();
|
||||
|
||||
activeSession?.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(activeSession?.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
it("does not mark a timeout after SDK compaction has completed as active compaction", async () => {
|
||||
const sdk = makeFakeSdk({
|
||||
onCreateSession: (session) => {
|
||||
session.sendAndWait.mockImplementationOnce(async () => {
|
||||
session.emit("session.compaction_start", {});
|
||||
session.emit("session.compaction_complete", { success: true });
|
||||
session.emit("session.idle", {});
|
||||
return undefined;
|
||||
});
|
||||
},
|
||||
@@ -1935,6 +1983,7 @@ describe("runCopilotAttempt", () => {
|
||||
expect(dualWriteMock.dualWriteCopilotTranscriptBestEffort).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
sdk.sessions[0]?.emit("session.compaction_complete", { success: true });
|
||||
sdk.sessions[0]?.emit("session.idle", {});
|
||||
mirror.resolve();
|
||||
|
||||
const result = await attempt;
|
||||
@@ -1978,6 +2027,10 @@ describe("runCopilotAttempt", () => {
|
||||
// replay-shim incorrectly treated the attempt as side-effect-safe.
|
||||
expect(result.replayMetadata?.hadPotentialSideEffects).toBe(true);
|
||||
expect(result.replayMetadata?.replaySafe).toBe(false);
|
||||
sdk.sessions[0]?.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(sdk.sessions[0]?.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
it("G1: SDK timeout flushes the in-flight delta chain before snapshot so assistant text is preserved", async () => {
|
||||
@@ -2019,6 +2072,10 @@ describe("runCopilotAttempt", () => {
|
||||
expect(result.timedOut).toBe(true);
|
||||
expect(onAssistantDelta).toHaveBeenCalledTimes(1);
|
||||
expect(result.assistantTexts?.join("")).toContain("partial-");
|
||||
session.emit("session.idle", {});
|
||||
await vi.waitFor(() => {
|
||||
expect(session.disconnect).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
it("model translation: unsupported provider", async () => {
|
||||
|
||||
@@ -186,34 +186,39 @@ async function finalizeCopilotAttempt(
|
||||
return result;
|
||||
}
|
||||
|
||||
async function awaitCompactionCompletionOrAbort(
|
||||
bridge: ReturnType<typeof attachEventBridge>,
|
||||
abortSignal: AbortSignal | undefined,
|
||||
): Promise<"aborted" | "completed"> {
|
||||
if (!abortSignal) {
|
||||
await bridge.awaitCompactionCompletion();
|
||||
async function awaitDeferredCleanupCompletionOrAbort(params: {
|
||||
abortSignal: AbortSignal | undefined;
|
||||
awaitSessionIdle: boolean;
|
||||
bridge: ReturnType<typeof attachEventBridge>;
|
||||
}): Promise<"aborted" | "completed"> {
|
||||
const awaitCompletion = async () => {
|
||||
if (params.awaitSessionIdle) {
|
||||
await params.bridge.awaitSessionIdle();
|
||||
}
|
||||
await params.bridge.awaitCompactionCompletion();
|
||||
};
|
||||
if (!params.abortSignal) {
|
||||
await awaitCompletion();
|
||||
return "completed";
|
||||
}
|
||||
if (abortSignal.aborted) {
|
||||
if (params.abortSignal.aborted) {
|
||||
return "aborted";
|
||||
}
|
||||
let resolveAbort: () => void = () => undefined;
|
||||
const aborted = new Promise<"aborted">((resolve) => {
|
||||
resolveAbort = () => resolve("aborted");
|
||||
});
|
||||
abortSignal.addEventListener("abort", resolveAbort, { once: true });
|
||||
params.abortSignal.addEventListener("abort", resolveAbort, { once: true });
|
||||
try {
|
||||
return await Promise.race([
|
||||
bridge.awaitCompactionCompletion().then(() => "completed" as const),
|
||||
aborted,
|
||||
]);
|
||||
return await Promise.race([awaitCompletion().then(() => "completed" as const), aborted]);
|
||||
} finally {
|
||||
abortSignal.removeEventListener("abort", resolveAbort);
|
||||
params.abortSignal.removeEventListener("abort", resolveAbort);
|
||||
}
|
||||
}
|
||||
|
||||
function deferBackgroundCompactionCleanup(params: {
|
||||
abortSignal: AbortSignal | undefined;
|
||||
awaitSessionIdle: boolean;
|
||||
bridge: ReturnType<typeof attachEventBridge>;
|
||||
handle: PooledClient;
|
||||
pool: CopilotClientPool;
|
||||
@@ -226,8 +231,9 @@ function deferBackgroundCompactionCleanup(params: {
|
||||
return (async () => {
|
||||
let outcome: "aborted" | "completed" | "deadline" = "deadline";
|
||||
try {
|
||||
outcome = await awaitCompactionCompletionBeforeDeadline({
|
||||
outcome = await awaitDeferredCleanupBeforeDeadline({
|
||||
abortSignal: params.abortSignal,
|
||||
awaitSessionIdle: params.awaitSessionIdle,
|
||||
bridge: params.bridge,
|
||||
timeoutMs: params.timeoutMs,
|
||||
});
|
||||
@@ -284,8 +290,9 @@ async function cancelBackgroundCompactionBeforeTeardown(session: SessionLike): P
|
||||
}
|
||||
}
|
||||
|
||||
async function awaitCompactionCompletionBeforeDeadline(params: {
|
||||
async function awaitDeferredCleanupBeforeDeadline(params: {
|
||||
abortSignal: AbortSignal | undefined;
|
||||
awaitSessionIdle: boolean;
|
||||
bridge: ReturnType<typeof attachEventBridge>;
|
||||
timeoutMs: number;
|
||||
}): Promise<"aborted" | "completed" | "deadline"> {
|
||||
@@ -294,10 +301,7 @@ async function awaitCompactionCompletionBeforeDeadline(params: {
|
||||
timeoutId = setTimeout(() => resolve("deadline"), params.timeoutMs);
|
||||
});
|
||||
try {
|
||||
return await Promise.race([
|
||||
awaitCompactionCompletionOrAbort(params.bridge, params.abortSignal),
|
||||
deadline,
|
||||
]);
|
||||
return await Promise.race([awaitDeferredCleanupCompletionOrAbort(params), deadline]);
|
||||
} finally {
|
||||
if (timeoutId !== undefined) {
|
||||
clearTimeout(timeoutId);
|
||||
@@ -802,7 +806,9 @@ export async function runCopilotAttempt(
|
||||
}
|
||||
} finally {
|
||||
settled = true;
|
||||
if (bridge?.hasObservedCompaction() && session && handle) {
|
||||
const retainSessionForDeferredCleanup =
|
||||
bridge?.hasObservedCompaction() || (timedOut && bridge?.hasObservedSessionIdle() === false);
|
||||
if (retainSessionForDeferredCleanup && bridge && session && handle) {
|
||||
const cleanupAbort = new AbortController();
|
||||
const abortCleanup = () => cleanupAbort.abort();
|
||||
if (params.abortSignal?.aborted) {
|
||||
@@ -812,6 +818,7 @@ export async function runCopilotAttempt(
|
||||
}
|
||||
const cleanup = deferBackgroundCompactionCleanup({
|
||||
abortSignal: cleanupAbort.signal,
|
||||
awaitSessionIdle: !bridge.hasObservedSessionIdle(),
|
||||
bridge,
|
||||
handle,
|
||||
pool: deps.pool,
|
||||
@@ -837,8 +844,9 @@ export async function runCopilotAttempt(
|
||||
}
|
||||
params.abortSignal?.removeEventListener("abort", onAbort);
|
||||
} else {
|
||||
// `sendAndWait` resolves on `session.idle`, which the SDK defines as
|
||||
// no background agents in flight. Only an observed compaction needs retention.
|
||||
// A normal sendAndWait result has observed session.idle, which the SDK
|
||||
// defines as no background agents in flight. Timeouts retain the bridge
|
||||
// until that event so compaction that starts after the timer still completes.
|
||||
await bridge?.awaitCompactionChain();
|
||||
bridge?.detach();
|
||||
params.abortSignal?.removeEventListener("abort", onAbort);
|
||||
|
||||
@@ -17,6 +17,7 @@ const REGISTERED_EVENT_TYPES = [
|
||||
"tool.execution_complete",
|
||||
"session.compaction_start",
|
||||
"session.compaction_complete",
|
||||
"session.idle",
|
||||
"session.error",
|
||||
"abort",
|
||||
] as const;
|
||||
@@ -662,6 +663,21 @@ describe("attachEventBridge", () => {
|
||||
expect(bridge.isCompacting()).toBe(false);
|
||||
});
|
||||
|
||||
it("waits for the SDK terminal idle event", async () => {
|
||||
const session = createFakeSession();
|
||||
const bridge = attachEventBridge(session, {
|
||||
getSdkSessionId: () => "sdk-session-id",
|
||||
isAborted: () => false,
|
||||
});
|
||||
|
||||
const idle = bridge.awaitSessionIdle();
|
||||
await flushAsync();
|
||||
session.emit("session.idle", makeEvent("session.idle", {}));
|
||||
await idle;
|
||||
|
||||
expect(bridge.hasObservedSessionIdle()).toBe(true);
|
||||
});
|
||||
|
||||
it("keeps compaction pending after an abort until the SDK reports completion", async () => {
|
||||
const session = createFakeSession();
|
||||
const bridge = attachEventBridge(session, {
|
||||
|
||||
@@ -69,9 +69,11 @@ export interface EventBridgeController {
|
||||
recordSendResult(result: SessionEvent | undefined): boolean;
|
||||
awaitCompactionChain(): Promise<void>;
|
||||
awaitCompactionCompletion(): Promise<void>;
|
||||
awaitSessionIdle(): Promise<void>;
|
||||
settleCompactionWait(): void;
|
||||
awaitDeltaChain(): Promise<void>;
|
||||
hasObservedCompaction(): boolean;
|
||||
hasObservedSessionIdle(): boolean;
|
||||
isCompacting(): boolean;
|
||||
snapshot(): EventBridgeSnapshot;
|
||||
buildAssistantMessage(args: BuildAssistantMessageArgs): AssistantMessage | undefined;
|
||||
@@ -104,6 +106,11 @@ export function attachEventBridge(
|
||||
let compactionChain = Promise.resolve();
|
||||
let compactionIdle = Promise.resolve();
|
||||
let resolveCompactionIdle: (() => void) | undefined;
|
||||
let observedSessionIdle = false;
|
||||
let resolveSessionIdle: (() => void) | undefined;
|
||||
const sessionIdle = new Promise<void>((resolve) => {
|
||||
resolveSessionIdle = resolve;
|
||||
});
|
||||
let firstDeltaError: unknown;
|
||||
let detached = false;
|
||||
const unsubscribeFns: Array<() => void> = [];
|
||||
@@ -217,6 +224,12 @@ export function attachEventBridge(
|
||||
}
|
||||
});
|
||||
|
||||
registerListener(session, unsubscribeFns, "session.idle", () => {
|
||||
observedSessionIdle = true;
|
||||
resolveSessionIdle?.();
|
||||
resolveSessionIdle = undefined;
|
||||
});
|
||||
|
||||
registerListener(session, unsubscribeFns, "session.error", (event) => {
|
||||
if (!options.isAborted()) {
|
||||
streamError = createPromptError(
|
||||
@@ -254,6 +267,9 @@ export function attachEventBridge(
|
||||
}
|
||||
await compactionChain;
|
||||
},
|
||||
awaitSessionIdle() {
|
||||
return observedSessionIdle ? Promise.resolve() : sessionIdle;
|
||||
},
|
||||
settleCompactionWait() {
|
||||
activeCompactionCount = 0;
|
||||
resolveCompactionIdle?.();
|
||||
@@ -265,6 +281,9 @@ export function attachEventBridge(
|
||||
hasObservedCompaction() {
|
||||
return observedCompaction;
|
||||
},
|
||||
hasObservedSessionIdle() {
|
||||
return observedSessionIdle;
|
||||
},
|
||||
isCompacting() {
|
||||
return activeCompactionCount > 0;
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user