99import { resolveBootstrapWarningSignaturesSeen } from "../../agents/bootstrap-budget.js" ;
1010import { estimateMessagesTokens } from "../../agents/compaction.js" ;
1111import { classifyCompactionReason } from "../../agents/embedded-agent-runner/compact-reasons.js" ;
12+ import { isSignalTimeoutReason } from "../../agents/failover-error.js" ;
1213import { resolveAgentHarnessPolicy } from "../../agents/harness/policy.js" ;
1314import { ensureSelectedAgentHarnessPlugin } from "../../agents/harness/runtime-plugin.js" ;
1415import { runWithModelFallback } from "../../agents/model-fallback.js" ;
@@ -63,6 +64,10 @@ import type { ReplyOperation } from "./reply-run-registry.js";
6364import { incrementCompactionCount } from "./session-updates.js" ;
6465
6566type EmbeddedAgentRuntime = typeof import ( "../../agents/embedded-agent.js" ) ;
67+ type ScopedAbortSignal = {
68+ abortSignal : AbortSignal ;
69+ cleanup : ( ) => void ;
70+ } ;
6671type UpdateSessionEntryParams = {
6772 storePath : string ;
6873 sessionKey : string ;
@@ -85,6 +90,44 @@ function loadEmbeddedAgentRuntime(): Promise<EmbeddedAgentRuntime> {
8590 return embeddedAgentRuntimeLoader . load ( ) ;
8691}
8792
93+ function createPreflightCompactionAbortSignal ( operation : ReplyOperation ) : ScopedAbortSignal {
94+ const controller = new AbortController ( ) ;
95+ const cleanupHandlers : Array < ( ) => void > = [ ] ;
96+ const abortFrom = ( signal : AbortSignal , opts ?: { ignoreTimeoutReason ?: boolean } ) => {
97+ if ( controller . signal . aborted ) {
98+ return ;
99+ }
100+ if ( opts ?. ignoreTimeoutReason && isSignalTimeoutReason ( signal . reason ) ) {
101+ return ;
102+ }
103+ controller . abort ( signal . reason ) ;
104+ } ;
105+ const watch = ( signal : AbortSignal | undefined , opts ?: { ignoreTimeoutReason ?: boolean } ) => {
106+ if ( ! signal ) {
107+ return ;
108+ }
109+ if ( signal . aborted ) {
110+ abortFrom ( signal , opts ) ;
111+ return ;
112+ }
113+ const onAbort = ( ) => abortFrom ( signal , opts ) ;
114+ signal . addEventListener ( "abort" , onAbort , { once : true } ) ;
115+ cleanupHandlers . push ( ( ) => signal . removeEventListener ( "abort" , onAbort ) ) ;
116+ } ;
117+
118+ watch ( operation . abortSignal , { ignoreTimeoutReason : true } ) ;
119+ watch ( operation . explicitAbortSignal ) ;
120+
121+ return {
122+ abortSignal : controller . signal ,
123+ cleanup : ( ) => {
124+ for ( const cleanup of cleanupHandlers . splice ( 0 ) ) {
125+ cleanup ( ) ;
126+ }
127+ } ,
128+ } ;
129+ }
130+
88131async function compactEmbeddedAgentSessionDefault (
89132 ...args : Parameters < typeof import ( "../../agents/embedded-agent.js" ) . compactEmbeddedAgentSession >
90133) : Promise <
@@ -943,45 +986,49 @@ export async function runPreflightCompactionIfNeeded(params: {
943986 params . sessionKey ?? params . followupRun . run . sessionKey ,
944987 { storePath : params . storePath } ,
945988 ) ;
946- const result = await deps . compactEmbeddedAgentSession ( {
947- sessionId : entry . sessionId ,
948- sessionKey : params . sessionKey ,
949- sandboxSessionKey : params . runtimePolicySessionKey ,
950- allowGatewaySubagentBinding : true ,
951- messageChannel : params . followupRun . run . messageProvider ,
952- groupId : entry . groupId ?? params . followupRun . run . groupId ,
953- groupChannel : entry . groupChannel ?? params . followupRun . run . groupChannel ,
954- groupSpace : entry . space ?? params . followupRun . run . groupSpace ,
955- senderId : params . followupRun . run . senderId ,
956- senderName : params . followupRun . run . senderName ,
957- senderUsername : params . followupRun . run . senderUsername ,
958- senderE164 : params . followupRun . run . senderE164 ,
959- sessionFile : sessionFile ?? params . followupRun . run . sessionFile ,
960- workspaceDir : params . followupRun . run . workspaceDir ,
961- cwd : params . followupRun . run . cwd ,
962- agentDir : params . followupRun . run . agentDir ,
963- config : params . cfg ,
964- skillsSnapshot : entry . skillsSnapshot ?? params . followupRun . run . skillsSnapshot ,
965- provider : params . followupRun . run . provider ,
966- model : params . followupRun . run . model ,
967- authProfileId : params . followupRun . run . authProfileId ,
968- agentHarnessId :
969- entry . sessionId === params . followupRun . run . sessionId ? entry . agentHarnessId : undefined ,
970- thinkLevel : params . followupRun . run . thinkLevel ,
971- bashElevated : params . followupRun . run . bashElevated ,
972- ...( params . replyOperation . explicitAbortSignal
973- ? { abortSignal : params . replyOperation . explicitAbortSignal }
974- : { } ) ,
975- trigger : "budget" ,
976- force : true ,
977- forcePreflight : true ,
978- preflightRequired : true ,
979- preflightCompactionTrigger : compactionTrigger ,
980- deferOwningContextEngineCompaction : false ,
981- contextTokenBudget : contextWindowTokens ,
982- currentTokenCount : tokenCountForCompaction ?? freshPersistedTokens ,
983- ownerNumbers : params . followupRun . run . ownerNumbers ,
984- } ) ;
989+ const preflightAbort = createPreflightCompactionAbortSignal ( params . replyOperation ) ;
990+ let result : Awaited < ReturnType < typeof deps . compactEmbeddedAgentSession > > | undefined ;
991+ try {
992+ result = await deps . compactEmbeddedAgentSession ( {
993+ sessionId : entry . sessionId ,
994+ sessionKey : params . sessionKey ,
995+ sandboxSessionKey : params . runtimePolicySessionKey ,
996+ allowGatewaySubagentBinding : true ,
997+ messageChannel : params . followupRun . run . messageProvider ,
998+ groupId : entry . groupId ?? params . followupRun . run . groupId ,
999+ groupChannel : entry . groupChannel ?? params . followupRun . run . groupChannel ,
1000+ groupSpace : entry . space ?? params . followupRun . run . groupSpace ,
1001+ senderId : params . followupRun . run . senderId ,
1002+ senderName : params . followupRun . run . senderName ,
1003+ senderUsername : params . followupRun . run . senderUsername ,
1004+ senderE164 : params . followupRun . run . senderE164 ,
1005+ sessionFile : sessionFile ?? params . followupRun . run . sessionFile ,
1006+ workspaceDir : params . followupRun . run . workspaceDir ,
1007+ cwd : params . followupRun . run . cwd ,
1008+ agentDir : params . followupRun . run . agentDir ,
1009+ config : params . cfg ,
1010+ skillsSnapshot : entry . skillsSnapshot ?? params . followupRun . run . skillsSnapshot ,
1011+ provider : params . followupRun . run . provider ,
1012+ model : params . followupRun . run . model ,
1013+ authProfileId : params . followupRun . run . authProfileId ,
1014+ agentHarnessId :
1015+ entry . sessionId === params . followupRun . run . sessionId ? entry . agentHarnessId : undefined ,
1016+ thinkLevel : params . followupRun . run . thinkLevel ,
1017+ bashElevated : params . followupRun . run . bashElevated ,
1018+ abortSignal : preflightAbort . abortSignal ,
1019+ trigger : "budget" ,
1020+ force : true ,
1021+ forcePreflight : true ,
1022+ preflightRequired : true ,
1023+ preflightCompactionTrigger : compactionTrigger ,
1024+ deferOwningContextEngineCompaction : false ,
1025+ contextTokenBudget : contextWindowTokens ,
1026+ currentTokenCount : tokenCountForCompaction ?? freshPersistedTokens ,
1027+ ownerNumbers : params . followupRun . run . ownerNumbers ,
1028+ } ) ;
1029+ } finally {
1030+ preflightAbort . cleanup ( ) ;
1031+ }
9851032
9861033 if ( ! result ?. ok ) {
9871034 const reason = result ?. reason ?? "not_compacted" ;
0 commit comments