2026-02-08 13:57:26 +09:00
import type { ToolContextWithMetadata , OpencodeClient } from "./types"
2026-02-09 15:37:19 +09:00
import type { SessionMessage } from "./executor-types"
2026-03-03 00:33:10 +09:00
import { getDefaultSyncPollTimeoutMs , getTimingConfig } from "./timing"
2026-02-08 18:03:15 +09:00
import { log } from "../../shared/logger"
2026-02-16 18:20:19 +09:00
import { normalizeSDKResponse } from "../../shared"
2026-04-28 15:28:20 +09:00
import { extractErrorMessage } from "../../features/background-agent/error-classifier"
2026-02-08 13:57:26 +09:00
2026-02-09 15:37:19 +09:00
const NON_TERMINAL_FINISH_REASONS = new Set ( [ "tool-calls" , "unknown" ] )
2026-04-05 09:30:19 +09:00
const PENDING_TOOL_PART_TYPES = new Set ( [ "tool" , "tool_use" , "tool-call" ] )
2026-05-10 13:52:22 +09:00
const ACTIVE_SESSION_STATUSES = new Set ( [ "busy" , "retry" , "running" ] )
2026-02-09 15:37:19 +09:00
2026-03-09 18:55:12 +01:00
function wait ( milliseconds : number ) : Promise < void > {
const sharedBuffer = new SharedArrayBuffer ( Int32Array . BYTES_PER_ELEMENT )
const typedArray = new Int32Array ( sharedBuffer )
const result = Atomics . waitAsync ( typedArray , 0 , 0 , milliseconds )
return result . async ? result . value . then ( ( ) = > undefined ) : Promise . resolve ( )
}
function abortSyncSession ( client : OpencodeClient , sessionID : string , reason : string ) : void {
log ( "[task] Aborting sync session" , { sessionID , reason } )
void client . session . abort ( {
path : { id : sessionID } ,
} ) . catch ( ( error : unknown ) = > {
log ( "[task] Failed to abort sync session" , { sessionID , reason , error : String ( error ) } )
} )
}
2026-05-10 13:52:22 +09:00
function isActiveSessionStatus ( status : { type : string } | undefined ) : boolean {
return status !== undefined && ACTIVE_SESSION_STATUSES . has ( status . type )
}
2026-04-12 02:30:22 +09:00
async function fetchSessionMessages (
client : OpencodeClient ,
sessionID : string
) : Promise < SessionMessage [ ] > {
const messagesResult = await client . session . messages ( { path : { id : sessionID } } )
const rawData = ( messagesResult as { data? : unknown } ) ? . data ? ? messagesResult
return Array . isArray ( rawData ) ? ( rawData as SessionMessage [ ] ) : [ ]
}
2026-04-28 15:28:20 +09:00
function getTerminalSessionError ( messages : SessionMessage [ ] ) : string | null {
const lastAssistant = [ . . . messages ] . reverse ( ) . find ( ( msg ) = > msg . info ? . role === "assistant" )
2026-04-28 21:43:10 +09:00
const lastUser = [ . . . messages ] . reverse ( ) . find ( ( msg ) = > msg . info ? . role === "user" )
if ( lastUser ? . info ? . id && lastAssistant ? . info ? . id && lastAssistant . info . id <= lastUser . info . id ) {
return null
}
2026-04-28 15:28:20 +09:00
if ( ! lastAssistant ? . info || ! ( "error" in lastAssistant . info ) ) {
return null
}
const errorMessage = extractErrorMessage ( ( lastAssistant . info as { error? : unknown } ) . error )
return errorMessage && errorMessage . length > 0 ? errorMessage : "Session error"
}
2026-02-09 15:37:19 +09:00
export function isSessionComplete ( messages : SessionMessage [ ] ) : boolean {
let lastUser : SessionMessage | undefined
let lastAssistant : SessionMessage | undefined
for ( let i = messages . length - 1 ; i >= 0 ; i -- ) {
const msg = messages [ i ]
if ( ! lastAssistant && msg . info ? . role === "assistant" ) lastAssistant = msg
if ( ! lastUser && msg . info ? . role === "user" ) lastUser = msg
if ( lastUser && lastAssistant ) break
}
if ( ! lastAssistant ? . info ? . finish ) return false
if ( NON_TERMINAL_FINISH_REASONS . has ( lastAssistant . info . finish ) ) return false
2026-04-05 09:30:19 +09:00
if ( lastAssistant . parts ? . some ( ( part ) = > part . type && PENDING_TOOL_PART_TYPES . has ( part . type ) ) ) return false
2026-02-09 15:37:19 +09:00
if ( ! lastUser ? . info ? . id || ! lastAssistant ? . info ? . id ) return false
return lastUser . info . id < lastAssistant . info . id
}
2026-03-15 12:05:42 +08:00
const DEFAULT_MAX_ASSISTANT_TURNS = 300
2026-02-08 13:57:26 +09:00
export async function pollSyncSession (
ctx : ToolContextWithMetadata ,
client : OpencodeClient ,
input : {
sessionID : string
agentToUse : string
toastManager : { removeTask : ( id : string ) = > void } | null | undefined
taskId : string | undefined
2026-02-10 16:17:24 +09:00
anchorMessageCount? : number
2026-03-15 12:05:42 +08:00
maxAssistantTurns? : number
2026-02-27 13:54:50 +08:00
} ,
timeoutMs? : number
2026-02-08 13:57:26 +09:00
) : Promise < string | null > {
const syncTiming = getTimingConfig ( )
2026-03-03 00:33:10 +09:00
const maxPollTimeMs = Math . max ( timeoutMs ? ? getDefaultSyncPollTimeoutMs ( ) , 50 )
2026-03-15 12:05:42 +08:00
const maxTurns = input . maxAssistantTurns ? ? DEFAULT_MAX_ASSISTANT_TURNS
2026-02-08 13:57:26 +09:00
const pollStart = Date . now ( )
2026-05-10 13:52:22 +09:00
let inactiveStart = pollStart
2026-02-08 13:57:26 +09:00
let pollCount = 0
2026-02-10 19:09:22 +09:00
let timedOut = false
2026-03-15 12:05:42 +08:00
let assistantTurnCount = 0
let lastSeenAssistantId : string | undefined
2026-02-08 13:57:26 +09:00
2026-03-15 12:05:42 +08:00
log ( "[task] Starting poll loop" , { sessionID : input.sessionID , agentToUse : input.agentToUse , maxTurns } )
2026-02-08 13:57:26 +09:00
2026-05-10 13:52:22 +09:00
while ( true ) {
const inactiveElapsedMs = Date . now ( ) - inactiveStart
if ( inactiveElapsedMs >= maxPollTimeMs ) {
timedOut = true
break
}
2026-02-08 13:57:26 +09:00
if ( ctx . abort ? . aborted ) {
2026-05-10 14:55:06 +09:00
let finalMessages : SessionMessage [ ] | null = null
const abortFetchAttempts = 3
for ( let attempt = 1 ; attempt <= abortFetchAttempts ; attempt ++ ) {
try {
finalMessages = await fetchSessionMessages ( client , input . sessionID )
break
} catch ( error ) {
log ( "[task] Final messages fetch failed after abort, retrying" , {
sessionID : input.sessionID ,
attempt ,
maxAttempts : abortFetchAttempts ,
error : String ( error ) ,
} )
if ( attempt < abortFetchAttempts ) {
await wait ( syncTiming . POLL_INTERVAL_MS )
}
}
}
if ( finalMessages ) {
2026-04-12 02:30:22 +09:00
const hasNewMessages =
2026-05-10 14:55:06 +09:00
input . anchorMessageCount === undefined || finalMessages . length > input . anchorMessageCount
if ( hasNewMessages && isSessionComplete ( finalMessages ) ) {
2026-04-12 02:30:22 +09:00
log ( "[task] Abort detected after session already completed" , { sessionID : input.sessionID } )
return null
}
}
2026-02-08 13:57:26 +09:00
log ( "[task] Aborted by user" , { sessionID : input.sessionID } )
2026-03-09 18:55:12 +01:00
abortSyncSession ( client , input . sessionID , "parent_abort" )
2026-02-08 13:57:26 +09:00
if ( input . toastManager && input . taskId ) input . toastManager . removeTask ( input . taskId )
return ` Task aborted. \ n \ nSession ID: ${ input . sessionID } `
}
2026-03-09 18:55:12 +01:00
await wait ( syncTiming . POLL_INTERVAL_MS )
2026-02-08 13:57:26 +09:00
pollCount ++
2026-02-10 16:17:24 +09:00
let statusResult : { data? : Record < string , { type : string } > }
try {
statusResult = await client . session . status ( )
} catch ( error ) {
log ( "[task] Poll status fetch failed, retrying" , { sessionID : input.sessionID , error : String ( error ) } )
continue
}
2026-02-16 18:20:19 +09:00
const allStatuses = normalizeSDKResponse ( statusResult , { } as Record < string , { type : string } > )
2026-02-08 13:57:26 +09:00
const sessionStatus = allStatuses [ input . sessionID ]
if ( pollCount % 10 === 0 ) {
2026-02-09 15:37:19 +09:00
log ( "[task] Poll status" , {
2026-02-08 13:57:26 +09:00
sessionID : input.sessionID ,
pollCount ,
elapsed : Math.floor ( ( Date . now ( ) - pollStart ) / 1000 ) + "s" ,
2026-05-10 13:52:22 +09:00
inactiveElapsed : Math.floor ( inactiveElapsedMs / 1000 ) + "s" ,
2026-02-08 13:57:26 +09:00
sessionStatus : sessionStatus?.type ? ? "not_in_status" ,
} )
}
2026-05-10 13:52:22 +09:00
if ( isActiveSessionStatus ( sessionStatus ) ) {
inactiveStart = Date . now ( )
2026-02-08 13:57:26 +09:00
continue
}
2026-04-12 02:30:22 +09:00
let messages : SessionMessage [ ]
2026-02-10 16:17:24 +09:00
try {
2026-04-12 02:30:22 +09:00
messages = await fetchSessionMessages ( client , input . sessionID )
2026-02-10 16:17:24 +09:00
} catch ( error ) {
log ( "[task] Poll messages fetch failed, retrying" , { sessionID : input.sessionID , error : String ( error ) } )
continue
}
2026-04-12 02:30:22 +09:00
if ( input . anchorMessageCount !== undefined && messages . length <= input . anchorMessageCount ) {
2026-02-10 16:17:24 +09:00
continue
}
2026-02-08 13:57:26 +09:00
2026-04-28 15:28:20 +09:00
const sessionError = getTerminalSessionError ( messages )
if ( sessionError ) {
log ( "[task] Poll detected terminal session error" , { sessionID : input.sessionID , sessionError } )
return sessionError
}
2026-04-12 02:30:22 +09:00
if ( isSessionComplete ( messages ) ) {
2026-02-09 15:37:19 +09:00
log ( "[task] Poll complete - terminal finish detected" , { sessionID : input.sessionID , pollCount } )
break
2026-02-08 13:57:26 +09:00
}
2026-02-10 22:52:17 +09:00
2026-04-28 21:43:10 +09:00
// Count new assistant turns to circuit-break infinite loops
2026-04-12 02:30:22 +09:00
const lastAssistant = [ . . . messages ] . reverse ( ) . find ( ( m ) = > m . info ? . role === "assistant" )
2026-03-15 12:05:42 +08:00
if ( lastAssistant ? . info ? . id && lastAssistant . info . id !== lastSeenAssistantId ) {
lastSeenAssistantId = lastAssistant . info . id
assistantTurnCount ++
if ( assistantTurnCount >= maxTurns ) {
log ( "[task] Max assistant turns reached, aborting to prevent infinite loop" , {
sessionID : input.sessionID ,
assistantTurnCount ,
maxTurns ,
} )
abortSyncSession ( client , input . sessionID , "max_turns_exceeded" )
if ( input . toastManager && input . taskId ) input . toastManager . removeTask ( input . taskId )
return ` Task aborted: subagent exceeded ${ maxTurns } assistant turns without completing. This usually indicates an infinite tool-call loop. Session ID: ${ input . sessionID } `
}
}
2026-04-12 02:30:22 +09:00
const hasAssistantText = messages . some ( ( m ) = > {
2026-02-10 22:52:17 +09:00
if ( m . info ? . role !== "assistant" ) return false
const parts = m . parts ? ? [ ]
return parts . some ( ( p ) = > {
if ( p . type !== "text" && p . type !== "reasoning" ) return false
const text = ( p . text ? ? "" ) . trim ( )
return text . length > 0
} )
} )
if ( ! lastAssistant ? . info ? . finish && hasAssistantText ) {
log ( "[task] Poll complete - assistant text detected (fallback)" , {
sessionID : input.sessionID ,
pollCount ,
} )
break
}
2026-02-08 13:57:26 +09:00
}
2026-05-10 13:52:22 +09:00
if ( timedOut ) {
2026-05-10 15:54:29 +09:00
log ( "[task] Poll inactivity timeout reached" , { sessionID : input.sessionID , pollCount } )
2026-03-09 18:55:12 +01:00
abortSyncSession ( client , input . sessionID , "poll_timeout" )
2026-02-08 13:57:26 +09:00
}
2026-05-10 15:54:29 +09:00
return timedOut
? ` Poll inactivity timeout reached after ${ maxPollTimeMs } ms without active OpenCode status for session ${ input . sessionID } `
: null
2026-02-08 13:57:26 +09:00
}