Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
aaa99c6
fix(Toolkit): route callTool through executeWithInfrastructure so id/…
KIM406-CMD Sep 13, 2026
c362855
fix(test): add ToolUseBlock.content to pass upstream schema validation
KIM406-CMD Sep 25, 2026
cbdfcfe
docs(Toolkit): document callTool(ToolCallParam) execution semantics
KIM406-CMD Sep 25, 2026
292a7da
docs(Toolkit): fix callTool javadoc — correct shutdown guard attribut…
KIM406-CMD Sep 25, 2026
2bcea60
docs(Toolkit): add scheduling hop note to callTool javadoc
KIM406-CMD Sep 25, 2026
63b85bd
docs(ToolExecutor): fix stale javadoc and align execute* family infra…
KIM406-CMD Sep 28, 2026
36a5738
docs(ToolExecutor): address PR #3130 review comments on javadoc
KIM406-CMD Sep 30, 2026
ffac69a
fix(tool): address review feedback — timeout opt-out, retry doc accur…
KIM406-CMD Oct 1, 2026
84440a5
Merge branch 'main' into fix
KIM406-CMD Oct 1, 2026
a826bfd
test(tool): add regression test for content-validation split
KIM406-CMD Oct 1, 2026
0cf254f
fix(tool): opt-out timeout sentinel, ToolRegistry atomicity, merge he…
KIM406-CMD Oct 1, 2026
d4a081d
fix(tool): fix retry/shutdown javadoc, add per-call config tests
KIM406-CMD Oct 1, 2026
549af6c
fix(tool): final javadoc and test: retry claim, callTool per-call cov…
KIM406-CMD Oct 1, 2026
6ce5af4
docs(tool): pin identity semantics in removeToolIfSame javadoc
KIM406-CMD Oct 1, 2026
d47a113
docs(test): flag that negative validation test pins current gap, not …
KIM406-CMD Oct 1, 2026
28f7484
fix(tool): add isTimeoutDisabled() helper and unify NO_TIMEOUT checks…
KIM406-CMD Oct 3, 2026
621faef
Merge branch 'main' into fix
KIM406-CMD Oct 3, 2026
772f1b6
fix(tool): tighten isTimeoutDisabled to sentinel, add builder validat…
KIM406-CMD Oct 3, 2026
b7ad1cc
fix: guard NO_TIMEOUT sentinel in openai-official SDK path and javado…
KIM406-CMD Oct 4, 2026
a9f7be8
Merge branch 'main' into fix
KIM406-CMD Oct 4, 2026
e21cbd6
fix: address all four info-level review notes from round 7
KIM406-CMD Oct 4, 2026
dcaf911
fix: address round-9 review — error message, log, mergeConfigs tests,…
KIM406-CMD Oct 4, 2026
bd1e1c6
test: add Duration.ZERO rejection test for Builder.timeout()
KIM406-CMD Oct 4, 2026
1896833
fix: align Builder.timeout() javadoc — ZERO is rejected, not immediat…
KIM406-CMD Oct 4, 2026
016626a
test: verify mergeConfigs inherits NO_TIMEOUT when primary timeout is…
KIM406-CMD Oct 4, 2026
8b1d6aa
fix: address round-10 review — error message includes value, producti…
KIM406-CMD Oct 5, 2026
6bed472
style: fix spotless formatting violation in ModelTimeoutRetryTest
KIM406-CMD Oct 5, 2026
e63fede
docs: add changelog entry for ExecutionConfig.Builder.timeout() valid…
KIM406-CMD Oct 5, 2026
dc85c84
Merge branch 'main' into fix
KIM406-CMD Oct 7, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,33 @@ private static boolean isRetryableError(Throwable error) {
.retryOn(RETRYABLE_ERRORS)
.build();

/**
* Sentinel value for {@link #timeout} meaning "no timeout". A negative duration is never
* produced by normal usage and is recognised by {@code ToolExecutor.applyTimeout} and
* {@link ModelUtils#applyTimeoutAndRetry ModelUtils.applyTimeoutAndRetry}
* as "skip the timeout operator entirely".
*
* <p>This is the only way to opt out of the timeout that {@link #TOOL_DEFAULTS} and {@link
* #MODEL_DEFAULTS} always carry, because {@link #mergeConfigs} treats {@code null} as
* "inherit from fallback".
*
* <p><b>Supported paths</b>: The sentinel is honoured on the <em>tool</em> path
* ({@code ToolExecutor.applyTimeout}) and the <em>model-flux</em> path
* ({@code ModelUtils.applyTimeoutAndRetry}), where it means genuinely unbounded
* (no timeout operator is applied). Extension consumers that read
* {@link #getTimeout()} directly and pass the value to a framework timeout operator
* must apply the same guard via {@link #isTimeoutDisabled()}:
* <ul>
* <li>{@code EmbeddingUtils.applyTimeoutAndRetry} in {@code rag-simple} — honours the
* sentinel by skipping the timeout operator, consistent with tool/model-flux.</li>
* <li>{@code openai-official} provider — the sentinel is recognised but degrades to
* the OpenAI SDK's own default request timeout rather than being truly unbounded,
* because the SDK client does not accept an unbounded timeout value. A warning is
* logged when this occurs.</li>
* </ul>
*/
public static final Duration NO_TIMEOUT = Duration.ofNanos(-1);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Critical] NO_TIMEOUT is a public constant on ExecutionConfig, but only the tool path recognises it. ToolExecutor.applyTimeout now guards with isNegative(), while ModelUtils.applyTimeoutAndRetry (agentscope-core/src/main/java/io/agentscope/core/model/ModelUtils.java:85) still only checks timeout != null and hands the value straight to responseFlux.timeout(...). So ExecutionConfig.builder().noTimeout().build() on a model/embedding call does not mean "no timeout" — a negative duration reaches Reactor and either fails validation or fires immediately.

