From 99712e4f06756d3736dd1a0c8909ba8dc7ce958e Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 08:41:07 +0530 Subject: [PATCH 01/11] update opentoolcall closer --- .../capabilities/AgentContextProcessor.ts | 5 +- .../src/core/runtime/AgentThread.ts | 19 ++- .../src/core/runtime/OpenToolCallCloser.ts | 134 ++++++++++++++---- 3 files changed, 127 insertions(+), 31 deletions(-) diff --git a/packages/trueforge-core/src/core/capabilities/AgentContextProcessor.ts b/packages/trueforge-core/src/core/capabilities/AgentContextProcessor.ts index 24b37bf01..ae3b0ff44 100644 --- a/packages/trueforge-core/src/core/capabilities/AgentContextProcessor.ts +++ b/packages/trueforge-core/src/core/capabilities/AgentContextProcessor.ts @@ -31,7 +31,10 @@ export type AgentContextProcessorOutput = | AgentContextProcessorCapabilityState; export interface PreSendContextProcessor { - processPreSend(execution: Readonly): AsyncIterable; + processPreSend( + execution: Readonly, + options: { userMessageIncoming: boolean }, + ): AsyncIterable; } // NOTE: Saved in Redis, Saved in AgentThread in memory. Persisted accross Agent Loop. diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.ts b/packages/trueforge-core/src/core/runtime/AgentThread.ts index 3e2c6785b..4402fac4c 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.ts @@ -184,7 +184,7 @@ function lastAssistantInContext(context: ContextMessage[]): InternalEnrichedAssi // since it mutates (deletes) the closable ids. function getUnclosableOpenToolCallIds(context: ContextMessage[], openToolCallIds: Set): Set { const blockingOpenToolCallIds = new Set(openToolCallIds); - for (const id of getClosableOpenToolCallIds(context)) { + for (const id of getClosableOpenToolCallIds({ context, userMessageIncoming: false })) { blockingOpenToolCallIds.delete(id); } return blockingOpenToolCallIds; @@ -599,7 +599,9 @@ export class AgentThread { this.contextBusy = true; try { - for await (const event of this.executeContextProcessors('preSend')) { + for await (const event of this.executeContextProcessors('preSend', { + userMessageIncoming: messages.some(isInputUserMessage), + })) { yield event; } this.preSendRanThisTurn = true; @@ -876,7 +878,10 @@ export class AgentThread { }; } - private executeContextProcessors(hook: 'preSend'): AsyncGenerator; + private executeContextProcessors( + hook: 'preSend', + options: { userMessageIncoming: boolean }, + ): AsyncGenerator; private executeContextProcessors( hook: 'preLLM' | 'postToolCall', ): AsyncGenerator< @@ -889,6 +894,7 @@ export class AgentThread { >; private async *executeContextProcessors( hook: 'preSend' | 'preLLM' | 'postToolCall', + options?: { userMessageIncoming: boolean }, ): AsyncGenerator< | ThreadOverwriteContextEvent | AgentThreadAppendContext @@ -902,7 +908,10 @@ export class AgentThread { ) => AsyncIterable)[]; switch (hook) { case 'preSend': - processors = this.preSendContextProcessors.map(p => p.processPreSend.bind(p)); + processors = this.preSendContextProcessors.map( + p => (execution: Readonly) => + p.processPreSend(execution, { userMessageIncoming: options?.userMessageIncoming === true }), + ); break; case 'preLLM': processors = this.preLLMContextProcessors.map(p => p.processPreLLM.bind(p)); @@ -1316,7 +1325,7 @@ export class AgentThread { } if (!this.preSendRanThisTurn) { - for await (const event of this.executeContextProcessors('preSend')) { + for await (const event of this.executeContextProcessors('preSend', { userMessageIncoming: false })) { yield event; } } diff --git a/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts b/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts index 972717912..3f81ae9f7 100644 --- a/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts +++ b/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts @@ -8,80 +8,164 @@ * conversation can continue. * * Uses is_thread_creation from InternalToolCallInfo to identify sub-agent - * tool calls (which should remain open); does NOT use SUB_AGENTS_SERVER_ID. - * Older persisted contexts without is_thread_creation follow the ordinary - * non-thread-creation path (no compatibility fallback). + * tool calls. Older persisted contexts without is_thread_creation follow the + * ordinary non-thread-creation path (no compatibility fallback). + * + * Resume / user-action batches close only dangling regular calls. A user + * message also closes pending approval, client-side, and thread-creation + * calls so the new message can sit after a complete tool-call/response pair. */ import type { AgentContextProcessorAppendContext, AgentThreadExecutionContext, PreSendContextProcessor, } from '../capabilities/AgentContextProcessor'; -import type { InternalEnrichedAssistantMessage, LLMToolMessage } from '../llm/LLMTypes'; +import { EventType, newEventId, type ToolResponseEvent } from '../events/schema'; +import type { InternalEnrichedAssistantMessage, InternalEnrichedToolCall, LLMToolMessage } from '../llm/LLMTypes'; import type { ContextMessage } from './AgentThread.types'; import { InternalEventType } from './AgentThread.types'; import { mergeCurrentContextUsage } from './contextUsage'; import { estimateTokensForContextMessages, isLLMContextMessage } from './contextUtils'; -const DUMMY_TOOL_MESSAGE_CONTENT = JSON.stringify({ +const DANGLING_TOOL_MESSAGE_CONTENT = JSON.stringify({ error: 'Tool call was not executed. Please retry this tool call.', }); -export function getClosableOpenToolCallIds(context: ContextMessage[]): Set { - const lastIdx = context.findLastIndex( +const USER_ACTION_TOOL_MESSAGE_CONTENT = JSON.stringify({ + error: + 'Tool call was not executed: the user sent a new message before it was resolved. Do not retry unless the user asks.', +}); + +const THREAD_CREATION_TOOL_MESSAGE_CONTENT = 'Sub-agent was cancelled because the user sent a new message.'; + +export type ClosableOpenToolCallKind = 'dangling' | 'user_action' | 'thread_creation'; + +export interface ClosableOpenToolCall { + tool_call_id: string; + close_kind: ClosableOpenToolCallKind; +} + +function closeKindForToolCall(toolCall: InternalEnrichedToolCall): ClosableOpenToolCallKind { + if (toolCall.tool_info.is_thread_creation === true) { + return 'thread_creation'; + } + if (toolCall.tool_info.is_approval_required === true || toolCall.tool_info.is_client_side === true) { + return 'user_action'; + } + return 'dangling'; +} + +function contentForCloseKind(closeKind: ClosableOpenToolCallKind): string { + switch (closeKind) { + case 'dangling': + return DANGLING_TOOL_MESSAGE_CONTENT; + case 'user_action': + return USER_ACTION_TOOL_MESSAGE_CONTENT; + case 'thread_creation': + return THREAD_CREATION_TOOL_MESSAGE_CONTENT; + } +} + +export function getClosableOpenToolCalls(input: { + context: ContextMessage[]; + userMessageIncoming: boolean; +}): ClosableOpenToolCall[] { + const lastIdx = input.context.findLastIndex( (msg): msg is InternalEnrichedAssistantMessage => isLLMContextMessage(msg) && msg.role === 'assistant' && !!msg.tool_calls?.length, ); if (lastIdx === -1) { - return new Set(); + return []; } - const lastAssistant = context[lastIdx]; + const lastAssistant = input.context[lastIdx]; if (lastAssistant === undefined || !isLLMContextMessage(lastAssistant) || lastAssistant.role !== 'assistant') { throw new Error('Unreachable'); } if (!lastAssistant.tool_calls) { - return new Set(); + return []; } if ( + !input.userMessageIncoming && lastAssistant.tool_calls.some( tc => tc.tool_info.is_approval_required === true || tc.tool_info.is_client_side === true, ) ) { - return new Set(); + return []; } - const requestedIds = new Set( - lastAssistant.tool_calls.filter(tc => tc.tool_info.is_thread_creation !== true).map(tc => tc.id), - ); - - const messagesAfterAssistant = context.slice(lastIdx + 1); - for (const msg of messagesAfterAssistant) { + const resolvedIds = new Set(); + for (const msg of input.context.slice(lastIdx + 1)) { if (isLLMContextMessage(msg) && msg.role === 'tool') { - requestedIds.delete(msg.tool_call_id); + resolvedIds.add(msg.tool_call_id); + } + } + + const closable: ClosableOpenToolCall[] = []; + for (const toolCall of lastAssistant.tool_calls) { + if (resolvedIds.has(toolCall.id)) { + continue; + } + const close_kind = closeKindForToolCall(toolCall); + if (!input.userMessageIncoming && close_kind === 'thread_creation') { + continue; } + closable.push({ tool_call_id: toolCall.id, close_kind }); } + return closable; +} - return requestedIds; +export function getClosableOpenToolCallIds(input: { + context: ContextMessage[]; + userMessageIncoming: boolean; +}): Set { + return new Set(getClosableOpenToolCalls(input).map(call => call.tool_call_id)); +} + +// we are closing open tool calls synthetically, the subscriber needs to understand +// the tool calls were closed. +function toToolResponseEvent(input: { threadId: string; toolCallId: string; content: string }): ToolResponseEvent { + return { + type: EventType.TOOL_RESPONSE, + id: newEventId(), + created_at: new Date().toISOString(), + thread_id: input.threadId, + tool_call_id: input.toolCallId, + content: input.content, + }; } export class OpenToolCallCloser implements PreSendContextProcessor { // eslint-disable-next-line @typescript-eslint/require-await -- async *: AsyncIterable contract; body is sync async *processPreSend( execution: Readonly, + options: { userMessageIncoming: boolean }, ): AsyncGenerator { - const requestedIds = getClosableOpenToolCallIds(execution.context); - if (requestedIds.size === 0) { + const closable = getClosableOpenToolCalls({ + context: execution.context, + userMessageIncoming: options.userMessageIncoming, + }); + if (closable.length === 0) { return; } - const dummyToolMessages: LLMToolMessage[] = [...requestedIds].map(toolCallId => ({ + const dummyToolMessages: LLMToolMessage[] = closable.map(call => ({ role: 'tool', - tool_call_id: toolCallId, - content: DUMMY_TOOL_MESSAGE_CONTENT, + tool_call_id: call.tool_call_id, + content: contentForCloseKind(call.close_kind), })); + const output: ToolResponseEvent[] = closable + .filter(call => call.close_kind !== 'dangling') + .map(call => + toToolResponseEvent({ + threadId: execution.threadId, + toolCallId: call.tool_call_id, + content: contentForCloseKind(call.close_kind), + }), + ); + const currentContextUsage = mergeCurrentContextUsage( execution.currentContextUsage, estimateTokensForContextMessages(dummyToolMessages), @@ -90,7 +174,7 @@ export class OpenToolCallCloser implements PreSendContextProcessor { yield { type: InternalEventType.AGENT_CONTEXT_APPEND, context: dummyToolMessages, - output: [], + output, current_context_usage: currentContextUsage, }; } From f020d5017e8dd5ca3f3179f2e6ec40a317147e6a Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 08:51:49 +0530 Subject: [PATCH 02/11] relax validation --- .../src/core/runtime/AgentThread.ts | 37 ++++--------------- 1 file changed, 7 insertions(+), 30 deletions(-) diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.ts b/packages/trueforge-core/src/core/runtime/AgentThread.ts index 4402fac4c..7f500adc2 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.ts @@ -94,7 +94,7 @@ import { } from './contextUtils'; import { DeferredTool } from './DeferredTool'; import { createEmptyAgentThreadMetrics, updateMetricsFromUsage, type AgentThreadMetrics } from './metrics'; -import { getClosableOpenToolCallIds, OpenToolCallCloser } from './OpenToolCallCloser'; +import { OpenToolCallCloser } from './OpenToolCallCloser'; import { isEmptyMessageContent, processAgentUserInput, type AgentInputUserMessage } from './UserInputMessage'; const DEFAULT_ITERATION_LIMIT = 25; @@ -177,19 +177,6 @@ function lastAssistantInContext(context: ContextMessage[]): InternalEnrichedAssi ); } -// Open tool calls that block a new user message: the open set minus those OpenToolCallCloser will -// auto-close during preSend. Lets a user message resume a thread whose only open calls are dangling -// regular tool calls (which the closer repairs), while still blocking on approval/client-side/ -// sub-agent calls that genuinely need resolution. Takes the already-computed open set; copies it -// since it mutates (deletes) the closable ids. -function getUnclosableOpenToolCallIds(context: ContextMessage[], openToolCallIds: Set): Set { - const blockingOpenToolCallIds = new Set(openToolCallIds); - for (const id of getClosableOpenToolCallIds({ context, userMessageIncoming: false })) { - blockingOpenToolCallIds.delete(id); - } - return blockingOpenToolCallIds; -} - function buildMCPInitializeEvent(initInfo: MCPServerInitInfo[], threadId: string): MCPInitializeEvent { return { type: EventType.MCP_INITIALIZE, @@ -260,17 +247,10 @@ function buildModelMessageEvent({ return event; } -function validateUserMessage( - message: { content: AgentInputUserMessage['content'] }, - blockingOpenToolCallIds: Set, - index: number, -): void { +function validateUserMessage(message: { content: AgentInputUserMessage['content'] }, index: number): void { if (isEmptyMessageContent(message.content)) { throw new InvalidAgentSendInputError(`messages[${String(index)}] user message has empty content`); } - if (blockingOpenToolCallIds.size > 0) { - throw new InvalidAgentSendInputError('user message cannot be sent while approvals or questions are pending'); - } } function validateToolMessage( @@ -316,11 +296,9 @@ function validateInputMessageTypesGivenContext( context: ContextMessage[], messages: AgentThreadRuntimeSendInput[], ): void { + const hasUserMessage = messages.some(isInputUserMessage); // Full open set: validates incoming tool responses and dedupes within the batch. const openToolCallIds = getOpenToolCallIds(context); - // Subset that blocks a fresh user message: excludes calls OpenToolCallCloser will auto-close - // during preSend, so a dangling regular tool call doesn't reject a user message it will repair. - const blockingOpenToolCallIds = getUnclosableOpenToolCallIds(context, openToolCallIds); const pendingApprovalIds = new Set(getPendingApprovalToolCalls(context).map(tc => tc.id)); const pendingClientSideIds = new Set(getPendingClientSideToolCalls(context).map(tc => tc.id)); @@ -333,11 +311,10 @@ function validateInputMessageTypesGivenContext( validateApprovalMessage(m, pendingApprovalIds, i); pendingApprovalIds.delete(m.tool_call_id); } else if (isInputUserMessage(m)) { - validateUserMessage(m, blockingOpenToolCallIds, i); + validateUserMessage(m, i); } else if (isClientSideToolResponseMessage(m) || isLLMToolMessage(m)) { validateToolMessage(m, openToolCallIds, i); openToolCallIds.delete(m.tool_call_id); - blockingOpenToolCallIds.delete(m.tool_call_id); pendingClientSideIds.delete(m.tool_call_id); } else { const _exhaustive: never = m; @@ -347,9 +324,9 @@ function validateInputMessageTypesGivenContext( } } - // A send for a thread awaiting user input must resolve every pending approval and client-side - // tool call in the same batch; any left unresolved (including an empty batch) is a blocker. - if (pendingApprovalIds.size > 0 || pendingClientSideIds.size > 0) { + // User messages interrupt pending approvals / client-side calls (OpenToolCallCloser + // synthesizes the missing responses). Empty and action batches still must resolve them all. + if (!hasUserMessage && (pendingApprovalIds.size > 0 || pendingClientSideIds.size > 0)) { const missing = [...pendingApprovalIds, ...pendingClientSideIds]; throw new InvalidAgentSendInputError( `Send batch must resolve all pending tool calls awaiting user input. Missing: ${missing.join(', ')}`, From a305ca29501f8b26e13dedd1be8a36077febdd96 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 09:02:58 +0530 Subject: [PATCH 03/11] nit --- .../src/agent-session/TurnHandle.ts | 4 +- .../trueforge-core/src/core/events/schema.ts | 9 +- .../src/core/runtime/AgentThread.types.ts | 6 +- .../core/runtime/AgentThreadOrchestrator.ts | 89 +++++++++++++------ .../core/runtime/openToolCallCloser.test.ts | 20 +++-- 5 files changed, 90 insertions(+), 38 deletions(-) diff --git a/packages/trueforge-core/src/agent-session/TurnHandle.ts b/packages/trueforge-core/src/agent-session/TurnHandle.ts index f0bb2be03..5932c3cbd 100644 --- a/packages/trueforge-core/src/agent-session/TurnHandle.ts +++ b/packages/trueforge-core/src/agent-session/TurnHandle.ts @@ -48,7 +48,9 @@ function toThreadDoneEvent(event: InternalThreadDoneEvent): ThreadDoneEvent { const state = event.status === 'error' ? { status: 'error' as const, error: event.error, ...(event.output && { output: event.output }) } - : { status: 'done' as const, output: event.output }; + : event.status === 'cancelled' + ? { status: 'cancelled' as const } + : { status: 'done' as const, output: event.output }; return { type: HarnessEventType.THREAD_DONE, id: newEventId(), diff --git a/packages/trueforge-core/src/core/events/schema.ts b/packages/trueforge-core/src/core/events/schema.ts index bdced8eda..2ee8e2714 100644 --- a/packages/trueforge-core/src/core/events/schema.ts +++ b/packages/trueforge-core/src/core/events/schema.ts @@ -219,8 +219,14 @@ export const ThreadStateErrorSchema = z }) .openapi('ThreadStateError'); +export const ThreadStateCancelledSchema = z + .object({ + status: z.literal('cancelled').describe('Thread was cancelled before completion.'), + }) + .openapi('ThreadStateCancelled'); + export const ThreadStateSchema = z - .discriminatedUnion('status', [ThreadStateDoneSchema, ThreadStateErrorSchema]) + .discriminatedUnion('status', [ThreadStateDoneSchema, ThreadStateErrorSchema, ThreadStateCancelledSchema]) .openapi('ThreadState'); export const BaseThreadDoneEventSchema = z @@ -381,6 +387,7 @@ export type ModelMessageDeltaEvent = z.infer; export type ThreadCreatedEvent = z.infer; export type ThreadStateError = z.infer; +export type ThreadStateCancelled = z.infer; export type ThreadState = z.infer; export type BaseThreadDoneEvent = z.infer; export type ThreadDoneEvent = z.infer; diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.types.ts b/packages/trueforge-core/src/core/runtime/AgentThread.types.ts index 5a5cf79f1..be1bc6747 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.types.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.types.ts @@ -77,7 +77,11 @@ export type InternalMCPAuthRequiredEvent = BaseMCPAuthRequiredEvent & { export type InternalThreadDoneEvent = BaseThreadDoneEvent & { type: typeof InternalEventType.AGENT_DONE; send_to_parent: LLMToolMessage | undefined; -} & ({ status: 'done'; output: ModelMessageEvent } | { status: 'error'; error: string; output?: ModelMessageEvent }); +} & ( + | { status: 'done'; output: ModelMessageEvent } + | { status: 'error'; error: string; output?: ModelMessageEvent } + | { status: 'cancelled' } + ); export type LLMContextMessage = LLMUserMessage | InternalEnrichedAssistantMessage | LLMToolMessage; diff --git a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts index a76c5203b..cc83a8d4f 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts @@ -45,7 +45,8 @@ function agentThreadEventToTerminalFields(event: AgentThreadExecutionEvent): { if ('parent' in event && event.parent) { return {}; } - if (event.status === 'error') { + // only a successful root completion becomes turn output + if (event.status !== 'done') { return {}; } return { output: event.output }; @@ -165,7 +166,7 @@ function createRootAgentSpan(mainThread: AgentThread, tracing: AgentTracing): Ro if (isInternalThreadDoneError(chunk)) { rootAgentErrorMessage = chunk.error; trace.setOutput(JSON.stringify({ error: chunk.error })); - } else { + } else if (chunk.status === 'done') { const content = assistantMessageContentToStringForSubAgent(chunk.output.content); trace.setOutput(JSON.stringify({ result: content })); } @@ -207,7 +208,7 @@ export async function* wrapWithSubAgentSpan( subTrace.setOutput(JSON.stringify({ error: event.error })); subTrace.setMetrics(currentThread.getAgentThreadMetrics()); subTrace.setError(event.error); - } else { + } else if (event.status === 'done') { const content = assistantMessageContentToStringForSubAgent(event.output.content); subTrace.setOutput(JSON.stringify({ result: content })); subTrace.setMetrics(currentThread.getAgentThreadMetrics()); @@ -239,6 +240,7 @@ export class AgentThreadOrchestrator { private readonly logger: Logger; // Finished sub-agents removed from `agentThreads`; kept so totals still include them. private finishedSubAgentMetrics: AgentThreadMetrics = createEmptyAgentThreadMetrics(); + private readonly cancelledThreadIds = new Set(); constructor(params: AgentThreadOrchestratorInput) { this.agentThreads = params.agentThreads; @@ -261,9 +263,24 @@ export class AgentThreadOrchestrator { return total; } + private markNonRootThreadsCancelled(): void { + for (const thread of this.agentThreads.values()) { + if (thread.parent) { + this.cancelledThreadIds.add(thread.threadId); + } + } + } + public async *send(messages: AgentThreadSendBatch): AsyncGenerator { + if (messages.length > 0 && !isUserToolApprovalOrResponseBatch(messages)) { + this.markNonRootThreadsCancelled(); + } + const byThread = new Map(); for (const thread of this.agentThreads.values()) { + if (this.cancelledThreadIds.has(thread.threadId)) { + continue; + } byThread.set(thread.threadId, []); } @@ -282,11 +299,6 @@ export class AgentThreadOrchestrator { byThread.set(threadId, batch); } } else if (messages.length > 0) { - if (this.agentThreads.size > 1) { - throw new InvalidAgentSendInputError( - 'Cannot process user messages while sub agents are running, please send empty input for previous conversation to complete', - ); - } byThread.set(getMainThreadId(this.agentThreads), messages); } @@ -352,25 +364,27 @@ export class AgentThreadOrchestrator { return { shouldStopExecution: true }; case InternalEventType.AGENT_DONE: { if (chunk.parent) { - if (!chunk.send_to_parent) { - throw new Error('unreachable'); - } - const parentThread = this.agentThreads.get(chunk.parent.thread_id); - if (!parentThread) { - throw new Error('unreachable: parent thread missing'); - } - const subAgentToolIsOpen = parentThread.hasOpenToolCallId(chunk.parent.tool_call_id); - if (subAgentToolIsOpen) { - const parentToolResponse: ToolResponseEvent = { - type: EventType.TOOL_RESPONSE, - id: newEventId(), - created_at: new Date().toISOString(), - thread_id: chunk.parent.thread_id, - tool_call_id: chunk.send_to_parent.tool_call_id, - content: '', - }; - yield parentToolResponse; - yield* this.sendToThread(chunk.parent.thread_id, [chunk.send_to_parent]); + if (chunk.status !== 'cancelled') { + if (!chunk.send_to_parent) { + throw new Error('unreachable'); + } + const parentThread = this.agentThreads.get(chunk.parent.thread_id); + if (!parentThread) { + throw new Error('unreachable: parent thread missing'); + } + const subAgentToolIsOpen = parentThread.hasOpenToolCallId(chunk.parent.tool_call_id); + if (subAgentToolIsOpen) { + const parentToolResponse: ToolResponseEvent = { + type: EventType.TOOL_RESPONSE, + id: newEventId(), + created_at: new Date().toISOString(), + thread_id: chunk.parent.thread_id, + tool_call_id: chunk.send_to_parent.tool_call_id, + content: '', + }; + yield parentToolResponse; + yield* this.sendToThread(chunk.parent.thread_id, [chunk.send_to_parent]); + } } yield chunk; // Move metrics to the finished bucket and drop the live entry with no `yield` between, @@ -443,6 +457,7 @@ export class AgentThreadOrchestrator { }); try { + yield* this.flushCancelledThreads(signal); while (agentThreads.size > 0) { if (shouldStopExecution) { break; @@ -525,4 +540,24 @@ export class AgentThreadOrchestrator { root_agent_error: rootAgentError, }; } + + private async *flushCancelledThreads(signal: AbortSignal): AsyncGenerator { + for (const threadId of this.cancelledThreadIds) { + const thread = this.agentThreads.get(threadId); + if (!thread?.parent) { + continue; + } + yield* this.processAgentStreamChunk( + { + type: InternalEventType.AGENT_DONE, + status: 'cancelled', + thread_id: thread.threadId, + title: thread.title, + parent: thread.parent, + send_to_parent: undefined, + }, + signal, + ); + } + } } diff --git a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts index 75b57b69c..79ae01676 100644 --- a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts +++ b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts @@ -35,7 +35,7 @@ function assistantWithToolCalls(toolCalls: InternalEnrichedToolCall[]): ContextM describe('getClosableOpenToolCallIds', () => { it('closes ordinary dangling tool calls on the last assistant message', () => { const context = assistantWithToolCalls([makeToolCall('tc-1'), makeToolCall('tc-2')]); - expect(getClosableOpenToolCallIds(context)).toEqual(new Set(['tc-1', 'tc-2'])); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set(['tc-1', 'tc-2'])); }); it('excludes tool calls that already have a matching tool response', () => { @@ -43,7 +43,7 @@ describe('getClosableOpenToolCallIds', () => { ...assistantWithToolCalls([makeToolCall('tc-1'), makeToolCall('tc-2')]), { role: 'tool', tool_call_id: 'tc-1', content: 'done' }, ]; - expect(getClosableOpenToolCallIds(context)).toEqual(new Set(['tc-2'])); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set(['tc-2'])); }); it('returns empty when any open call requires approval', () => { @@ -51,7 +51,7 @@ describe('getClosableOpenToolCallIds', () => { makeToolCall('tc-ordinary'), makeToolCall('tc-approval', { is_approval_required: true }), ]); - expect(getClosableOpenToolCallIds(context)).toEqual(new Set()); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set()); }); it('returns empty when any open call is client-side', () => { @@ -59,7 +59,7 @@ describe('getClosableOpenToolCallIds', () => { makeToolCall('tc-ordinary'), makeToolCall('tc-client', { is_client_side: true }), ]); - expect(getClosableOpenToolCallIds(context)).toEqual(new Set()); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set()); }); it('excludes thread-creation tool calls (is_thread_creation)', () => { @@ -67,7 +67,7 @@ describe('getClosableOpenToolCallIds', () => { makeToolCall('tc-regular'), makeToolCall('tc-sub-agent', { is_thread_creation: true }), ]); - expect(getClosableOpenToolCallIds(context)).toEqual(new Set(['tc-regular'])); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set(['tc-regular'])); }); it('treats legacy sub-agent calls without is_thread_creation as ordinary closable calls', () => { @@ -78,7 +78,9 @@ describe('getClosableOpenToolCallIds', () => { original_tool_name: 'create_sub_agent', }), ]); - expect(getClosableOpenToolCallIds(context)).toEqual(new Set(['tc-legacy-sub-agent'])); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual( + new Set(['tc-legacy-sub-agent']), + ); }); it('only inspects tool calls on the last assistant message', () => { @@ -95,10 +97,12 @@ describe('getClosableOpenToolCallIds', () => { tool_calls: [makeToolCall('new-tc')], }, ]; - expect(getClosableOpenToolCallIds(context)).toEqual(new Set(['new-tc'])); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: false })).toEqual(new Set(['new-tc'])); }); it('returns empty when there is no assistant message with tool calls', () => { - expect(getClosableOpenToolCallIds([{ role: 'user', content: 'hello' }])).toEqual(new Set()); + expect( + getClosableOpenToolCallIds({ context: [{ role: 'user', content: 'hello' }], userMessageIncoming: false }), + ).toEqual(new Set()); }); }); From 8146a694d2ca88b6861853dcad40e48a76cd0d7b Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 09:12:25 +0530 Subject: [PATCH 04/11] nit --- .../src/agent-session/SessionHandle.ts | 5 ++++- .../trueforge-core/src/core/runtime/AgentThread.ts | 14 +++++++++++++- 2 files changed, 17 insertions(+), 2 deletions(-) diff --git a/packages/trueforge-core/src/agent-session/SessionHandle.ts b/packages/trueforge-core/src/agent-session/SessionHandle.ts index 1aad0297a..0382f5c00 100644 --- a/packages/trueforge-core/src/agent-session/SessionHandle.ts +++ b/packages/trueforge-core/src/agent-session/SessionHandle.ts @@ -1,6 +1,7 @@ /** * Bound session handle: starts turns via {@link SessionHandle.createTurn}. */ +import { InvalidAgentSendInputError } from '../core/errors'; import { newEventId } from '../core/events/schema'; import type { AgentDefinition } from '../core/runtime/AgentDefinition'; import { AgentThread } from '../core/runtime/AgentThread'; @@ -55,7 +56,9 @@ function toSendBatch(input: TurnInputItem[] | undefined): AgentThreadSendBatch { if (input.every(isInputUserMessage)) { return input; } - throw new Error('input must be homogeneous: all user messages, or all approval/tool-response messages'); + throw new InvalidAgentSendInputError( + 'input must be homogeneous: all user messages, or all approval/tool-response messages', + ); } function toNewThreadInit(snapshot: AgentThreadSnapshot): NewThreadInit { diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.ts b/packages/trueforge-core/src/core/runtime/AgentThread.ts index 7f500adc2..0b8466f38 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.ts @@ -480,6 +480,7 @@ export class AgentThread { private deferredTool?: DeferredTool | undefined; private convertedTools: ConvertToolsResult | undefined; private pendingSandboxCreatedEvents: SandboxCreatedEvent[] = []; + private pendingPreSendOutputEvents: ToolResponseEvent[] = []; private sandbox?: Sandbox | undefined; private readonly tracing: AgentTracing; private readonly logger: Logger; @@ -579,7 +580,14 @@ export class AgentThread { for await (const event of this.executeContextProcessors('preSend', { userMessageIncoming: messages.some(isInputUserMessage), })) { - yield event; + // createTurn drains send() for context only; surface closer tool.response + // events at execute start so they persist after turn.created. + for (const item of event.output) { + if (item.type === EventType.TOOL_RESPONSE) { + this.pendingPreSendOutputEvents.push(item); + } + } + yield { ...event, output: [] }; } this.preSendRanThisTurn = true; @@ -1307,6 +1315,10 @@ export class AgentThread { } } this.preSendRanThisTurn = false; + for (const event of this.pendingPreSendOutputEvents) { + yield event; + } + this.pendingPreSendOutputEvents = []; const { initializationInfo, authRequirementInfo } = await this.tracing.withInitSpan(() => this.init()); if (initializationInfo.length > 0) { From be08848b25d67372edd6d5e7d1f398d4cbb22ef8 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 09:21:17 +0530 Subject: [PATCH 05/11] add tests --- .changeset/pre/steer-session-anytime.md | 5 + .../tests/agent-session/sessions.test.ts | 3 +- .../core/runtime/openToolCallCloser.test.ts | 90 +++++- .../orchestrationUserMessage.test.ts | 278 ++++++++++++++++++ 4 files changed, 374 insertions(+), 2 deletions(-) create mode 100644 .changeset/pre/steer-session-anytime.md create mode 100644 packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts diff --git a/.changeset/pre/steer-session-anytime.md b/.changeset/pre/steer-session-anytime.md new file mode 100644 index 000000000..74ce1d9f9 --- /dev/null +++ b/.changeset/pre/steer-session-anytime.md @@ -0,0 +1,5 @@ +--- +"@truefoundry/trueforge-core": patch +--- + +Allow a user message to start a turn while approvals, client-side tools, or sub-agent threads are pending: close those calls synthetically and cancel open sub-agents with `thread.done` status `cancelled`. diff --git a/packages/trueforge-core/tests/agent-session/sessions.test.ts b/packages/trueforge-core/tests/agent-session/sessions.test.ts index 39e7861be..327e93ddd 100644 --- a/packages/trueforge-core/tests/agent-session/sessions.test.ts +++ b/packages/trueforge-core/tests/agent-session/sessions.test.ts @@ -4,6 +4,7 @@ import { Sessions } from '../../src/agent-session/Sessions'; import { InMemorySessionStore } from '../../src/agent-session/store/InMemorySessionStore'; import { TurnNotFoundError } from '../../src/agent-session/store/SessionStoreErrors'; import { TurnHandle } from '../../src/agent-session/TurnHandle'; +import { InvalidAgentSendInputError } from '../../src/core/errors'; import { makeAgentSpec, makeTestResolver, mintTestTurnId } from './testHelpers'; describe('Sessions / SessionHandle / TurnHandle (storage + createTurn)', () => { @@ -314,7 +315,7 @@ describe('Sessions / SessionHandle / TurnHandle (storage + createTurn)', () => { signal: new AbortController().signal, resolver: makeTestResolver(), }), - ).rejects.toThrow(); + ).rejects.toThrow(InvalidAgentSendInputError); const turns = await store.listTurns({ session_id: 's1', limit: 10, diff --git a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts index 79ae01676..57229fcf6 100644 --- a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts +++ b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts @@ -1,6 +1,12 @@ +import { EventType } from '../../../src/core/events/schema'; import type { InternalEnrichedAssistantMessage, InternalEnrichedToolCall } from '../../../src/core/llm/LLMTypes'; import type { ContextMessage } from '../../../src/core/runtime/AgentThread.types'; -import { getClosableOpenToolCallIds } from '../../../src/core/runtime/OpenToolCallCloser'; +import { getEmptyCurrentContextUsage } from '../../../src/core/runtime/contextUsage'; +import { + getClosableOpenToolCallIds, + getClosableOpenToolCalls, + OpenToolCallCloser, +} from '../../../src/core/runtime/OpenToolCallCloser'; import '../harnessMocks'; function makeToolCall( @@ -105,4 +111,86 @@ describe('getClosableOpenToolCallIds', () => { getClosableOpenToolCallIds({ context: [{ role: 'user', content: 'hello' }], userMessageIncoming: false }), ).toEqual(new Set()); }); + + it('closes approval, client-side, and thread-creation calls when a user message is incoming', () => { + const context = assistantWithToolCalls([ + makeToolCall('tc-regular'), + makeToolCall('tc-approval', { is_approval_required: true }), + makeToolCall('tc-client', { is_client_side: true }), + makeToolCall('tc-sub-agent', { is_thread_creation: true }), + ]); + expect(getClosableOpenToolCalls({ context, userMessageIncoming: true })).toEqual([ + { tool_call_id: 'tc-regular', close_kind: 'dangling' }, + { tool_call_id: 'tc-approval', close_kind: 'user_action' }, + { tool_call_id: 'tc-client', close_kind: 'user_action' }, + { tool_call_id: 'tc-sub-agent', close_kind: 'thread_creation' }, + ]); + }); + + it('still excludes already-resolved calls when a user message is incoming', () => { + const context: ContextMessage[] = [ + ...assistantWithToolCalls([ + makeToolCall('tc-approval', { is_approval_required: true }), + makeToolCall('tc-sub-agent', { is_thread_creation: true }), + ]), + { role: 'tool', tool_call_id: 'tc-approval', content: 'already closed' }, + ]; + expect(getClosableOpenToolCalls({ context, userMessageIncoming: true })).toEqual([ + { tool_call_id: 'tc-sub-agent', close_kind: 'thread_creation' }, + ]); + }); +}); + +describe('OpenToolCallCloser.processPreSend', () => { + async function collectPreSend(context: ContextMessage[], userMessageIncoming: boolean) { + const closer = new OpenToolCallCloser(); + const yielded = []; + for await (const event of closer.processPreSend( + { + threadId: 'main', + currentContextUsage: getEmptyCurrentContextUsage(), + context, + }, + { userMessageIncoming }, + )) { + yielded.push(event); + } + return yielded; + } + + it('appends dangling dummy tool messages with no output events on resume', async () => { + const yielded = await collectPreSend(assistantWithToolCalls([makeToolCall('tc-1')]), false); + expect(yielded).toHaveLength(1); + expect(yielded[0]?.context).toEqual([ + { + role: 'tool', + tool_call_id: 'tc-1', + content: JSON.stringify({ error: 'Tool call was not executed. Please retry this tool call.' }), + }, + ]); + expect(yielded[0]?.output).toEqual([]); + }); + + it('emits tool.response events for user-action and thread-creation closures', async () => { + const yielded = await collectPreSend( + assistantWithToolCalls([ + makeToolCall('tc-approval', { is_approval_required: true }), + makeToolCall('tc-sub-agent', { is_thread_creation: true }), + ]), + true, + ); + expect(yielded).toHaveLength(1); + expect(yielded[0]?.output).toEqual([ + expect.objectContaining({ type: EventType.TOOL_RESPONSE, tool_call_id: 'tc-approval', thread_id: 'main' }), + expect.objectContaining({ type: EventType.TOOL_RESPONSE, tool_call_id: 'tc-sub-agent', thread_id: 'main' }), + ]); + }); + + it('is idempotent after dummy responses are in context', async () => { + const context = assistantWithToolCalls([makeToolCall('tc-1')]); + const first = await collectPreSend(context, false); + const closed = [...context, ...(first[0]?.context ?? [])]; + expect(await collectPreSend(closed, false)).toEqual([]); + expect(getClosableOpenToolCalls({ context: closed, userMessageIncoming: true })).toEqual([]); + }); }); diff --git a/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts b/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts new file mode 100644 index 000000000..f518dd904 --- /dev/null +++ b/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts @@ -0,0 +1,278 @@ +import { InvalidAgentSendInputError } from '../../src/core/errors'; +import { EventType } from '../../src/core/events/schema'; +import type { InternalEnrichedAssistantMessage } from '../../src/core/llm/LLMTypes'; +import { AgentThread } from '../../src/core/runtime/AgentThread'; +import { InternalEventType } from '../../src/core/runtime/AgentThread.types'; +import { AgentThreadOrchestrator } from '../../src/core/runtime/AgentThreadOrchestrator'; +import { NOOP_AGENT_TRACING } from '../../src/core/tracing/NoopAgentTracing'; +import { makeSilentLogger } from '../core/harnessMocks'; +import { + llmCreateInputs, + makeApprovalGatedWriteNoteToolSet, + runTurn, + textReplyStream, + WRITE_NOTE_CALL_ID, + writeNoteToolCallStream, +} from './helpers/helpers'; + +const ROOT_ID = 'thread_root'; +const CHILD_ID = 'thread_child'; +const SUB_AGENT_CALL_ID = 'call-sub'; +const INSTRUCTION = 'You are running in a test setup.'; +const STEER_REPLY = 'acknowledged the new instruction'; +const CHILD_REPLY = 'child still running'; + +describe('orchestration: user message while work is pending', () => { + it('rejects empty and incomplete action batches while approval is pending, then steers with a user message without executing the tool', async () => { + const { orchestrator, thread, callTool } = makeApprovalHarness(STEER_REPLY); + + await runTurn({ + orchestrator, + sendBatch: [{ type: EventType.USER_MESSAGE, content: 'hello' }], + }); + expect(callTool).not.toHaveBeenCalled(); + + await expect(runTurn({ orchestrator, sendBatch: [] })).rejects.toThrow(InvalidAgentSendInputError); + await expect( + runTurn({ + orchestrator, + sendBatch: [ + { + type: EventType.USER_TOOL_APPROVAL, + thread_id: ROOT_ID, + tool_call_id: 'unknown-call', + approval: { status: 'allow' }, + }, + ], + }), + ).rejects.toThrow(/no pending approval/); + + const steered = await runTurn({ + orchestrator, + sendBatch: [{ type: EventType.USER_MESSAGE, content: 'never mind, do this instead' }], + }); + expect(callTool).not.toHaveBeenCalled(); + expect(steered.events).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + thread_id: ROOT_ID, + tool_call_id: WRITE_NOTE_CALL_ID, + }), + expect.objectContaining({ + type: InternalEventType.AGENT_DONE, + thread_id: ROOT_ID, + status: 'done', + }), + ]), + ); + expect(steered.result.required_actions).toEqual([]); + expect(llmCreateInputs(thread.definition.modelClient).at(-1)).toMatchObject({ + messages: expect.arrayContaining([ + { role: 'tool', tool_call_id: WRITE_NOTE_CALL_ID, content: expect.stringContaining('new message') }, + { role: 'user', content: 'never mind, do this instead' }, + ]), + }); + + await expect( + runTurn({ + orchestrator, + sendBatch: [ + { + type: EventType.USER_TOOL_APPROVAL, + thread_id: ROOT_ID, + tool_call_id: WRITE_NOTE_CALL_ID, + approval: { status: 'allow' }, + }, + ], + }), + ).rejects.toThrow(/no pending approval/); + }); + + it('cancels an open sub-agent when a user message arrives', async () => { + const { orchestrator, childCreate } = makeOpenSubAgentHarness(); + const steered = await runTurn({ + orchestrator, + sendBatch: [{ type: EventType.USER_MESSAGE, content: 'stop the worker' }], + }); + expect(childCreate).not.toHaveBeenCalled(); + expect(steered.events[0]).toMatchObject({ + type: InternalEventType.AGENT_DONE, + thread_id: CHILD_ID, + status: 'cancelled', + }); + expect(steered.events).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + thread_id: ROOT_ID, + tool_call_id: SUB_AGENT_CALL_ID, + }), + expect.objectContaining({ type: InternalEventType.AGENT_DONE, thread_id: ROOT_ID, status: 'done' }), + ]), + ); + }); + + it('resumes an open sub-agent on empty input', async () => { + const { orchestrator, childCreate } = makeOpenSubAgentHarness(); + const resumed = await runTurn({ orchestrator, sendBatch: [] }); + expect(childCreate).toHaveBeenCalled(); + expect(resumed.events).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + thread_id: ROOT_ID, + tool_call_id: SUB_AGENT_CALL_ID, + }), + expect.objectContaining({ type: InternalEventType.AGENT_DONE, thread_id: CHILD_ID, status: 'done' }), + expect.objectContaining({ type: InternalEventType.AGENT_DONE, thread_id: ROOT_ID, status: 'done' }), + ]), + ); + expect(resumed.events.some(e => e.type === InternalEventType.AGENT_DONE && e.status === 'cancelled')).toBe(false); + }); +}); + +function makeApprovalHarness(finalReply: string): { + orchestrator: AgentThreadOrchestrator; + thread: AgentThread; + callTool: jest.Mock; +} { + const { toolSet, callTool } = makeApprovalGatedWriteNoteToolSet(); + const thread = new AgentThread({ + definition: { + modelClient: { + create: jest + .fn() + .mockImplementationOnce(() => writeNoteToolCallStream()) + .mockImplementation(() => textReplyStream(finalReply)), + createNonStream: jest.fn(), + }, + instruction: INSTRUCTION, + messages: undefined, + modelParams: undefined, + responseFormat: undefined, + iterationLimit: undefined, + toolSets: undefined, + }, + threadId: ROOT_ID, + title: 'orchestration-user-message', + parent: undefined, + agentInfo: undefined, + context: undefined, + currentContextUsage: undefined, + preComputedCompletion: undefined, + sandbox: undefined, + capabilities: [ + { + systemToolSets: [toolSet], + preSendProcessors: undefined, + preLLMProcessors: undefined, + preLLMEphemeralProcessors: undefined, + postToolCallProcessors: undefined, + toolResponseProcessors: undefined, + instructionBuilders: undefined, + }, + ], + capabilityState: undefined, + tracing: NOOP_AGENT_TRACING, + logger: makeSilentLogger(), + }); + return { + orchestrator: new AgentThreadOrchestrator({ + agentThreads: new Map([[thread.threadId, thread]]), + createDynamicSubAgentThread: () => Promise.reject(new Error('unexpected sub-agent')), + tracing: NOOP_AGENT_TRACING, + logger: makeSilentLogger(), + }), + thread, + callTool, + }; +} + +function makeOpenSubAgentHarness(): { orchestrator: AgentThreadOrchestrator; childCreate: jest.Mock } { + const assistant: InternalEnrichedAssistantMessage = { + role: 'assistant', + content: null, + tool_calls: [ + { + id: SUB_AGENT_CALL_ID, + type: 'function', + function: { name: 'create_sub_agent', arguments: '{"name":"worker","input":"task"}' }, + tool_info: { + type: 'mcp', + mcp_server_id: 'sub-agents', + mcp_server_name: 'sub-agents', + original_tool_name: 'create_sub_agent', + is_thread_creation: true, + is_approval_required: false, + is_client_side: false, + }, + }, + ], + }; + const root = new AgentThread({ + definition: { + modelClient: { + create: jest.fn().mockImplementation(() => textReplyStream(STEER_REPLY)), + createNonStream: jest.fn(), + }, + instruction: INSTRUCTION, + messages: undefined, + modelParams: undefined, + responseFormat: undefined, + iterationLimit: undefined, + toolSets: undefined, + }, + threadId: ROOT_ID, + title: 'main', + parent: undefined, + agentInfo: undefined, + context: [assistant], + currentContextUsage: undefined, + preComputedCompletion: undefined, + sandbox: undefined, + capabilities: undefined, + capabilityState: undefined, + tracing: NOOP_AGENT_TRACING, + logger: makeSilentLogger(), + }); + const childCreate = jest.fn().mockImplementation(() => textReplyStream(CHILD_REPLY)); + const child = new AgentThread({ + definition: { + modelClient: { + create: childCreate, + createNonStream: jest.fn(), + }, + instruction: undefined, + messages: [{ role: 'user', content: 'task' }], + modelParams: undefined, + responseFormat: undefined, + iterationLimit: undefined, + toolSets: undefined, + }, + threadId: CHILD_ID, + title: 'worker', + parent: { thread_id: ROOT_ID, tool_call_id: SUB_AGENT_CALL_ID }, + agentInfo: { type: 'dynamic', name: 'worker', input: 'task' }, + context: undefined, + currentContextUsage: undefined, + preComputedCompletion: undefined, + sandbox: undefined, + capabilities: undefined, + capabilityState: undefined, + tracing: NOOP_AGENT_TRACING, + logger: makeSilentLogger(), + }); + return { + orchestrator: new AgentThreadOrchestrator({ + agentThreads: new Map([ + [root.threadId, root], + [child.threadId, child], + ]), + createDynamicSubAgentThread: () => Promise.reject(new Error('unexpected sub-agent spawn')), + tracing: NOOP_AGENT_TRACING, + logger: makeSilentLogger(), + }), + childCreate, + }; +} From 2dcf006104abca7f5738da7579ab564fc1cecb4e Mon Sep 17 00:00:00 2001 From: "trueforge-dev-bot[bot]" Date: Thu, 17 Sep 2026 03:54:55 +0000 Subject: [PATCH 06/11] Regenerate OpenAPI document and SDKs --- ...60917035455-regenerate-sdk-from-openapi.md | 5 +++++ .github/fern/openapi/openapi.json | 21 ++++++++++++++++++- docs/openapi.json | 21 ++++++++++++++++++- .../src/api/types/ThreadState.ts | 2 +- .../src/api/types/ThreadStateCancelled.ts | 5 +++++ packages/trueforge-sdk/src/api/types/index.ts | 1 + .../src/serialization/types/ThreadState.ts | 5 +++-- .../types/ThreadStateCancelled.ts | 18 ++++++++++++++++ .../src/serialization/types/index.ts | 1 + python/trueforge_sdk/.fern/metadata.json | 2 +- .../src/trueforge_sdk/__init__.py | 3 +++ .../src/trueforge_sdk/types/__init__.py | 3 +++ .../src/trueforge_sdk/types/thread_state.py | 3 ++- .../types/thread_state_cancelled.py | 19 +++++++++++++++++ 14 files changed, 102 insertions(+), 7 deletions(-) create mode 100644 .changeset/20260917035455-regenerate-sdk-from-openapi.md create mode 100644 packages/trueforge-sdk/src/api/types/ThreadStateCancelled.ts create mode 100644 packages/trueforge-sdk/src/serialization/types/ThreadStateCancelled.ts create mode 100644 python/trueforge_sdk/src/trueforge_sdk/types/thread_state_cancelled.py diff --git a/.changeset/20260917035455-regenerate-sdk-from-openapi.md b/.changeset/20260917035455-regenerate-sdk-from-openapi.md new file mode 100644 index 000000000..efd8ff00f --- /dev/null +++ b/.changeset/20260917035455-regenerate-sdk-from-openapi.md @@ -0,0 +1,5 @@ +--- +"@truefoundry/trueforge-sdk": patch +--- + +Regenerate SDK from updated OpenAPI spec. diff --git a/.github/fern/openapi/openapi.json b/.github/fern/openapi/openapi.json index e9d988bef..b333a454d 100644 --- a/.github/fern/openapi/openapi.json +++ b/.github/fern/openapi/openapi.json @@ -4564,6 +4564,7 @@ "ThreadState": { "discriminator": { "mapping": { + "cancelled": "#/components/schemas/ThreadStateCancelled", "done": "#/components/schemas/ThreadStateDone", "error": "#/components/schemas/ThreadStateError" }, @@ -4575,9 +4576,27 @@ }, { "$ref": "#/components/schemas/ThreadStateError" + }, + { + "$ref": "#/components/schemas/ThreadStateCancelled" } ] }, + "ThreadStateCancelled": { + "properties": { + "status": { + "description": "Thread was cancelled before completion.", + "enum": [ + "cancelled" + ], + "type": "string" + } + }, + "required": [ + "status" + ], + "type": "object" + }, "ThreadStateDone": { "properties": { "output": { @@ -5663,7 +5682,7 @@ "info": { "description": "HTTP API for the TrueForge agent server (`/api/v1`). Interactive docs are served at `/api/v1/docs` (OpenAPI JSON at `/api/v1/openapi.json`).\n\n**Authentication:** Standalone auth accepts requests without credentials — middleware stamps a local default user. When OIDC or TrueFoundry auth is configured, protected routes require a valid cookie or `Authorization: Bearer` token. There is no built-in API-key scheme; pass custom headers only if your reverse proxy or IdP layer requires them.\n\nCovers DB-backed sessions, the agent registry, settings catalogs, and model/MCP/skill/sandbox providers.", "title": "TrueForge API", - "version": "0.2.0-rc.11" + "version": "0.2.0-rc.12" }, "openapi": "3.1.0", "paths": { diff --git a/docs/openapi.json b/docs/openapi.json index e9d988bef..b333a454d 100644 --- a/docs/openapi.json +++ b/docs/openapi.json @@ -4564,6 +4564,7 @@ "ThreadState": { "discriminator": { "mapping": { + "cancelled": "#/components/schemas/ThreadStateCancelled", "done": "#/components/schemas/ThreadStateDone", "error": "#/components/schemas/ThreadStateError" }, @@ -4575,9 +4576,27 @@ }, { "$ref": "#/components/schemas/ThreadStateError" + }, + { + "$ref": "#/components/schemas/ThreadStateCancelled" } ] }, + "ThreadStateCancelled": { + "properties": { + "status": { + "description": "Thread was cancelled before completion.", + "enum": [ + "cancelled" + ], + "type": "string" + } + }, + "required": [ + "status" + ], + "type": "object" + }, "ThreadStateDone": { "properties": { "output": { @@ -5663,7 +5682,7 @@ "info": { "description": "HTTP API for the TrueForge agent server (`/api/v1`). Interactive docs are served at `/api/v1/docs` (OpenAPI JSON at `/api/v1/openapi.json`).\n\n**Authentication:** Standalone auth accepts requests without credentials — middleware stamps a local default user. When OIDC or TrueFoundry auth is configured, protected routes require a valid cookie or `Authorization: Bearer` token. There is no built-in API-key scheme; pass custom headers only if your reverse proxy or IdP layer requires them.\n\nCovers DB-backed sessions, the agent registry, settings catalogs, and model/MCP/skill/sandbox providers.", "title": "TrueForge API", - "version": "0.2.0-rc.11" + "version": "0.2.0-rc.12" }, "openapi": "3.1.0", "paths": { diff --git a/packages/trueforge-sdk/src/api/types/ThreadState.ts b/packages/trueforge-sdk/src/api/types/ThreadState.ts index 76ec3e26c..6e8a37976 100644 --- a/packages/trueforge-sdk/src/api/types/ThreadState.ts +++ b/packages/trueforge-sdk/src/api/types/ThreadState.ts @@ -2,4 +2,4 @@ import type * as TrueForge from "../index.js"; -export type ThreadState = TrueForge.ThreadStateDone | TrueForge.ThreadStateError; +export type ThreadState = TrueForge.ThreadStateCancelled | TrueForge.ThreadStateDone | TrueForge.ThreadStateError; diff --git a/packages/trueforge-sdk/src/api/types/ThreadStateCancelled.ts b/packages/trueforge-sdk/src/api/types/ThreadStateCancelled.ts new file mode 100644 index 000000000..d740759e0 --- /dev/null +++ b/packages/trueforge-sdk/src/api/types/ThreadStateCancelled.ts @@ -0,0 +1,5 @@ +// This file was auto-generated by Fern from our API Definition. + +export interface ThreadStateCancelled { + status: "cancelled"; +} diff --git a/packages/trueforge-sdk/src/api/types/index.ts b/packages/trueforge-sdk/src/api/types/index.ts index c88fa727c..bdde6e618 100644 --- a/packages/trueforge-sdk/src/api/types/index.ts +++ b/packages/trueforge-sdk/src/api/types/index.ts @@ -192,6 +192,7 @@ export * from "./TextContent.js"; export * from "./ThreadCreatedEvent.js"; export * from "./ThreadDoneEvent.js"; export * from "./ThreadState.js"; +export * from "./ThreadStateCancelled.js"; export * from "./ThreadStateDone.js"; export * from "./ThreadStateError.js"; export * from "./Timezone.js"; diff --git a/packages/trueforge-sdk/src/serialization/types/ThreadState.ts b/packages/trueforge-sdk/src/serialization/types/ThreadState.ts index e59ed63a0..c8fae3122 100644 --- a/packages/trueforge-sdk/src/serialization/types/ThreadState.ts +++ b/packages/trueforge-sdk/src/serialization/types/ThreadState.ts @@ -3,12 +3,13 @@ import type * as TrueForge from "../../api/index.js"; import * as core from "../../core/index.js"; import type * as serializers from "../index.js"; +import { ThreadStateCancelled } from "./ThreadStateCancelled.js"; import { ThreadStateDone } from "./ThreadStateDone.js"; import { ThreadStateError } from "./ThreadStateError.js"; export const ThreadState: core.serialization.Schema = - core.serialization.undiscriminatedUnion([ThreadStateDone, ThreadStateError]); + core.serialization.undiscriminatedUnion([ThreadStateCancelled, ThreadStateDone, ThreadStateError]); export declare namespace ThreadState { - export type Raw = ThreadStateDone.Raw | ThreadStateError.Raw; + export type Raw = ThreadStateCancelled.Raw | ThreadStateDone.Raw | ThreadStateError.Raw; } diff --git a/packages/trueforge-sdk/src/serialization/types/ThreadStateCancelled.ts b/packages/trueforge-sdk/src/serialization/types/ThreadStateCancelled.ts new file mode 100644 index 000000000..b7ab5e92b --- /dev/null +++ b/packages/trueforge-sdk/src/serialization/types/ThreadStateCancelled.ts @@ -0,0 +1,18 @@ +// This file was auto-generated by Fern from our API Definition. + +import type * as TrueForge from "../../api/index.js"; +import * as core from "../../core/index.js"; +import type * as serializers from "../index.js"; + +export const ThreadStateCancelled: core.serialization.ObjectSchema< + serializers.ThreadStateCancelled.Raw, + TrueForge.ThreadStateCancelled +> = core.serialization.object({ + status: core.serialization.stringLiteral("cancelled"), +}); + +export declare namespace ThreadStateCancelled { + export interface Raw { + status: "cancelled"; + } +} diff --git a/packages/trueforge-sdk/src/serialization/types/index.ts b/packages/trueforge-sdk/src/serialization/types/index.ts index c88fa727c..bdde6e618 100644 --- a/packages/trueforge-sdk/src/serialization/types/index.ts +++ b/packages/trueforge-sdk/src/serialization/types/index.ts @@ -192,6 +192,7 @@ export * from "./TextContent.js"; export * from "./ThreadCreatedEvent.js"; export * from "./ThreadDoneEvent.js"; export * from "./ThreadState.js"; +export * from "./ThreadStateCancelled.js"; export * from "./ThreadStateDone.js"; export * from "./ThreadStateError.js"; export * from "./Timezone.js"; diff --git a/python/trueforge_sdk/.fern/metadata.json b/python/trueforge_sdk/.fern/metadata.json index 0c193c627..942a2f2af 100644 --- a/python/trueforge_sdk/.fern/metadata.json +++ b/python/trueforge_sdk/.fern/metadata.json @@ -21,7 +21,7 @@ }, "pyproject_python_version": ">=3.10" }, - "originGitCommit": "fd1bf7fedb0663f5a5cda4f5842984b7a60bf933", + "originGitCommit": "be08848b25d67372edd6d5e7d1f398d4cbb22ef8", "originGitCommitIsDirty": true, "invokedBy": "ci", "requestedVersion": "0.2.0-rc.8", diff --git a/python/trueforge_sdk/src/trueforge_sdk/__init__.py b/python/trueforge_sdk/src/trueforge_sdk/__init__.py index f336deeec..396cf809f 100644 --- a/python/trueforge_sdk/src/trueforge_sdk/__init__.py +++ b/python/trueforge_sdk/src/trueforge_sdk/__init__.py @@ -201,6 +201,7 @@ ThreadCreatedEvent, ThreadDoneEvent, ThreadState, + ThreadStateCancelled, ThreadStateDone, ThreadStateError, Timezone, @@ -464,6 +465,7 @@ "ThreadCreatedEvent": ".types", "ThreadDoneEvent": ".types", "ThreadState": ".types", + "ThreadStateCancelled": ".types", "ThreadStateDone": ".types", "ThreadStateError": ".types", "Timezone": ".types", @@ -747,6 +749,7 @@ def __dir__(): "ThreadCreatedEvent", "ThreadDoneEvent", "ThreadState", + "ThreadStateCancelled", "ThreadStateDone", "ThreadStateError", "Timezone", diff --git a/python/trueforge_sdk/src/trueforge_sdk/types/__init__.py b/python/trueforge_sdk/src/trueforge_sdk/types/__init__.py index c242e7818..150ec8b81 100644 --- a/python/trueforge_sdk/src/trueforge_sdk/types/__init__.py +++ b/python/trueforge_sdk/src/trueforge_sdk/types/__init__.py @@ -200,6 +200,7 @@ from .thread_created_event import ThreadCreatedEvent from .thread_done_event import ThreadDoneEvent from .thread_state import ThreadState + from .thread_state_cancelled import ThreadStateCancelled from .thread_state_done import ThreadStateDone from .thread_state_error import ThreadStateError from .timezone import Timezone @@ -431,6 +432,7 @@ "ThreadCreatedEvent": ".thread_created_event", "ThreadDoneEvent": ".thread_done_event", "ThreadState": ".thread_state", + "ThreadStateCancelled": ".thread_state_cancelled", "ThreadStateDone": ".thread_state_done", "ThreadStateError": ".thread_state_error", "Timezone": ".timezone", @@ -686,6 +688,7 @@ def __dir__(): "ThreadCreatedEvent", "ThreadDoneEvent", "ThreadState", + "ThreadStateCancelled", "ThreadStateDone", "ThreadStateError", "Timezone", diff --git a/python/trueforge_sdk/src/trueforge_sdk/types/thread_state.py b/python/trueforge_sdk/src/trueforge_sdk/types/thread_state.py index d69da3db5..0cb2d299e 100644 --- a/python/trueforge_sdk/src/trueforge_sdk/types/thread_state.py +++ b/python/trueforge_sdk/src/trueforge_sdk/types/thread_state.py @@ -2,7 +2,8 @@ import typing +from .thread_state_cancelled import ThreadStateCancelled from .thread_state_done import ThreadStateDone from .thread_state_error import ThreadStateError -ThreadState = typing.Union[ThreadStateDone, ThreadStateError] +ThreadState = typing.Union[ThreadStateCancelled, ThreadStateDone, ThreadStateError] diff --git a/python/trueforge_sdk/src/trueforge_sdk/types/thread_state_cancelled.py b/python/trueforge_sdk/src/trueforge_sdk/types/thread_state_cancelled.py new file mode 100644 index 000000000..1ef0f38e5 --- /dev/null +++ b/python/trueforge_sdk/src/trueforge_sdk/types/thread_state_cancelled.py @@ -0,0 +1,19 @@ +# This file was auto-generated by Fern from our API Definition. + +import typing + +import pydantic +from ..core.pydantic_utilities import IS_PYDANTIC_V2 +from ..core.unchecked_base_model import UncheckedBaseModel + + +class ThreadStateCancelled(UncheckedBaseModel): + status: typing.Literal["cancelled"] = "cancelled" + + if IS_PYDANTIC_V2: + model_config: typing.ClassVar[pydantic.ConfigDict] = pydantic.ConfigDict(extra="allow") # type: ignore # Pydantic v2 + else: + + class Config: + smart_union = True + extra = pydantic.Extra.allow From 8d1ed0decf62ea87bd290320e6eafe1b8e7b8402 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 09:38:43 +0530 Subject: [PATCH 07/11] nit --- .../src/core/runtime/OpenToolCallCloser.ts | 103 ++++++------------ .../core/runtime/openToolCallCloser.test.ts | 50 ++++++--- .../orchestrationUserMessage.test.ts | 6 +- 3 files changed, 71 insertions(+), 88 deletions(-) diff --git a/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts b/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts index 3f81ae9f7..d7cb4b991 100644 --- a/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts +++ b/packages/trueforge-core/src/core/runtime/OpenToolCallCloser.ts @@ -31,51 +31,26 @@ const DANGLING_TOOL_MESSAGE_CONTENT = JSON.stringify({ error: 'Tool call was not executed. Please retry this tool call.', }); -const USER_ACTION_TOOL_MESSAGE_CONTENT = JSON.stringify({ - error: - 'Tool call was not executed: the user sent a new message before it was resolved. Do not retry unless the user asks.', -}); - -const THREAD_CREATION_TOOL_MESSAGE_CONTENT = 'Sub-agent was cancelled because the user sent a new message.'; - -export type ClosableOpenToolCallKind = 'dangling' | 'user_action' | 'thread_creation'; - -export interface ClosableOpenToolCall { - tool_call_id: string; - close_kind: ClosableOpenToolCallKind; -} +const CANCELLED_TOOL_MESSAGE_CONTENT = 'Tool call was cancelled: a new turn was started.'; -function closeKindForToolCall(toolCall: InternalEnrichedToolCall): ClosableOpenToolCallKind { - if (toolCall.tool_info.is_thread_creation === true) { - return 'thread_creation'; - } - if (toolCall.tool_info.is_approval_required === true || toolCall.tool_info.is_client_side === true) { - return 'user_action'; - } - return 'dangling'; +function isPendingUserAction(toolCall: InternalEnrichedToolCall): boolean { + return toolCall.tool_info.is_approval_required === true || toolCall.tool_info.is_client_side === true; } -function contentForCloseKind(closeKind: ClosableOpenToolCallKind): string { - switch (closeKind) { - case 'dangling': - return DANGLING_TOOL_MESSAGE_CONTENT; - case 'user_action': - return USER_ACTION_TOOL_MESSAGE_CONTENT; - case 'thread_creation': - return THREAD_CREATION_TOOL_MESSAGE_CONTENT; - } +function isThreadCreation(toolCall: InternalEnrichedToolCall): boolean { + return toolCall.tool_info.is_thread_creation === true; } -export function getClosableOpenToolCalls(input: { +export function getClosableOpenToolCallIds(input: { context: ContextMessage[]; userMessageIncoming: boolean; -}): ClosableOpenToolCall[] { +}): Set { const lastIdx = input.context.findLastIndex( (msg): msg is InternalEnrichedAssistantMessage => isLLMContextMessage(msg) && msg.role === 'assistant' && !!msg.tool_calls?.length, ); if (lastIdx === -1) { - return []; + return new Set(); } const lastAssistant = input.context[lastIdx]; @@ -83,16 +58,11 @@ export function getClosableOpenToolCalls(input: { throw new Error('Unreachable'); } if (!lastAssistant.tool_calls) { - return []; + return new Set(); } - if ( - !input.userMessageIncoming && - lastAssistant.tool_calls.some( - tc => tc.tool_info.is_approval_required === true || tc.tool_info.is_client_side === true, - ) - ) { - return []; + if (!input.userMessageIncoming && lastAssistant.tool_calls.some(isPendingUserAction)) { + return new Set(); } const resolvedIds = new Set(); @@ -102,27 +72,19 @@ export function getClosableOpenToolCalls(input: { } } - const closable: ClosableOpenToolCall[] = []; + const closable = new Set(); for (const toolCall of lastAssistant.tool_calls) { if (resolvedIds.has(toolCall.id)) { continue; } - const close_kind = closeKindForToolCall(toolCall); - if (!input.userMessageIncoming && close_kind === 'thread_creation') { + if (!input.userMessageIncoming && isThreadCreation(toolCall)) { continue; } - closable.push({ tool_call_id: toolCall.id, close_kind }); + closable.add(toolCall.id); } return closable; } -export function getClosableOpenToolCallIds(input: { - context: ContextMessage[]; - userMessageIncoming: boolean; -}): Set { - return new Set(getClosableOpenToolCalls(input).map(call => call.tool_call_id)); -} - // we are closing open tool calls synthetically, the subscriber needs to understand // the tool calls were closed. function toToolResponseEvent(input: { threadId: string; toolCallId: string; content: string }): ToolResponseEvent { @@ -142,29 +104,32 @@ export class OpenToolCallCloser implements PreSendContextProcessor { execution: Readonly, options: { userMessageIncoming: boolean }, ): AsyncGenerator { - const closable = getClosableOpenToolCalls({ - context: execution.context, - userMessageIncoming: options.userMessageIncoming, - }); - if (closable.length === 0) { + const closableIds = [ + ...getClosableOpenToolCallIds({ + context: execution.context, + userMessageIncoming: options.userMessageIncoming, + }), + ]; + if (closableIds.length === 0) { return; } - const dummyToolMessages: LLMToolMessage[] = closable.map(call => ({ + const content = options.userMessageIncoming ? CANCELLED_TOOL_MESSAGE_CONTENT : DANGLING_TOOL_MESSAGE_CONTENT; + const dummyToolMessages: LLMToolMessage[] = closableIds.map(toolCallId => ({ role: 'tool', - tool_call_id: call.tool_call_id, - content: contentForCloseKind(call.close_kind), + tool_call_id: toolCallId, + content, })); - const output: ToolResponseEvent[] = closable - .filter(call => call.close_kind !== 'dangling') - .map(call => - toToolResponseEvent({ - threadId: execution.threadId, - toolCallId: call.tool_call_id, - content: contentForCloseKind(call.close_kind), - }), - ); + const output: ToolResponseEvent[] = options.userMessageIncoming + ? closableIds.map(toolCallId => + toToolResponseEvent({ + threadId: execution.threadId, + toolCallId, + content, + }), + ) + : []; const currentContextUsage = mergeCurrentContextUsage( execution.currentContextUsage, diff --git a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts index 57229fcf6..b73187d8f 100644 --- a/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts +++ b/packages/trueforge-core/tests/core/runtime/openToolCallCloser.test.ts @@ -2,11 +2,7 @@ import { EventType } from '../../../src/core/events/schema'; import type { InternalEnrichedAssistantMessage, InternalEnrichedToolCall } from '../../../src/core/llm/LLMTypes'; import type { ContextMessage } from '../../../src/core/runtime/AgentThread.types'; import { getEmptyCurrentContextUsage } from '../../../src/core/runtime/contextUsage'; -import { - getClosableOpenToolCallIds, - getClosableOpenToolCalls, - OpenToolCallCloser, -} from '../../../src/core/runtime/OpenToolCallCloser'; +import { getClosableOpenToolCallIds, OpenToolCallCloser } from '../../../src/core/runtime/OpenToolCallCloser'; import '../harnessMocks'; function makeToolCall( @@ -119,12 +115,9 @@ describe('getClosableOpenToolCallIds', () => { makeToolCall('tc-client', { is_client_side: true }), makeToolCall('tc-sub-agent', { is_thread_creation: true }), ]); - expect(getClosableOpenToolCalls({ context, userMessageIncoming: true })).toEqual([ - { tool_call_id: 'tc-regular', close_kind: 'dangling' }, - { tool_call_id: 'tc-approval', close_kind: 'user_action' }, - { tool_call_id: 'tc-client', close_kind: 'user_action' }, - { tool_call_id: 'tc-sub-agent', close_kind: 'thread_creation' }, - ]); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: true })).toEqual( + new Set(['tc-regular', 'tc-approval', 'tc-client', 'tc-sub-agent']), + ); }); it('still excludes already-resolved calls when a user message is incoming', () => { @@ -135,9 +128,7 @@ describe('getClosableOpenToolCallIds', () => { ]), { role: 'tool', tool_call_id: 'tc-approval', content: 'already closed' }, ]; - expect(getClosableOpenToolCalls({ context, userMessageIncoming: true })).toEqual([ - { tool_call_id: 'tc-sub-agent', close_kind: 'thread_creation' }, - ]); + expect(getClosableOpenToolCallIds({ context, userMessageIncoming: true })).toEqual(new Set(['tc-sub-agent'])); }); }); @@ -171,18 +162,41 @@ describe('OpenToolCallCloser.processPreSend', () => { expect(yielded[0]?.output).toEqual([]); }); - it('emits tool.response events for user-action and thread-creation closures', async () => { + it('cancels every unmatched last-assistant call with the same content and tool.response events', async () => { const yielded = await collectPreSend( assistantWithToolCalls([ + makeToolCall('tc-regular'), makeToolCall('tc-approval', { is_approval_required: true }), makeToolCall('tc-sub-agent', { is_thread_creation: true }), ]), true, ); + const cancelled = 'Tool call was cancelled: a new turn was started.'; expect(yielded).toHaveLength(1); + expect(yielded[0]?.context).toEqual([ + { role: 'tool', tool_call_id: 'tc-regular', content: cancelled }, + { role: 'tool', tool_call_id: 'tc-approval', content: cancelled }, + { role: 'tool', tool_call_id: 'tc-sub-agent', content: cancelled }, + ]); expect(yielded[0]?.output).toEqual([ - expect.objectContaining({ type: EventType.TOOL_RESPONSE, tool_call_id: 'tc-approval', thread_id: 'main' }), - expect.objectContaining({ type: EventType.TOOL_RESPONSE, tool_call_id: 'tc-sub-agent', thread_id: 'main' }), + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + tool_call_id: 'tc-regular', + thread_id: 'main', + content: cancelled, + }), + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + tool_call_id: 'tc-approval', + thread_id: 'main', + content: cancelled, + }), + expect.objectContaining({ + type: EventType.TOOL_RESPONSE, + tool_call_id: 'tc-sub-agent', + thread_id: 'main', + content: cancelled, + }), ]); }); @@ -191,6 +205,6 @@ describe('OpenToolCallCloser.processPreSend', () => { const first = await collectPreSend(context, false); const closed = [...context, ...(first[0]?.context ?? [])]; expect(await collectPreSend(closed, false)).toEqual([]); - expect(getClosableOpenToolCalls({ context: closed, userMessageIncoming: true })).toEqual([]); + expect(getClosableOpenToolCallIds({ context: closed, userMessageIncoming: true })).toEqual(new Set()); }); }); diff --git a/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts b/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts index f518dd904..8f1edb972 100644 --- a/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts +++ b/packages/trueforge-core/tests/orchestration/orchestrationUserMessage.test.ts @@ -69,7 +69,11 @@ describe('orchestration: user message while work is pending', () => { expect(steered.result.required_actions).toEqual([]); expect(llmCreateInputs(thread.definition.modelClient).at(-1)).toMatchObject({ messages: expect.arrayContaining([ - { role: 'tool', tool_call_id: WRITE_NOTE_CALL_ID, content: expect.stringContaining('new message') }, + { + role: 'tool', + tool_call_id: WRITE_NOTE_CALL_ID, + content: 'Tool call was cancelled: a new turn was started.', + }, { role: 'user', content: 'never mind, do this instead' }, ]), }); From 09bba0889c007e938b821a1057d5ef52e5af3752 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 09:46:55 +0530 Subject: [PATCH 08/11] nit --- .changeset/{pre => }/steer-session-anytime.md | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename .changeset/{pre => }/steer-session-anytime.md (100%) diff --git a/.changeset/pre/steer-session-anytime.md b/.changeset/steer-session-anytime.md similarity index 100% rename from .changeset/pre/steer-session-anytime.md rename to .changeset/steer-session-anytime.md From ebdd489b6933e8d7d6eeb28e5e824918e2ef5127 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 17:21:13 +0530 Subject: [PATCH 09/11] nit --- .../trueforge-core/src/core/runtime/AgentThread.ts | 10 ++++++---- .../src/core/runtime/AgentThreadOrchestrator.ts | 7 +++---- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/packages/trueforge-core/src/core/runtime/AgentThread.ts b/packages/trueforge-core/src/core/runtime/AgentThread.ts index 0b8466f38..b36d827a8 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThread.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThread.ts @@ -296,7 +296,6 @@ function validateInputMessageTypesGivenContext( context: ContextMessage[], messages: AgentThreadRuntimeSendInput[], ): void { - const hasUserMessage = messages.some(isInputUserMessage); // Full open set: validates incoming tool responses and dedupes within the batch. const openToolCallIds = getOpenToolCallIds(context); const pendingApprovalIds = new Set(getPendingApprovalToolCalls(context).map(tc => tc.id)); @@ -324,9 +323,12 @@ function validateInputMessageTypesGivenContext( } } - // User messages interrupt pending approvals / client-side calls (OpenToolCallCloser - // synthesizes the missing responses). Empty and action batches still must resolve them all. - if (!hasUserMessage && (pendingApprovalIds.size > 0 || pendingClientSideIds.size > 0)) { + // User messages interrupt pending work (OpenToolCallCloser synthesizes responses). + if (messages.some(isInputUserMessage)) { + return; + } + + if (pendingApprovalIds.size > 0 || pendingClientSideIds.size > 0) { const missing = [...pendingApprovalIds, ...pendingClientSideIds]; throw new InvalidAgentSendInputError( `Send batch must resolve all pending tool calls awaiting user input. Missing: ${missing.join(', ')}`, diff --git a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts index cc83a8d4f..b562ac482 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts @@ -45,11 +45,10 @@ function agentThreadEventToTerminalFields(event: AgentThreadExecutionEvent): { if ('parent' in event && event.parent) { return {}; } - // only a successful root completion becomes turn output - if (event.status !== 'done') { - return {}; + if (event.status === 'done') { + return { output: event.output }; } - return { output: event.output }; + return {}; } case EventType.TOOL_APPROVAL_REQUIRED: case EventType.TOOL_RESPONSE_REQUIRED: From 4963fd6c46eb0420202298bdd4c12253c8449dc9 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 17:32:33 +0530 Subject: [PATCH 10/11] fix spans --- packages/trueforge-core/src/core/runtime/contextUtils.ts | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/packages/trueforge-core/src/core/runtime/contextUtils.ts b/packages/trueforge-core/src/core/runtime/contextUtils.ts index d5cdffca4..63f0982c2 100644 --- a/packages/trueforge-core/src/core/runtime/contextUtils.ts +++ b/packages/trueforge-core/src/core/runtime/contextUtils.ts @@ -80,6 +80,12 @@ export function isInternalThreadDoneError( return event.status === 'error'; } +export function isInternalThreadDoneCancelled( + event: InternalThreadDoneEvent, +): event is InternalThreadDoneEvent & { status: 'cancelled' } { + return event.status === 'cancelled'; +} + export function scanApprovalDecisions(context: ContextMessage[]): Map { const decisions = new Map(); for (const msg of context) { From 564171dbfd3f28881f1051effc22760a2e1f2a20 Mon Sep 17 00:00:00 2001 From: Srajan Asthana Date: Thu, 17 Sep 2026 17:32:40 +0530 Subject: [PATCH 11/11] nit --- .../src/core/runtime/AgentThreadOrchestrator.ts | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts index b562ac482..c0a894a2c 100644 --- a/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts +++ b/packages/trueforge-core/src/core/runtime/AgentThreadOrchestrator.ts @@ -29,6 +29,7 @@ import { getThreadId, isApprovalDecisionMessage, isClientSideToolResponseMessage, + isInternalThreadDoneCancelled, isInternalThreadDoneError, } from './contextUtils'; import type { CreateDynamicSubAgentThread } from './CreateDynamicSubAgentThread'; @@ -165,7 +166,10 @@ function createRootAgentSpan(mainThread: AgentThread, tracing: AgentTracing): Ro if (isInternalThreadDoneError(chunk)) { rootAgentErrorMessage = chunk.error; trace.setOutput(JSON.stringify({ error: chunk.error })); - } else if (chunk.status === 'done') { + } else if (isInternalThreadDoneCancelled(chunk)) { + rootAgentErrorMessage = 'Agent thread cancelled'; + trace.setOutput(JSON.stringify({ cancelled: true })); + } else { const content = assistantMessageContentToStringForSubAgent(chunk.output.content); trace.setOutput(JSON.stringify({ result: content })); } @@ -207,7 +211,11 @@ export async function* wrapWithSubAgentSpan( subTrace.setOutput(JSON.stringify({ error: event.error })); subTrace.setMetrics(currentThread.getAgentThreadMetrics()); subTrace.setError(event.error); - } else if (event.status === 'done') { + } else if (isInternalThreadDoneCancelled(event)) { + subTrace.setOutput(JSON.stringify({ cancelled: true })); + subTrace.setMetrics(currentThread.getAgentThreadMetrics()); + subTrace.setError('Agent thread cancelled'); + } else { const content = assistantMessageContentToStringForSubAgent(event.output.content); subTrace.setOutput(JSON.stringify({ result: content })); subTrace.setMetrics(currentThread.getAgentThreadMetrics());