/*--------------------------------------------------------------------------------------------- * Copyright (c) Microsoft Corporation. All rights reserved. *--------------------------------------------------------------------------------------------*/ import type { ChildProcess } from "node:child_process"; import { AssistantMessageEvent, CopilotSession, SessionEvent } from "../../../src"; const CHILD_SHUTDOWN_TIMEOUT_MS = 1_000; export async function stopChildProcess(child: ChildProcess): Promise { if (hasChildExited(child)) { return; } child.kill("SIGTERM"); if (await waitForChildExit(child, CHILD_SHUTDOWN_TIMEOUT_MS)) { return; } child.kill("SIGKILL"); if (!(await waitForChildExit(child, CHILD_SHUTDOWN_TIMEOUT_MS))) { throw new Error("Child process did not exit after SIGKILL"); } } function hasChildExited(child: ChildProcess): boolean { return child.exitCode !== null || child.signalCode !== null; } function waitForChildExit(child: ChildProcess, timeoutMs: number): Promise { if (hasChildExited(child)) { return Promise.resolve(true); } return new Promise((resolvePromise) => { let settled = false; const finish = (exited: boolean) => { if (settled) { return; } settled = true; clearTimeout(timeout); child.off("exit", onExit); resolvePromise(exited); }; const onExit = () => finish(true); const timeout = setTimeout(() => finish(false), timeoutMs); child.once("exit", onExit); if (hasChildExited(child)) { onExit(); } }); } export async function getFinalAssistantMessage( session: CopilotSession, { alreadyIdle = false }: { alreadyIdle?: boolean } = {} ): Promise { // Install the live subscription (via getFutureFinalResponse) before issuing the // existing-messages RPC so we don't miss events that arrive while that RPC is in flight. const futurePromise = getFutureFinalResponse(session); // We may end up returning from the existing-messages path; attach a noop handler so // the unawaited future-response rejection doesn't surface as an unhandled rejection. futurePromise.catch(() => {}); const existing = await getExistingFinalResponse(session, alreadyIdle); if (existing) { return existing; } return futurePromise; } async function getExistingFinalResponse( session: CopilotSession, alreadyIdle: boolean = false ): Promise { const messages = await session.getEvents(); const finalUserMessageIndex = messages.findLastIndex((m) => m.type === "user.message"); const currentTurnMessages = finalUserMessageIndex < 0 ? messages : messages.slice(finalUserMessageIndex); const currentTurnError = currentTurnMessages.find((m) => m.type === "session.error"); if (currentTurnError) { const error = new Error(currentTurnError.data.message); error.stack = currentTurnError.data.stack; throw error; } const sessionIdleMessageIndex = alreadyIdle ? currentTurnMessages.length : currentTurnMessages.findIndex((m) => m.type === "session.idle"); if (sessionIdleMessageIndex !== -1) { return currentTurnMessages .slice(0, sessionIdleMessageIndex) .findLast((m) => m.type === "assistant.message") as AssistantMessageEvent | undefined; } return undefined; } function getFutureFinalResponse(session: CopilotSession): Promise { return new Promise((resolve, reject) => { let finalAssistantMessage: AssistantMessageEvent | undefined; session.on((event) => { if (event.type === "assistant.message") { finalAssistantMessage = event; } else if (event.type === "session.idle") { if (!finalAssistantMessage) { reject( new Error("Received session.idle without a preceding assistant.message") ); } else { resolve(finalAssistantMessage); } } else if (event.type === "session.error") { const error = new Error(event.data.message); error.stack = event.data.stack; reject(error); } }); }); } export async function retry( message: string, fn: () => Promise, maxTries: number = 100, delay: number = 100 ) { let failedAttempts = 0; while (true) { try { await fn(); return; } catch (error: unknown) { failedAttempts++; if (failedAttempts >= maxTries) { throw new Error( `Failed to ${message} after ${maxTries} attempts\n${formatError(error)}` ); } await new Promise((resolve) => setTimeout(resolve, delay)); } } } export function formatError(error: unknown): string { if (error instanceof Error) { return String(error); } else if (typeof error === "object" && error !== null) { try { return JSON.stringify(error); } catch { return "[object with circular reference]"; } } else { return String(error); } } export function getNextEventOfType( session: CopilotSession, eventType: SessionEvent["type"] ): Promise { return new Promise((resolve, reject) => { const unsubscribe = session.on((event) => { if (event.type === eventType) { unsubscribe(); resolve(event); } else if (event.type === "session.error") { unsubscribe(); reject(new Error(`${event.data.message}\n${event.data.stack}`)); } }); }); } export async function waitForCondition( predicate: () => boolean | Promise, { timeoutMs = 30_000, intervalMs = 100, timeoutMessage = "Timed out waiting for condition.", }: { timeoutMs?: number; intervalMs?: number; timeoutMessage?: string } = {} ): Promise { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { if (await predicate()) { return; } await new Promise((resolve) => setTimeout(resolve, intervalMs)); } throw new Error(timeoutMessage); }