feat(background-agent): add wait-for-task-session helper
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -1,2 +1,4 @@
|
||||
export * from "./types"
|
||||
export { BackgroundManager, type SubagentSessionCreatedEvent, type OnSubagentSessionCreated } from "./manager"
|
||||
export { waitForTaskSessionID } from "./wait-for-task-session"
|
||||
export type { WaitForTaskSessionIDOptions } from "./wait-for-task-session"
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
import { describe, expect, test } from "bun:test"
|
||||
|
||||
import type { BackgroundTaskStatus } from "./types"
|
||||
import { waitForTaskSessionID } from "./wait-for-task-session"
|
||||
|
||||
interface TaskSnapshot {
|
||||
sessionID?: string
|
||||
status?: BackgroundTaskStatus
|
||||
}
|
||||
|
||||
function createManager(responses: TaskSnapshot[]) {
|
||||
let index = 0
|
||||
|
||||
return {
|
||||
getTask(_taskID: string): TaskSnapshot {
|
||||
const response = responses[Math.min(index, responses.length - 1)]
|
||||
index += 1
|
||||
return response
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
describe("waitForTaskSessionID", () => {
|
||||
test("#given task already has a session id #when waiting #then it returns immediately", async () => {
|
||||
// given
|
||||
const manager = createManager([{ sessionID: "ses_ready_123", status: "running" }])
|
||||
|
||||
// when
|
||||
const sessionID = await waitForTaskSessionID(manager, "bg_ready")
|
||||
|
||||
// then
|
||||
expect(sessionID).toBe("ses_ready_123")
|
||||
})
|
||||
|
||||
test("#given session appears later #when waiting #then it polls until resolved", async () => {
|
||||
// given
|
||||
const manager = createManager([
|
||||
{ status: "running" },
|
||||
{ status: "running" },
|
||||
{ sessionID: "ses_late_123", status: "running" },
|
||||
])
|
||||
|
||||
// when
|
||||
const sessionID = await waitForTaskSessionID(manager, "bg_late", {
|
||||
intervalMs: 1,
|
||||
timeoutMs: 20,
|
||||
})
|
||||
|
||||
// then
|
||||
expect(sessionID).toBe("ses_late_123")
|
||||
})
|
||||
|
||||
test("#given aborted signal #when waiting #then it returns undefined", async () => {
|
||||
// given
|
||||
const controller = new AbortController()
|
||||
controller.abort()
|
||||
const manager = createManager([{ status: "running" }])
|
||||
|
||||
// when
|
||||
const sessionID = await waitForTaskSessionID(manager, "bg_abort", {
|
||||
signal: controller.signal,
|
||||
})
|
||||
|
||||
// then
|
||||
expect(sessionID).toBeUndefined()
|
||||
})
|
||||
|
||||
test("#given task never resolves #when waiting past timeout #then it returns undefined", async () => {
|
||||
// given
|
||||
const manager = createManager([{ status: "running" }, { status: "running" }, { status: "running" }])
|
||||
|
||||
// when
|
||||
const sessionID = await waitForTaskSessionID(manager, "bg_timeout", {
|
||||
intervalMs: 1,
|
||||
timeoutMs: 3,
|
||||
})
|
||||
|
||||
// then
|
||||
expect(sessionID).toBeUndefined()
|
||||
})
|
||||
|
||||
test.each(["error", "cancelled", "interrupt"] satisfies BackgroundTaskStatus[])(
|
||||
"#given %s task state #when waiting #then it returns undefined",
|
||||
async (status: BackgroundTaskStatus) => {
|
||||
// given
|
||||
const manager = createManager([{ status }])
|
||||
|
||||
// when
|
||||
const sessionID = await waitForTaskSessionID(manager, `bg_${status}`)
|
||||
|
||||
// then
|
||||
expect(sessionID).toBeUndefined()
|
||||
}
|
||||
)
|
||||
})
|
||||
@@ -0,0 +1,68 @@
|
||||
import { getTimingConfig } from "../../tools/delegate-task/timing"
|
||||
import type { BackgroundTaskStatus } from "./types"
|
||||
|
||||
type SessionWaitTerminalStatus = Extract<BackgroundTaskStatus, "error" | "cancelled" | "interrupt">
|
||||
type AbortSignalLike = { aborted: boolean }
|
||||
|
||||
interface TaskReader {
|
||||
getTask(taskID: string): { sessionID?: string; status?: BackgroundTaskStatus } | undefined
|
||||
}
|
||||
|
||||
export interface WaitForTaskSessionIDOptions {
|
||||
timeoutMs?: number
|
||||
intervalMs?: number
|
||||
signal?: AbortSignalLike
|
||||
}
|
||||
|
||||
function isTerminalStatus(status: BackgroundTaskStatus | undefined): status is SessionWaitTerminalStatus {
|
||||
return status === "error" || status === "cancelled" || status === "interrupt"
|
||||
}
|
||||
|
||||
function waitForInterval(intervalMs: number): Promise<void> {
|
||||
return new Promise(resolve => {
|
||||
const scheduler = globalThis as { setTimeout: (handler: () => void, timeout?: number) => unknown }
|
||||
scheduler.setTimeout(resolve, intervalMs)
|
||||
})
|
||||
}
|
||||
|
||||
export async function waitForTaskSessionID(
|
||||
manager: TaskReader,
|
||||
taskID: string,
|
||||
options: WaitForTaskSessionIDOptions = {}
|
||||
): Promise<string | undefined> {
|
||||
const timing = getTimingConfig()
|
||||
const timeoutMs = options.timeoutMs ?? timing.WAIT_FOR_SESSION_TIMEOUT_MS
|
||||
const intervalMs = options.intervalMs ?? timing.WAIT_FOR_SESSION_INTERVAL_MS
|
||||
|
||||
if (options.signal?.aborted) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
const initialTask = manager.getTask(taskID)
|
||||
if (initialTask?.sessionID) {
|
||||
return initialTask.sessionID
|
||||
}
|
||||
if (isTerminalStatus(initialTask?.status)) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
const deadline = Date.now() + timeoutMs
|
||||
|
||||
while (Date.now() < deadline) {
|
||||
if (options.signal?.aborted) {
|
||||
return undefined
|
||||
}
|
||||
|
||||
await waitForInterval(intervalMs)
|
||||
|
||||
const task = manager.getTask(taskID)
|
||||
if (task?.sessionID) {
|
||||
return task.sessionID
|
||||
}
|
||||
if (isTerminalStatus(task?.status)) {
|
||||
return undefined
|
||||
}
|
||||
}
|
||||
|
||||
return undefined
|
||||
}
|
||||
Reference in New Issue
Block a user