Merge remote-tracking branch 'origin/master' into codex/simp-drop-unused-schema-default

This commit is contained in:
Tianyi Cui
2026-07-19 02:28:15 +08:00
95 changed files with 3232 additions and 876 deletions

View File

@@ -117,9 +117,6 @@ declare module 'cordis' {
}
}
// TODO(concurrency): revisit these shapes when concurrency metadata becomes useful
// (for example, a read-only hint that would permit safe parallel execution).
/** Tool output, optionally with lossless-JSON presentation metadata persisted for replay. */
export type ToolExecuteReturn = ContentBlock[] | { content: ContentBlock[]; meta?: unknown }
@@ -134,6 +131,20 @@ export interface ToolDefinition extends ToolSchema {
* cooperative implementation that can reach quiescence when the signal aborts.
*/
timeoutMs?: number
/**
* Pure synchronous classifier for overlap with sibling tool calls. Only
* `true` opts in; omission, exceptions, non-`true` returns, and invalid
* `defineTool` arguments are exclusive. This metadata is never model-visible.
*
* Opted-in executions must not mutate parent-owned state. Shared state must
* tolerate concurrent dispatch; recorder races are permitted only when they
* commute or fail closed. See the
* [parallel-tool-call RFC](../../../../docs/rfc/implemented/feature/2026-07-10-parallel-tool-call-execution.md)
* for the full contract.
* @param args - parsed arguments; `defineTool` validates before calling.
* @returns Whether this call may join a parallel group.
*/
isConcurrencySafe?(args: unknown): boolean
/**
* Optional: how to present the PENDING state of one call in a UI, derived from
* the call's `args` (parsed arguments, `unknown` — the tool validates/narrows
@@ -195,6 +206,14 @@ export interface ToolExecutionInput {
signal?: AbortSignal
}
/**
* Scheduling mode for one pending call. `parallel` may overlap with siblings;
* `exclusive` runs alone and forms an ordering barrier.
*/
export type ToolExecutionMode =
| { kind: 'parallel' }
| { kind: 'exclusive' }
/**
* One pending tool call inside the registry pipeline. Parsed arguments cross
* one lossless-JSON materialization boundary before policy and are deep-frozen;
@@ -222,6 +241,47 @@ export interface ToolRunContext extends ToolExecution {
deferContext(context: HookContext): void
}
/**
* Scheduler-only result after ordered pre-execute and guards. A `post-result`
* still receives post-execute; a `final-result` bypasses it.
* @internal
*/
export type ScheduledToolPreparation =
| { kind: 'dispatch'; exec: ToolRunContext }
| { kind: 'post-result'; exec: ToolRunContext; result: ToolExecutionResult }
| { kind: 'final-result'; exec: ToolRunContext; result: ToolExecutionResult }
/**
* Scheduler-only dispatch result. A `post-result` still receives post-execute;
* a `final-result` already matches {@link ToolRegistry.execute} failure semantics.
* @internal
*/
export type ScheduledToolDispatch =
| { kind: 'post-result'; result: ToolExecutionResult }
| { kind: 'final-result'; result: ToolExecutionResult }
/**
* Symbol-keyed scheduler view that keeps pre/post policy ordered while
* overlapping dispatch. Ordinary callers use {@link ToolRegistry.execute};
* this is not a plugin seam.
* @internal
*/
export interface ToolRegistryScheduler {
/** Materialize input, run the ordered pre-execute/guard gate, and decide what stage follows. */
prepare(exec: ToolExecutionInput): Promise<ScheduledToolPreparation>
/** Run only the around-dispatch/body stage. */
dispatch(exec: ToolRunContext): Promise<ScheduledToolDispatch>
/** Run ordered post-execute finalization, then materialize and notify the final outcome. */
finalize(exec: ToolRunContext, result: ToolExecutionResult): Promise<ToolExecutionResult>
/** Materialize and notify a final outcome that must bypass post-execute. */
finish(exec: ToolRunContext, result: ToolExecutionResult): ToolExecutionResult
}
/**
* Scheduler entry point omitted from the generated named service API.
* @internal
*/
export const TOOL_REGISTRY_SCHEDULER: unique symbol = Symbol('@deepseek-ai/dsh-tools.scheduler')
/** Structured error metadata for a failed tool call (alongside the model-facing text). */
export interface ToolErrorInfo {
name: string
@@ -382,6 +442,16 @@ export class ToolRegistry extends Service {
mode: z.union(['native', 'code', 'both'] as const).default('native'),
})
/** Internal staged view consumed by `dsh-agent-loop`'s parallel scheduler. */
readonly [TOOL_REGISTRY_SCHEDULER]: ToolRegistryScheduler = {
prepare: exec => this.prepareScheduledExecution(exec),
dispatch: exec => this.dispatchScheduledExecution(exec),
finalize: (exec, result) => this.finalizeScheduledExecution(exec, result),
finish: (exec, result) => this.finishScheduledExecution(exec, result),
}
/** Context deferred by a running tool body, keyed by its scheduler-owned execution. */
private deferredContexts = new WeakMap<ToolRunContext, HookContext[]>()
private global = new Map<string, ToolDefinition>()
private scoped = new Map<ScopeKey, Map<string, ToolDefinition>>()
/** Compiled restriction filters, per scope (see {@link restrict}). */
@@ -682,6 +752,24 @@ export class ToolRegistry extends Service {
}
}
/**
* Classify a pending call through the caller's visible tool definition. Only
* an exact `true` is parallel; unknown, hidden, undeclared, invalid, or
* throwing classifiers are exclusive.
* @param exec - call name, parsed arguments, and optional agent scope.
* @returns the fail-closed scheduling mode.
*/
executionMode(exec: ToolExecutionInput): ToolExecutionMode {
const tool = this.get(exec.name, exec.agent)
if (!tool?.isConcurrencySafe) return { kind: 'exclusive' }
try {
const concurrencySafe: unknown = tool.isConcurrencySafe(exec.arguments)
return concurrencySafe === true ? { kind: 'parallel' } : { kind: 'exclusive' }
} catch {
return { kind: 'exclusive' }
}
}
/**
* Execute through pre-policy, guards, around-dispatch, post-policy, and final
* notification. Tool and listener failures resolve as materialized error
@@ -692,6 +780,28 @@ export class ToolRegistry extends Service {
* @returns the materialized final result.
*/
async execute(exec: ToolExecutionInput): Promise<ToolExecutionResult> {
return this.prepareExecution(exec, prepared => this.completeScheduledExecution(prepared))
}
private async completeScheduledExecution(prepared: ScheduledToolPreparation): Promise<ToolExecutionResult> {
switch (prepared.kind) {
case 'dispatch': {
const dispatched = await this.dispatchScheduledExecution(prepared.exec)
return dispatched.kind === 'post-result'
? await this.finalizeScheduledExecution(prepared.exec, dispatched.result)
: this.finishScheduledExecution(prepared.exec, dispatched.result)
}
case 'post-result':
return await this.finalizeScheduledExecution(prepared.exec, prepared.result)
case 'final-result':
return this.finishScheduledExecution(prepared.exec, prepared.result)
/* v8 ignore next -- closed-union exhaustiveness guard */
default:
return assertNever(prepared, 'scheduled tool preparation')
}
}
private createExecution(exec: ToolExecutionInput): ScheduledToolPreparation | { kind: 'ready'; exec: ToolRunContext } {
const deferredContexts: HookContext[] = []
const token = createExecutionToken()
const callId = exec.callId
@@ -710,105 +820,143 @@ export class ToolRegistry extends Service {
deferredContexts.push(context)
},
}
let execution: ToolRunContext
try {
const detached = snapshotJsonValue(exec.arguments)
if (detached === undefined) {
throw new TypeError('tool execution arguments must be losslessly JSON-serializable')
}
execution = {
...base,
arguments: deepFreeze(detached),
}
const execution: ToolRunContext = { ...base, arguments: deepFreeze(detached) }
this.deferredContexts.set(execution, deferredContexts)
return { kind: 'ready', exec: execution }
} catch (error: unknown) {
execution = { ...base, arguments: undefined }
const result = this.materializeFinalResult(toolErrorResult(error))
this.notifyResult(execution, result)
return result
const execution: ToolRunContext = { ...base, arguments: undefined }
return { kind: 'final-result', exec: execution, result: toolErrorResult(error) }
}
let result: ToolExecutionResult
}
/**
* Run the ordered pre-execute and monotonic guard stages for the scheduler.
* @param input - the caller-supplied execution input.
* @returns the prepared execution plus the next scheduler stage.
* @internal
*/
private async prepareScheduledExecution(input: ToolExecutionInput): Promise<ScheduledToolPreparation> {
return this.prepareExecution(input, prepared => prepared)
}
private async prepareExecution<T>(
input: ToolExecutionInput,
next: (prepared: ScheduledToolPreparation) => T | PromiseLike<T>,
): Promise<T> {
const created = this.createExecution(input)
if (created.kind !== 'ready') return next(created)
const exec = created.exec
try {
result = this.materializeFinalResult(await this.executePipeline(execution, deferredContexts))
const carrier = scopeTarget(this, exec.agent)
const gate = await this.ctx.waterfall(
carrier, 'tools/pre-execute', exec,
() => Promise.resolve<PreToolDecision>({ kind: 'allow' }),
)
const decision = gate.kind === 'ask' ? await this.serviceAsk(exec, gate) : gate
const denialReason = decision.kind === 'allow'
? this.guardReason(exec)
: decision.reason
if (denialReason !== undefined) {
return await next({
kind: 'post-result',
exec,
result: {
content: [{ type: 'text', text: `Error: ${denialReason}` }],
isError: true,
},
})
}
return await next({ kind: 'dispatch', exec })
} catch (error: unknown) {
// Outer backstop: a throwing pre/post-execute listener, guard, or the
// waterfall machinery becomes an isError result, never a turn failure.
result = this.materializeFinalResult(toolErrorResult(error))
return next({ kind: 'final-result', exec, result: toolErrorResult(error) })
}
this.notifyResult(execution, result)
return result
}
/** Run the transformable pipeline; {@link execute} owns final normalization and notification. */
private async executePipeline(exec: ToolRunContext, deferredContexts: HookContext[]): Promise<ToolExecutionResult> {
// --- Gate: tools/pre-execute. An `ask` resolves through the optional
// approval seam (or degrades to deny) before the monotonic guards run. The
// carrier keys dispatch by exec.agent, so an `agent.ctx` listener gates only
// its own agent's calls (agent-less calls are subject-less).
const carrier = scopeTarget(this, exec.agent)
const gate = await this.ctx.waterfall(
carrier, 'tools/pre-execute', exec,
() => Promise.resolve<PreToolDecision>({ kind: 'allow' }),
)
const decision = gate.kind === 'ask' ? await this.serviceAsk(exec, gate) : gate
const denialReason = decision.kind === 'allow'
? this.guardReason(exec)
: decision.reason
if (denialReason !== undefined) {
// Every non-grant, including a failed/unavailable approval request, takes
// the same deny path and still reaches post-policy plus result observers.
const denied: ToolExecutionResult = {
content: [{ type: 'text', text: `Error: ${denialReason}` }],
isError: true,
}
return await this.postExecute(exec, denied)
}
// --- Around-dispatch: tools/execute. The base `next` is the dispatch-
// with-normalization thunk — the tool body's own try/catch turns a throw
// into an isError result so a wrapper (and post-execute) can inspect it;
// an unknown tool routes through the same catch. A `tools/execute` listener
// (e.g. a timeout plugin) wraps this thunk: it may replace `exec.signal`
// before delegating and inspect the normalized result after. Dispatched with the
// same carrier as the gate, so an `agent.ctx` wrapper wraps only its own
// agent's calls. ---
const result = await this.ctx.waterfall(
carrier, 'tools/execute', exec,
async (): Promise<ToolExecutionResult> => {
try {
// Resolve through the CALLER's visible view ({@link get}): a scoped
// tool shadows its global name-twin for that agent, and a
// restricted-away global tool is exactly as absent as a nonexistent
// one — same UNKNOWN_TOOL result, no capability leak in the error.
const tool = this.get(exec.name, exec.agent)
if (!tool) throw new ToolNotFoundError(exec.name)
// Normalize the two `execute` return shapes: a bare ContentBlock[] (no
// meta) or a { content, meta } object (a tool attaching a private
// presentation payload). An array IS the content; the object carries it.
const returned = await tool.execute(exec.arguments, exec)
const content = Array.isArray(returned) ? returned : returned.content
const meta = Array.isArray(returned) ? undefined : returned.meta
return { content, isError: false, ...meta !== undefined ? { meta } : {} }
} catch (error: unknown) {
return toolErrorResult(error)
/**
* Run around-dispatch and the tool body. Tool and unknown-tool failures still
* receive post-execute; pipeline failures are already final.
* @param exec - the prepared execution.
* @returns whether the result still needs post-execute.
* @internal
*/
private async dispatchScheduledExecution(exec: ToolRunContext): Promise<ScheduledToolDispatch> {
try {
const carrier = scopeTarget(this, exec.agent)
const result = await this.ctx.waterfall(
carrier, 'tools/execute', exec,
async (): Promise<ToolExecutionResult> => {
try {
const tool = this.get(exec.name, exec.agent)
if (!tool) throw new ToolNotFoundError(exec.name)
const returned = await tool.execute(exec.arguments, exec)
const content = Array.isArray(returned) ? returned : returned.content
const meta = Array.isArray(returned) ? undefined : returned.meta
return { content, isError: false, ...meta !== undefined ? { meta } : {} }
} catch (error: unknown) {
return toolErrorResult(error)
}
},
)
const deferredContexts = this.deferredContexts.get(exec)
/* v8 ignore next -- dispatch only receives executions minted by this registry's prepare stage */
if (deferredContexts === undefined) throw new Error('tool registry scheduler invariant violated: unprepared execution')
const resultWithDeferredContexts: ToolExecutionResult = deferredContexts.length === 0
? result
: {
...result,
additionalContexts: [
...deferredContexts,
...result.additionalContexts ?? [],
],
}
},
)
const resultWithDeferredContexts: ToolExecutionResult = deferredContexts.length === 0
? result
: {
...result,
additionalContexts: [
...deferredContexts,
...result.additionalContexts ?? [],
],
}
return await this.postExecute(exec, resultWithDeferredContexts)
return { kind: 'post-result', result: resultWithDeferredContexts }
} catch (error: unknown) {
return { kind: 'final-result', result: toolErrorResult(error) }
}
}
/** Notify final-result observers without giving them a mutation/error channel into the outcome. */
/**
* Run ordered post-execute, then materialize and notify the final outcome.
* @param exec - the prepared execution.
* @param result - dispatch/pre result that still needs post-execute.
* @returns the materialized final result.
* @internal
*/
private async finalizeScheduledExecution(exec: ToolRunContext, result: ToolExecutionResult): Promise<ToolExecutionResult> {
try {
return this.finishScheduledExecution(exec, await this.postExecute(exec, result))
} catch (error: unknown) {
return this.finishScheduledExecution(exec, toolErrorResult(error))
}
}
/**
* Materialize and notify a final result that must bypass post-execute.
* @param exec - the prepared execution.
* @param result - final result.
* @returns the materialized final result.
* @internal
*/
private finishScheduledExecution(exec: ToolRunContext, result: ToolExecutionResult): ToolExecutionResult {
let finalResult: ToolExecutionResult
try {
finalResult = this.materializeFinalResult(result)
} catch (error: unknown) {
finalResult = this.materializeFinalResult(toolErrorResult(error))
}
this.notifyResult(exec, finalResult)
return finalResult
}
/** Notify observers without exposing a mutation or error channel into the outcome. */
private notifyResult(exec: ToolExecution, result: ToolExecutionResult): void {
// The pipeline is over: freeze the remaining mutable signal slot so every
// observer sees the SAME WeakMap-keyable execution without a mutation race.
// Freeze the remaining mutable signal slot before observers receive the
// shared WeakMap-keyable execution object.
Object.freeze(exec)
const callbacks = this.ctx.events.dispatch('emit', [
scopeTarget(this, exec.agent), 'tools/result', exec, result,

View File

@@ -279,6 +279,14 @@ export interface DefineToolOptions<S extends SchemaSpec> {
* is never sent to the model.
*/
readonly timeoutMs?: number
/**
* Optional pure synchronous classifier for sibling overlap. It receives typed
* arguments after soft validation; invalid input returns `false` without
* invoking it. See {@link ToolDefinition.isConcurrencySafe}.
* @param args - typed validated arguments.
* @returns whether this call may join a parallel group.
*/
isConcurrencySafe?(args: InferArgs<S>): boolean
/**
* Tool execution function. `args` is typed as {@link InferArgs<S>} — zero
* casts needed. Returns either a bare {@link ContentBlock}`[]` (model-facing
@@ -311,7 +319,7 @@ export interface DefineToolOptions<S extends SchemaSpec> {
* @param options - the tool's name, description, typed parameter schema,
* execute body, and optional presenters.
* @returns a registry-ready definition with strict execution validation and
* soft presenter validation for replay compatibility.
* soft presenter and classifier validation for replay compatibility.
*/
export function defineTool<S extends SchemaSpec>(options: DefineToolOptions<S>): ToolDefinition {
// Object-literal execute methods don't use `this`; the reference is safe.
@@ -321,6 +329,8 @@ export function defineTool<S extends SchemaSpec>(options: DefineToolOptions<S>):
const userPresentCall = options.presentCall
// eslint-disable-next-line @typescript-eslint/unbound-method
const userPresentResult = options.presentResult
// eslint-disable-next-line @typescript-eslint/unbound-method
const userIsConcurrencySafe = options.isConcurrencySafe
if (options.timeoutMs !== undefined && (!Number.isFinite(options.timeoutMs) || options.timeoutMs <= 0)) {
throw new Error(`defineTool(${options.name}): timeoutMs must be a positive finite number`)
}
@@ -355,5 +365,12 @@ export function defineTool<S extends SchemaSpec>(options: DefineToolOptions<S>):
return userPresentResult(args as InferArgs<S>, result)
}
}
// Invalid arguments fail closed without invoking the typed classifier.
if (userIsConcurrencySafe) {
tool.isConcurrencySafe = (args: unknown): boolean => {
if (validateArgs(options.parameters, args).length > 0) return false
return userIsConcurrencySafe(args as InferArgs<S>)
}
}
return tool
}