fix(background-agent): apply smart circuit breaker to manager events
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -0,0 +1,239 @@
|
|||||||
|
import { describe, expect, test } from "bun:test"
|
||||||
|
import type { PluginInput } from "@opencode-ai/plugin"
|
||||||
|
import { tmpdir } from "node:os"
|
||||||
|
import type { BackgroundTaskConfig } from "../../config/schema"
|
||||||
|
import { BackgroundManager } from "./manager"
|
||||||
|
import type { BackgroundTask } from "./types"
|
||||||
|
|
||||||
|
function createManager(config?: BackgroundTaskConfig): BackgroundManager {
|
||||||
|
const client = {
|
||||||
|
session: {
|
||||||
|
prompt: async () => ({}),
|
||||||
|
promptAsync: async () => ({}),
|
||||||
|
abort: async () => ({}),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
const manager = new BackgroundManager({ client, directory: tmpdir() } as unknown as PluginInput, config)
|
||||||
|
const testManager = manager as unknown as {
|
||||||
|
enqueueNotificationForParent: (sessionID: string, fn: () => Promise<void>) => Promise<void>
|
||||||
|
notifyParentSession: (task: BackgroundTask) => Promise<void>
|
||||||
|
tasks: Map<string, BackgroundTask>
|
||||||
|
}
|
||||||
|
|
||||||
|
testManager.enqueueNotificationForParent = async (_sessionID, fn) => {
|
||||||
|
await fn()
|
||||||
|
}
|
||||||
|
testManager.notifyParentSession = async () => {}
|
||||||
|
|
||||||
|
return manager
|
||||||
|
}
|
||||||
|
|
||||||
|
function getTaskMap(manager: BackgroundManager): Map<string, BackgroundTask> {
|
||||||
|
return (manager as unknown as { tasks: Map<string, BackgroundTask> }).tasks
|
||||||
|
}
|
||||||
|
|
||||||
|
async function flushAsyncWork() {
|
||||||
|
await new Promise(resolve => setTimeout(resolve, 0))
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("BackgroundManager circuit breaker", () => {
|
||||||
|
describe("#given the same tool dominates the recent window", () => {
|
||||||
|
test("#when tool events arrive #then the task is cancelled early", async () => {
|
||||||
|
const manager = createManager({
|
||||||
|
circuitBreaker: {
|
||||||
|
windowSize: 20,
|
||||||
|
repetitionThresholdPercent: 80,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
const task: BackgroundTask = {
|
||||||
|
id: "task-loop-1",
|
||||||
|
sessionID: "session-loop-1",
|
||||||
|
parentSessionID: "parent-1",
|
||||||
|
parentMessageID: "msg-1",
|
||||||
|
description: "Looping task",
|
||||||
|
prompt: "loop",
|
||||||
|
agent: "explore",
|
||||||
|
status: "running",
|
||||||
|
startedAt: new Date(Date.now() - 60_000),
|
||||||
|
progress: {
|
||||||
|
toolCalls: 0,
|
||||||
|
lastUpdate: new Date(Date.now() - 60_000),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
getTaskMap(manager).set(task.id, task)
|
||||||
|
|
||||||
|
for (const toolName of [
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"grep",
|
||||||
|
"read",
|
||||||
|
"edit",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"bash",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"glob",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
"read",
|
||||||
|
]) {
|
||||||
|
manager.handleEvent({
|
||||||
|
type: "message.part.updated",
|
||||||
|
properties: { sessionID: task.sessionID, type: "tool", tool: toolName },
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
await flushAsyncWork()
|
||||||
|
|
||||||
|
expect(task.status).toBe("cancelled")
|
||||||
|
expect(task.error).toContain("repeatedly called read 16/20 times")
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
describe("#given recent tool calls are diverse", () => {
|
||||||
|
test("#when the window fills #then the task keeps running", async () => {
|
||||||
|
const manager = createManager({
|
||||||
|
circuitBreaker: {
|
||||||
|
windowSize: 10,
|
||||||
|
repetitionThresholdPercent: 80,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
const task: BackgroundTask = {
|
||||||
|
id: "task-diverse-1",
|
||||||
|
sessionID: "session-diverse-1",
|
||||||
|
parentSessionID: "parent-1",
|
||||||
|
parentMessageID: "msg-1",
|
||||||
|
description: "Healthy task",
|
||||||
|
prompt: "work",
|
||||||
|
agent: "explore",
|
||||||
|
status: "running",
|
||||||
|
startedAt: new Date(Date.now() - 60_000),
|
||||||
|
progress: {
|
||||||
|
toolCalls: 0,
|
||||||
|
lastUpdate: new Date(Date.now() - 60_000),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
getTaskMap(manager).set(task.id, task)
|
||||||
|
|
||||||
|
for (const toolName of [
|
||||||
|
"read",
|
||||||
|
"grep",
|
||||||
|
"edit",
|
||||||
|
"bash",
|
||||||
|
"glob",
|
||||||
|
"read",
|
||||||
|
"lsp_diagnostics",
|
||||||
|
"grep",
|
||||||
|
"edit",
|
||||||
|
"read",
|
||||||
|
]) {
|
||||||
|
manager.handleEvent({
|
||||||
|
type: "message.part.updated",
|
||||||
|
properties: { sessionID: task.sessionID, type: "tool", tool: toolName },
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
await flushAsyncWork()
|
||||||
|
|
||||||
|
expect(task.status).toBe("running")
|
||||||
|
expect(task.progress?.toolCalls).toBe(10)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
describe("#given the absolute cap is configured lower than the repetition detector needs", () => {
|
||||||
|
test("#when the raw tool-call cap is reached #then the backstop still cancels the task", async () => {
|
||||||
|
const manager = createManager({
|
||||||
|
maxToolCalls: 3,
|
||||||
|
circuitBreaker: {
|
||||||
|
windowSize: 10,
|
||||||
|
repetitionThresholdPercent: 95,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
const task: BackgroundTask = {
|
||||||
|
id: "task-cap-1",
|
||||||
|
sessionID: "session-cap-1",
|
||||||
|
parentSessionID: "parent-1",
|
||||||
|
parentMessageID: "msg-1",
|
||||||
|
description: "Backstop task",
|
||||||
|
prompt: "work",
|
||||||
|
agent: "explore",
|
||||||
|
status: "running",
|
||||||
|
startedAt: new Date(Date.now() - 60_000),
|
||||||
|
progress: {
|
||||||
|
toolCalls: 0,
|
||||||
|
lastUpdate: new Date(Date.now() - 60_000),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
getTaskMap(manager).set(task.id, task)
|
||||||
|
|
||||||
|
for (const toolName of ["read", "grep", "edit"]) {
|
||||||
|
manager.handleEvent({
|
||||||
|
type: "message.part.updated",
|
||||||
|
properties: { sessionID: task.sessionID, type: "tool", tool: toolName },
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
await flushAsyncWork()
|
||||||
|
|
||||||
|
expect(task.status).toBe("cancelled")
|
||||||
|
expect(task.error).toContain("maximum tool call limit (3)")
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
describe("#given the same running tool part emits multiple updates", () => {
|
||||||
|
test("#when duplicate running updates arrive #then it only counts the tool once", async () => {
|
||||||
|
const manager = createManager({
|
||||||
|
maxToolCalls: 2,
|
||||||
|
circuitBreaker: {
|
||||||
|
windowSize: 5,
|
||||||
|
repetitionThresholdPercent: 80,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
const task: BackgroundTask = {
|
||||||
|
id: "task-dedupe-1",
|
||||||
|
sessionID: "session-dedupe-1",
|
||||||
|
parentSessionID: "parent-1",
|
||||||
|
parentMessageID: "msg-1",
|
||||||
|
description: "Dedupe task",
|
||||||
|
prompt: "work",
|
||||||
|
agent: "explore",
|
||||||
|
status: "running",
|
||||||
|
startedAt: new Date(Date.now() - 60_000),
|
||||||
|
progress: {
|
||||||
|
toolCalls: 0,
|
||||||
|
lastUpdate: new Date(Date.now() - 60_000),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
getTaskMap(manager).set(task.id, task)
|
||||||
|
|
||||||
|
for (let index = 0; index < 3; index += 1) {
|
||||||
|
manager.handleEvent({
|
||||||
|
type: "message.part.updated",
|
||||||
|
properties: {
|
||||||
|
part: {
|
||||||
|
id: "tool-1",
|
||||||
|
sessionID: task.sessionID,
|
||||||
|
type: "tool",
|
||||||
|
tool: "bash",
|
||||||
|
state: { status: "running" },
|
||||||
|
},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
await flushAsyncWork()
|
||||||
|
|
||||||
|
expect(task.status).toBe("running")
|
||||||
|
expect(task.progress?.toolCalls).toBe(1)
|
||||||
|
expect(task.progress?.countedToolPartIDs).toEqual(["tool-1"])
|
||||||
|
})
|
||||||
|
})
|
||||||
|
})
|
||||||
@@ -51,6 +51,11 @@ import { join } from "node:path"
|
|||||||
import { pruneStaleTasksAndNotifications } from "./task-poller"
|
import { pruneStaleTasksAndNotifications } from "./task-poller"
|
||||||
import { checkAndInterruptStaleTasks } from "./task-poller"
|
import { checkAndInterruptStaleTasks } from "./task-poller"
|
||||||
import { removeTaskToastTracking } from "./remove-task-toast-tracking"
|
import { removeTaskToastTracking } from "./remove-task-toast-tracking"
|
||||||
|
import {
|
||||||
|
detectRepetitiveToolUse,
|
||||||
|
recordToolCall,
|
||||||
|
resolveCircuitBreakerSettings,
|
||||||
|
} from "./loop-detector"
|
||||||
import {
|
import {
|
||||||
createSubagentDepthLimitError,
|
createSubagentDepthLimitError,
|
||||||
createSubagentDescendantLimitError,
|
createSubagentDescendantLimitError,
|
||||||
@@ -64,9 +69,11 @@ type OpencodeClient = PluginInput["client"]
|
|||||||
|
|
||||||
|
|
||||||
interface MessagePartInfo {
|
interface MessagePartInfo {
|
||||||
|
id?: string
|
||||||
sessionID?: string
|
sessionID?: string
|
||||||
type?: string
|
type?: string
|
||||||
tool?: string
|
tool?: string
|
||||||
|
state?: { status?: string }
|
||||||
}
|
}
|
||||||
|
|
||||||
interface EventProperties {
|
interface EventProperties {
|
||||||
@@ -80,6 +87,19 @@ interface Event {
|
|||||||
properties?: EventProperties
|
properties?: EventProperties
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function resolveMessagePartInfo(properties: EventProperties | undefined): MessagePartInfo | undefined {
|
||||||
|
if (!properties || typeof properties !== "object") {
|
||||||
|
return undefined
|
||||||
|
}
|
||||||
|
|
||||||
|
const nestedPart = properties.part
|
||||||
|
if (nestedPart && typeof nestedPart === "object") {
|
||||||
|
return nestedPart as MessagePartInfo
|
||||||
|
}
|
||||||
|
|
||||||
|
return properties as MessagePartInfo
|
||||||
|
}
|
||||||
|
|
||||||
interface Todo {
|
interface Todo {
|
||||||
content: string
|
content: string
|
||||||
status: string
|
status: string
|
||||||
@@ -720,6 +740,8 @@ export class BackgroundManager {
|
|||||||
|
|
||||||
existingTask.progress = {
|
existingTask.progress = {
|
||||||
toolCalls: existingTask.progress?.toolCalls ?? 0,
|
toolCalls: existingTask.progress?.toolCalls ?? 0,
|
||||||
|
toolCallWindow: existingTask.progress?.toolCallWindow,
|
||||||
|
countedToolPartIDs: existingTask.progress?.countedToolPartIDs,
|
||||||
lastUpdate: new Date(),
|
lastUpdate: new Date(),
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -852,8 +874,7 @@ export class BackgroundManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (event.type === "message.part.updated" || event.type === "message.part.delta") {
|
if (event.type === "message.part.updated" || event.type === "message.part.delta") {
|
||||||
if (!props || typeof props !== "object" || !("sessionID" in props)) return
|
const partInfo = resolveMessagePartInfo(props)
|
||||||
const partInfo = props as unknown as MessagePartInfo
|
|
||||||
const sessionID = partInfo?.sessionID
|
const sessionID = partInfo?.sessionID
|
||||||
if (!sessionID) return
|
if (!sessionID) return
|
||||||
|
|
||||||
@@ -876,10 +897,50 @@ export class BackgroundManager {
|
|||||||
task.progress.lastUpdate = new Date()
|
task.progress.lastUpdate = new Date()
|
||||||
|
|
||||||
if (partInfo?.type === "tool" || partInfo?.tool) {
|
if (partInfo?.type === "tool" || partInfo?.tool) {
|
||||||
|
const countedToolPartIDs = task.progress.countedToolPartIDs ?? []
|
||||||
|
const shouldCountToolCall =
|
||||||
|
!partInfo.id ||
|
||||||
|
partInfo.state?.status !== "running" ||
|
||||||
|
!countedToolPartIDs.includes(partInfo.id)
|
||||||
|
|
||||||
|
if (!shouldCountToolCall) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if (partInfo.id && partInfo.state?.status === "running") {
|
||||||
|
task.progress.countedToolPartIDs = [...countedToolPartIDs, partInfo.id]
|
||||||
|
}
|
||||||
|
|
||||||
task.progress.toolCalls += 1
|
task.progress.toolCalls += 1
|
||||||
task.progress.lastTool = partInfo.tool
|
task.progress.lastTool = partInfo.tool
|
||||||
|
const circuitBreaker = resolveCircuitBreakerSettings(this.config)
|
||||||
|
if (partInfo.tool) {
|
||||||
|
task.progress.toolCallWindow = recordToolCall(
|
||||||
|
task.progress.toolCallWindow,
|
||||||
|
partInfo.tool,
|
||||||
|
circuitBreaker
|
||||||
|
)
|
||||||
|
|
||||||
const maxToolCalls = this.config?.maxToolCalls ?? 200
|
const loopDetection = detectRepetitiveToolUse(task.progress.toolCallWindow)
|
||||||
|
if (loopDetection.triggered) {
|
||||||
|
log("[background-agent] Circuit breaker: repetitive tool usage detected", {
|
||||||
|
taskId: task.id,
|
||||||
|
agent: task.agent,
|
||||||
|
sessionID,
|
||||||
|
toolName: loopDetection.toolName,
|
||||||
|
repeatedCount: loopDetection.repeatedCount,
|
||||||
|
sampleSize: loopDetection.sampleSize,
|
||||||
|
thresholdPercent: loopDetection.thresholdPercent,
|
||||||
|
})
|
||||||
|
void this.cancelTask(task.id, {
|
||||||
|
source: "circuit-breaker",
|
||||||
|
reason: `Subagent repeatedly called ${loopDetection.toolName} ${loopDetection.repeatedCount}/${loopDetection.sampleSize} times in the recent tool-call window (${loopDetection.thresholdPercent}% threshold). This usually indicates an infinite loop. The task was automatically cancelled to prevent excessive token usage.`,
|
||||||
|
})
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const maxToolCalls = circuitBreaker.maxToolCalls
|
||||||
if (task.progress.toolCalls >= maxToolCalls) {
|
if (task.progress.toolCalls >= maxToolCalls) {
|
||||||
log("[background-agent] Circuit breaker: tool call limit reached", {
|
log("[background-agent] Circuit breaker: tool call limit reached", {
|
||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
|
|||||||
Reference in New Issue
Block a user