Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CONTEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -297,3 +297,12 @@ cancellation, output, and command results remain invocation-scoped.

See [ADR 0002](docs/adr/0002-stable-execution-worlds.md) for the adapter
identity and compatibility boundary.

## Prepared Subagent Invocations

A **Prepared Subagent Invocation** is a foreground built-in subagent invocation
started after its provider seals the arguments but before the parent model
attempt completes. It is owned by that attempt and the executing agent, bound
to the tool-call ID and canonical arguments, and adopted once by ToolNode for
normal result processing. It is process-local and is not a Durable Subagent
Execution replay mechanism. See [ADR 0009](docs/adr/0009-prepare-sealed-subagent-invocations.md).
80 changes: 80 additions & 0 deletions docs/adr/0009-prepare-sealed-subagent-invocations.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
# ADR 0009: Prepare Sealed Foreground Subagent Invocations in ToolNode

## Status

Accepted

## Context

Event-driven tools already start from provider stream completion signals, but
all direct graph tools are excluded. A parent therefore finishes generating
every subagent prompt before any child starts. The graph-tool exclusion also
preserves a real ordering guarantee: direct tools can interrupt or redirect
the batch before event-driven tools reach the host.

Starting SubagentExecutor directly from the stream would duplicate ToolNode's
runtime setup and bypass its hooks, replay, output limits and reference
registration. Its Execution Record only deduplicates pending child work; it
is not a replacement for the complete tool lifecycle. Durable Subagent
Execution additionally requires the finalized parent checkpoint/batch identity.

## Decision

A **Prepared Subagent Invocation** is the raw invocation of a built-in
foreground subagent started by ToolNode during one explicitly open model
attempt. It is bound to its owning agent, tool-call ID and canonical arguments.
A graph-owned module reserves bounded work, adopts the raw result once, and
cancels abandoned work. The stream uses the existing provider sealing and
canonical-argument parsing implementation; it never infers completion from a
closing JSON brace alone.

The first adapter is the built-in subagent wrapper. It invokes the existing
SubagentExecutor through the same runtime construction as normal tools. Normal
ToolNode execution adopts its raw output, then applies existing output limits,
references and completion handling. Ordinary event tools retain the existing
direct-before-event ordering. There is no generic eager-direct-tool interface
or second subagent scheduler.

Early starts require event-driven execution with eager execution enabled, a
single built-in subagent graph tool, no checkpointer, no human-in-the-loop
mode, no interrupting tool names, and no parent tool lifecycle hooks. Background
calls, continued child threads, output-reference arguments and excluded tools
retain normal execution. Child lifecycle hooks still execute in the child.
An observation-only PostToolBatch registry remains compatible.

The provider-attempt scope opens admission and closes it before retry or
fallback. Completion reconciles reservations with the final tool calls. Late
buffered stream events cannot reopen a closed attempt. If an attempt fails,
discards or changes a call after delegation begins, the run fails closed rather
than automatically retrying work whose effects cannot be undone. Reset and
terminal cleanup cancel unadopted and still-running adopted invocations.

`eagerEventToolExecution.maxPendingSubagents` bounds outstanding early work
and retained raw results per graph; the default is four and zero disables this
path. Excess calls execute normally with the completed batch. An unsettled tool invocation keeps its admission slot until it settles. This is an early-work
limit, not a new process-wide limit for all foreground subagents.

## Consequences

- Long sibling prompts overlap with the first child's execution.
- ToolNode remains the module owning tool runtime and result processing;
streaming does not acquire graph lifecycle responsibilities.
- The reservation interface concentrates identity, cancellation, admission,
failure containment and result adoption in one testable module.
- Checkpointed/HITL runs keep their established replay semantics and do not
receive the latency improvement in this first pass.
- Earlier child side effects are possible before the parent response commits.
Cancellation is best effort, never rollback. Automatic provider fallback is
deliberately unavailable after an early child invocation starts.
- Normal result order, references and completion events remain batch-owned;
child activity can appear before the parent finishes generating sibling calls.
- Tests exercise actual streamed model and child graphs, plus cancellation,
argument changes, duplicate admission, capacity and ToolNode adoption.

### Trace ownership

Each early invocation opens a singleton tool-dispatch chain before starting its tool and child graph. The normal completed batch collects that result and dispatches deferred calls. This represents the actual streaming timeline without reparenting spans or transferring ownership of a future graph-node span. The dispatch uses the attempt-owned callbacks and the executing agent’s tracing scope, and its input contains only that call.

The early dispatch chain emits only completion status. Raw results stay in the preparation closure and tool observation, preserving tool-specific output redaction even for scalar results that cannot be recursively identified as tool messages.

