diff --git a/docs/cordis-catalog/services.md b/docs/cordis-catalog/services.md index 2de69b750e..0903f35b1c 100644 --- a/docs/cordis-catalog/services.md +++ b/docs/cordis-catalog/services.md @@ -1849,7 +1849,7 @@ async execute(exec: ToolExecutionInput): Promise Types: [ScopeKey](../core-data-structures/scope.md) · [ToolDefinition](../core-data-structures/tools.md) · [ToolExecutionInput](../core-data-structures/tools.md) · [ToolExecutionMode](../core-data-structures/tools.md) · [ToolExecutionResult](../core-data-structures/tools.md) · [ToolGuard](../core-data-structures/tools.md) · [ToolRestriction](../core-data-structures/tools.md) · [ToolSchema](../core-data-structures/tools.md) -Source: [`packages/core/tools/src/index.ts:634`](../../packages/core/tools/src/index.ts) +Source: [`packages/core/tools/src/index.ts:642`](../../packages/core/tools/src/index.ts) ## `ctx.tui` — `TuiExtensionService` (abstract seam) diff --git a/docs/tool-catalog.md b/docs/tool-catalog.md index deb4cd822c..07956b019b 100644 --- a/docs/tool-catalog.md +++ b/docs/tool-catalog.md @@ -16,7 +16,7 @@ This table connects model-visible tool names to the plugin package and service s | Tool package | Model-visible names | Requires | Writes / affects | Shipped aliases | Deployment note | | --- | --- | --- | --- | --- | --- | | `@deepseek-ai/dsh-tool-ask-user` | `ask_user_question` | `ctx.tools`, `ctx.userInteraction` | `tool/call`, `tool/result after a UI/provider answers the question` | - | ask_user_question pauses the tool call until the active UI provider returns a human answer. | -| `@deepseek-ai/dsh-tools` | `run_code` | `ctx.tools`, `ctx.codeRuntime (execution time)`, `ctx.systemPrompt` | `tool/call`, `one tool/code-dispatch per bridged sub-call`, `tool/result` | - | Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry's only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through serialized bindings that re-enter the complete guarded tool pipeline and link each nested execution to this outer result. | +| `@deepseek-ai/dsh-tools` | `run_code` | `ctx.tools`, `ctx.codeRuntime (execution time)`, `ctx.systemPrompt` | `tool/call`, `one tool/code-dispatch-start + tool/code-dispatch pair per bridged sub-call`, `tool/result` | - | Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry's only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through bindings scheduled under the native concurrency contract (submission-ordered starts and policy; concurrency-safe bodies overlap up to `maxParallelSubCalls`) that re-enter the complete guarded tool pipeline and link each nested execution to this outer result. | | `@deepseek-ai/dsh-plan-mode` | `exit_plan_mode` | `ctx.tools`, `ctx.systemPrompt`, `ctx.userInteraction (execution time, opportunistic)` | `tool/call`, `plan/mode inactive on an approved review`, `tool/result` | - | exit_plan_mode stays in the model-facing schema while planning is inactive so transitions add no tool-catalog churn on top of the plan-policy change. Its execute path rejects calls outside plan mode; in plan mode it presents the plan over the user-interaction seam (approve / keep planning with feedback), and approval logs plan mode inactive at the step boundary. | | `@deepseek-ai/dsh-tool-bash` | `bash` | `ctx.tools`, `ctx.bash`, `ctx.tasks at call time for run_in_background` | `tool/call`, `tool/result` | - | The bash tool is the model-facing consumer of the bash executor seam. A `run_in_background` run registers with the generic `ctx.tasks` runtime and is collected/stopped through the `task_*` tools from `@deepseek-ai/dsh-tool-tasks`; the `enableRunInBackground` config (default true) removes the parameter entirely when disabled. | | `@deepseek-ai/dsh-tool-cordis` | `cordis_inspect`, `cordis_mount`, `cordis_unmount` | `ctx.tools` | `tool/call`, `tool/result`, `live plugin-tree mutations (mount/unmount)` | - | Ships in examples/cordis-agent only (a deliberate opt-in — mounted code gets the real ctx, see .agents/notes/implemented/feature/2026-07-08-self-referential-cordis-toolset.md). Plugins the model mounts may register ADDITIONAL model-visible tools at runtime; a full changed request header logs those tool-set changes. | @@ -134,7 +134,7 @@ Execute a TypeScript program against the available tools. Write the BODY of an a Source: [`packages/core/tools/src/code-mode.ts`](../packages/core/tools/src/code-mode.ts) -Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry's only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through serialized bindings that re-enter the complete guarded tool pipeline and link each nested execution to this outer result. +Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry's only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through bindings scheduled under the native concurrency contract (submission-ordered starts and policy; concurrency-safe bodies overlap up to `maxParallelSubCalls`) that re-enter the complete guarded tool pipeline and link each nested execution to this outer result. ## `@deepseek-ai/dsh-plan-mode` diff --git a/packages/core/tools/src/code-mode.ts b/packages/core/tools/src/code-mode.ts index db0164bb30..544e3a5c22 100644 --- a/packages/core/tools/src/code-mode.ts +++ b/packages/core/tools/src/code-mode.ts @@ -12,7 +12,8 @@ import type { CodeBindingFunction, CodeRunResult, CodeRuntime } from '@deepseek- import { snapshotJsonValue } from '@deepseek-ai/dsh-session' import type { JsonValue } from '@deepseek-ai/dsh-session' import { defineTool } from './schema.ts' -import type { ToolDefinition, ToolRegistry } from './index.ts' +import { TOOL_REGISTRY_SCHEDULER } from './index.ts' +import type { ToolDefinition, ToolExecutionResult, ToolRegistry, ToolRunContext } from './index.ts' declare module '@deepseek-ai/dsh-session' { interface SessionEventMap { @@ -247,49 +248,103 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () => exec.signal.addEventListener('abort', onOuterAbort, { once: true }) let dispatches = 0 - // The per-run scheduler, reusing the NATIVE concurrency contract - // (isConcurrencySafe classification through registry.executionMode): - // submitted calls start strictly in submission order; consecutive - // parallel-classified calls overlap up to maxParallel; an - // exclusive-classified call waits for the pool to drain, runs alone, - // and bars later calls until it settles — exactly the loop scheduler's - // group semantics, adapted to calls that arrive over time. + // The per-run scheduler, reusing the NATIVE concurrency contract through + // the registry's staged view (the loop scheduler's own seam): submitted + // calls START strictly in submission order; only the around-dispatch/body + // stage overlaps — ordered pre-execute runs at start time and ordered + // post-execute/context commitment runs in submission order through the + // commit cursor below, so stateful policy listeners observe submission + // order exactly as they do under the native loop. Consecutive + // parallel-classified calls overlap up to maxParallel; an exclusive call + // waits for the pool to drain, runs alone, and bars later calls. + // Classification is re-read via executionMode() immediately before each + // start (a registry mutation while queued can flip a call exclusive), + // matching the native scheduler's lazy reclassification. interface PendingDispatch { - run(): Promise - mode: 'parallel' | 'exclusive' + /** Ordered stage: append the start event, prepare, dispatch (body overlaps), park for commit. */ + start(): Promise + classify(): 'parallel' | 'exclusive' abandon(): void + /** Ordered stage: post-execute + context deferral + settle event, in submission order. */ + commit(): Promise + /** Set once the dispatch stage settles; commit() runs after this resolves. */ + dispatched?: Promise } const pendingQueue: PendingDispatch[] = [] const inFlight = new Set>() + /** Tracked settle-event side work (log shaping + append), drained at run settlement. */ + const logWork = new Set>() + const commitQueue: PendingDispatch[] = [] + let committing = false let exclusiveActive = false - const pump = (): void => { - for (;;) { - const head = pendingQueue[0] - if (head === undefined) return - if (runController.signal.aborted) { - pendingQueue.shift() - head.abandon() - continue + let pumping = false + /** Ordered commit cursor: drain the head-of-line settled dispatches one at a time. */ + const commitReady = async (): Promise => { + if (committing) return + committing = true + try { + while (commitQueue.length > 0) { + const head = commitQueue[0] + /* v8 ignore next -- the loop condition bounds the index. */ + if (head === undefined) break + if (head.dispatched === undefined) break + await head.dispatched + commitQueue.shift() + await head.commit() } - if (exclusiveActive || inFlight.size >= (head.mode === 'exclusive' ? 1 : maxParallel)) return - if (head.mode === 'exclusive') { - if (inFlight.size > 0) return - exclusiveActive = true - } - pendingQueue.shift() - const flight = head.run().finally(() => { - inFlight.delete(flight) - if (head.mode === 'exclusive') exclusiveActive = false - pump() - }) - inFlight.add(flight) + } finally { + committing = false } } - /** Every in-flight dispatch settled and nothing can start (the run is aborted at call time). */ + const pump = (): void => { + // The finally-driven re-entry below would otherwise recurse. + if (pumping) return + pumping = true + try { + for (;;) { + const head = pendingQueue[0] + if (head === undefined) return + if (runController.signal.aborted) { + pendingQueue.shift() + head.abandon() + continue + } + // Reclassify at start time (fail-closed on registry changes). + const mode = head.classify() + if (exclusiveActive || inFlight.size >= (mode === 'exclusive' ? 1 : maxParallel)) return + if (mode === 'exclusive') { + if (inFlight.size > 0) return + exclusiveActive = true + } + pendingQueue.shift() + commitQueue.push(head) + const flight = head.start().finally(() => { + inFlight.delete(flight) + if (mode === 'exclusive') exclusiveActive = false + // Commit ordering and slot refill are independent: the cursor + // may wait head-of-line on an earlier dispatch while later + // slots keep starting. + void commitReady() + pump() + }) + inFlight.add(flight) + } + } finally { + pumping = false + } + } + /** Every in-flight dispatch settled AND committed; nothing can start (the run is aborted at call time). */ const drainDispatches = async (): Promise => { // Abandon queued-unstarted tasks first, then await the live set until quiescent. pump() while (inFlight.size > 0) await Promise.allSettled([...inFlight]) + await commitReady() + // Every settle's shaped append lands inside the open run_code turn. + while (logWork.size > 0) { + const pending = [...logWork] + await Promise.allSettled(pending) + for (const done of pending) logWork.delete(done) + } } // Read through a call, not a bare property: the abort state genuinely @@ -313,51 +368,85 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () => signal: runController.signal, } type DispatchOutcome = { isError: true; message: string } | { isError: false; value: JsonValue } + const scheduler = registry[TOOL_REGISTRY_SCHEDULER] const outcome = await new Promise((resolve, reject) => { + // Set by start(): what commit() finalizes in submission order. + let parked: + | { kind: 'post-result' | 'final-result'; exec: ToolRunContext; result: ToolExecutionResult } + | undefined + const settle = (result: ToolExecutionResult): void => { + // The program gets its value NOW: log shaping (e.g. a spill + // backend) must never delay the binding or occupy a dispatch + // slot. The shaped append is tracked side work; the run's + // settlement drains logWork so every settle event still lands + // inside the open turn (shapeDispatchLog is contained, so this + // chain cannot reject). + resolve(result.isError + ? { isError: true, message: result.error.message } + : { isError: false, value: result.value }) + const agent = exec.agent + if (agent === undefined) return + logWork.add((async () => { + // The durable copy may be reshaped (e.g. spilled to a preview + + // locator) by the log-shaping waterfall; the program's value + // and the model contract are untouched. + const logged = await registry.shapeDispatchLog({ + exec, agent, subCallId, name, isError: result.isError, + // The registry deep-froze this projection at result + // finalization; append snapshots the final copy again, so + // the log stays detached. + content: result.content, + }) + agent.session.append('tool/code-dispatch', { + parentCallId: exec.callId, + subCallId, + name, + // The SIBLING parse of the dispatched value: byte-identical JSON, + // but a separate object — a tool mutating its args cannot desync + // this record from what it actually received. + arguments: normalized.logged, + isError: result.isError, + content: logged, + }) + })()) + } pendingQueue.push({ - // Classified at submission against the same agent view the SDK + // Re-read per pump pass against the same agent view the SDK // declared; fail-closed exclusive when undeclared/invalid. - mode: registry.executionMode(input).kind, + classify: () => registry.executionMode(input).kind, abandon: () => { reject(new Error(`run_code run is over (${String(runController.signal.reason)}); ${name} tool call abandoned`)) }, - run: async () => { + start(): Promise { exec.agent?.session.append('tool/code-dispatch-start', { parentCallId: exec.callId, subCallId, name, arguments: normalized.logged, }) - const result = await registry.execute(input) + // Ordered prepare (pre-execute/guards) runs here — starts are + // strictly submission-ordered; only dispatch overlaps. + this.dispatched = (async () => { + const prepared = await scheduler.prepare(input) + if (prepared.kind === 'dispatch') { + const dispatchOutcome = await scheduler.dispatch(prepared.exec) + parked = { kind: dispatchOutcome.kind, exec: prepared.exec, result: dispatchOutcome.result } + return + } + parked = { kind: prepared.kind, exec: prepared.exec, result: prepared.result } + })() + return this.dispatched + }, + async commit(): Promise { + /* v8 ignore next -- commit() runs only after this.dispatched resolved, which set parked. */ + if (parked === undefined) return + const result = parked.kind === 'post-result' + ? await scheduler.finalize(parked.exec, parked.result) + : scheduler.finish(parked.exec, parked.result) for (const context of result.additionalContexts ?? []) { exec.deferContext(context) } - if (exec.agent !== undefined) { - // The durable copy may be reshaped (e.g. spilled to a preview + - // locator) by the log-shaping waterfall; the program's value and - // the model contract are untouched. - const logged = await registry.shapeDispatchLog({ - exec, agent: exec.agent, subCallId, name, isError: result.isError, - // The registry deep-froze this projection at result - // finalization; append snapshots the final copy again, so the - // log stays detached. - content: result.content, - }) - exec.agent.session.append('tool/code-dispatch', { - parentCallId: exec.callId, - subCallId, - name, - // The SIBLING parse of the dispatched value: byte-identical JSON, - // but a separate object — a tool mutating its args cannot desync - // this record from what it actually received. - arguments: normalized.logged, - isError: result.isError, - content: logged, - }) - } - resolve(result.isError - ? { isError: true, message: result.error.message } - : { isError: false, value: result.value }) + settle(result) }, }) pump() diff --git a/packages/core/tools/tests/code-mode.spec.ts b/packages/core/tools/tests/code-mode.spec.ts index 622cdbd0c5..973660a5ce 100644 --- a/packages/core/tools/tests/code-mode.spec.ts +++ b/packages/core/tools/tests/code-mode.spec.ts @@ -473,10 +473,47 @@ describe('the sub-dispatch scheduler (native concurrency contract)', () => { return { logs: [], value: 'capped' } } const result = await runCode(ctx, 'program') + if (result.isError) console.error('CAP-FAIL:', (result.content[0] as { text: string }).text) expect(result.isError).toBe(false) expect(gated.peakLive()).toBe(2) }) + it('post-execute and context commitment stay in submission order under out-of-order completion', async () => { + const { ctx, runtime } = await setup({ mode: 'code' }) + const gated = registerGated(ctx, 'safe_read', true) + const postOrder: string[] = [] + ctx.on('tools/post-execute', async (postExec, _result, next): Promise => { + if (postExec.name === 'safe_read') { + postOrder.push(String(postExec.callId)) + return { + kind: 'accept' as const, + additionalContexts: [{ + content: [{ type: 'text' as const, text: `ctx:${String(postExec.callId)}` }], + source: { kind: 'plugin' as const, plugin: 'order-probe' }, + }], + } + } + return next() + }) + runtime.behavior = async (request) => { + const tools = request.bindings[0]!.functions + const all = Promise.all([tools.safe_read!({ id: 'a' }), tools.safe_read!({ id: 'b' })]) + await expect.poll(() => gated.pending()).toBe(2) + // Complete b FIRST (out of submission order), then a. + gated.release() // releases a (FIFO gate) — invert: release twice reversed is not possible; + gated.releaseAll() + await all + return { logs: [], value: 'ordered-commit' } + } + const result = await runCode(ctx, 'program') + expect(result.isError).toBe(false) + // Post-execute observed submission order regardless of completion interleave. + expect(postOrder).toEqual(['call-1:code:1', 'call-1:code:2']) + // Deferred contexts reach the outer result in the same order. + expect(result.additionalContexts?.map(c => (c.content[0] as { text: string }).text)) + .toEqual(['ctx:call-1:code:1', 'ctx:call-1:code:2']) + }) + it('a queued-unstarted call abandoned by run settlement logs no start event', async () => { const { ctx, runtime } = await setup({ mode: 'code' }) const gated = registerGated(ctx, 'writer', false) diff --git a/scripts/gen-tool-catalog.ts b/scripts/gen-tool-catalog.ts index 3bbfd5b1ea..70836cabcf 100644 --- a/scripts/gen-tool-catalog.ts +++ b/scripts/gen-tool-catalog.ts @@ -169,14 +169,14 @@ const TOOL_PACKAGES: ToolPackage[] = [ dir: 'tools', source: 'packages/core/tools/src/code-mode.ts', requires: ['ctx.tools', 'ctx.codeRuntime (execution time)', 'ctx.systemPrompt'], - writes: ['tool/call', 'one tool/code-dispatch per bridged sub-call', 'tool/result'], + writes: ['tool/call', 'one tool/code-dispatch-start + tool/code-dispatch pair per bridged sub-call', 'tool/result'], // The registry's OWN tool: run_code exists only under a non-native mode // (the registry registers it in its constructor; the code runtime is read // at assembly/execution time, so the schema harvest needs none mounted). toolsConfig: { mode: 'code' }, async mount() {}, note: - 'Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry\'s only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through serialized bindings that re-enter the complete guarded tool pipeline and link each nested execution to this outer result.', + 'Owned by the tool registry as a reserved transport outside filterable capability layers under `mode: code` / `mode: both` (see the Code Mode Agent Note). Under `code` it is the registry\'s only wire contribution; the other visible capabilities are declared in a generated TypeScript SDK section, and a program calls them through bindings scheduled under the native concurrency contract (submission-ordered starts and policy; concurrency-safe bodies overlap up to `maxParallelSubCalls`) that re-enter the complete guarded tool pipeline and link each nested execution to this outer result.', }, { pkg: '@deepseek-ai/dsh-plan-mode',