@@ -12,6 +12,7 @@ import { formatErrorMessage } from "../../infra/errors.js";
1212import { getGlobalHookRunner } from "../../plugins/hook-runner-global.js" ;
1313import { resolveProviderAuthProfileId } from "../../plugins/provider-runtime.js" ;
1414import { enqueueCommandInLane } from "../../process/command-queue.js" ;
15+ import type { CommandQueueEnqueueOptions } from "../../process/command-queue.types.js" ;
1516import { normalizeOptionalString } from "../../shared/string-coerce.js" ;
1617import { sanitizeForLog } from "../../terminal/ansi.js" ;
1718import { resolveUserPath } from "../../utils.js" ;
@@ -76,6 +77,7 @@ import {
7677 pickFallbackThinkingLevel ,
7778} from "../pi-embedded-helpers.js" ;
7879import { resolveProviderIdForAuth } from "../provider-auth-aliases.js" ;
80+ import { runAgentCleanupStep } from "../run-cleanup-timeout.js" ;
7981import { buildAgentRuntimeAuthPlan } from "../runtime-plan/auth.js" ;
8082import { buildAgentRuntimePlan } from "../runtime-plan/build.js" ;
8183import { ensureRuntimePluginsLoaded } from "../runtime-plugins.js" ;
@@ -159,8 +161,26 @@ import { createUsageAccumulator, mergeUsageIntoAccumulator } from "./usage-accum
159161type ApiKeyInfo = ResolvedProviderAuth ;
160162
161163const MAX_SAME_MODEL_IDLE_TIMEOUT_RETRIES = 1 ;
164+ const EMBEDDED_RUN_LANE_TIMEOUT_GRACE_MS = 30_000 ;
162165type EmbeddedRunAttemptForRunner = Awaited < ReturnType < typeof runEmbeddedAttemptWithBackend > > ;
163166
167+ function resolveEmbeddedRunLaneTimeoutMs ( timeoutMs : number ) : number | undefined {
168+ if ( ! Number . isFinite ( timeoutMs ) || timeoutMs <= 0 ) {
169+ return undefined ;
170+ }
171+ return Math . floor ( timeoutMs ) + EMBEDDED_RUN_LANE_TIMEOUT_GRACE_MS ;
172+ }
173+
174+ function withEmbeddedRunLaneTimeout (
175+ opts : CommandQueueEnqueueOptions | undefined ,
176+ laneTaskTimeoutMs : number | undefined ,
177+ ) : CommandQueueEnqueueOptions | undefined {
178+ if ( laneTaskTimeoutMs === undefined || opts ?. taskTimeoutMs !== undefined ) {
179+ return opts ;
180+ }
181+ return { ...opts , taskTimeoutMs : laneTaskTimeoutMs } ;
182+ }
183+
164184function normalizeEmbeddedRunAttemptResult (
165185 attempt : EmbeddedRunAttemptForRunner ,
166186) : EmbeddedRunAttemptForRunner {
@@ -292,10 +312,15 @@ export async function runEmbeddedPiAgent(
292312 }
293313 const sessionLane = resolveSessionLane ( params . sessionKey ?. trim ( ) || params . sessionId ) ;
294314 const globalLane = resolveGlobalLane ( params . lane ) ;
295- const enqueueGlobal =
296- params . enqueue ?? ( ( task , opts ) => enqueueCommandInLane ( globalLane , task , opts ) ) ;
297- const enqueueSession =
298- params . enqueue ?? ( ( task , opts ) => enqueueCommandInLane ( sessionLane , task , opts ) ) ;
315+ const laneTaskTimeoutMs = resolveEmbeddedRunLaneTimeoutMs ( params . timeoutMs ) ;
316+ const withLaneTimeout = ( opts ?: CommandQueueEnqueueOptions ) =>
317+ withEmbeddedRunLaneTimeout ( opts , laneTaskTimeoutMs ) ;
318+ const enqueueGlobal = < T > ( task : ( ) => Promise < T > , opts ?: CommandQueueEnqueueOptions ) =>
319+ params . enqueue
320+ ? params . enqueue ( task , withLaneTimeout ( opts ) )
321+ : enqueueCommandInLane ( globalLane , task , withLaneTimeout ( opts ) ) ;
322+ const enqueueSession = < T > ( task : ( ) => Promise < T > , opts ?: CommandQueueEnqueueOptions ) =>
323+ params . enqueue ? params . enqueue ( task , opts ) : enqueueCommandInLane ( sessionLane , task , opts ) ;
299324 const channelHint = params . messageChannel ?? params . messageProvider ;
300325 const resolvedToolResultFormat =
301326 params . toolResultFormat ??
@@ -2489,26 +2514,42 @@ export async function runEmbeddedPiAgent(
24892514 }
24902515 } finally {
24912516 forgetPromptBuildDrainCacheForRun ( params . runId ) ;
2492- await contextEngine . dispose ?.( ) ;
24932517 stopRuntimeAuthRefreshTimer ( ) ;
2518+ await runAgentCleanupStep ( {
2519+ runId : params . runId ,
2520+ sessionId : params . sessionId ,
2521+ step : "context-engine-dispose" ,
2522+ log,
2523+ cleanup : async ( ) => {
2524+ await contextEngine . dispose ?.( ) ;
2525+ } ,
2526+ } ) ;
24942527 if ( params . cleanupBundleMcpOnRunEnd === true ) {
2495- const onError = ( error : unknown , sessionId : string ) => {
2496- log . warn (
2497- `bundle-mcp cleanup failed after run for ${ sessionId } : ${ formatErrorMessage ( error ) } ` ,
2498- ) ;
2499- } ;
2500- const retiredBySessionKey = await retireSessionMcpRuntimeForSessionKey ( {
2501- sessionKey : params . sessionKey ,
2502- reason : "embedded-run-end" ,
2503- onError,
2528+ await runAgentCleanupStep ( {
2529+ runId : params . runId ,
2530+ sessionId : params . sessionId ,
2531+ step : "bundle-mcp-retire" ,
2532+ log,
2533+ cleanup : async ( ) => {
2534+ const onError = ( error : unknown , sessionId : string ) => {
2535+ log . warn (
2536+ `bundle-mcp cleanup failed after run for ${ sessionId } : ${ formatErrorMessage ( error ) } ` ,
2537+ ) ;
2538+ } ;
2539+ const retiredBySessionKey = await retireSessionMcpRuntimeForSessionKey ( {
2540+ sessionKey : params . sessionKey ,
2541+ reason : "embedded-run-end" ,
2542+ onError,
2543+ } ) ;
2544+ if ( ! retiredBySessionKey ) {
2545+ await retireSessionMcpRuntime ( {
2546+ sessionId : params . sessionId ,
2547+ reason : "embedded-run-end" ,
2548+ onError,
2549+ } ) ;
2550+ }
2551+ } ,
25042552 } ) ;
2505- if ( ! retiredBySessionKey ) {
2506- await retireSessionMcpRuntime ( {
2507- sessionId : params . sessionId ,
2508- reason : "embedded-run-end" ,
2509- onError,
2510- } ) ;
2511- }
25122553 }
25132554 }
25142555 } ) ;
0 commit comments