The javadoc right above makes this worse: it explicitly advertises NO_TIMEOUT as the way to opt out of "the timeout that TOOL_DEFAULTS and MODEL_DEFAULTS always carry", which invites exactly the model-path usage that is not handled. Since this repo's own MODEL_DEFAULTS is merged into every model call, callers will hit this.

Suggest centralising the predicate so no consumer can forget it, e.g.

public boolean hasTimeout() {
    return timeout != null && !timeout.isNegative();
}

and using it in both ToolExecutor.applyTimeout and ModelUtils.applyTimeoutAndRetry (plus EmbeddingUtils, which routes through the same helper). A model-path test mirroring ModelTimeoutRetryTest with noTimeout() would pin it down.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Carry-over from the previous round, still open at d47a113d: NO_TIMEOUT is honoured only on the tool path. ToolExecutor.applyTimeout now skips the operator when config.getTimeout().isNegative(), but ModelUtils.applyTimeoutAndRetry still does if (timeout != null) { responseFlux.timeout(timeout, ...) }, so a model/embedding call configured with ExecutionConfig.builder().noTimeout() hands Duration.ofNanos(-1) directly to Reactor — which either rejects the negative duration or fires the timeout immediately. Either way the call does not run unbounded, the opposite of what this javadoc promises, and the javadoc explicitly cites MODEL_DEFAULTS (which always carries a timeout) as the reason the sentinel exists.

Two ways to close it: mirror the isNegative() guard in ModelUtils, or — preferred, because every consumer otherwise has to remember the sentinel — put the decision on the value itself:

/** @return true when the configured timeout is the {@link #NO_TIMEOUT} sentinel. */
public boolean isTimeoutDisabled() {
    return timeout != null && timeout.isNegative();
}

and use it in both ToolExecutor.applyTimeout and ModelUtils.applyTimeoutAndRetry, plus one model-path regression test so the two call sites cannot drift apart again. Non-blocking for this PR if you prefer to scope it out: it is still a public-API promise that the runtime does not keep.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer Good call — added ExecutionConfig#isTimeoutDisabled() and switched both
ToolExecutor.applyTimeout and ModelUtils.applyTimeoutAndRetry to use it.

  • ExecutionConfig.isTimeoutDisabled() returns true when timeout is the NO_TIMEOUT sentinel
  • ToolExecutor.applyTimeout: config.getTimeout().isNegative() → config.isTimeoutDisabled()
  • ModelUtils.applyTimeoutAndRetry: if (timeout != null) → if (timeout != null && !execConfig.isTimeoutDisabled())

This way both call sites share the same guard and won't drift apart.
See: 28f7484


/**
* Standard defaults for tool executions.
*
Expand Down Expand Up @@ -184,6 +211,20 @@ public Duration getTimeout() {
return timeout;
}

/**
* Returns true when the configured timeout is the {@link #NO_TIMEOUT} sentinel,
* meaning consumers should skip applying any timeout operator.
*
* <p>Only the exact {@link #NO_TIMEOUT} sentinel is recognised; stray negative
* durations (which should be rejected by {@link Builder#timeout(Duration)}) are
* not treated as "no timeout".
*
* @return true if timeout is disabled via {@link #NO_TIMEOUT}
*/
public boolean isTimeoutDisabled() {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] The javadoc says this returns true when the timeout is the NO_TIMEOUT sentinel, but the implementation matches any negative Duration, and Builder.timeout(...) accepts negative values without validation.

So ExecutionConfig.builder().timeout(Duration.ofSeconds(-30)).build() silently becomes "no timeout" instead of a config error — the same trap this sentinel exists to avoid. Two ways to close the gap:

