fix(background-agent): defer stale cancel on activity lookup errors

This commit is contained in:
YeonGyu-Kim
2026-05-21 12:43:51 +09:00
parent 036e04c15b
commit 6e6f300c7d
5 changed files with 159 additions and 45 deletions
@@ -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<PollingManager>(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()
})
})
@@ -1,7 +1,12 @@
import { isRecord, log } from "../../shared"
import type { OpencodeClient } from "./opencode-client"
export type SessionActivityResolver = (sessionID: string) => Promise<Date | undefined>
export type SessionActivityLookup =
| { readonly type: "activity"; readonly activity: Date }
| { readonly type: "missing" }
| { readonly type: "unavailable" }
export type SessionActivityResolver = (sessionID: string) => Promise<SessionActivityLookup>
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<Date | undefined> {
): Promise<SessionActivityLookup> {
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" }
}
}
@@ -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<TaskActivityRefreshResult> {
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)
}
@@ -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()
})
})
+8 -36
View File
@@ -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<BackgroundTask["status"]>([
"completed",
@@ -112,38 +113,6 @@ export function pruneStaleTasksAndNotifications(args: {
export type SessionStatusMap = Record<string, { type: string }>
async function refreshTaskActivityFromSession(
task: BackgroundTask,
getSessionActivity: SessionActivityResolver,
): Promise<number | undefined> {
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<BackgroundTask>
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