From a064e1367622cf2429e98de6a3f1087000bb46de Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Sun, 10 May 2026 14:55:06 +0900 Subject: [PATCH] fix(delegate-task): recover sync results on abort race --- .../delegate-task/sync-continuation.test.ts | 118 +++++++++++++++++- src/tools/delegate-task/sync-continuation.ts | 20 ++- .../delegate-task/sync-session-poller.test.ts | 31 +++++ .../delegate-task/sync-session-poller.ts | 31 +++-- src/tools/delegate-task/sync-task.test.ts | 57 ++++++++- src/tools/delegate-task/sync-task.ts | 30 +++++ 6 files changed, 272 insertions(+), 15 deletions(-) diff --git a/src/tools/delegate-task/sync-continuation.test.ts b/src/tools/delegate-task/sync-continuation.test.ts index 9f2a690ef..b33b7f721 100644 --- a/src/tools/delegate-task/sync-continuation.test.ts +++ b/src/tools/delegate-task/sync-continuation.test.ts @@ -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 () => { diff --git a/src/tools/delegate-task/sync-continuation.ts b/src/tools/delegate-task/sync-continuation.ts index add3afd0d..36679981f 100644 --- a/src/tools/delegate-task/sync-continuation.ts +++ b/src/tools/delegate-task/sync-continuation.ts @@ -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) diff --git a/src/tools/delegate-task/sync-session-poller.test.ts b/src/tools/delegate-task/sync-session-poller.test.ts index 004fa0cb8..552b543bc 100644 --- a/src/tools/delegate-task/sync-session-poller.test.ts +++ b/src/tools/delegate-task/sync-session-poller.test.ts @@ -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", () => { diff --git a/src/tools/delegate-task/sync-session-poller.ts b/src/tools/delegate-task/sync-session-poller.ts index 9d69e4157..5c3d5e5b9 100644 --- a/src/tools/delegate-task/sync-session-poller.ts +++ b/src/tools/delegate-task/sync-session-poller.ts @@ -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 }) diff --git a/src/tools/delegate-task/sync-task.test.ts b/src/tools/delegate-task/sync-task.test.ts index bd346b8f0..11d2c4177 100644 --- a/src/tools/delegate-task/sync-task.test.ts +++ b/src/tools/delegate-task/sync-task.test.ts @@ -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 = { diff --git a/src/tools/delegate-task/sync-task.ts b/src/tools/delegate-task/sync-task.ts index 5601247b7..762146442 100644 --- a/src/tools/delegate-task/sync-task.ts +++ b/src/tools/delegate-task/sync-task.ts @@ -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