Capacity tracks outstanding tool invocation promises. Cancellation is still governed by the existing tool and executor contracts: a provider request that outlives a settled cancellation response is not a physical resource this process-local registry can account for. The tracing wrapper must not add an earlier settlement race of its own.
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@librechat/agents",
"version": "3.7.21",
"version": "3.7.22",
"reova": {
"enabled": true,
"endpoint": "https://telemetry.reo.dev/data"
Expand Down
68 changes: 63 additions & 5 deletions src/graphs/Graph.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,7 @@ import {
findCallback,
type CallbackEntry,
} from '@/utils/callbacks';
import { PreparedSubagents, PreparedSubagentError } from '@/tools/preparedSubagents';
import { ToolNode as CustomToolNode, toolsCondition } from '@/tools/ToolNode';
import { shouldTraceToolNodeForLangfuse } from '@/langfuseToolOutputTracing';
import { createLocalCodingToolBundle } from '@/tools/local/LocalCodingTools';
Expand Down Expand Up @@ -898,6 +899,7 @@ export abstract class Graph<
* Call after a run completes and content has been extracted.
*/
clearHeavyState(): void {
this.preparedSubagents.clear();
this.config = undefined;
this.signal = undefined;
this.callerSignal = undefined;
Expand Down Expand Up @@ -1239,6 +1241,8 @@ export abstract class Graph<
restoreSubagentResumeState(state: SubagentToolNodeResumeState): void;
}> = new Set();
private _subagentExecutors = new Set<SubagentExecutor>();
readonly preparedSubagents = new PreparedSubagents();
protected readonly subagentToolNodes = new Map<string, CustomToolNode<t.BaseGraphState>>();
public getOrCreateFileCheckpointer(): t.LocalFileCheckpointer | undefined {
// Return the cached instance unconditionally if one exists. The
// toolExecution check below decides whether to *create* a new
Expand Down Expand Up @@ -1397,16 +1401,16 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
* paths must not run in that state. */
protected resolveTrippedBreakerReason(
breakerSignal: AbortSignal = this.breakerAbort.signal
): StreamLimitExceededError | undefined {
): StreamLimitExceededError | PreparedSubagentError | undefined {
if (
breakerSignal.aborted &&
breakerSignal.reason instanceof StreamLimitExceededError
(breakerSignal.reason instanceof StreamLimitExceededError || breakerSignal.reason instanceof PreparedSubagentError)
) {
return breakerSignal.reason;
}
if (
this.signal?.aborted === true &&
this.signal.reason instanceof StreamLimitExceededError
(this.signal.reason instanceof StreamLimitExceededError || this.signal.reason instanceof PreparedSubagentError)
) {
return this.signal.reason;
}
Expand Down Expand Up @@ -1581,6 +1585,7 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
/* Init */

resetValues(keepContent?: boolean, checkpointScope?: string): void {
this.preparedSubagents.clear();
this.resetSubagentCheckpointThreadIds();
this.messages = [];
this.hasRestoredCheckpointMessages = false;
Expand Down Expand Up @@ -2717,6 +2722,49 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {

/* Graph */

canPrestartSubagents(agentContext?: AgentContext): boolean {
if (
this.eagerEventToolExecution?.enabled !== true ||
(this.eagerEventToolExecution.maxPendingSubagents ?? 4) <= 0 ||
this.compileOptions?.checkpointer != null ||
this.humanInTheLoop?.enabled === true ||
(this.interruptingToolNames?.length ?? 0) > 0 ||
agentContext == null ||
!this.subagentToolNodes.has(agentContext.agentId)
) {
return false;
}
const graphTools = agentContext.graphTools as t.GenericTool[] | undefined;
return (
graphTools?.length === 1 &&
graphTools[0].name === Constants.SUBAGENT &&
SUBAGENT_REPLAY_CONTROLLER in graphTools[0]
);
}

prestartSubagent(
call: ToolCall,
attempt: string,
agentContext?: AgentContext
): void {
const config = this.preparedSubagents.getConfig(attempt);
if (
!this.canPrestartSubagents(agentContext) ||
config == null ||
agentContext == null
) {
return;
}
this.subagentToolNodes
.get(agentContext.agentId)
?.prestartSubagent(
call,
attempt,
config,
this.eagerEventToolExecution?.maxPendingSubagents ?? 4
);
}

initializeTools({
currentTools,
currentToolMap,
Expand Down Expand Up @@ -2791,6 +2839,7 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
hookRegistry: this.hookRegistry,
humanInTheLoop: this.humanInTheLoop,
eagerEventToolExecution: this.eagerEventToolExecution,
preparedSubagents: this.preparedSubagents,
codeSessionToolNames: this.codeSessionToolNames,
eagerEventToolExecutions: this.eagerEventToolExecutions,
eagerEventToolUsageCount: this.getEagerEventToolUsageCount(
Expand All @@ -2816,6 +2865,9 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
StandardGraph.handleToolCallErrorStatic(this, data, metadata),
});
this.registerCompiledToolNode(node);
if (agentContext != null) {
this.subagentToolNodes.set(agentContext.agentId, node);
}
return node;
}

Expand Down Expand Up @@ -4092,7 +4144,10 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
* succeeding fallback would resolve a run the public contract says
* must reject. Rethrow before any recovery path.
*/
if (primaryError instanceof StreamLimitExceededError) {
if (
primaryError instanceof StreamLimitExceededError ||
primaryError instanceof PreparedSubagentError
) {
/** Tripped before rethrowing so parallel agent nodes' in-flight
* model calls and subagents stop while the rejection propagates. */
attemptBreaker.abort(primaryError);
Expand Down Expand Up @@ -4342,7 +4397,10 @@ export class StandardGraph extends Graph<t.BaseGraphState, t.GraphNode> {
})
);
} catch (fallbackError) {
if (fallbackError instanceof StreamLimitExceededError) {
if (
fallbackError instanceof StreamLimitExceededError ||
fallbackError instanceof PreparedSubagentError
) {
/** Same treatment as the primary path: a fallback stream that
* trips the breaker must stop parallel agent nodes' model calls
* and subagents before the rejection propagates. */
Expand Down
43 changes: 35 additions & 8 deletions src/llm/invoke.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ import {
} from '@/llm/providers';
import { ChatModelStreamHandler, dispatchesChatModelStream } from '@/stream';
import { Constants, ContentTypes, GraphEvents, Providers } from '@/common';
import { PreparedSubagentError } from '@/tools/preparedSubagents';
import { assertNotTruncatedToolCall } from '@/llm/truncation';
import { resolveClientOptionsModel } from '@/llm/request';
import { safeDispatchCustomEvent } from '@/utils/events';
Expand Down Expand Up @@ -702,9 +703,8 @@ export async function attemptInvoke(
},
};
const request = resolveAttemptRequest(params, providerStampedConfig);
const configuredModel = providerStampedConfig.metadata?.[
Constants.INVOKED_MODEL
];
const configuredModel =
providerStampedConfig.metadata?.[Constants.INVOKED_MODEL];
const modelId =
request.modelId ??
(typeof configuredModel === 'string' ? configuredModel : undefined);
Expand Down Expand Up @@ -745,8 +745,14 @@ export async function attemptInvoke(
if (leaseTarget != null && generationKey != null) {
registerActiveStreamLimitGeneration(leaseTarget, generationKey);
}
const prepared =
params.context?.eagerEventToolExecution?.enabled === true
? params.context.preparedSubagents
: undefined;
const preparedAttempt = resolveGenerationKey(stampedConfig.metadata);
prepared?.begin(preparedAttempt, stampedConfig);
try {
return await attemptInvokeBody(
const result = await attemptInvokeBody(
{
request,
context: params.context,
Expand All @@ -755,7 +761,27 @@ export async function attemptInvoke(
},
stampedConfig
);
const calls =
result.messages?.flatMap((message) =>
'tool_calls' in message
? ((message.tool_calls as ToolCall[] | undefined) ?? [])
: []
) ?? [];
prepared?.finish(preparedAttempt, calls);
return result;
} catch (error) {
prepared?.finish(preparedAttempt, undefined, error);
throw error;
} finally {
const chunks = params.context?.eagerEventToolCallChunks;
if (prepared != null && chunks != null) {
const prefix = `prepared:${preparedAttempt}\u0000`;
for (const key of chunks.keys()) {
if (key.startsWith(prefix)) {
chunks.delete(key);
}
}
}
if (leaseTarget != null && generationKey != null) {
releaseStreamLimitGeneration(leaseTarget, generationKey);
}
Expand Down Expand Up @@ -1053,7 +1079,7 @@ async function attemptInvokeBody(
const signal = config.signal;
if (
signal?.aborted === true &&
signal.reason instanceof StreamLimitExceededError
(signal.reason instanceof StreamLimitExceededError || signal.reason instanceof PreparedSubagentError)
) {
throw signal.reason;
}
Expand Down Expand Up @@ -1290,6 +1316,7 @@ async function attemptInvokeBody(
const ownAbort =
restartRoute === 'aborted' &&
!(error instanceof StreamLimitExceededError) &&
!(error instanceof PreparedSubagentError) &&
config.signal?.aborted !== true;
if (!ownAbort) {
throw error;
Expand Down Expand Up @@ -1639,7 +1666,7 @@ export async function tryFallbackProviders({
* a run that must reject. Check before every fallback invocation. */
if (
config?.signal?.aborted === true &&
config.signal.reason instanceof StreamLimitExceededError
(config.signal.reason instanceof StreamLimitExceededError || config.signal.reason instanceof PreparedSubagentError)
) {
throw config.signal.reason;
}
Expand Down Expand Up @@ -1673,7 +1700,7 @@ export async function tryFallbackProviders({
* provider failure. Continuing would try the remaining fallbacks and a
* succeeding one would resolve a run that must reject.
*/
if (e instanceof StreamLimitExceededError) {
if (e instanceof StreamLimitExceededError || e instanceof PreparedSubagentError) {
throw e;
}
/** A parallel sibling's trip aborts this branch's composed signal, and
Expand All @@ -1682,7 +1709,7 @@ export async function tryFallbackProviders({
* abort. Rethrow the breaker's own reason instead. */
if (
config?.signal?.aborted === true &&
config.signal.reason instanceof StreamLimitExceededError
(config.signal.reason instanceof StreamLimitExceededError || config.signal.reason instanceof PreparedSubagentError)
) {
throw config.signal.reason;
}
Expand Down
Loading
Loading