Merge pull request #3912 from code-yeongyu/fix/sync-poller-abort-handling
fix(delegate-task): handle abort race in sync session polling
This commit is contained in:
@@ -186,6 +186,231 @@ describe("executeSyncContinuation - toast cleanup error paths", () => {
|
|||||||
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("recovers from MessageAbortedError poll error when result already exists", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => ({
|
||||||
|
data: [
|
||||||
|
{ info: { id: "msg_001", role: "user", time: { created: 1000 } } },
|
||||||
|
{
|
||||||
|
info: { id: "msg_002", role: "assistant", time: { created: 2000 }, finish: "end_turn" },
|
||||||
|
parts: [{ type: "text", text: "Response" }],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
promptAsync: async () => ({}),
|
||||||
|
status: async () => ({
|
||||||
|
data: { ses_test: { type: "idle" } },
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncContinuation } = require("./sync-continuation")
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
pollSyncSession: async () => "MessageAbortedError: aborted by user",
|
||||||
|
fetchSyncResult: async () => ({ ok: true as const, textContent: "Recovered result" }),
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
task_id: "ses_test_12345678",
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncContinuation(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
messageID: "parent-message",
|
||||||
|
}, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toContain("Task continued and completed in")
|
||||||
|
expect(result).toContain("Recovered result")
|
||||||
|
expect(removeTaskCalls.length).toBe(1)
|
||||||
|
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("recovers from canonical aborted-operation message", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => ({
|
||||||
|
data: [
|
||||||
|
{ info: { id: "msg_001", role: "user", time: { created: 1000 } } },
|
||||||
|
{
|
||||||
|
info: { id: "msg_002", role: "assistant", time: { created: 2000 }, finish: "end_turn" },
|
||||||
|
parts: [{ type: "text", text: "Response" }],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
promptAsync: async () => ({}),
|
||||||
|
status: async () => ({
|
||||||
|
data: { ses_test: { type: "idle" } },
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncContinuation } = require("./sync-continuation")
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
pollSyncSession: async () => "The operation was aborted.",
|
||||||
|
fetchSyncResult: async () => ({ ok: true as const, textContent: "Recovered result" }),
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
task_id: "ses_test_12345678",
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncContinuation(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
messageID: "parent-message",
|
||||||
|
}, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toContain("Task continued and completed in")
|
||||||
|
expect(result).toContain("Recovered result")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("returns MessageAbortedError poll error when recovery fetch has no result", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => ({
|
||||||
|
data: [
|
||||||
|
{ info: { id: "msg_001", role: "user", time: { created: 1000 } } },
|
||||||
|
{
|
||||||
|
info: { id: "msg_002", role: "assistant", time: { created: 2000 }, finish: "end_turn" },
|
||||||
|
parts: [{ type: "text", text: "Response" }],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
promptAsync: async () => ({}),
|
||||||
|
status: async () => ({
|
||||||
|
data: { ses_test: { type: "idle" } },
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncContinuation } = require("./sync-continuation")
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
pollSyncSession: async () => "MessageAbortedError: aborted by user",
|
||||||
|
fetchSyncResult: async () => ({ ok: false as const, error: "No assistant response found" }),
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
task_id: "ses_test_12345678",
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncContinuation(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
messageID: "parent-message",
|
||||||
|
}, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toBe("MessageAbortedError: aborted by user")
|
||||||
|
expect(removeTaskCalls.length).toBe(1)
|
||||||
|
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("does not recover abort poll error when anchor cannot be established", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => {
|
||||||
|
throw new Error("messages unavailable")
|
||||||
|
},
|
||||||
|
promptAsync: async () => ({}),
|
||||||
|
status: async () => ({
|
||||||
|
data: { ses_test: { type: "idle" } },
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncContinuation } = require("./sync-continuation")
|
||||||
|
let fetchSyncResultCalled = false
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
pollSyncSession: async () => "The operation was aborted.",
|
||||||
|
fetchSyncResult: async () => {
|
||||||
|
fetchSyncResultCalled = true
|
||||||
|
return { ok: true as const, textContent: "Recovered result" }
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
task_id: "ses_test_12345678",
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncContinuation(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
messageID: "parent-message",
|
||||||
|
}, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toBe("The operation was aborted.")
|
||||||
|
expect(fetchSyncResultCalled).toBe(false)
|
||||||
|
})
|
||||||
|
|
||||||
test("removes toast on successful completion", async () => {
|
test("removes toast on successful completion", async () => {
|
||||||
//#given - mock successful completion with messages growing after anchor
|
//#given - mock successful completion with messages growing after anchor
|
||||||
const mockClient = {
|
const mockClient = {
|
||||||
@@ -306,7 +531,7 @@ describe("executeSyncContinuation - toast cleanup error paths", () => {
|
|||||||
//#then - removeTask should be called at least once (poller and finally may both call it)
|
//#then - removeTask should be called at least once (poller and finally may both call it)
|
||||||
expect(removeTaskCalls.length).toBeGreaterThanOrEqual(1)
|
expect(removeTaskCalls.length).toBeGreaterThanOrEqual(1)
|
||||||
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||||
expect(result).toContain("Task aborted")
|
expect(result).toBe("Task aborted.\n\nSession ID: ses_test_12345678")
|
||||||
})
|
})
|
||||||
|
|
||||||
test("no crash when toastManager is null", async () => {
|
test("no crash when toastManager is null", async () => {
|
||||||
|
|||||||
@@ -24,6 +24,32 @@ type ResumeContext = {
|
|||||||
anchorMessageCount?: number
|
anchorMessageCount?: number
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function shouldAttemptPollErrorRecovery(pollError: string): boolean {
|
||||||
|
const trimmed = pollError.trim()
|
||||||
|
|
||||||
|
if (trimmed.length === 0) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bMessageAbortedError\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bDOMException\b/u.test(trimmed) && /\bAbortError\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bAbortError\b/u.test(trimmed) && !/\bTask aborted\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/^the operation was aborted\.?$/iu.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
async function resolveResumeContext(
|
async function resolveResumeContext(
|
||||||
client: ExecutorContext["client"],
|
client: ExecutorContext["client"],
|
||||||
continuationID: string
|
continuationID: string
|
||||||
@@ -161,7 +187,32 @@ export async function executeSyncContinuation(
|
|||||||
taskId,
|
taskId,
|
||||||
anchorMessageCount,
|
anchorMessageCount,
|
||||||
}, syncPollTimeoutMs)
|
}, syncPollTimeoutMs)
|
||||||
if (pollError) {
|
if (pollError && shouldAttemptPollErrorRecovery(pollError)) {
|
||||||
|
if (anchorMessageCount === undefined) {
|
||||||
|
return pollError
|
||||||
|
}
|
||||||
|
const recoveredResult = await deps.fetchSyncResult(client, continuationID, anchorMessageCount, {
|
||||||
|
strictAbortRecovery: true,
|
||||||
|
})
|
||||||
|
if (!recoveredResult.ok) {
|
||||||
|
return pollError
|
||||||
|
}
|
||||||
|
|
||||||
|
const duration = formatDuration(startTime)
|
||||||
|
|
||||||
|
return `Task continued and completed in ${duration}.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
${recoveredResult.textContent || "(No text output)"}
|
||||||
|
|
||||||
|
${buildTaskMetadataBlock({
|
||||||
|
sessionId: continuationID,
|
||||||
|
taskId: continuationID,
|
||||||
|
agent: resumeAgent,
|
||||||
|
category: args.category,
|
||||||
|
})}`
|
||||||
|
} else if (pollError) {
|
||||||
return pollError
|
return pollError
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -141,4 +141,65 @@ describe("fetchSyncResult", () => {
|
|||||||
expect(result.ok).toBe(false)
|
expect(result.ok).toBe(false)
|
||||||
expect(result.error).toContain("No assistant response found")
|
expect(result.error).toContain("No assistant response found")
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("strict abort recovery: does not fall back to older text when latest assistant is error", async () => {
|
||||||
|
//#given
|
||||||
|
const { fetchSyncResult } = require("./sync-result-fetcher")
|
||||||
|
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => ({
|
||||||
|
data: [
|
||||||
|
{ info: { id: "msg_001", role: "user", time: { created: 1000 } } },
|
||||||
|
{
|
||||||
|
info: { id: "msg_002", role: "assistant", time: { created: 2000 } },
|
||||||
|
parts: [{ type: "text", text: "Older text" }],
|
||||||
|
},
|
||||||
|
{
|
||||||
|
info: {
|
||||||
|
id: "msg_003",
|
||||||
|
role: "assistant",
|
||||||
|
time: { created: 3000 },
|
||||||
|
error: { name: "MessageAbortedError", message: "The operation was aborted." },
|
||||||
|
},
|
||||||
|
parts: [],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await fetchSyncResult(mockClient, "ses_test", 1, { strictAbortRecovery: true })
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result.ok).toBe(false)
|
||||||
|
expect(result.error).toContain("Latest assistant message is an error")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("strict abort recovery: requires latest assistant text output", async () => {
|
||||||
|
//#given
|
||||||
|
const { fetchSyncResult } = require("./sync-result-fetcher")
|
||||||
|
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
messages: async () => ({
|
||||||
|
data: [
|
||||||
|
{ info: { id: "msg_001", role: "user", time: { created: 1000 } } },
|
||||||
|
{
|
||||||
|
info: { id: "msg_002", role: "assistant", time: { created: 2000 } },
|
||||||
|
parts: [{ type: "tool", toolCallId: "t1", toolName: "x", state: "output-available", input: {}, output: {} }],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await fetchSyncResult(mockClient, "ses_test", 0, { strictAbortRecovery: true })
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result.ok).toBe(false)
|
||||||
|
expect(result.error).toContain("No assistant text output found in latest response")
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -5,7 +5,8 @@ import { normalizeSDKResponse } from "../../shared"
|
|||||||
export async function fetchSyncResult(
|
export async function fetchSyncResult(
|
||||||
client: OpencodeClient,
|
client: OpencodeClient,
|
||||||
sessionID: string,
|
sessionID: string,
|
||||||
anchorMessageCount?: number
|
anchorMessageCount?: number,
|
||||||
|
options?: { strictAbortRecovery?: boolean }
|
||||||
): Promise<{ ok: true; textContent: string } | { ok: false; error: string }> {
|
): Promise<{ ok: true; textContent: string } | { ok: false; error: string }> {
|
||||||
const messagesResult = await client.session.messages({
|
const messagesResult = await client.session.messages({
|
||||||
path: { id: sessionID },
|
path: { id: sessionID },
|
||||||
@@ -44,6 +45,26 @@ export async function fetchSyncResult(
|
|||||||
return { ok: false, error: `No assistant response found.\n\nSession ID: ${sessionID}` }
|
return { ok: false, error: `No assistant response found.\n\nSession ID: ${sessionID}` }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (options?.strictAbortRecovery) {
|
||||||
|
if (lastMessage.info && "error" in lastMessage.info) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
error: `Latest assistant message is an error; refusing abort recovery.\n\nSession ID: ${sessionID}`,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const lastTextParts = lastMessage.parts?.filter((p) => p.type === "text" || p.type === "reasoning") ?? []
|
||||||
|
const lastContent = lastTextParts.map((p) => p.text ?? "").filter(Boolean).join("\n")
|
||||||
|
if (!lastContent) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
error: `No assistant text output found in latest response.\n\nSession ID: ${sessionID}`,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return { ok: true, textContent: lastContent }
|
||||||
|
}
|
||||||
|
|
||||||
// Search assistant messages (newest first) for one with text/reasoning content.
|
// Search assistant messages (newest first) for one with text/reasoning content.
|
||||||
// The last assistant message may only contain tool calls with no text.
|
// The last assistant message may only contain tool calls with no text.
|
||||||
let textContent = ""
|
let textContent = ""
|
||||||
@@ -56,5 +77,12 @@ export async function fetchSyncResult(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (!textContent) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
error: `No assistant text output found in completed response.\n\nSession ID: ${sessionID}`,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
return { ok: true, textContent }
|
return { ok: true, textContent }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -421,6 +421,37 @@ describe("pollSyncSession", () => {
|
|||||||
expect(result).toContain("ses_abort")
|
expect(result).toContain("ses_abort")
|
||||||
expect(abortCount).toBe(1)
|
expect(abortCount).toBe(1)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("retries final message fetch on abort before returning aborted", async () => {
|
||||||
|
// given: abort signal set and message fetch keeps failing
|
||||||
|
const { pollSyncSession } = require("./sync-session-poller")
|
||||||
|
let abortCount = 0
|
||||||
|
let messageCallCount = 0
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
abort: async () => {
|
||||||
|
abortCount++
|
||||||
|
},
|
||||||
|
messages: async () => {
|
||||||
|
messageCallCount++
|
||||||
|
throw new Error("temporary fetch failure")
|
||||||
|
},
|
||||||
|
status: async () => ({ data: {} }),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const result = await pollSyncSession(createMockCtx(true), mockClient, {
|
||||||
|
sessionID: "ses_abort_retry",
|
||||||
|
agentToUse: "test-agent",
|
||||||
|
toastManager: { removeTask: () => {} },
|
||||||
|
taskId: "task_123",
|
||||||
|
})
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(result).toContain("Task aborted")
|
||||||
|
expect(messageCallCount).toBe(3)
|
||||||
|
expect(abortCount).toBe(1)
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
describe("timeout handling", () => {
|
describe("timeout handling", () => {
|
||||||
|
|||||||
@@ -105,19 +105,32 @@ export async function pollSyncSession(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (ctx.abort?.aborted) {
|
if (ctx.abort?.aborted) {
|
||||||
try {
|
let finalMessages: SessionMessage[] | null = null
|
||||||
const messages = await fetchSessionMessages(client, input.sessionID)
|
const abortFetchAttempts = 3
|
||||||
|
for (let attempt = 1; attempt <= abortFetchAttempts; attempt++) {
|
||||||
|
try {
|
||||||
|
finalMessages = await fetchSessionMessages(client, input.sessionID)
|
||||||
|
break
|
||||||
|
} catch (error) {
|
||||||
|
log("[task] Final messages fetch failed after abort, retrying", {
|
||||||
|
sessionID: input.sessionID,
|
||||||
|
attempt,
|
||||||
|
maxAttempts: abortFetchAttempts,
|
||||||
|
error: String(error),
|
||||||
|
})
|
||||||
|
if (attempt < abortFetchAttempts) {
|
||||||
|
await wait(syncTiming.POLL_INTERVAL_MS)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if (finalMessages) {
|
||||||
const hasNewMessages =
|
const hasNewMessages =
|
||||||
input.anchorMessageCount === undefined || messages.length > input.anchorMessageCount
|
input.anchorMessageCount === undefined || finalMessages.length > input.anchorMessageCount
|
||||||
if (hasNewMessages && isSessionComplete(messages)) {
|
if (hasNewMessages && isSessionComplete(finalMessages)) {
|
||||||
log("[task] Abort detected after session already completed", { sessionID: input.sessionID })
|
log("[task] Abort detected after session already completed", { sessionID: input.sessionID })
|
||||||
return null
|
return null
|
||||||
}
|
}
|
||||||
} catch (error) {
|
|
||||||
log("[task] Final messages fetch failed after abort, continuing with abort", {
|
|
||||||
sessionID: input.sessionID,
|
|
||||||
error: String(error),
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
|
|
||||||
log("[task] Aborted by user", { sessionID: input.sessionID })
|
log("[task] Aborted by user", { sessionID: input.sessionID })
|
||||||
|
|||||||
@@ -173,7 +173,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
|||||||
expect(rollback).toHaveBeenCalledTimes(1)
|
expect(rollback).toHaveBeenCalledTimes(1)
|
||||||
})
|
})
|
||||||
|
|
||||||
test("cleans up toast and subagentSessions when pollSyncSession returns error", async () => {
|
test("recovers from MessageAbortedError poll error when result already exists", async () => {
|
||||||
const mockClient = {
|
const mockClient = {
|
||||||
session: {
|
session: {
|
||||||
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
||||||
@@ -185,7 +185,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
|||||||
const deps = {
|
const deps = {
|
||||||
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
||||||
sendSyncPrompt: async () => null,
|
sendSyncPrompt: async () => null,
|
||||||
pollSyncSession: async () => "Poll error",
|
pollSyncSession: async () => "MessageAbortedError: aborted by user",
|
||||||
fetchSyncResult: async () => ({ ok: true as const, textContent: "Result" }),
|
fetchSyncResult: async () => ({ ok: true as const, textContent: "Result" }),
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -210,19 +210,172 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
|||||||
command: null,
|
command: null,
|
||||||
}
|
}
|
||||||
|
|
||||||
//#when - executeSyncTask with pollSyncSession failing
|
//#when - executeSyncTask with MessageAbortedError poll error
|
||||||
const result = await executeSyncTask(args, mockCtx, mockExecutorCtx, {
|
const result = await executeSyncTask(args, mockCtx, mockExecutorCtx, {
|
||||||
sessionID: "parent-session",
|
sessionID: "parent-session",
|
||||||
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
||||||
|
|
||||||
//#then - should return error and cleanup resources
|
//#then - should recover via fetchSyncResult and cleanup resources
|
||||||
expect(result).toBe("Poll error")
|
expect(result).toContain("Task completed in")
|
||||||
|
expect(result).toContain("Result")
|
||||||
expect(removeTaskCalls.length).toBe(1)
|
expect(removeTaskCalls.length).toBe(1)
|
||||||
expect(removeTaskCalls[0]).toBe("sync_ses_test")
|
expect(removeTaskCalls[0]).toBe("sync_ses_test")
|
||||||
expect(deleteCalls.length).toBe(1)
|
expect(deleteCalls.length).toBe(1)
|
||||||
expect(deleteCalls[0]).toBe("ses_test_12345678")
|
expect(deleteCalls[0]).toBe("ses_test_12345678")
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("recovers from canonical aborted-operation message", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncTask } = require("./sync-task")
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
||||||
|
sendSyncPrompt: async () => null,
|
||||||
|
pollSyncSession: async () => "The operation was aborted.",
|
||||||
|
fetchSyncResult: async () => ({ ok: true as const, textContent: "Recovered result" }),
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
directory: "/tmp",
|
||||||
|
onSyncSessionCreated: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
command: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncTask(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toContain("Task completed in")
|
||||||
|
expect(result).toContain("Recovered result")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("does not recover from non-abort poll error containing abort-like words", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncTask } = require("./sync-task")
|
||||||
|
let fetchSyncResultCalled = false
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
||||||
|
sendSyncPrompt: async () => null,
|
||||||
|
pollSyncSession: async () => "Task aborted: subagent exceeded 5 assistant turns without completing",
|
||||||
|
fetchSyncResult: async () => {
|
||||||
|
fetchSyncResultCalled = true
|
||||||
|
return { ok: true as const, textContent: "unexpected" }
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
directory: "/tmp",
|
||||||
|
onSyncSessionCreated: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
command: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncTask(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toBe("Task aborted: subagent exceeded 5 assistant turns without completing")
|
||||||
|
expect(fetchSyncResultCalled).toBe(false)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("returns abort poll error when recovery fetch has no result", async () => {
|
||||||
|
const mockClient = {
|
||||||
|
session: {
|
||||||
|
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const { executeSyncTask } = require("./sync-task")
|
||||||
|
|
||||||
|
let fetchSyncResultCalled = false
|
||||||
|
|
||||||
|
const deps = {
|
||||||
|
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
||||||
|
sendSyncPrompt: async () => null,
|
||||||
|
pollSyncSession: async () => "MessageAbortedError: aborted by user",
|
||||||
|
fetchSyncResult: async () => {
|
||||||
|
fetchSyncResultCalled = true
|
||||||
|
return { ok: false as const, error: "No assistant response found" }
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockCtx = {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
callID: "call-123",
|
||||||
|
metadata: () => {},
|
||||||
|
}
|
||||||
|
|
||||||
|
const mockExecutorCtx = {
|
||||||
|
client: mockClient,
|
||||||
|
directory: "/tmp",
|
||||||
|
onSyncSessionCreated: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
const args = {
|
||||||
|
prompt: "test prompt",
|
||||||
|
description: "test task",
|
||||||
|
category: "test",
|
||||||
|
load_skills: [],
|
||||||
|
run_in_background: false,
|
||||||
|
command: null,
|
||||||
|
}
|
||||||
|
|
||||||
|
//#when
|
||||||
|
const result = await executeSyncTask(args, mockCtx, mockExecutorCtx, {
|
||||||
|
sessionID: "parent-session",
|
||||||
|
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
||||||
|
|
||||||
|
//#then
|
||||||
|
expect(result).toBe("MessageAbortedError: aborted by user")
|
||||||
|
expect(fetchSyncResultCalled).toBe(true)
|
||||||
|
expect(removeTaskCalls.length).toBe(1)
|
||||||
|
expect(deleteCalls.length).toBe(1)
|
||||||
|
})
|
||||||
|
|
||||||
test("#given fallback chain set #when sendSyncPrompt fails #then retries with next model", async () => {
|
test("#given fallback chain set #when sendSyncPrompt fails #then retries with next model", async () => {
|
||||||
//#given
|
//#given
|
||||||
const mockClient = {
|
const mockClient = {
|
||||||
@@ -501,7 +654,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
|||||||
expect(result).toContain("Result from ses_second")
|
expect(result).toContain("Result from ses_second")
|
||||||
expect(deleteCalls).toContain("ses_first")
|
expect(deleteCalls).toContain("ses_first")
|
||||||
|
|
||||||
const finalMetadata = metadataCalls.at(-1)
|
const finalMetadata = metadataCalls[metadataCalls.length - 1]
|
||||||
expect(finalMetadata.metadata.sessionId).toBe("ses_second")
|
expect(finalMetadata.metadata.sessionId).toBe("ses_second")
|
||||||
expect(finalMetadata.metadata.taskId).toBe("ses_second")
|
expect(finalMetadata.metadata.taskId).toBe("ses_second")
|
||||||
expect(finalMetadata.metadata.model).toEqual({
|
expect(finalMetadata.metadata.model).toEqual({
|
||||||
@@ -652,7 +805,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
|||||||
}, "sisyphus-junior", initialModel, undefined, undefined, fallbackChain, deps)
|
}, "sisyphus-junior", initialModel, undefined, undefined, fallbackChain, deps)
|
||||||
|
|
||||||
expect(result).toBe("Final retry failed")
|
expect(result).toBe("Final retry failed")
|
||||||
const finalMetadata = metadataCalls.at(-1)
|
const finalMetadata = metadataCalls[metadataCalls.length - 1]
|
||||||
expect(finalMetadata.metadata.sessionId).toBe("ses_second")
|
expect(finalMetadata.metadata.sessionId).toBe("ses_second")
|
||||||
expect(finalMetadata.metadata.taskId).toBe("ses_second")
|
expect(finalMetadata.metadata.taskId).toBe("ses_second")
|
||||||
expect(finalMetadata.metadata.model).toEqual({
|
expect(finalMetadata.metadata.model).toEqual({
|
||||||
|
|||||||
@@ -15,6 +15,32 @@ import { resolveMetadataModel } from "./resolve-metadata-model"
|
|||||||
import { shouldRetryError } from "../../shared/model-error-classifier"
|
import { shouldRetryError } from "../../shared/model-error-classifier"
|
||||||
import type { ModelFallbackState } from "../../hooks/model-fallback/hook"
|
import type { ModelFallbackState } from "../../hooks/model-fallback/hook"
|
||||||
|
|
||||||
|
function shouldAttemptPollErrorRecovery(pollError: string): boolean {
|
||||||
|
const trimmed = pollError.trim()
|
||||||
|
|
||||||
|
if (trimmed.length === 0) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bMessageAbortedError\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bDOMException\b/u.test(trimmed) && /\bAbortError\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/\bAbortError\b/u.test(trimmed) && !/\bTask aborted\b/u.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
if (/^the operation was aborted\.?$/iu.test(trimmed)) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
export async function executeSyncTask(
|
export async function executeSyncTask(
|
||||||
args: DelegateTaskArgs,
|
args: DelegateTaskArgs,
|
||||||
ctx: ToolContextWithMetadata,
|
ctx: ToolContextWithMetadata,
|
||||||
@@ -213,6 +239,33 @@ export async function executeSyncTask(
|
|||||||
taskId,
|
taskId,
|
||||||
}, syncPollTimeoutMs)
|
}, syncPollTimeoutMs)
|
||||||
if (pollError) {
|
if (pollError) {
|
||||||
|
if (shouldAttemptPollErrorRecovery(pollError)) {
|
||||||
|
const recoveredResult = await deps.fetchSyncResult(client, activeSessionID, undefined, {
|
||||||
|
strictAbortRecovery: true,
|
||||||
|
})
|
||||||
|
if (recoveredResult.ok) {
|
||||||
|
const duration = formatDuration(startTime)
|
||||||
|
|
||||||
|
const actualModelStr = effectiveCategoryModel
|
||||||
|
? `${effectiveCategoryModel.providerID}/${effectiveCategoryModel.modelID}`
|
||||||
|
: undefined
|
||||||
|
const parentModelStr = parentContext.model
|
||||||
|
? `${parentContext.model.providerID}/${parentContext.model.modelID}`
|
||||||
|
: undefined
|
||||||
|
let modelRoutingNote = ""
|
||||||
|
if (actualModelStr && parentModelStr && actualModelStr !== parentModelStr) {
|
||||||
|
modelRoutingNote = `\n⚠️ Model fallback used: requested ${parentModelStr}, executed ${actualModelStr}`
|
||||||
|
}
|
||||||
|
|
||||||
|
return `Task completed in ${duration}.\n\n---\n\n${recoveredResult.textContent || "(No text output)"}${modelRoutingNote}\n\n${buildTaskMetadataBlock({
|
||||||
|
sessionId: activeSessionID,
|
||||||
|
taskId: activeSessionID,
|
||||||
|
agent: agentToUse,
|
||||||
|
category: args.category,
|
||||||
|
})}`
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
const nextFallbackModel = shouldRetryError({ message: pollError })
|
const nextFallbackModel = shouldRetryError({ message: pollError })
|
||||||
? getNextSyncFallbackModel(activeSessionID, fallbackState)
|
? getNextSyncFallbackModel(activeSessionID, fallbackState)
|
||||||
: null
|
: null
|
||||||
|
|||||||
Reference in New Issue
Block a user