From 6e6f300c7dde09904c9fdcfd6b59ce7fbcd64032 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Thu, 21 May 2026 12:43:51 +0900 Subject: [PATCH] fix(background-agent): defer stale cancel on activity lookup errors --- .../manager-session-activity.test.ts | 49 +++++++++++++++++++ .../background-agent/session-activity.ts | 27 +++++++--- .../background-agent/task-activity-refresh.ts | 49 +++++++++++++++++++ .../task-poller-session-activity.test.ts | 35 +++++++++++-- src/features/background-agent/task-poller.ts | 44 +++-------------- 5 files changed, 159 insertions(+), 45 deletions(-) create mode 100644 src/features/background-agent/task-activity-refresh.ts diff --git a/src/features/background-agent/manager-session-activity.test.ts b/src/features/background-agent/manager-session-activity.test.ts index 69ec9a4e6..4a7bbf3fe 100644 --- a/src/features/background-agent/manager-session-activity.test.ts +++ b/src/features/background-agent/manager-session-activity.test.ts @@ -102,4 +102,53 @@ describe("BackgroundManager persisted session activity stale checks", () => { await manager.shutdown() }) + + test("keeps a busy task running when session.get returns an error response", async () => { + //#given - live event progress is stale and OpenCode session lookup fails without throwing + spyOn(globalThis.Date, "now").mockReturnValue(fixedTime) + let abortCallCount = 0 + const sessionGet = mock(async () => ({ + error: "lookup failed", + data: undefined, + })) + const client = { + session: { + status: async () => ({ data: { "ses-active": { type: "busy" } } }), + get: sessionGet, + prompt: async () => ({}), + promptAsync: async () => ({}), + abort: async () => { + abortCallCount += 1 + return {} + }, + todo: async () => ({ data: [] }), + messages: async () => ({ data: [] }), + }, + } + const manager = new BackgroundManager({ + pluginContext: createPluginContext(client), + config: { staleTimeoutMs: 180_000 }, + enableParentSessionNotifications: false, + }) + const task = createRunningTask({ + startedAt: new Date(Date.now() - 45 * 60 * 1000), + progress: { + toolCalls: 3, + lastUpdate: new Date(Date.now() - 45 * 60 * 1000), + }, + }) + const pollingManager = unsafeTestValue(manager) + pollingManager.tasks.set(task.id, task) + + //#when - polling tries to confirm stale activity through the SDK response + await pollingManager.pollRunningTasks() + + //#then - the lookup failure defers cancellation for the active child session + expect(task.status).toBe("running") + expect(task.error).toBeUndefined() + expect(abortCallCount).toBe(0) + expect(sessionGet).toHaveBeenCalledTimes(1) + + await manager.shutdown() + }) }) diff --git a/src/features/background-agent/session-activity.ts b/src/features/background-agent/session-activity.ts index 463f6216d..10979452c 100644 --- a/src/features/background-agent/session-activity.ts +++ b/src/features/background-agent/session-activity.ts @@ -1,7 +1,12 @@ import { isRecord, log } from "../../shared" import type { OpencodeClient } from "./opencode-client" -export type SessionActivityResolver = (sessionID: string) => Promise +export type SessionActivityLookup = + | { readonly type: "activity"; readonly activity: Date } + | { readonly type: "missing" } + | { readonly type: "unavailable" } + +export type SessionActivityResolver = (sessionID: string) => Promise function dateFromMillis(value: unknown): Date | undefined { if (typeof value !== "number") return undefined @@ -15,27 +20,37 @@ export function extractSessionActivityDate(sessionInfo: unknown): Date | undefin return dateFromMillis(time?.updated) ?? dateFromMillis(sessionInfo.time_updated) } +function sessionActivityLookupFromInfo(sessionInfo: unknown): SessionActivityLookup { + const activity = extractSessionActivityDate(sessionInfo) + return activity ? { type: "activity", activity } : { type: "missing" } +} + export async function getSessionActivityFromClient( client: OpencodeClient, sessionID: string, directory?: string, -): Promise { +): Promise { const sessionGet = client.session.get - if (typeof sessionGet !== "function") return undefined + if (typeof sessionGet !== "function") return { type: "missing" } try { const response = await sessionGet({ path: { id: sessionID }, ...(directory ? { query: { directory } } : {}), }) + if (isRecord(response) && response.error !== undefined && response.error !== null) { + log("[background-agent] Failed to read session activity:", { sessionID, error: response.error }) + return { type: "unavailable" } + } + const sessionInfo = isRecord(response) && "data" in response ? response.data : response - return extractSessionActivityDate(sessionInfo) + return sessionActivityLookupFromInfo(sessionInfo) } catch (error) { if (error instanceof Error) { log("[background-agent] Failed to read session activity:", { sessionID, error: error.message }) - return undefined + return { type: "unavailable" } } log("[background-agent] Failed to read session activity:", { sessionID, error }) - return undefined + return { type: "unavailable" } } } diff --git a/src/features/background-agent/task-activity-refresh.ts b/src/features/background-agent/task-activity-refresh.ts new file mode 100644 index 000000000..ca364649d --- /dev/null +++ b/src/features/background-agent/task-activity-refresh.ts @@ -0,0 +1,49 @@ +import { log } from "../../shared" +import type { SessionActivityLookup, SessionActivityResolver } from "./session-activity" +import type { BackgroundTask } from "./types" + +export type TaskActivityRefreshResult = + | { readonly type: "activity"; readonly activityTime: number } + | { readonly type: "missing" } + | { readonly type: "unavailable" } + +function updateTaskActivityFromLookup( + task: BackgroundTask, + lookup: SessionActivityLookup, +): TaskActivityRefreshResult { + if (lookup.type !== "activity") return lookup + + const activityTime = lookup.activity.getTime() + if (!Number.isFinite(activityTime)) return { type: "missing" } + + const baseline = task.progress?.lastUpdate.getTime() ?? task.startedAt?.getTime() + if (baseline !== undefined && activityTime <= baseline) return { type: "activity", activityTime } + + if (!task.progress) { + task.progress = { toolCalls: 0, lastUpdate: new Date(activityTime) } + } else { + task.progress.lastUpdate = new Date(activityTime) + } + return { type: "activity", activityTime } +} + +export async function refreshTaskActivityFromSession( + task: BackgroundTask, + getSessionActivity: SessionActivityResolver, +): Promise { + if (!task.sessionId) return { type: "missing" } + + let lookup: SessionActivityLookup + try { + lookup = await getSessionActivity(task.sessionId) + } catch (error) { + if (error instanceof Error) { + log("[background-agent] Error refreshing task session activity:", { taskId: task.id, error: error.message }) + return { type: "unavailable" } + } + log("[background-agent] Error refreshing task session activity:", { taskId: task.id, error }) + return { type: "unavailable" } + } + + return updateTaskActivityFromLookup(task, lookup) +} diff --git a/src/features/background-agent/task-poller-session-activity.test.ts b/src/features/background-agent/task-poller-session-activity.test.ts index de4b687d6..2cd76b9b9 100644 --- a/src/features/background-agent/task-poller-session-activity.test.ts +++ b/src/features/background-agent/task-poller-session-activity.test.ts @@ -61,7 +61,7 @@ describe("checkAndInterruptStaleTasks persisted session activity", () => { concurrencyManager: mockConcurrencyManager, notifyParentSession: mockNotify, sessionStatuses: { "ses-1": { type: "busy" } }, - getSessionActivity: async () => freshActivity, + getSessionActivity: async () => ({ type: "activity", activity: freshActivity }), }) //#then - task stays running and future stale checks use the persisted activity timestamp @@ -88,7 +88,7 @@ describe("checkAndInterruptStaleTasks persisted session activity", () => { concurrencyManager: mockConcurrencyManager, notifyParentSession: mockNotify, sessionStatuses: { "ses-1": { type: "busy" } }, - getSessionActivity: async () => freshActivity, + getSessionActivity: async () => ({ type: "activity", activity: freshActivity }), }) //#then - the task stays running and receives a progress timestamp for future stale checks @@ -118,7 +118,7 @@ describe("checkAndInterruptStaleTasks persisted session activity", () => { concurrencyManager: mockConcurrencyManager, notifyParentSession: mockNotify, sessionStatuses: { "ses-1": { type: "busy" } }, - getSessionActivity: async () => stalePersistedActivity, + getSessionActivity: async () => ({ type: "activity", activity: stalePersistedActivity }), }) //#then - cancellation still happens, but the stale age reflects persisted activity @@ -128,4 +128,33 @@ describe("checkAndInterruptStaleTasks persisted session activity", () => { expect(task.progress?.lastUpdate).toEqual(stalePersistedActivity) expect(mockNotify).toHaveBeenCalledWith(task) }) + + test("keeps a busy task running when persisted session activity lookup is unavailable", async () => { + //#given - the child session is active, but session.get returned an error response during confirmation + spyOn(globalThis.Date, "now").mockReturnValue(fixedTime) + const staleActivity = new Date(Date.now() - 45 * 60 * 1000) + const task = createRunningTask({ + startedAt: staleActivity, + progress: { + toolCalls: 2, + lastUpdate: staleActivity, + }, + }) + + //#when - stale checking cannot verify persisted activity for this poll + await checkAndInterruptStaleTasks({ + tasks: [task], + client: mockClient, + config: { staleTimeoutMs: 180_000 }, + concurrencyManager: mockConcurrencyManager, + notifyParentSession: mockNotify, + sessionStatuses: { "ses-1": { type: "busy" } }, + getSessionActivity: async () => ({ type: "unavailable" }), + }) + + //#then - active-session cancellation is deferred instead of treating lookup failure as inactivity + expect(task.status).toBe("running") + expect(task.progress?.lastUpdate).toEqual(staleActivity) + expect(mockNotify).not.toHaveBeenCalled() + }) }) diff --git a/src/features/background-agent/task-poller.ts b/src/features/background-agent/task-poller.ts index 6b6a42dec..9f0e0362a 100644 --- a/src/features/background-agent/task-poller.ts +++ b/src/features/background-agent/task-poller.ts @@ -20,6 +20,7 @@ import { MIN_SESSION_GONE_POLLS, verifySessionExists } from "./session-existence import { isActiveSessionStatus } from "./session-status-classifier" import { getSessionActivityFromClient, type SessionActivityResolver } from "./session-activity" +import { refreshTaskActivityFromSession } from "./task-activity-refresh" const TERMINAL_TASK_STATUSES = new Set([ "completed", @@ -112,38 +113,6 @@ export function pruneStaleTasksAndNotifications(args: { export type SessionStatusMap = Record -async function refreshTaskActivityFromSession( - task: BackgroundTask, - getSessionActivity: SessionActivityResolver, -): Promise { - if (!task.sessionId) return undefined - - let activity: Date | undefined - try { - activity = await getSessionActivity(task.sessionId) - } catch (error) { - if (error instanceof Error) { - log("[background-agent] Error refreshing task session activity:", { taskId: task.id, error: error.message }) - return undefined - } - log("[background-agent] Error refreshing task session activity:", { taskId: task.id, error }) - return undefined - } - - const activityTime = activity?.getTime() - if (activityTime === undefined || !Number.isFinite(activityTime)) return undefined - - const baseline = task.progress?.lastUpdate.getTime() ?? task.startedAt?.getTime() - if (baseline !== undefined && activityTime <= baseline) return activityTime - - if (!task.progress) { - task.progress = { toolCalls: 0, lastUpdate: new Date(activityTime) } - } else { - task.progress.lastUpdate = new Date(activityTime) - } - return activityTime -} - export async function checkAndInterruptStaleTasks(args: { tasks: Iterable client: OpencodeClient @@ -204,8 +173,9 @@ export async function checkAndInterruptStaleTasks(args: { if (runtime <= effectiveTimeout) continue if (shouldRefreshFromSessionActivity) { - const activityTime = await refreshTaskActivityFromSession(task, getSessionActivity) - if (activityTime !== undefined && now - activityTime <= effectiveTimeout) continue + const activityRefresh = await refreshTaskActivityFromSession(task, getSessionActivity) + if (activityRefresh.type === "unavailable") continue + if (activityRefresh.type === "activity" && now - activityRefresh.activityTime <= effectiveTimeout) continue } if (sessionGone && await verifySessionExists(client, sessionID, directory)) { @@ -246,8 +216,10 @@ export async function checkAndInterruptStaleTasks(args: { if (timeSinceLastUpdate <= effectiveStaleTimeout) continue if (shouldRefreshFromSessionActivity) { - const activityTime = await refreshTaskActivityFromSession(task, getSessionActivity) - const refreshedLastUpdate = task.progress?.lastUpdate.getTime() ?? activityTime + const activityRefresh = await refreshTaskActivityFromSession(task, getSessionActivity) + if (activityRefresh.type === "unavailable") continue + const refreshedLastUpdate = task.progress?.lastUpdate.getTime() + ?? (activityRefresh.type === "activity" ? activityRefresh.activityTime : undefined) if (refreshedLastUpdate !== undefined && now - refreshedLastUpdate <= effectiveStaleTimeout) continue if (refreshedLastUpdate !== undefined) { timeSinceLastUpdate = now - refreshedLastUpdate