/** * Integration tests for decision D12: a background step's `needs_attention` flag clears when its cause resolves * (the open tool ends, a failed steer is delivered, the failure streak is broken, activity resumes) and can be raised again, * while stop/interrupt/terminal still clear it unconditionally. * * Requires pi packages to be importable. Skips gracefully if unavailable. */ import { describe, it } from "node:test"; import assert from "node:assert/strict"; import * as fs from "node:fs"; import * as path from "node:path"; import { events, makeAgent } from "../support/helpers.ts"; import { deliverStopRequest, requestAsyncSteer } from "../../src/runs/background/control-channel.ts"; import type { AsyncStatusPayload } from "../support/async-execution-fixture.ts"; import { installAsyncExecutionHooks, available, isAsyncAvailable, executeAsyncSingle, ASYNC_DIR, waitForAsyncResultFile, waitForAsyncState, tempDir, mockPi, } from "../support/async-execution-fixture.ts"; const skip = !isAsyncAvailable() ? "jiti not available" : undefined; const flagged = (status: AsyncStatusPayload) => status.steps?.[0]?.activityState === "needs_attention" && status.activityState === "needs_attention"; const unflagged = (status: AsyncStatusPayload) => status.steps?.[0]?.activityState !== "needs_attention" && status.activityState !== "needs_attention"; type RecordedEvent = { type?: string; reason?: string; from?: string; to?: string | null; index?: number; event?: { type?: string; reason?: string } }; function readEvents(id: string): RecordedEvent[] { const eventsPath = path.join(ASYNC_DIR, id, "events.jsonl"); if (!fs.existsSync(eventsPath)) return []; return fs.readFileSync(eventsPath, "utf-8").split("\n").filter(Boolean).map((line) => JSON.parse(line) as RecordedEvent); } const clearedEvents = (id: string, reason: string) => readEvents(id).filter((event) => event.type === "subagent.attention.cleared" && event.reason === reason); const attentionEvents = (id: string, reason: string) => readEvents(id).filter((event) => event.type === "subagent.control" && event.event?.type === "needs_attention" && event.event.reason === reason); /** An assistant turn that ends in a tool call, so the child keeps running (a `stop` turn would settle it). */ function midTurnMessage(text: string): object { const message = events.assistantMessage(text) as { message: { stopReason: string } }; message.message.stopReason = "toolUse"; return message; } function launch(id: string, controlConfig: Record) { return executeAsyncSingle(id, { agent: "worker", task: "Do the work", agentConfig: makeAgent("worker"), ctx: { pi: { events: { emit() {} } }, cwd: tempDir, currentSessionId: "session-1" }, artifactConfig: { enabled: false, includeInput: false, includeOutput: false, includeJsonl: false, includeMetadata: false, cleanupDays: 7 }, shareEnabled: false, sessionRoot: path.join(tempDir, "sessions"), maxSubagentDepth: 2, controlConfig: { enabled: true, needsAttentionAfterMs: 999_999, activeNoticeAfterMs: 100, failedToolAttemptsBeforeAttention: 3, notifyOn: ["needs_attention"], notifyChannels: ["event", "async", "intercom"], ...controlConfig, }, }); } describe("async needs_attention clearing (D12)", { skip: !available ? "pi packages not available" : undefined }, () => { installAsyncExecutionHooks(); it("clears when the open tool ends and flags again for a second long tool", { timeout: 45_000, skip }, async () => { const id = `async-attention-tool-${Date.now().toString(36)}`; const release = [1, 2, 3].map((n) => path.join(tempDir, `${id}.${n}`)); const release3 = release[2]!; mockPi.onCall({ steps: [ { jsonl: [{ type: "tool_execution_start", toolCallId: "bash-1", toolName: "bash", args: { command: "sleep 600" } }] }, { waitForPath: release[0]!, jsonl: [{ type: "tool_execution_end", toolCallId: "bash-1", toolName: "bash" }, events.toolResult("bash", "first done")] }, { waitForPath: release[1]!, jsonl: [{ type: "tool_execution_start", toolCallId: "bash-2", toolName: "bash", args: { command: "sleep 600" } }] }, { waitForPath: release3, jsonl: [{ type: "tool_execution_end", toolCallId: "bash-2", toolName: "bash" }, events.toolResult("bash", "second done"), events.assistantMessage("Done")] }, ], }); launch(id, {}); await waitForAsyncState(id, flagged); assert.equal(attentionEvents(id, "tool_open_threshold").length, 1); assert.equal(clearedEvents(id, "tool_open_threshold").length, 0, "still flagged while the tool is open"); fs.writeFileSync(release[0]!, ""); const recovered = await waitForAsyncState(id, unflagged); assert.notEqual(recovered.steps?.[0]?.activityState, "needs_attention"); assert.notEqual(recovered.activityState, "needs_attention"); const cleared = clearedEvents(id, "tool_open_threshold"); assert.equal(cleared.length, 1); assert.equal(cleared[0]?.from, "needs_attention"); assert.equal(cleared[0]?.index, 0); fs.writeFileSync(release[1]!, ""); await waitForAsyncState(id, flagged); assert.equal(attentionEvents(id, "tool_open_threshold").length, 2, "a recurrence is flagged and notified again"); fs.writeFileSync(release3, ""); await waitForAsyncResultFile(id); assert.equal(clearedEvents(id, "tool_open_threshold").length, 2); const finalStatus = JSON.parse(fs.readFileSync(path.join(ASYNC_DIR, id, "status.json"), "utf-8")) as AsyncStatusPayload; assert.equal(finalStatus.activityState, undefined); assert.notEqual(finalStatus.steps?.[0]?.activityState, "needs_attention"); }); it("still clears unconditionally on stop while the flag is raised", { timeout: 30_000, skip }, async () => { const id = `async-attention-stop-${Date.now().toString(36)}`; mockPi.onCall({ steps: [{ jsonl: [{ type: "tool_execution_start", toolCallId: "bash-1", toolName: "bash", args: { command: "sleep 600" } }] }, { waitForPath: path.join(tempDir, `${id}.never`), jsonl: [] }] }); launch(id, {}); await waitForAsyncState(id, flagged); deliverStopRequest({ asyncDir: path.join(ASYNC_DIR, id), source: "test" }); const stopped = await waitForAsyncState(id, (status) => status.state === "stopped"); assert.equal(stopped.activityState, undefined); assert.equal(stopped.steps?.[0]?.activityState, undefined); assert.equal(clearedEvents(id, "tool_open_threshold").length, 0, "terminal handling does not announce a recovery"); }); it("clears steering attention once a steer is delivered, but not when one is merely queued", { timeout: 30_000, skip }, async () => { const id = `async-attention-steer-${Date.now().toString(36)}`; const release = path.join(tempDir, `${id}.release`); mockPi.onCall({ failSteers: 1, steps: [{ waitForPath: release, jsonl: [events.assistantMessage("Done")] }] }); launch(id, {}); const asyncDir = path.join(ASYNC_DIR, id); await waitForAsyncState(id, (status) => status.steps?.[0]?.status === "running"); requestAsyncSteer(asyncDir, { message: "Please focus on the tests.", targetIndex: 0, source: "test" }); await waitForAsyncState(id, flagged); assert.equal(readEvents(id).filter((event) => event.type === "subagent.steer.failed").length, 1); requestAsyncSteer(asyncDir, { message: "Queue this for later.", mode: "follow_up", targetIndex: 0, source: "test" }); await waitForAsyncState(id, () => readEvents(id).some((event) => event.type === "subagent.steer.queued")); assert.equal(clearedEvents(id, "steering").length, 0, "a queued follow-up is not a delivery"); const stillFlagged = JSON.parse(fs.readFileSync(path.join(asyncDir, "status.json"), "utf-8")) as AsyncStatusPayload; assert.equal(stillFlagged.steps?.[0]?.activityState, "needs_attention"); requestAsyncSteer(asyncDir, { message: "Please focus on the tests, again.", targetIndex: 0, source: "test" }); await waitForAsyncState(id, unflagged); const cleared = clearedEvents(id, "steering"); assert.equal(cleared.length, 1); assert.equal(cleared[0]?.index, 0); fs.writeFileSync(release, ""); await waitForAsyncResultFile(id); }); it("clears a mutating-failure flag when a mutating tool succeeds again", { timeout: 30_000, skip }, async () => { const id = `async-attention-failures-${Date.now().toString(36)}`; const release = path.join(tempDir, `${id}.release`); const failingEdit = (n: number) => [ { type: "tool_execution_start", toolCallId: `edit-${n}`, toolName: "edit", args: { path: "src/a.ts" } }, { type: "tool_execution_end", toolCallId: `edit-${n}`, toolName: "edit" }, events.toolResult("edit", "Error: could not find the text to replace", true), ]; mockPi.onCall({ steps: [ { jsonl: [...failingEdit(1), ...failingEdit(2), ...failingEdit(3)] }, { waitForPath: release, jsonl: [ { type: "tool_execution_start", toolCallId: "edit-4", toolName: "edit", args: { path: "src/a.ts" } }, { type: "tool_execution_end", toolCallId: "edit-4", toolName: "edit" }, events.toolResult("edit", "Successfully replaced text in src/a.ts"), ] }, { jsonl: [events.assistantMessage("Done")] }, ], }); launch(id, { activeNoticeAfterMs: 999_999 }); await waitForAsyncState(id, flagged); assert.equal(attentionEvents(id, "tool_failures").length, 1); fs.writeFileSync(release, ""); await waitForAsyncState(id, unflagged); assert.equal(clearedEvents(id, "tool_failures").length, 1); await waitForAsyncResultFile(id); }); it("clears idle attention when the child becomes active again", { timeout: 45_000, skip }, async () => { const id = `async-attention-idle-${Date.now().toString(36)}`; const release = path.join(tempDir, `${id}.release`); mockPi.onCall({ steps: [ { jsonl: [midTurnMessage("Thinking about it")] }, { waitForPath: release, jsonl: [{ type: "tool_execution_start", toolCallId: "read-1", toolName: "read", args: { path: "README.md" } }, { type: "tool_execution_end", toolCallId: "read-1", toolName: "read" }] }, { delay: 3_500, jsonl: [events.assistantMessage("Done")] }, ], }); launch(id, { needsAttentionAfterMs: 2_000, activeNoticeAfterMs: 999_999 }); await waitForAsyncState(id, flagged, 20_000); assert.equal(attentionEvents(id, "idle").length, 1); fs.writeFileSync(release, ""); await waitForAsyncState(id, unflagged, 20_000); assert.equal(clearedEvents(id, "idle").length, 1); await waitForAsyncResultFile(id, 30_000); }); });