feat(tools): require cancellation signal on every invocation
This commit is contained in:
77 files changed
+1129
-446
No files matched your search
@@ -160,9 +160,8 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
// (its executor kills on this signal) instead of orphaned, and
|
||||
// queued-unstarted dispatches are abandoned.
|
||||
const runController = new AbortController()
|
||||
const onOuterAbort = (): void => { runController.abort(exec.signal?.reason) }
|
||||
if (exec.signal?.aborted) onOuterAbort()
|
||||
exec.signal?.addEventListener('abort', onOuterAbort, { once: true })
|
||||
const onOuterAbort = (): void => { runController.abort(exec.signal.reason) }
|
||||
exec.signal.addEventListener('abort', onOuterAbort, { once: true })
|
||||
|
||||
let dispatches = 0
|
||||
// The per-run serialization queue: every binding call chains onto the tail, so even
|
||||
@@ -273,7 +272,7 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
meta,
|
||||
}
|
||||
} finally {
|
||||
exec.signal?.removeEventListener('abort', onOuterAbort)
|
||||
exec.signal.removeEventListener('abort', onOuterAbort)
|
||||
}
|
||||
},
|
||||
// ACP execute cards use the program as their visible title.
|
||||
|
||||
@@ -90,12 +90,13 @@ declare module 'cordis' {
|
||||
* @param exec - the allowed call about to dispatch (name, parsed arguments, caller agent, signal).
|
||||
* @mode waterfall
|
||||
*/
|
||||
'tools/execute'(this: Scoped<ToolRegistry>, exec: ToolExecution, next: () => Promise<ToolExecutionResult>): Promise<ToolExecutionResult>
|
||||
'tools/execute'(this: Scoped<ToolRegistry>, exec: ToolDispatchExecution, next: () => Promise<ToolExecutionResult>): Promise<ToolExecutionResult>
|
||||
/**
|
||||
* Accept, replace, enrich, or block a normalized dispatch result. `next()`
|
||||
* accepts it unchanged; thrown tools still reach this seam as errors. Async
|
||||
* listeners must observe `exec.signal`; after they settle, caller
|
||||
* cancellation replaces only a successful accepted outcome with `ABORTED`.
|
||||
* cancellation replaces only a successful accepted outcome with the code
|
||||
* selected by whether the tool body was invoked.
|
||||
* Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent's calls.
|
||||
* @param exec - the call that just ran (name, parsed arguments, caller agent).
|
||||
* @param result - the dispatch outcome a listener may accept, replace, or block.
|
||||
@@ -218,7 +219,8 @@ export interface ToolExecutionInput {
|
||||
* the outer `run_code` outcome without receiving its live mutable execution.
|
||||
*/
|
||||
readonly parent?: ToolExecutionToken
|
||||
signal?: AbortSignal
|
||||
/** Required caller-owned cancellation for this invocation. */
|
||||
readonly signal: AbortSignal
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -232,17 +234,25 @@ export type ToolExecutionMode =
|
||||
/**
|
||||
* One pending tool call inside the registry pipeline. Parsed arguments cross
|
||||
* one lossless-JSON materialization boundary before policy and are deep-frozen;
|
||||
* call identity and the registry-assigned {@link token} are readonly. An
|
||||
* around-dispatch wrapper may set, replace, or remove `signal`; immediately
|
||||
* before the body, the registry re-fuses the original caller signal so a
|
||||
* wrapper cannot detach caller cancellation. The registry freezes the complete
|
||||
* object before `tools/result` observers run.
|
||||
* call identity, the caller signal, and the registry-assigned {@link token} are
|
||||
* readonly. The registry freezes the complete object before `tools/result`
|
||||
* observers run.
|
||||
*/
|
||||
export interface ToolExecution extends ToolExecutionInput {
|
||||
/** Registry-assigned identity shared with nested calls only as their opaque `parent` token. */
|
||||
readonly token: ToolExecutionToken
|
||||
}
|
||||
|
||||
/**
|
||||
* Around-dispatch view of a {@link ToolExecution}. A `tools/execute` wrapper
|
||||
* may replace the signal for its delegated lifetime, but it cannot remove it.
|
||||
* The registry fuses every replacement with the captured caller signal.
|
||||
*/
|
||||
export interface ToolDispatchExecution extends Omit<ToolExecution, 'signal'> {
|
||||
/** Cancellation signal visible to the next wrapper or tool body. */
|
||||
signal: AbortSignal
|
||||
}
|
||||
|
||||
/**
|
||||
* Runtime context handed to a tool implementation after the registry has
|
||||
* accepted a {@link ToolExecution}. A composite tool uses
|
||||
@@ -258,6 +268,9 @@ export interface ToolRunContext extends ToolExecution {
|
||||
deferContext(context: HookContext): void
|
||||
}
|
||||
|
||||
/** Registry-owned live execution object; public pipeline views stay readonly. */
|
||||
type MutableToolRunContext = Omit<ToolRunContext, 'signal'> & { signal: AbortSignal }
|
||||
|
||||
/**
|
||||
* Scheduler-only result after ordered pre-execute and guards. A `post-result`
|
||||
* still receives post-execute; a `final-result` bypasses it.
|
||||
@@ -299,6 +312,13 @@ export interface ToolRegistryScheduler {
|
||||
* @internal
|
||||
*/
|
||||
export const TOOL_REGISTRY_SCHEDULER: unique symbol = Symbol('@deepseek-ai/dsh-tools.scheduler')
|
||||
|
||||
/** Canonical error code for cancellation after a tool body was invoked. */
|
||||
export const TOOL_ABORTED = 'ABORTED'
|
||||
|
||||
/** Canonical error code for cancellation before a tool body was invoked. */
|
||||
export const TOOL_ABORTED_BEFORE_DISPATCH = 'ABORTED_BEFORE_DISPATCH'
|
||||
|
||||
/** Structured error metadata for a failed tool call (alongside the model-facing text). */
|
||||
export interface ToolErrorInfo {
|
||||
name: string
|
||||
@@ -448,15 +468,21 @@ interface ToolGuardRegistration {
|
||||
guard: ToolGuard
|
||||
}
|
||||
|
||||
/** Caller cancellation captured before around-dispatch wrappers may replace the public signal slot. */
|
||||
/** Approval decision plus whether the approval channel reported cancellation. */
|
||||
interface ToolAskResolution {
|
||||
readonly decision: Extract<PreToolDecision, { kind: 'allow' | 'deny' }>
|
||||
readonly approvalCancelled: boolean
|
||||
}
|
||||
|
||||
/** Caller cancellation and dispatch state kept outside the around-wrapper view. */
|
||||
interface ToolCancellationState {
|
||||
readonly callerSignal: AbortSignal | undefined
|
||||
readonly abortedAtEntry: boolean
|
||||
readonly callerSignal: AbortSignal
|
||||
bodyInvoked: boolean
|
||||
}
|
||||
|
||||
/** One dispatch-scoped fused signal plus listener cleanup after the body settles. */
|
||||
interface FusedToolSignal {
|
||||
readonly signal: AbortSignal | undefined
|
||||
readonly signal: AbortSignal
|
||||
dispose(): void
|
||||
}
|
||||
|
||||
@@ -807,9 +833,9 @@ export class ToolRegistry extends Service {
|
||||
* results; an invisible tool reports `UNKNOWN_TOOL`. The returned outcome is
|
||||
* the same lossless, frozen snapshot final observers receive. Cancellation
|
||||
* arriving after entry and before final result materialization skips a
|
||||
* not-yet-started body or replaces a successful pipeline outcome with
|
||||
* `ABORTED`; already-started work is still drained and may retain a
|
||||
* tool-owned structured error.
|
||||
* not-yet-started body with `ABORTED_BEFORE_DISPATCH` or replaces a
|
||||
* successful started outcome with `ABORTED`; already-started work is still
|
||||
* drained and may retain a tool-owned structured error.
|
||||
* @param exec - the typed same-process call input. The registry assigns its
|
||||
* correlation token before policy begins.
|
||||
* @returns the materialized final result.
|
||||
@@ -836,7 +862,7 @@ export class ToolRegistry extends Service {
|
||||
}
|
||||
}
|
||||
|
||||
private createExecution(exec: ToolExecutionInput): ScheduledToolPreparation | { kind: 'ready'; exec: ToolRunContext } {
|
||||
private createExecution(exec: ToolExecutionInput): ScheduledToolPreparation | { kind: 'ready'; exec: MutableToolRunContext } {
|
||||
const deferredContexts: HookContext[] = []
|
||||
const token = createExecutionToken()
|
||||
const callId = exec.callId
|
||||
@@ -848,9 +874,9 @@ export class ToolRegistry extends Service {
|
||||
token,
|
||||
callId,
|
||||
name,
|
||||
signal,
|
||||
...agent !== undefined ? { agent } : {},
|
||||
...parent !== undefined ? { parent } : {},
|
||||
...signal !== undefined ? { signal } : {},
|
||||
deferContext(context: HookContext): void {
|
||||
deferredContexts.push(context)
|
||||
},
|
||||
@@ -860,15 +886,15 @@ export class ToolRegistry extends Service {
|
||||
if (detached === undefined) {
|
||||
throw new TypeError('tool execution arguments must be losslessly JSON-serializable')
|
||||
}
|
||||
const execution: ToolRunContext = { ...base, arguments: deepFreeze(detached) }
|
||||
const execution: MutableToolRunContext = { ...base, arguments: deepFreeze(detached) }
|
||||
this.deferredContexts.set(execution, deferredContexts)
|
||||
this.cancellationStates.set(execution, {
|
||||
callerSignal: signal,
|
||||
abortedAtEntry: signal?.aborted === true,
|
||||
bodyInvoked: false,
|
||||
})
|
||||
return { kind: 'ready', exec: execution }
|
||||
} catch (error: unknown) {
|
||||
const execution: ToolRunContext = { ...base, arguments: undefined }
|
||||
const execution: MutableToolRunContext = { ...base, arguments: undefined }
|
||||
return { kind: 'final-result', exec: execution, result: toolErrorResult(error) }
|
||||
}
|
||||
}
|
||||
@@ -890,15 +916,21 @@ export class ToolRegistry extends Service {
|
||||
const created = this.createExecution(input)
|
||||
if (created.kind !== 'ready') return next(created)
|
||||
const exec = created.exec
|
||||
if (this.callerCancelled(exec)) {
|
||||
return next({ kind: 'final-result', exec, result: toolAbortedBeforeDispatchResult() })
|
||||
}
|
||||
try {
|
||||
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
|
||||
if (this.callerCancelledAfterEntry(exec)) {
|
||||
return await next({ kind: 'post-result', exec, result: toolAbortedResult() })
|
||||
const askResolution: ToolAskResolution = gate.kind === 'ask'
|
||||
? await this.serviceAsk(exec, gate)
|
||||
: { decision: gate, approvalCancelled: false }
|
||||
const { decision } = askResolution
|
||||
if (this.callerCancelled(exec) && askResolution.approvalCancelled) {
|
||||
return await next({ kind: 'post-result', exec, result: toolAbortedBeforeDispatchResult() })
|
||||
}
|
||||
const denialReason = decision.kind === 'allow'
|
||||
? this.guardReason(exec)
|
||||
@@ -913,20 +945,31 @@ export class ToolRegistry extends Service {
|
||||
},
|
||||
})
|
||||
}
|
||||
if (this.callerCancelled(exec)) {
|
||||
return await next({ kind: 'post-result', exec, result: toolAbortedBeforeDispatchResult() })
|
||||
}
|
||||
return await next({ kind: 'dispatch', exec })
|
||||
} catch (error: unknown) {
|
||||
return this.callerCancelledAfterEntry(exec)
|
||||
? await next({ kind: 'post-result', exec, result: toolAbortedResult() })
|
||||
: next({ kind: 'final-result', exec, result: toolErrorResult(error) })
|
||||
return next({ kind: 'final-result', exec, result: toolErrorResult(error) })
|
||||
}
|
||||
}
|
||||
|
||||
/** Whether the original live caller signal aborted after this execution entered the registry. */
|
||||
private callerCancelledAfterEntry(exec: ToolRunContext): boolean {
|
||||
/** Whether the original caller signal is currently aborted. */
|
||||
private callerCancelled(exec: ToolRunContext): boolean {
|
||||
const state = this.cancellationStates.get(exec)
|
||||
/* v8 ignore next -- only registry-minted executions reach the staged scheduler methods */
|
||||
if (state === undefined) throw new Error('tool registry scheduler invariant violated: missing cancellation state')
|
||||
return !state.abortedAtEntry && state.callerSignal?.aborted === true
|
||||
return state.callerSignal.aborted
|
||||
}
|
||||
|
||||
/** Canonical cancellation outcome selected by whether the tool body started. */
|
||||
private cancellationResult(exec: ToolRunContext, prior?: ToolExecutionResult): ToolExecutionResult {
|
||||
const state = this.cancellationStates.get(exec)
|
||||
/* v8 ignore next -- only registry-minted executions reach the staged scheduler methods */
|
||||
if (state === undefined) throw new Error('tool registry scheduler invariant violated: missing cancellation state')
|
||||
return state.bodyInvoked
|
||||
? toolAbortedResult(prior)
|
||||
: toolAbortedBeforeDispatchResult(prior)
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -934,24 +977,23 @@ export class ToolRegistry extends Service {
|
||||
* into any around-wrapper replacement. Cancellation never abandons the body:
|
||||
* a started promise reaches quiescence before its outcome becomes `ABORTED`.
|
||||
*/
|
||||
private async dispatchToolBody(exec: ToolRunContext): Promise<ToolExecutionResult> {
|
||||
private async dispatchToolBody(exec: MutableToolRunContext): Promise<ToolExecutionResult> {
|
||||
const state = this.cancellationStates.get(exec)
|
||||
/* v8 ignore next -- only registry-minted executions reach the staged scheduler methods */
|
||||
if (state === undefined) throw new Error('tool registry scheduler invariant violated: missing cancellation state')
|
||||
const wrapperSignal = exec.signal
|
||||
const fused = fuseToolSignals(state.callerSignal, wrapperSignal)
|
||||
const signal = fused.signal
|
||||
const abortedBeforeBody = isAborted(signal)
|
||||
|
||||
if (!state.abortedAtEntry && abortedBeforeBody) {
|
||||
if (isAborted(signal)) {
|
||||
fused.dispose()
|
||||
return toolAbortedResult()
|
||||
return toolAbortedBeforeDispatchResult()
|
||||
}
|
||||
if (signal === undefined) delete exec.signal
|
||||
else exec.signal = signal
|
||||
exec.signal = signal
|
||||
try {
|
||||
const tool = this.get(exec.name, exec.agent)
|
||||
if (!tool) throw new ToolNotFoundError(exec.name)
|
||||
state.bodyInvoked = true
|
||||
const returned = await tool.execute(exec.arguments, exec)
|
||||
const content = Array.isArray(returned) ? returned : returned.content
|
||||
const meta = Array.isArray(returned) ? undefined : returned.meta
|
||||
@@ -960,15 +1002,14 @@ export class ToolRegistry extends Service {
|
||||
isError: false,
|
||||
...meta !== undefined ? { meta } : {},
|
||||
}
|
||||
return !abortedBeforeBody && isAborted(signal)
|
||||
return isAborted(signal)
|
||||
? toolAbortedResult(result)
|
||||
: result
|
||||
} catch (error: unknown) {
|
||||
return toolErrorResult(error)
|
||||
} finally {
|
||||
fused.dispose()
|
||||
if (wrapperSignal === undefined) delete exec.signal
|
||||
else exec.signal = wrapperSignal
|
||||
exec.signal = wrapperSignal
|
||||
}
|
||||
}
|
||||
|
||||
@@ -981,10 +1022,11 @@ export class ToolRegistry extends Service {
|
||||
*/
|
||||
private async dispatchScheduledExecution(exec: ToolRunContext): Promise<ScheduledToolDispatch> {
|
||||
try {
|
||||
const mutableExec = exec as MutableToolRunContext
|
||||
const carrier = scopeTarget(this, exec.agent)
|
||||
const result = await this.ctx.waterfall(
|
||||
carrier, 'tools/execute', exec,
|
||||
() => this.dispatchToolBody(exec),
|
||||
carrier, 'tools/execute', mutableExec,
|
||||
() => this.dispatchToolBody(mutableExec),
|
||||
)
|
||||
const deferredContexts = this.deferredContexts.get(exec)
|
||||
/* v8 ignore next -- dispatch only receives executions minted by this registry's prepare stage */
|
||||
@@ -1000,8 +1042,8 @@ export class ToolRegistry extends Service {
|
||||
}
|
||||
return {
|
||||
kind: 'post-result',
|
||||
result: this.callerCancelledAfterEntry(exec) && !resultWithDeferredContexts.isError
|
||||
? toolAbortedResult(resultWithDeferredContexts)
|
||||
result: this.callerCancelled(exec) && !resultWithDeferredContexts.isError
|
||||
? this.cancellationResult(exec, resultWithDeferredContexts)
|
||||
: resultWithDeferredContexts,
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
@@ -1021,8 +1063,8 @@ export class ToolRegistry extends Service {
|
||||
const postResult = await this.postExecute(exec, result)
|
||||
return this.finishScheduledExecution(
|
||||
exec,
|
||||
this.callerCancelledAfterEntry(exec) && !postResult.isError
|
||||
? toolAbortedResult(postResult)
|
||||
this.callerCancelled(exec) && !postResult.isError
|
||||
? this.cancellationResult(exec, postResult)
|
||||
: postResult,
|
||||
)
|
||||
} catch (error: unknown) {
|
||||
@@ -1050,8 +1092,8 @@ export class ToolRegistry extends Service {
|
||||
|
||||
/** Notify observers without exposing a mutation or error channel into the outcome. */
|
||||
private notifyResult(exec: ToolExecution, result: ToolExecutionResult): void {
|
||||
// Freeze the remaining mutable signal slot before observers receive the
|
||||
// shared WeakMap-keyable execution object.
|
||||
// Freeze the registry's live object before observers receive its readonly
|
||||
// WeakMap-keyable view.
|
||||
Object.freeze(exec)
|
||||
const callbacks = this.ctx.events.dispatch('emit', [
|
||||
scopeTarget(this, exec.agent), 'tools/result', exec, result,
|
||||
@@ -1079,26 +1121,41 @@ export class ToolRegistry extends Service {
|
||||
private async serviceAsk(
|
||||
exec: ToolExecution,
|
||||
ask: Extract<PreToolDecision, { kind: 'ask' }>,
|
||||
): Promise<Extract<PreToolDecision, { kind: 'allow' | 'deny' }>> {
|
||||
): Promise<ToolAskResolution> {
|
||||
const approval = this.ctx.get('approval')
|
||||
if (approval === undefined) {
|
||||
return { kind: 'deny', reason: ask.reason ?? `tool "${exec.name}" requires approval (not yet supported)` }
|
||||
return {
|
||||
decision: { kind: 'deny', reason: ask.reason ?? `tool "${exec.name}" requires approval (not yet supported)` },
|
||||
approvalCancelled: false,
|
||||
}
|
||||
}
|
||||
if (exec.agent === undefined) {
|
||||
return { kind: 'deny', reason: `tool "${exec.name}" requires approval, but the call has no agent to route it through` }
|
||||
return {
|
||||
decision: { kind: 'deny', reason: `tool "${exec.name}" requires approval, but the call has no agent to route it through` },
|
||||
approvalCancelled: false,
|
||||
}
|
||||
}
|
||||
const outcome = await approval.request({
|
||||
agent: exec.agent,
|
||||
toolName: exec.name,
|
||||
callId: exec.callId,
|
||||
...ask.reason !== undefined ? { reason: ask.reason } : {},
|
||||
...exec.signal !== undefined ? { signal: exec.signal } : {},
|
||||
signal: exec.signal,
|
||||
})
|
||||
switch (outcome) {
|
||||
case 'allowed-once': return { kind: 'allow' }
|
||||
case 'rejected': return { kind: 'deny', reason: `the user rejected tool "${exec.name}"` }
|
||||
case 'cancelled': return { kind: 'deny', reason: `approval for tool "${exec.name}" was cancelled` }
|
||||
case 'unavailable': return { kind: 'deny', reason: `tool "${exec.name}" requires approval, but no approval channel is available` }
|
||||
case 'allowed-once': return { decision: { kind: 'allow' }, approvalCancelled: false }
|
||||
case 'rejected': return {
|
||||
decision: { kind: 'deny', reason: `the user rejected tool "${exec.name}"` },
|
||||
approvalCancelled: false,
|
||||
}
|
||||
case 'cancelled': return {
|
||||
decision: { kind: 'deny', reason: `approval for tool "${exec.name}" was cancelled` },
|
||||
approvalCancelled: true,
|
||||
}
|
||||
case 'unavailable': return {
|
||||
decision: { kind: 'deny', reason: `tool "${exec.name}" requires approval, but no approval channel is available` },
|
||||
approvalCancelled: false,
|
||||
}
|
||||
default: return assertNever(outcome, 'ApprovalOutcome')
|
||||
}
|
||||
}
|
||||
@@ -1165,19 +1222,16 @@ function toolErrorResult(error: unknown): ToolExecutionResult {
|
||||
}
|
||||
|
||||
/** Read live abort state across an await without treating it as synchronously immutable. */
|
||||
function isAborted(signal: AbortSignal | undefined): boolean {
|
||||
return signal?.aborted === true
|
||||
function isAborted(signal: AbortSignal): boolean {
|
||||
return signal.aborted
|
||||
}
|
||||
|
||||
/**
|
||||
* Fuse caller and wrapper cancellation without nesting `AbortSignal.any`.
|
||||
* Keeping the relay dispatch-scoped also removes listeners when work settles.
|
||||
*/
|
||||
function fuseToolSignals(caller: AbortSignal | undefined, wrapper: AbortSignal | undefined): FusedToolSignal {
|
||||
if (caller === undefined || caller === wrapper) {
|
||||
return { signal: wrapper ?? caller, dispose() {} }
|
||||
}
|
||||
if (wrapper === undefined) return { signal: caller, dispose() {} }
|
||||
function fuseToolSignals(caller: AbortSignal, wrapper: AbortSignal): FusedToolSignal {
|
||||
if (caller === wrapper) return { signal: caller, dispose() {} }
|
||||
|
||||
const controller = new AbortController()
|
||||
let listening = false
|
||||
@@ -1205,13 +1259,24 @@ function fuseToolSignals(caller: AbortSignal | undefined, wrapper: AbortSignal |
|
||||
return { signal: controller.signal, dispose }
|
||||
}
|
||||
|
||||
/** Canonical result when cancellation prevents dispatch or supersedes a successful outcome. */
|
||||
/** Canonical result when cancellation supersedes success after body invocation. */
|
||||
function toolAbortedResult(prior?: ToolExecutionResult): ToolExecutionResult {
|
||||
const additionalContexts = prior?.additionalContexts ?? []
|
||||
return {
|
||||
content: [{ type: 'text', text: 'Error: tool call aborted' }],
|
||||
isError: true,
|
||||
error: { name: 'AbortError', code: 'ABORTED' },
|
||||
error: { name: 'AbortError', code: TOOL_ABORTED },
|
||||
...additionalContexts.length > 0 ? { additionalContexts } : {},
|
||||
}
|
||||
}
|
||||
|
||||
/** Canonical result when cancellation prevents tool body invocation. */
|
||||
function toolAbortedBeforeDispatchResult(prior?: ToolExecutionResult): ToolExecutionResult {
|
||||
const additionalContexts = prior?.additionalContexts ?? []
|
||||
return {
|
||||
content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],
|
||||
isError: true,
|
||||
error: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
|
||||
...additionalContexts.length > 0 ? { additionalContexts } : {},
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user