fix: use pi settled lifecycle event
This commit is contained in:
parent
1f2487554b
commit
44f2211608
|
|
@ -66,14 +66,17 @@ type Handler = (event: unknown, context: unknown) => unknown;
|
|||
|
||||
function createExtensionHarness() {
|
||||
const handlers = new Map<string, Handler>();
|
||||
const eventHandlers = new Map<string, Handler>();
|
||||
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, Handler>): 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<void> {
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -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<void> {
|
|||
}
|
||||
}
|
||||
|
||||
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<typeof setTimeout> | undefined;
|
||||
let retryTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
let rootSession = false;
|
||||
|
||||
function clearTimer(timer: ReturnType<typeof setTimeout> | 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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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"));
|
||||
|
|
|
|||
Loading…
Reference in New Issue