Merge branch 'code-mode-ui/dispatch-spill' into code-mode-ui/shiki
This commit is contained in:
@@ -1849,7 +1849,7 @@ async execute(exec: ToolExecutionInput): Promise<ToolExecutionResult>
|
|||||||
|
|
||||||
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)
|
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)
|
## `ctx.tui` — `TuiExtensionService` (abstract seam)
|
||||||
|
|
||||||
|
|||||||
@@ -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 |
|
| 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-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-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-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. |
|
| `@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)
|
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`
|
## `@deepseek-ai/dsh-plan-mode`
|
||||||
|
|
||||||
|
|||||||
@@ -12,7 +12,8 @@ import type { CodeBindingFunction, CodeRunResult, CodeRuntime } from '@deepseek-
|
|||||||
import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
|
import { snapshotJsonValue } from '@deepseek-ai/dsh-session'
|
||||||
import type { JsonValue } from '@deepseek-ai/dsh-session'
|
import type { JsonValue } from '@deepseek-ai/dsh-session'
|
||||||
import { defineTool } from './schema.ts'
|
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' {
|
declare module '@deepseek-ai/dsh-session' {
|
||||||
interface SessionEventMap {
|
interface SessionEventMap {
|
||||||
@@ -247,49 +248,103 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
|||||||
exec.signal.addEventListener('abort', onOuterAbort, { once: true })
|
exec.signal.addEventListener('abort', onOuterAbort, { once: true })
|
||||||
|
|
||||||
let dispatches = 0
|
let dispatches = 0
|
||||||
// The per-run scheduler, reusing the NATIVE concurrency contract
|
// The per-run scheduler, reusing the NATIVE concurrency contract through
|
||||||
// (isConcurrencySafe classification through registry.executionMode):
|
// the registry's staged view (the loop scheduler's own seam): submitted
|
||||||
// submitted calls start strictly in submission order; consecutive
|
// calls START strictly in submission order; only the around-dispatch/body
|
||||||
// parallel-classified calls overlap up to maxParallel; an
|
// stage overlaps — ordered pre-execute runs at start time and ordered
|
||||||
// exclusive-classified call waits for the pool to drain, runs alone,
|
// post-execute/context commitment runs in submission order through the
|
||||||
// and bars later calls until it settles — exactly the loop scheduler's
|
// commit cursor below, so stateful policy listeners observe submission
|
||||||
// group semantics, adapted to calls that arrive over time.
|
// 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 {
|
interface PendingDispatch {
|
||||||
run(): Promise<void>
|
/** Ordered stage: append the start event, prepare, dispatch (body overlaps), park for commit. */
|
||||||
mode: 'parallel' | 'exclusive'
|
start(): Promise<void>
|
||||||
|
classify(): 'parallel' | 'exclusive'
|
||||||
abandon(): void
|
abandon(): void
|
||||||
|
/** Ordered stage: post-execute + context deferral + settle event, in submission order. */
|
||||||
|
commit(): Promise<void>
|
||||||
|
/** Set once the dispatch stage settles; commit() runs after this resolves. */
|
||||||
|
dispatched?: Promise<void>
|
||||||
}
|
}
|
||||||
const pendingQueue: PendingDispatch[] = []
|
const pendingQueue: PendingDispatch[] = []
|
||||||
const inFlight = new Set<Promise<void>>()
|
const inFlight = new Set<Promise<void>>()
|
||||||
|
/** Tracked settle-event side work (log shaping + append), drained at run settlement. */
|
||||||
|
const logWork = new Set<Promise<void>>()
|
||||||
|
const commitQueue: PendingDispatch[] = []
|
||||||
|
let committing = false
|
||||||
let exclusiveActive = false
|
let exclusiveActive = false
|
||||||
const pump = (): void => {
|
let pumping = false
|
||||||
for (;;) {
|
/** Ordered commit cursor: drain the head-of-line settled dispatches one at a time. */
|
||||||
const head = pendingQueue[0]
|
const commitReady = async (): Promise<void> => {
|
||||||
if (head === undefined) return
|
if (committing) return
|
||||||
if (runController.signal.aborted) {
|
committing = true
|
||||||
pendingQueue.shift()
|
try {
|
||||||
head.abandon()
|
while (commitQueue.length > 0) {
|
||||||
continue
|
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
|
} finally {
|
||||||
if (head.mode === 'exclusive') {
|
committing = false
|
||||||
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)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
/** 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<void> => {
|
const drainDispatches = async (): Promise<void> => {
|
||||||
// Abandon queued-unstarted tasks first, then await the live set until quiescent.
|
// Abandon queued-unstarted tasks first, then await the live set until quiescent.
|
||||||
pump()
|
pump()
|
||||||
while (inFlight.size > 0) await Promise.allSettled([...inFlight])
|
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
|
// 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,
|
signal: runController.signal,
|
||||||
}
|
}
|
||||||
type DispatchOutcome = { isError: true; message: string } | { isError: false; value: JsonValue }
|
type DispatchOutcome = { isError: true; message: string } | { isError: false; value: JsonValue }
|
||||||
|
const scheduler = registry[TOOL_REGISTRY_SCHEDULER]
|
||||||
const outcome = await new Promise<DispatchOutcome>((resolve, reject) => {
|
const outcome = await new Promise<DispatchOutcome>((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({
|
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.
|
// declared; fail-closed exclusive when undeclared/invalid.
|
||||||
mode: registry.executionMode(input).kind,
|
classify: () => registry.executionMode(input).kind,
|
||||||
abandon: () => {
|
abandon: () => {
|
||||||
reject(new Error(`run_code run is over (${String(runController.signal.reason)}); ${name} tool call abandoned`))
|
reject(new Error(`run_code run is over (${String(runController.signal.reason)}); ${name} tool call abandoned`))
|
||||||
},
|
},
|
||||||
run: async () => {
|
start(): Promise<void> {
|
||||||
exec.agent?.session.append('tool/code-dispatch-start', {
|
exec.agent?.session.append('tool/code-dispatch-start', {
|
||||||
parentCallId: exec.callId,
|
parentCallId: exec.callId,
|
||||||
subCallId,
|
subCallId,
|
||||||
name,
|
name,
|
||||||
arguments: normalized.logged,
|
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<void> {
|
||||||
|
/* 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 ?? []) {
|
for (const context of result.additionalContexts ?? []) {
|
||||||
exec.deferContext(context)
|
exec.deferContext(context)
|
||||||
}
|
}
|
||||||
if (exec.agent !== undefined) {
|
settle(result)
|
||||||
// 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 })
|
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
pump()
|
pump()
|
||||||
|
|||||||
@@ -473,10 +473,47 @@ describe('the sub-dispatch scheduler (native concurrency contract)', () => {
|
|||||||
return { logs: [], value: 'capped' }
|
return { logs: [], value: 'capped' }
|
||||||
}
|
}
|
||||||
const result = await runCode(ctx, 'program')
|
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(result.isError).toBe(false)
|
||||||
expect(gated.peakLive()).toBe(2)
|
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<PostToolDecision> => {
|
||||||
|
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 () => {
|
it('a queued-unstarted call abandoned by run settlement logs no start event', async () => {
|
||||||
const { ctx, runtime } = await setup({ mode: 'code' })
|
const { ctx, runtime } = await setup({ mode: 'code' })
|
||||||
const gated = registerGated(ctx, 'writer', false)
|
const gated = registerGated(ctx, 'writer', false)
|
||||||
|
|||||||
@@ -290,6 +290,66 @@ describe('the durable dispatch-log arm', () => {
|
|||||||
expect(spill.saves.filter(entry => entry.source.label === 'dispatch')).toHaveLength(0)
|
expect(spill.saves.filter(entry => entry.source.label === 'dispatch')).toHaveLength(0)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('a slow spill backend never delays the program value or a later dispatch slot', async () => {
|
||||||
|
const ctx = new Context()
|
||||||
|
await ctx.plugin(SystemPrompt)
|
||||||
|
await ctx.plugin(ToolRegistry, { mode: 'code' })
|
||||||
|
await ctx.plugin(StubStore)
|
||||||
|
await ctx.plugin(SpillPolicy, { maxInlineBytes: 100 })
|
||||||
|
await ctx.plugin(WorkerCodeRuntime, {})
|
||||||
|
// A spill backend that hangs until released.
|
||||||
|
let releaseSave!: () => void
|
||||||
|
const gate = new Promise<void>((resolve) => { releaseSave = resolve })
|
||||||
|
const store = ctx.spillStore as StubStore
|
||||||
|
const realSave = store.saveText.bind(store)
|
||||||
|
store.saveText = async (input) => {
|
||||||
|
await gate
|
||||||
|
return realSave(input)
|
||||||
|
}
|
||||||
|
const events: { type: string; data: unknown }[] = []
|
||||||
|
const agent = {
|
||||||
|
session: {
|
||||||
|
header: { id: SessionId('dispatch-slow-spill'), cwd: '/workspace' },
|
||||||
|
append: (type: string, data: unknown) => { events.push({ type, data }) },
|
||||||
|
},
|
||||||
|
}
|
||||||
|
ctx.tools.register(textTool('huge_read', 'H'.repeat(2_000)))
|
||||||
|
ctx.tools.register(textTool('small_read', 'tiny'))
|
||||||
|
let smallAfterHuge = false
|
||||||
|
const runPromise = ctx.tools.execute({
|
||||||
|
signal: testToolSignal,
|
||||||
|
callId: CallId('parent-3'),
|
||||||
|
name: 'run_code',
|
||||||
|
arguments: {
|
||||||
|
// The program takes BOTH values while the spill backend hangs: the
|
||||||
|
// huge read's binding resolves immediately (its logged copy is side
|
||||||
|
// work), so the small read proceeds without waiting.
|
||||||
|
code: 'const big = await tools.huge_read({});\nconst small = await tools.small_read({});\nreturn big[0].text.length + small[0].text.length',
|
||||||
|
description: 'Prove log shaping is off the program path',
|
||||||
|
},
|
||||||
|
agent: agent as never,
|
||||||
|
}).then((result) => {
|
||||||
|
return result
|
||||||
|
})
|
||||||
|
// The run cannot COMPLETE while the settle append is gated (drain waits
|
||||||
|
// for logWork), but the program itself already ran both calls; release
|
||||||
|
// the backend and observe the settle events land inside the turn.
|
||||||
|
await vi.waitFor(() => {
|
||||||
|
// The second dispatch STARTED while the first one's spill hung.
|
||||||
|
smallAfterHuge = events.some(event => event.type === 'tool/code-dispatch-start'
|
||||||
|
&& (event.data as { name: string }).name === 'small_read')
|
||||||
|
if (!smallAfterHuge) throw new Error('small_read not started yet')
|
||||||
|
})
|
||||||
|
releaseSave()
|
||||||
|
const result = await runPromise
|
||||||
|
expect(result.isError).toBe(false)
|
||||||
|
if (result.isError) throw new Error('expected success')
|
||||||
|
expect(result.value).toMatchObject({ result: 2_004 })
|
||||||
|
const settles = events.filter(event => event.type === 'tool/code-dispatch')
|
||||||
|
expect(settles).toHaveLength(2)
|
||||||
|
expect(smallAfterHuge).toBe(true)
|
||||||
|
})
|
||||||
|
|
||||||
it('a saveText failure keeps the complete content in the durable log (best-effort)', async () => {
|
it('a saveText failure keeps the complete content in the durable log (best-effort)', async () => {
|
||||||
const ctx = new Context()
|
const ctx = new Context()
|
||||||
await ctx.plugin(SystemPrompt)
|
await ctx.plugin(SystemPrompt)
|
||||||
|
|||||||
@@ -169,14 +169,14 @@ const TOOL_PACKAGES: ToolPackage[] = [
|
|||||||
dir: 'tools',
|
dir: 'tools',
|
||||||
source: 'packages/core/tools/src/code-mode.ts',
|
source: 'packages/core/tools/src/code-mode.ts',
|
||||||
requires: ['ctx.tools', 'ctx.codeRuntime (execution time)', 'ctx.systemPrompt'],
|
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's OWN tool: run_code exists only under a non-native mode
|
||||||
// (the registry registers it in its constructor; the code runtime is read
|
// (the registry registers it in its constructor; the code runtime is read
|
||||||
// at assembly/execution time, so the schema harvest needs none mounted).
|
// at assembly/execution time, so the schema harvest needs none mounted).
|
||||||
toolsConfig: { mode: 'code' },
|
toolsConfig: { mode: 'code' },
|
||||||
async mount() {},
|
async mount() {},
|
||||||
note:
|
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',
|
pkg: '@deepseek-ai/dsh-plan-mode',
|
||||||
|
|||||||
Reference in New Issue
Block a user