fix(delegate-task): recover sync results on abort race
This commit is contained in:
@@ -186,6 +186,121 @@ describe("executeSyncContinuation - toast cleanup error paths", () => {
|
||||
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||
})
|
||||
|
||||
test("recovers from pollSyncSession 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 () => "Task aborted.\n\nSession ID: ses_test_12345678",
|
||||
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("returns 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 () => "Task aborted.\n\nSession ID: ses_test_12345678",
|
||||
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("Task aborted.\n\nSession ID: ses_test_12345678")
|
||||
expect(removeTaskCalls.length).toBe(1)
|
||||
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||
})
|
||||
|
||||
test("removes toast on successful completion", async () => {
|
||||
//#given - mock successful completion with messages growing after anchor
|
||||
const mockClient = {
|
||||
@@ -306,7 +421,8 @@ describe("executeSyncContinuation - toast cleanup error paths", () => {
|
||||
//#then - removeTask should be called at least once (poller and finally may both call it)
|
||||
expect(removeTaskCalls.length).toBeGreaterThanOrEqual(1)
|
||||
expect(removeTaskCalls[0]).toBe("resume_sync_ses_test")
|
||||
expect(result).toContain("Task aborted")
|
||||
expect(result).toContain("Task continued and completed in")
|
||||
expect(result).toContain("Result")
|
||||
})
|
||||
|
||||
test("no crash when toastManager is null", async () => {
|
||||
|
||||
@@ -162,7 +162,25 @@ export async function executeSyncContinuation(
|
||||
anchorMessageCount,
|
||||
}, syncPollTimeoutMs)
|
||||
if (pollError) {
|
||||
return pollError
|
||||
const recoveredResult = await deps.fetchSyncResult(client, continuationID, anchorMessageCount)
|
||||
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,
|
||||
})}`
|
||||
}
|
||||
|
||||
const result = await deps.fetchSyncResult(client, continuationID, anchorMessageCount)
|
||||
|
||||
@@ -421,6 +421,37 @@ describe("pollSyncSession", () => {
|
||||
expect(result).toContain("ses_abort")
|
||||
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", () => {
|
||||
|
||||
@@ -105,19 +105,32 @@ export async function pollSyncSession(
|
||||
}
|
||||
|
||||
if (ctx.abort?.aborted) {
|
||||
try {
|
||||
const messages = await fetchSessionMessages(client, input.sessionID)
|
||||
let finalMessages: SessionMessage[] | null = null
|
||||
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 =
|
||||
input.anchorMessageCount === undefined || messages.length > input.anchorMessageCount
|
||||
if (hasNewMessages && isSessionComplete(messages)) {
|
||||
input.anchorMessageCount === undefined || finalMessages.length > input.anchorMessageCount
|
||||
if (hasNewMessages && isSessionComplete(finalMessages)) {
|
||||
log("[task] Abort detected after session already completed", { sessionID: input.sessionID })
|
||||
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 })
|
||||
|
||||
@@ -173,7 +173,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
||||
expect(rollback).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
test("cleans up toast and subagentSessions when pollSyncSession returns error", async () => {
|
||||
test("recovers from pollSyncSession error when result already exists", async () => {
|
||||
const mockClient = {
|
||||
session: {
|
||||
create: async () => ({ data: { id: "ses_test_12345678" } }),
|
||||
@@ -185,7 +185,7 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
||||
const deps = {
|
||||
createSyncSession: async () => ({ ok: true, sessionID: "ses_test_12345678" }),
|
||||
sendSyncPrompt: async () => null,
|
||||
pollSyncSession: async () => "Poll error",
|
||||
pollSyncSession: async () => "Task aborted.\n\nSession ID: ses_test_12345678",
|
||||
fetchSyncResult: async () => ({ ok: true as const, textContent: "Result" }),
|
||||
}
|
||||
|
||||
@@ -215,14 +215,63 @@ describe("executeSyncTask - cleanup on error paths", () => {
|
||||
sessionID: "parent-session",
|
||||
}, "test-agent", undefined, undefined, undefined, undefined, deps)
|
||||
|
||||
//#then - should return error and cleanup resources
|
||||
expect(result).toBe("Poll error")
|
||||
//#then - should recover via fetchSyncResult and cleanup resources
|
||||
expect(result).toContain("Task completed in")
|
||||
expect(result).toContain("Result")
|
||||
expect(removeTaskCalls.length).toBe(1)
|
||||
expect(removeTaskCalls[0]).toBe("sync_ses_test")
|
||||
expect(deleteCalls.length).toBe(1)
|
||||
expect(deleteCalls[0]).toBe("ses_test_12345678")
|
||||
})
|
||||
|
||||
test("returns poll error when recovery fetch has no result", 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 () => "Poll error",
|
||||
fetchSyncResult: async () => ({ 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("Poll error")
|
||||
expect(removeTaskCalls.length).toBe(1)
|
||||
expect(deleteCalls.length).toBe(1)
|
||||
})
|
||||
|
||||
test("#given fallback chain set #when sendSyncPrompt fails #then retries with next model", async () => {
|
||||
//#given
|
||||
const mockClient = {
|
||||
|
||||
@@ -15,6 +15,11 @@ import { resolveMetadataModel } from "./resolve-metadata-model"
|
||||
import { shouldRetryError } from "../../shared/model-error-classifier"
|
||||
import type { ModelFallbackState } from "../../hooks/model-fallback/hook"
|
||||
|
||||
function shouldAttemptPollErrorRecovery(pollError: string): boolean {
|
||||
const normalized = pollError.toLowerCase()
|
||||
return normalized.includes("aborted") || normalized.includes("abort")
|
||||
}
|
||||
|
||||
export async function executeSyncTask(
|
||||
args: DelegateTaskArgs,
|
||||
ctx: ToolContextWithMetadata,
|
||||
@@ -213,6 +218,31 @@ export async function executeSyncTask(
|
||||
taskId,
|
||||
}, syncPollTimeoutMs)
|
||||
if (pollError) {
|
||||
if (shouldAttemptPollErrorRecovery(pollError)) {
|
||||
const recoveredResult = await deps.fetchSyncResult(client, activeSessionID)
|
||||
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 })
|
||||
? getNextSyncFallbackModel(activeSessionID, fallbackState)
|
||||
: null
|
||||
|
||||
Reference in New Issue
Block a user