public boolean isTimeoutDisabled() {
    return NO_TIMEOUT.equals(timeout);
}

…or keep the broad check but reject accidental negatives at construction time:

public Builder timeout(Duration timeout) {
    if (timeout != null && timeout.isNegative() && !NO_TIMEOUT.equals(timeout)) {
        throw new IllegalArgumentException(
                "timeout must be >= 0; use NO_TIMEOUT (or noTimeout()) to disable it");
    }
    this.timeout = timeout;
    return this;
}

The second option also matches how maxAttempts(...) already validates its argument, so it would be consistent with the rest of the builder.

return NO_TIMEOUT.equals(timeout);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] Exact-sentinel matching is the correct narrowing — "any negative duration means disabled" is exactly what let the isNegative() / "inf" style checks diverge across consumers, and with the builder validation in place a stray negative can no longer enter through the builder.

What remains is that getTimeout() still hands out -1ns, so each consumer has to remember to pair it with isTimeoutDisabled(). This PR had to patch four call sites (ToolExecutor, ModelUtils, EmbeddingUtils, and the openai-official client factory) and the new "Supported paths" javadoc is effectively a hand-maintained list that will drift as providers are added. To make the contract enforceable rather than documented, either:

  • add Duration effectiveTimeoutOrNull() returning null when disabled and have consumers use only that (leaving getTimeout() as the raw accessor for compatibility), or
  • centralize wrapping in one shared applyTimeoutAndRetry helper that core and extensions both call.

Fine as a follow-up, not a blocker for this PR.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed — effectiveTimeoutOrNull() (or a shared helper) is the right long-term solution. I'll track it as a follow-up; not in this PR.

}

/**
* Gets the maximum number of attempts.
*
Expand Down Expand Up @@ -305,14 +346,36 @@ public static class Builder {
/**
* Sets the timeout duration for a single execution.
*
* @param timeout the timeout duration, or null for no timeout
* @param timeout the timeout duration (must be &gt; 0, or {@link #NO_TIMEOUT}),
* or null to inherit from fallback; {@code Duration.ZERO} is not a synonym
* for {@link #noTimeout()} and is therefore rejected
* @return this builder instance
* @throws IllegalArgumentException if timeout is a negative duration other than
* {@link #NO_TIMEOUT}, or if timeout is {@code Duration.ZERO}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] One item from the previous round is still open, and it is documentation rather than logic: Builder.timeout(...) now rejects values it used to accept (Duration.ZERO and any negative duration other than the sentinel), and that rejection fires on the mergeConfigs path too — so code written against an earlier agentscope-core can start throwing IllegalArgumentException after an upgrade, at config-merge time rather than at the call site. The message improvement (got <value>) makes it diagnosable, which helps a lot. Could you add a line to docs/v2/en/docs/change-log.md + docs/v2/zh/docs/change-log.md noting the stricter validation and pointing at noTimeout()? That is the only thing I would still like to see here; I am not blocking the review on it.

*/
public Builder timeout(Duration timeout) {
if (timeout != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] Right place for the validation, but note that mergeConfigs funnels through this same builder:

builder.timeout(primary.timeout != null ? primary.timeout : fallback.timeout);

so a config produced before this change — or assembled by a path that bypasses the builder, e.g. deserialized session/agent state or an extension option map — carrying Duration.ZERO or a stray negative now raises IllegalArgumentException during a merge, i.e. per request inside the reactive chain rather than once at configuration time. Please confirm no persistence/rehydration path can feed a legacy value back in; if it can, the read side needs normalization in addition to builder validation.

Two test asks while you are here:

  • mergeConfigs(noTimeout, withTimeout) must keep the sentinel and mergeConfigs(withTimeout, noTimeout) must keep the positive timeout — the current tests pin the consumer-side skip but not the merge precedence.
  • a non-sentinel negative must be rejected, not treated as disabled, so the semantics cannot regress to the old isNegative() behaviour.

And a changelog line: Builder.timeout(Duration.ZERO) used to be accepted (instant-expiry) and now throws. Rejecting it is the better contract, but a config-driven timeout: 0 will hit it after upgrade.

Nit: the message says "timeout must be > 0" while the branch also rejects ZERO — "timeout must be positive; use NO_TIMEOUT (or noTimeout()) to disable it" names both cases.

@KIM406-CMD KIM406-CMD Oct 4, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All three items addressed in dcaf9112:

Error message changed to "timeout must be positive; use NO_TIMEOUT or noTimeout() to disable it" — covers both ZERO and negative cases.
mergeConfigs(noTimeout, withTimeout) test: asserts the sentinel survives when primary.
mergeConfigs(withTimeout, noTimeout) test: asserts the positive timeout overrides fallback NO_TIMEOUT.
shouldRejectNonSentinelNegativeDurations test: verifies a non-sentinel negative (e.g. -500ms) throws IllegalArgumentException — cannot regress to the old isNegative() behaviour.

Agreed on the changelog line for Duration.ZERO — I'll make sure it's called out.

&& (timeout.isNegative() || timeout.isZero())

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] Rejecting Duration.ZERO here is the right call semantically (a zero timeout would otherwise expire every call instantly), but it is a runtime-visible change on a public builder: code that today does .timeout(Duration.ZERO) and builds successfully will start throwing IllegalArgumentException after this lands. Since ExecutionConfig is public API and timeout(...) is documented in the v1/v2 task docs, this deserves a line in the changelog/release notes so downstream users are not surprised by a build-time exception where they previously had a working (if pointless) config. The mergeConfigs() path is fine either way — no instance can hold ZERO because every value goes through this builder. Not blocking; just call it out in the docs.

