diff --git a/src/integration/assets/herdr-agent-state.test.ts b/src/integration/assets/herdr-agent-state.test.ts index 50f1fc52..1dadbd5f 100644 --- a/src/integration/assets/herdr-agent-state.test.ts +++ b/src/integration/assets/herdr-agent-state.test.ts @@ -66,14 +66,17 @@ type Handler = (event: unknown, context: unknown) => unknown; function createExtensionHarness() { const handlers = new Map(); + const eventHandlers = new Map(); return { handlers, + eventHandlers, pi: { on(event: string, handler: Handler) { handlers.set(event, handler); }, events: { - on() { + on(event: string, handler: Handler) { + eventHandlers.set(event, handler); return () => {}; }, }, @@ -224,6 +227,62 @@ for (const integration of integrations) { }); } +test("Pi reports idle only after the agent settles", async () => { + const requests = await startRecordingServer("pi-settled"); + const { handlers, pi } = createExtensionHarness(); + const { default: install } = await importFresh("./pi/herdr-agent-state.ts"); + install(pi); + + expect(completionHandlers(handlers)).toEqual(["agent_settled"]); + let idle = true; + const context = piContext(() => idle); + await handlers.get("session_start")?.({ reason: "startup" }, context); + await waitFor(() => requestStates(requests).length === 1); + + idle = false; + handlers.get("agent_start")?.({}, context); + await waitFor(() => requestStates(requests).length === 2); + expect(requestStates(requests)).toEqual(["idle", "working"]); + expect(handlers.has("agent_end")).toBe(false); + + const requestCountBeforeStaleSettlement = requests.length; + handlers.get("agent_settled")?.({}, context); + await Bun.sleep(25); + expect(requests).toHaveLength(requestCountBeforeStaleSettlement); + expect(requestStates(requests)).toEqual(["idle", "working"]); + + idle = true; + handlers.get("agent_settled")?.({}, context); + await waitFor(() => requestStates(requests).length === 3); + expect(requestStates(requests)).toEqual(["idle", "working", "idle"]); +}); + +test("Pi settlement preserves explicit blocked-state precedence", async () => { + const requests = await startRecordingServer("pi-settled-blocked"); + const { eventHandlers, handlers, pi } = createExtensionHarness(); + const { default: install } = await importFresh("./pi/herdr-agent-state.ts"); + install(pi); + + let idle = true; + const context = piContext(() => idle); + await handlers.get("session_start")?.({ reason: "startup" }, context); + await waitFor(() => requestStates(requests).length === 1); + idle = false; + handlers.get("agent_start")?.({}, context); + await waitFor(() => requestStates(requests).length === 2); + eventHandlers.get("herdr:blocked")?.({ active: true, label: "approval" }, context); + await waitFor(() => requestStates(requests).length === 3); + + idle = true; + handlers.get("agent_settled")?.({}, context); + await Bun.sleep(25); + expect(requestStates(requests)).toEqual(["idle", "working", "blocked"]); + + eventHandlers.get("herdr:blocked")?.({ active: false }, context); + await waitFor(() => requestStates(requests).length === 4); + expect(requestStates(requests)).toEqual(["idle", "working", "blocked", "idle"]); +}); + test("Pi reports the session replacement source", async () => { const requests = await startRecordingServer("pi-session-source"); const { handlers, pi } = createExtensionHarness(); @@ -447,6 +506,35 @@ test("Pi retries working state after an unanswered socket attempt", async () => expect(reportedWorking()).toBe(true); }); +function completionHandlers(handlers: Map): string[] { + return ["agent_end", "agent_settled"].filter((event) => handlers.has(event)); +} + +function piContext(isIdle: () => boolean) { + return { + hasUI: true, + isIdle, + sessionManager: { + getSessionFile: () => undefined, + getSessionId: () => undefined, + }, + }; +} + +function requestStates(requests: unknown[]): unknown[] { + return requests + .filter((request) => isRecord(request) && request.method === "pane.report_agent") + .map(requestState); +} + +async function waitFor(predicate: () => boolean, timeoutMs = 1_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline && !predicate()) { + await Bun.sleep(5); + } + expect(predicate()).toBe(true); +} + function requestState(request: unknown): unknown { if (!isRecord(request) || !isRecord(request.params)) { return undefined; diff --git a/src/integration/assets/pi/herdr-agent-state.ts b/src/integration/assets/pi/herdr-agent-state.ts index 218cede9..caa50ed9 100644 --- a/src/integration/assets/pi/herdr-agent-state.ts +++ b/src/integration/assets/pi/herdr-agent-state.ts @@ -61,10 +61,6 @@ type QueuedState = { seq: number; }; -const idleDebounceMs = parseDurationEnv("HERDR_PI_IDLE_DEBOUNCE_MS", 250); -const retryGraceMs = parseDurationEnv("HERDR_PI_RETRY_GRACE_MS", 2500); -const retryableErrorPattern = - /overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|network.?error|connection.?error|connection.?refused|connection.?lost|websocket.?closed|websocket.?error|other side closed|fetch failed|upstream.?connect|reset before headers|socket hang up|ended without|http2 request did not get a response|timed? out|timeout|terminated|retry delay/i; let reportSeq = Date.now() * 1000; let currentAgentSessionId: string | undefined; let currentAgentSessionPath: string | undefined; @@ -74,18 +70,6 @@ function nextReportSeq(): number { return reportSeq; } -function parseDurationEnv(name: string, fallback: number): number { - const raw = process.env[name]; - if (!raw) { - return fallback; - } - const parsed = Number.parseInt(raw, 10); - if (!Number.isFinite(parsed) || parsed < 0) { - return fallback; - } - return parsed; -} - function updateSessionRef(ctx: any): void { try { const file = ctx?.sessionManager?.getSessionFile?.(); @@ -211,74 +195,23 @@ async function drainStateQueue(): Promise { } } -function lastAssistantMessage(messages: unknown[]): any | undefined { - for (let i = messages.length - 1; i >= 0; i -= 1) { - const message = messages[i] as any; - if (message?.role === "assistant") { - return message; - } - } - return undefined; -} - -function retryableErrorMessage(event: any): string | undefined { - const messages = Array.isArray(event?.messages) ? event.messages : []; - const assistant = lastAssistantMessage(messages); - if (assistant?.stopReason !== "error") { - return undefined; - } - - const errorMessage = String(assistant.errorMessage ?? ""); - if (!retryableErrorPattern.test(errorMessage)) { - return undefined; - } - return errorMessage || "retryable provider error"; -} - export default function (pi) { if (!enabled()) { return; } let agentActive = false; - let retryHoldActive = false; - let failureBlocked = false; - let failureMessage: string | undefined; let blockedCount = 0; let blockedMessage: string | undefined; let lastState: AgentState | undefined; let lastMessage: string | undefined; - let idleTimer: ReturnType | undefined; - let retryTimer: ReturnType | undefined; let rootSession = false; - function clearTimer(timer: ReturnType | undefined) { - if (timer) { - clearTimeout(timer); - } - } - - function clearPendingTimers() { - clearTimer(idleTimer); - clearTimer(retryTimer); - idleTimer = undefined; - retryTimer = undefined; - } - - function clearFailureState() { - retryHoldActive = false; - failureBlocked = false; - failureMessage = undefined; - } - function desiredState() { if (blockedCount > 0) { return { state: "blocked" as const, message: blockedMessage }; } - if (failureBlocked) { - return { state: "blocked" as const, message: failureMessage }; - } - if (agentActive || retryHoldActive) { + if (agentActive) { return { state: "working" as const, message: undefined }; } return { state: "idle" as const, message: undefined }; @@ -294,32 +227,6 @@ export default function (pi) { queueState(next.state, next.message); } - function scheduleIdle() { - clearPendingTimers(); - clearFailureState(); - idleTimer = setTimeout(() => { - idleTimer = undefined; - publishState(); - }, idleDebounceMs); - idleTimer.unref?.(); - } - - function holdForRetry(message: string) { - clearPendingTimers(); - retryHoldActive = true; - failureBlocked = false; - failureMessage = message; - publishState(); - - retryTimer = setTimeout(() => { - retryTimer = undefined; - retryHoldActive = false; - failureBlocked = true; - publishState(); - }, retryGraceMs); - retryTimer.unref?.(); - } - pi.events.on("herdr:blocked", (data) => { if (!rootSession) { return; @@ -333,7 +240,6 @@ export default function (pi) { return; } - clearPendingTimers(); blockedCount += 1; blockedMessage = data.label; publishState(); @@ -357,39 +263,23 @@ export default function (pi) { } updateSessionRef(ctx); void reportSession(); - clearPendingTimers(); - clearFailureState(); agentActive = true; publishState(); }); - pi.on("agent_end", (event) => { - if (!rootSession) { - return; - } - if (!agentActive) { - // Pi can emit duplicate/late end events while auto-retry is already - // holding the pane in Working. Do not let an unqualified duplicate end - // cancel the retry hold and publish a false Idle. + pi.on("agent_settled", (_event, ctx) => { + if (!rootSession || ctx?.isIdle?.() !== true) { return; } agentActive = false; - - const retryableMessage = retryableErrorMessage(event); - if (retryableMessage) { - holdForRetry(retryableMessage); - return; - } - - scheduleIdle(); + publishState(); }); pi.on("session_shutdown", async (event) => { if (!rootSession) { return; } - clearPendingTimers(); if (shouldReleaseOnSessionShutdown(event)) { await releaseAgent(); } diff --git a/src/integration/tests.rs b/src/integration/tests.rs index 786e67fb..e2cde978 100644 --- a/src/integration/tests.rs +++ b/src/integration/tests.rs @@ -2640,7 +2640,7 @@ fn bundled_integration_assets_report_session_refs() { assert!(PI_EXTENSION_ASSET.contains("pane.report_agent_session")); assert!(PI_EXTENSION_ASSET.contains("pane.report_agent\"")); assert!(PI_EXTENSION_ASSET.contains("pi.on(\"agent_start\"")); - assert!(PI_EXTENSION_ASSET.contains("pi.on(\"agent_end\"")); + assert!(PI_EXTENSION_ASSET.contains("pi.on(\"agent_settled\"")); assert!(PI_EXTENSION_ASSET.contains("pane.release_agent")); assert!(PI_EXTENSION_ASSET.contains("pi.on(\"session_shutdown\"")); assert!(OMP_EXTENSION_ASSET.contains("agent_session_path"));