@KIM406-CMD KIM406-CMD Oct 4, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer Acknowledged. The javadoc already documents the Duration.ZERO rejection. Maintainers can include this in the release notes when applicable. Thanks.

&& !NO_TIMEOUT.equals(timeout)) {
throw new IllegalArgumentException(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] Throwing from Builder.timeout(...) turns a previously accepted value into a hard failure, and this builder is reached through mergeConfigs (line 315: builder.timeout(primary.timeout != null ? primary.timeout : fallback.timeout)). Any config that already exists with a stray negative/zero timeout — built before this change, or assembled by third-party code that constructs ExecutionConfig reflectively/deserialised — will now throw at merge time rather than at configuration time, i.e. deeper inside the request path where the stack trace won't point at the offending builder call.

Two things worth doing before merge:

  1. Note the behaviour change in the PR description / changelog — users who today pass Duration.ZERO (or a negative) and effectively get "no timeout applied" now get an IllegalArgumentException instead.
  2. Include the rejected value in the message ("timeout must be positive, got " + timeout) so the failure is self-describing; maxAttempts/backoffMultiplier have the same weakness but this one is newly enforced on a previously-valid input, so it is the one users will hit.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done. The error message now includes the rejected value: "timeout must be positive, got " + timeout + "; use NO_TIMEOUT or noTimeout() to disable it".

"timeout must be positive, got "
+ timeout
+ "; use NO_TIMEOUT or noTimeout() to disable it");
}
this.timeout = timeout;
return this;
}

/**
* Opt out of timeout entirely for this call. Equivalent to {@code timeout(NO_TIMEOUT)}.
* This is the only way to prevent the timeout inherited from {@link #TOOL_DEFAULTS} /
* {@link #MODEL_DEFAULTS}, because {@link #mergeConfigs} treats {@code null} as inherit.
*/
public Builder noTimeout() {
this.timeout = NO_TIMEOUT;
return this;
}

/**
* Sets the maximum number of attempts (including the initial attempt).
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ public static Flux<ChatResponse> applyTimeoutAndRetry(
if (execConfig != null) {
// Apply timeout if configured
Duration timeout = execConfig.getTimeout();
if (timeout != null) {
if (timeout != null && !execConfig.isTimeoutDisabled()) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] This is a genuine behavior fix and worth calling out: Reactor clamps a negative delay to immediate expiry, so before this change timeout(NO_TIMEOUT) on the model path meant "fail instantly", not "no timeout" (verified against reactor-core 3.8.4 — a slow source with .timeout(Duration.ofNanos(-1), fallback) fires the fallback right away).

The unification is still partial, though. Other consumers that null-check getTimeout() only will hand the negative sentinel straight to a timeout operator:

  • agentscope-extensions-rag-simple — EmbeddingUtils.applyTimeoutAndRetry (if (timeout != null) { embeddingMono.timeout(timeout, …) }) takes the same ExecutionConfig, so a user setting noTimeout() on an embedding config gets an immediate timeout error.
  • agentscope-extensions-model-openai-official — OpenAIResponsesChatModel passes effectiveOptions.getExecutionConfig().getTimeout() into OpenAISdkClientFactory.createClient(…, timeout, …).

Since NO_TIMEOUT / noTimeout() become public API in this PR, could the same guard be applied in EmbeddingUtils here (it is one line), or should this PR at least state in the NO_TIMEOUT javadoc that the sentinel is honoured on the tool and model-flux paths only? A follow-up issue for the extension consumers is fine too, but worth linking from the javadoc so the limitation is discoverable.

responseFlux =
responseFlux.timeout(
timeout,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -155,7 +155,9 @@ private void invokeChunkCallback(
// ==================== Single Tool Execution ====================

/**
* Execute a single tool call with full infrastructure support.
* Execute a single tool call (core execution only: Tracer + {@link #executeCore};
* no scheduling, timeout, retry, shutdown guard, or id/name stamping). Use
* {@link #executeWithInfrastructure(ToolCallParam, ExecutionConfig)} for the full-infrastructure path.
*
* @param param Tool call parameters
* @return Mono containing execution result
Expand All @@ -171,8 +173,9 @@ Mono<ToolResultBlock> execute(ToolCallParam param) {

/**
* Execute a single tool call with a per-call tool request config and a per-call internal chunk
* callback. This is the single core entry point; the no-arg {@link #execute(ToolCallParam)}
* resolves the request config from its explicit runtime context and uses no internal callback.
* callback. This is the single core entry point; the 1-param {@link #execute(ToolCallParam)}
* overload resolves the request config from its explicit runtime context and uses no internal
* callback.
*/
Mono<ToolResultBlock> execute(
ToolCallParam param,
Expand Down Expand Up @@ -342,7 +345,10 @@ private Collection<String> resolveActiveGroups(ToolCallParam param) {
// ==================== Batch Tool Execution ====================

/**
* Execute multiple tool calls with concurrency control, timeout, and retry.
* Execute multiple tool calls with concurrency control plus full per-call infrastructure
* (scheduling, timeout, retry, shutdown guard, id/name stamping). Each single call is routed
* through {@link #executeWithInfrastructure(ToolUseBlock, ExecutionConfig, Agent,
* RuntimeContext, ToolRequestConfig, BiConsumer)}.
*
* @param toolCalls List of tool calls to execute
* @param parallel Whether to execute in parallel
Expand Down Expand Up @@ -449,33 +455,70 @@ private boolean isConcurrencySafe(ToolUseBlock toolCall, ToolRequestConfig reque
}

/**
* Execute a single tool call with infrastructure (scheduling, timeout, retry).
* Execute a single tool call with infrastructure (scheduling, timeout, retry, shutdown
* guard), and stamps the result with the tool call's id/name.
*
* <p>This overload is used by the batch path ({@link #executeAll(List, boolean,
* ExecutionConfig, Agent, RuntimeContext)}), which routes each {@link ToolUseBlock} with
* the infrastructure config, per-call request config, and chunk callback.
*/
private Mono<ToolResultBlock> executeWithInfrastructure(
Mono<ToolResultBlock> executeWithInfrastructure(
ToolUseBlock toolCall,
ExecutionConfig executionConfig,
Agent agent,
RuntimeContext agentRuntimeContext,
ToolRequestConfig requestConfig,
BiConsumer<ToolUseBlock, ToolResultBlock> internalChunkCallback) {
// Build tool call parameter
ToolCallParam param =
ToolCallParam.builder()
.toolUseBlock(toolCall)
.agent(agent)
.runtimeContext(agentRuntimeContext)
.build();

// Get core execution
Mono<ToolResultBlock> execution = execute(param, requestConfig, internalChunkCallback);

// Apply infrastructure layers
return applyInfrastructure(execution, executionConfig, toolCall);
}

/**
* Execute a single tool call with full infrastructure, preserving all fields from the
* original {@link ToolCallParam} (including input).
*
* <p>This overload is used by {@code Toolkit.callTool} so that user-supplied fields on the
* param object are not silently discarded before reaching {@link #executeCore}.
*/
Mono<ToolResultBlock> executeWithInfrastructure(
ToolCallParam param, ExecutionConfig executionConfig) {
ToolUseBlock toolCall = param.getToolUseBlock();

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] The two overloads are now asymmetric in field preservation. The 4-arg overload still rebuilds a fresh ToolCallParam from toolCall/agent/runtimeContext (so any input/emitter attached to an existing param is not carried through), while the new overload preserves the param verbatim. That is fine for today's call sites (executeAll only has the ToolUseBlock), but the doc comment here ("preserving all fields ... not silently discarded") reads as if the batch path had the same guarantee. Consider adding a note that the batch path intentionally constructs the param, or unify both to accept a ToolCallParam.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] Everything from here to the onErrorResume is a copy of :465-495 — the four infrastructure calls plus id/name stamping and the error-to-result conversion — and the only real difference is which execute(...) it starts from.

Duplicating exactly the block this PR identifies as "the thing the single path was missing" is how the two drift apart again, and the drift is invisible because each copy looks correct on its own. applyScheduling ignoring executionConfig while applyTimeout/applyRetry honor it is already only documented in one of the two copies.

Consider extracting private Mono<ToolResultBlock> applyInfrastructure(Mono<ToolResultBlock> execution, ExecutionConfig config, ToolUseBlock toolCall) and ending both overloads with return applyInfrastructure(execute(...), config, toolCall); — a future layer is then added once and both entry points get it.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer Done. I extracted the duplicated infrastructure pipeline — the four layers (applyScheduling → applyTimeout → applyRetry → applyShutdownGuard), the id/name stamping, and the error-to-result onErrorResume — into a single private applyInfrastructure(Mono, ExecutionConfig, ToolUseBlock) method. Both executeWithInfrastructure overloads now delegate to it in one line:

  Mono<ToolResultBlock> execution = execute(param);
  return applyInfrastructure(execution, executionConfig, toolCall);

This way, if a future layer (metrics tracing, rate limiting, etc.) needs to be added, it goes in one place and both entry points get it. The helper's javadoc also documents the retry-only-on-infrastructure semantics — which addresses the previous finding's doc gap as a bonus.


Mono<ToolResultBlock> execution = execute(param);

return applyInfrastructure(execution, executionConfig, toolCall);
}

/**
* Applies the shared infrastructure pipeline (scheduling, timeout, retry, shutdown guard)
* and stamps the result with the tool call's id/name. The four infrastructure layers and
* the error-to-result conversion live here so that both entry points (batch and single)
* stay in sync when a new layer is added.
*
* <p><b>Retry semantics</b>: {@link #applyRetry} only fires for the timeout
* {@code RuntimeException} emitted by {@link #applyTimeout}. Tool failures are converted
* to normal {@link ToolResultBlock#error} completions inside {@link #executeCore} before
* this pipeline runs, and {@link #applyShutdownGuard} runs <em>after</em> retry so
* shutdown signals are never seen by {@code retryWhen} either. "Retry" here means
* "retry on timeout", nothing else.
*/
private Mono<ToolResultBlock> applyInfrastructure(
Mono<ToolResultBlock> execution,
ExecutionConfig executionConfig,
ToolUseBlock toolCall) {
execution = applyScheduling(execution);
execution = applyTimeout(execution, executionConfig, toolCall);
execution = applyRetry(execution, executionConfig, toolCall);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] The callTool javadoc advertises retry ("non-idempotent tools may be re-invoked on timeout when ... maxAttempts > 1"), but applyRetry can never observe a tool that failed: executeCore ends with a blanket onErrorResume(e -> Mono.just(ToolResultBlock.error("Tool execution failed: " + errorMsg))) at :305-313, so an exception thrown by tool.callAsync is converted into a normal completion carrying an ERROR-state value before this line ever sees it.

What still works is retry over timeout and shutdown-guard signals, because those are produced by the infrastructure layers themselves (:546, and firstWithSignal at :606). So the accurate statement is "retried on timeout, never on tool failure".

Nothing regresses here — the batch path has always behaved this way — but this PR is the first place that promises retry semantics to callTool users, and the promise does not hold for the common case. Either tighten the wording, or, if retry-on-tool-failure is what's intended, move the exception-to-result conversion above the retry layer so retryWhen sees the error. A test with maxAttempts(2) plus an invocation counter would pin down which behavior is actually intended (and would also document the at-least-once caveat you call out in the javadoc).

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer This is a sharp catch — I didn't even realize half of my own test was a workaround 😅
The root cause is that executeCore validates and executes against two different payloads:
Validation reads toolCall.getContent() (:236) — schema validator gets the raw JSON string from ToolUseBlock.content
Execution reads the merged mergedInput (:275-278) — preferring param.input, falling back to toolCall.input
That's exactly why the tests in this PR needed .content(JsonUtils.getJsonCodec().toJson(input)) to pass — my own commit message called it out as "add ToolUseBlock.content to pass upstream schema validation" but didn't spell out why. The two directions are both broken: a caller who sets only param.input gets rejected outright for any tool with a non-empty schema, and a caller who stuffs out-of-schema values into param.input bypasses validation entirely.
For this PR I documented the contract in the callTool javadoc rather than changing behavior:
When you build a ToolCallParam with only input populated, also set ToolUseBlock.content to the JSON form of that input, or schema validation will reject the call.
A proper fix — validating against the effective merged input, or having ToolCallParam.Builder auto-serialize content from input — would also affect the batch path (where callTools passes through LLM-provided raw content that we don't want to silently rewrite). I suggest opening a follow-up issue for this; I'm happy to take it.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer Fully agree — the previous javadoc sentence ("non-idempotent tools may be re-invoked on timeout when maxAttempts > 1") promised more than the code delivers.
The chain is: executeCore has a blanket onErrorResume(e -> Mono.just(ToolResultBlock.error(...))) at :305-313, so any exception thrown by the tool is caught and converted into a normal ERROR-state completion before the infrastructure pipeline runs. applyRetry sits after applyTimeout, so its retryWhen only ever sees exceptions produced by the infrastructure layers themselves — timeout signals from applyTimeout and shutdown signals from applyShutdownGuard's firstWithSignal. Tool failures never reach retry.
I went with tightening the javadoc rather than changing the code (you noted this PR shouldn't fix executeCore's error conversion). The updated wording now says:
retry only fires on timeout or shutdown signals, never on tool failures. Tool exceptions are caught and converted into a normal ToolResultBlock.error completion before the retry layer runs, so maxAttempts > 1 has no effect on a failing tool — only on infrastructure-level aborts.
If we want retry to actually cover tool failures in the future, we'd need to move that blanket onErrorResume out of executeCore — that's a larger behavioral change and worth a separate PR.

execution = applyShutdownGuard(execution);

// Add tool metadata and error handling
return execution
.map(result -> result.withIdAndName(toolCall.getId(), toolCall.getName()))
.onErrorResume(
Expand All @@ -499,7 +542,8 @@ private Mono<ToolResultBlock> applyScheduling(Mono<ToolResultBlock> execution) {

private Mono<ToolResultBlock> applyTimeout(
Mono<ToolResultBlock> execution, ExecutionConfig config, ToolUseBlock toolCall) {
if (config == null || config.getTimeout() == null) {
// null = inherit from fallback, NO_TIMEOUT sentinel = explicitly disabled
if (config == null || config.getTimeout() == null || config.isTimeoutDisabled()) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] Nit on the extracted applyTimeout guard: config.getTimeout() == null is now redundant with isTimeoutDisabled() in the sense that readers will wonder which check does what. Not a bug — the null check must stay (the helper returns false for null). Suggest just folding the comment/javadoc to make the intent explicit, e.g. "null = inherit, negative = explicitly disabled", so the next person does not try to simplify it away.

Verified against reactor-core 3.8.4 that Mono/Flux.timeout(negativeDuration, …) does not throw but expires immediately — which is exactly why this guard is load-bearing and worth a test (see the ToolkitTest comment).

return execution;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,9 @@
* and retrieve tools.
*
* <p><b>Thread Safety:</b> This class is thread-safe, using {@link ConcurrentHashMap} for internal
* storage to support concurrent tool registration and lookup operations.
* storage to support concurrent tool registration and lookup operations. Tool instance and
* registration metadata are stored together in a single compound map entry, so put/remove of the
* two are a single atomic operation.
*
* <p><b>Key Responsibilities:</b>
* <ul>
Expand All @@ -39,8 +41,9 @@
*/
class ToolRegistry {

private final Map<String, AgentTool> tools = new ConcurrentHashMap<>();
private final Map<String, RegisteredToolFunction> registeredTools = new ConcurrentHashMap<>();
private record Entry(AgentTool tool, RegisteredToolFunction registered) {}

private final Map<String, Entry> entries = new ConcurrentHashMap<>();

/**
* Register a tool with its metadata.
Expand All @@ -53,8 +56,7 @@ void registerTool(String toolName, AgentTool tool, RegisteredToolFunction regist
if (toolName == null || toolName.isBlank()) {
throw new IllegalArgumentException("Tool name cannot be null or blank");
}
tools.put(toolName, tool);
registeredTools.put(toolName, registered);
entries.put(toolName, new Entry(tool, registered));
}

/**
Expand All @@ -67,7 +69,8 @@ AgentTool getTool(String name) {
if (name == null || name.isBlank()) {
return null;
}
return tools.get(name);
Entry e = entries.get(name);
return e != null ? e.tool() : null;
}

/**
Expand All @@ -80,7 +83,8 @@ RegisteredToolFunction getRegisteredTool(String name) {
if (name == null || name.isBlank()) {
return null;
}
return registeredTools.get(name);
Entry e = entries.get(name);
return e != null ? e.registered() : null;
}

/**
Expand All @@ -89,7 +93,7 @@ RegisteredToolFunction getRegisteredTool(String name) {
* @return Set of tool names
*/
Set<String> getToolNames() {
return new HashSet<>(tools.keySet());
return new HashSet<>(entries.keySet());
}

/**
Expand All @@ -98,7 +102,13 @@ Set<String> getToolNames() {
* @return Map of tool name to RegisteredToolFunction
*/
Map<String, RegisteredToolFunction> getAllRegisteredTools() {
return new ConcurrentHashMap<>(registeredTools);
Map<String, RegisteredToolFunction> result = new ConcurrentHashMap<>();
for (Map.Entry<String, Entry> e : entries.entrySet()) {
if (e.getValue().registered() != null) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Info] Nice win on atomicity — one compound map entry removes the window where tools and registeredTools disagreed. One behavioural wrinkle: getToolNames() (line 95) returns every key, while getAllRegisteredTools() now drops keys whose registered() is null. A caller that iterates names and then assumes a matching entry in getAllRegisteredTools() will silently miss those tools. Previously registerTool(name, tool, null) would have thrown NPE on registeredTools.put(...), so this state is new. Worth either filtering both views the same way or noting the asymmetry in the javadoc.

result.put(e.getKey(), e.getValue().registered());
}
}
return result;
}

/**
Expand All @@ -110,24 +120,35 @@ void removeTool(String toolName) {
if (toolName == null || toolName.isBlank()) {
throw new IllegalArgumentException("Tool name cannot be null or blank");
}
tools.remove(toolName);
registeredTools.remove(toolName);
entries.remove(toolName);
}

/**
* Atomically remove a tool only if the current instance matches the expected one.
* Uses {@link ConcurrentHashMap#remove(Object, Object)} to avoid TOCTOU races.
*
* <p><b>Identity semantics</b>: The guard check ({@code existing.tool() == expected})
* compares the expected tool by reference ({@code ==}), not via {@link Object#equals}.
* Two {@code AgentTool} instances that are {@link Object#equals equal} but not the same
* reference will not match — this guards against accidental removal of a tool that was
* re-registered under the same name by another caller. The CAS at
* {@link ConcurrentHashMap#remove(Object, Object)} additionally depends on the
* {@code Entry} record's {@link Object#equals}, which compares both the
* {@code AgentTool} and {@code RegisteredToolFunction} fields; callers that

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new identity-semantics paragraph matches the implementation (existing.tool() == expected, then entries.remove(toolName, existing)), thanks for pinning it down. One precision point: "Two AgentTool instances that are equal but not the same reference will not match" describes the guard, not the map CAS. The CAS compares the whole Entry record, so RegisteredToolFunction equality participates in it too. The last sentence already says this, so tying the first sentence to the guard (rather than to remove(key, value)) would keep a reader from concluding the two checks are the same check.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

reworded so the first sentence explicitly ties to the guard check
(existing.tool() == expected) rather than to the map CAS, keeping the two checks
clearly separated. The CAS detail (Entry.equals comparing both fields) remains in
the following sentence.

See: 28f7484

* rebuild {@code Entry} objects (e.g. via {@code copyTo}) must ensure
* {@code RegisteredToolFunction} equality remains stable across rebuilds.
*
* @param toolName Tool name to remove
* @param expected The expected AgentTool instance (identity comparison)
* @param expected The expected {@link AgentTool} instance, compared by reference
* ({@code ==}), not by {@link Object#equals}
* @return true if the tool was removed, false if it was already replaced or absent
*/
boolean removeToolIfSame(String toolName, AgentTool expected) {
boolean removed = tools.remove(toolName, expected);
if (removed) {
registeredTools.remove(toolName);
Entry existing = entries.get(toolName);
if (existing != null && existing.tool() == expected) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] Identity semantics changed here. The old code was tools.remove(toolName, expected), i.e. ConcurrentHashMap compares with AgentTool.equals. The new existing.tool() == expected compares by reference. For AgentTool implementations that override equals (proxies, wrappers, or anything copied across a module boundary — copyTo rebuilds Entry objects, so copies are common in this code path), a call that used to remove successfully now returns false. Also note the final entries.remove(toolName, existing) CAS compares the whole record, so it depends on RegisteredToolFunction.equals staying value-stable.

If reference-equality is the intent here (it probably is, for a "same instance" check), please say so explicitly in the javadoc so a future equals override on AgentTool doesn't quietly change the semantics.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@oss-maintainer Good catch. Added an explicit Identity semantics paragraph to removeToolIfSame javadoc documenting that the == comparison is by design — it guards against removal of a tool that was re-registered under the same name by another caller between the get() and the remove(). An equals()-based comparison would silently delete a replacement tool that happened to be equal.
Also documented the entries.remove(toolName, existing) CAS dependency on the Entry record's equals(), including the RegisteredToolFunction field — so callers that rebuild Entry objects through copyTo know that RegisteredToolFunction.equals() must remain stable.

return entries.remove(toolName, existing);
}
return removed;
return false;
}

/**
Expand All @@ -149,20 +170,20 @@ void removeTools(Set<String> toolNames) {
* @param target The target registry to copy tools to
*/
void copyTo(ToolRegistry target) {
for (Map.Entry<String, AgentTool> entry : tools.entrySet()) {
String toolName = entry.getKey();
AgentTool tool = entry.getValue();
RegisteredToolFunction registered = registeredTools.get(toolName);
target.registerTool(
for (Map.Entry<String, Entry> e : entries.entrySet()) {
String toolName = e.getKey();
Entry entry = e.getValue();
target.entries.put(
toolName,
tool,
registered == null
? null
: new RegisteredToolFunction(
tool,
registered.getExtendedModel(),
registered.getMcpClientName(),
registered.getPresetParameters()));
new Entry(
entry.tool(),
entry.registered() == null
? null
: new RegisteredToolFunction(
entry.tool(),
entry.registered().getExtendedModel(),
entry.registered().getMcpClientName(),
entry.registered().getPresetParameters())));
}
}
}
Loading
Loading