Merge remote-tracking branch 'origin/master' into xtr/agent-loop-message-machine
# Conflicts: # docs/cordis-catalog/services.md # examples/acp-agent/tests/snapshots/code-mode-workspace-context/session.jsonl # packages/core/tools/README.i18n.yaml # packages/core/tools/src/code-mode.ts
This commit is contained in:
77 files changed
+3404
-1456
No files matched your search
@@ -1,7 +1,8 @@
|
||||
/**
|
||||
* Code Mode `run_code` transport. Programs call the registry's agent-visible
|
||||
* tools through nested, sequential executions; each sub-dispatch is logged for
|
||||
* reconstruction, while only the outer curated result enters model history.
|
||||
* tools through nested executions scheduled under the native concurrency
|
||||
* contract; each sub-dispatch is logged for reconstruction, while only the
|
||||
* outer curated result enters model history.
|
||||
* @module @deepseek-ai/dsh-tools/src/code-mode
|
||||
*/
|
||||
|
||||
@@ -11,23 +12,39 @@ 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 {
|
||||
/**
|
||||
* One bridged sub-dispatch from a `run_code` program: the parent
|
||||
* `run_code` call id, the deterministic sub-call id
|
||||
* (`<parent>:code:<n>`), the tool `name` with its JSON-normalized
|
||||
* `arguments` — the exact value dispatched, normalized BEFORE dispatch,
|
||||
* so this append can never fail on payload shape — and the sub-call's
|
||||
* complete model-facing outcome in `tool/result`'s own vocabulary
|
||||
* One sub-dispatch STARTING inside a `run_code` program: the parent
|
||||
* `run_code` call id, the deterministic sub-call id (`<parent>:code:<n>`,
|
||||
* numbered in submission order), and the tool `name` with its
|
||||
* JSON-normalized `arguments` — the exact value dispatched, normalized
|
||||
* BEFORE dispatch, so this append can never fail on payload shape.
|
||||
* Appended when the scheduler actually starts the call (not at
|
||||
* submission), so a start means the tool body pipeline was entered; a
|
||||
* call abandoned in the queue logs nothing. Log-only: `deriveMessages()`
|
||||
* ignores it; UIs use it for live per-sub-call running state and pair it
|
||||
* with `tool/code-dispatch` by `subCallId` (timing = the two events'
|
||||
* `time` fields).
|
||||
*/
|
||||
'tool/code-dispatch-start': { parentCallId: CallId; subCallId: CallId; name: string; arguments: unknown }
|
||||
/**
|
||||
* One bridged sub-dispatch SETTLING: the pairing ids (matching the
|
||||
* `tool/code-dispatch-start` with the same `subCallId`), the tool `name`
|
||||
* with the same JSON-normalized `arguments`, and the sub-call's complete
|
||||
* model-facing outcome in `tool/result`'s own vocabulary
|
||||
* (`content` + `isError`), so UIs render a sub-call through the exact
|
||||
* code path that renders a native call.
|
||||
* code path that renders a native call. Every started sub-call settles
|
||||
* with exactly one of these (abort included: the aborted pipeline result
|
||||
* is an `isError` outcome).
|
||||
* Log-only: `deriveMessages()` ignores it, so sub-calls never re-enter
|
||||
* model context; persistence and UIs get every call. Appended inside the
|
||||
* parent `run_code`'s execution (the bridge drains its queue before
|
||||
* returning), so the turn-enclosure invariant holds by construction.
|
||||
* parent `run_code`'s execution (the bridge drains in-flight dispatches
|
||||
* before returning), so the turn-enclosure invariant holds by
|
||||
* construction.
|
||||
*/
|
||||
'tool/code-dispatch': { parentCallId: CallId; subCallId: CallId; name: string; arguments: unknown; isError: boolean; content: ContentBlock[] }
|
||||
}
|
||||
@@ -179,9 +196,11 @@ type RunCodeOutput = { logs: string[]; result?: JsonValue }
|
||||
* bindings cover its registered tools).
|
||||
* @param requireRuntime - resolves `ctx.codeRuntime` or throws the loud
|
||||
* misconfiguration error (shared with the registry's assembly-time checks).
|
||||
* @param maxParallel - the run's overlap cap for parallel-classified
|
||||
* sub-calls (the registry passes its validated `maxParallelSubCalls`).
|
||||
* @returns the registry-ready definition.
|
||||
*/
|
||||
export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () => CodeRuntime): ToolDefinition {
|
||||
export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () => CodeRuntime, maxParallel: number): ToolDefinition {
|
||||
return defineTool({
|
||||
name: RUN_CODE_NAME,
|
||||
description:
|
||||
@@ -229,19 +248,115 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
exec.signal.addEventListener('abort', onOuterAbort, { once: true })
|
||||
|
||||
let dispatches = 0
|
||||
// The per-run serialization queue: every binding call chains onto the tail, so even
|
||||
// `Promise.all` executes the underlying tool calls one at a time in submission order (the
|
||||
// tool contract carries no concurrency-safety metadata yet).
|
||||
let queue: Promise<void> = Promise.resolve()
|
||||
const enqueue = <T>(task: () => Promise<T>): Promise<T> => {
|
||||
const turn = queue.then(() => {
|
||||
if (runController.signal.aborted) {
|
||||
throw new Error(`run_code run is over (${String(runController.signal.reason)}); tool call abandoned`)
|
||||
// The per-run scheduler, reusing the NATIVE concurrency contract through
|
||||
// the registry's staged view (the loop scheduler's own seam) — and the
|
||||
// native loop's SEQUENCING: every ordered stage (the dispatch-start
|
||||
// append, prepare = pre-execute/guards, finalize/finish = post-execute,
|
||||
// context deferral, the settle append) runs inside ONE driver lane, so
|
||||
// ordered policy stages never overlap each other and only the
|
||||
// around-dispatch/body stage runs concurrently. Starts are strictly
|
||||
// submission-ordered; results commit in submission order through the
|
||||
// head-of-line cursor. Consecutive parallel-classified calls overlap up
|
||||
// to maxParallel; an exclusive call waits for the pool to drain, runs
|
||||
// alone, and holds its barrier until its COMMIT (post-execute included)
|
||||
// completes, exactly like a native exclusive group. 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 {
|
||||
/** Ordered stage: append the start event, await prepare (pre-execute/guards), launch the body into `flight`. */
|
||||
start(): Promise<void>
|
||||
classify(): 'parallel' | 'exclusive'
|
||||
abandon(): void
|
||||
/** Ordered stage: post-execute + context deferral + settle event, in submission order. */
|
||||
commit(): Promise<void>
|
||||
/** The launched around-dispatch/body stage; resolved until start() replaces it. */
|
||||
flight: Promise<void>
|
||||
/** True once the dispatch stage parked its outcome; the commit cursor waits on it. */
|
||||
settled: boolean
|
||||
/** The classification this entry started under; an exclusive holds its barrier through commit(). */
|
||||
mode?: 'parallel' | 'exclusive'
|
||||
}
|
||||
const pendingQueue: PendingDispatch[] = []
|
||||
const inFlight = new Set<Promise<void>>()
|
||||
const commitQueue: PendingDispatch[] = []
|
||||
let exclusiveActive = false
|
||||
let driving = false
|
||||
let driverRun: Promise<void> = Promise.resolve()
|
||||
let wake: (() => void) | undefined
|
||||
const wakeup = (): void => {
|
||||
const release = wake
|
||||
wake = undefined
|
||||
release?.()
|
||||
}
|
||||
/**
|
||||
* The single ordered lane. Each pass commits the head-of-line settled
|
||||
* dispatch (ordered post-execute), then starts the next queued entry if
|
||||
* its slot is free (ordered pre-execute), and otherwise sleeps until a
|
||||
* body settles or a new submission arrives. One run reaching the
|
||||
* empty-queues/empty-pool state is quiescence.
|
||||
*/
|
||||
const drive = (): Promise<void> => {
|
||||
if (driving) return driverRun
|
||||
driving = true
|
||||
driverRun = (async () => {
|
||||
try {
|
||||
for (;;) {
|
||||
// Arm before inspecting state so a settle or submission landing
|
||||
// between the checks and the await below cannot be lost.
|
||||
const signal = new Promise<void>((resolve) => { wake = resolve })
|
||||
const commitHead = commitQueue[0]
|
||||
if (commitHead !== undefined && commitHead.settled) {
|
||||
commitQueue.shift()
|
||||
await commitHead.commit()
|
||||
// The barrier covers post-execute: later starts wait for the
|
||||
// exclusive call's full pipeline, as under the native loop.
|
||||
if (commitHead.mode === 'exclusive') exclusiveActive = false
|
||||
continue
|
||||
}
|
||||
const head = pendingQueue[0]
|
||||
if (head !== undefined) {
|
||||
if (runController.signal.aborted) {
|
||||
pendingQueue.shift()
|
||||
head.abandon()
|
||||
continue
|
||||
}
|
||||
// Reclassify at start time (fail-closed on registry changes).
|
||||
const mode = head.classify()
|
||||
const capacity = !exclusiveActive
|
||||
&& (mode === 'exclusive' ? inFlight.size === 0 : inFlight.size < maxParallel)
|
||||
if (capacity) {
|
||||
if (mode === 'exclusive') exclusiveActive = true
|
||||
head.mode = mode
|
||||
pendingQueue.shift()
|
||||
// Joined before start() so the commit cursor sees submission
|
||||
// order; nothing commits it until `settled` flips.
|
||||
commitQueue.push(head)
|
||||
await head.start()
|
||||
const flight: Promise<void> = head.flight.finally(() => {
|
||||
inFlight.delete(flight)
|
||||
wakeup()
|
||||
})
|
||||
inFlight.add(flight)
|
||||
continue
|
||||
}
|
||||
}
|
||||
if (pendingQueue.length === 0 && commitQueue.length === 0 && inFlight.size === 0) return
|
||||
await signal
|
||||
}
|
||||
} finally {
|
||||
driving = false
|
||||
wake = undefined
|
||||
}
|
||||
return task()
|
||||
})
|
||||
queue = turn.then(() => undefined, () => undefined)
|
||||
return turn
|
||||
})()
|
||||
return driverRun
|
||||
}
|
||||
/** Every dispatch settled AND committed; nothing can start (the run is aborted at call time). */
|
||||
const drainDispatches = async (): Promise<void> => {
|
||||
// The abort already fired: the driver abandons queued-unstarted
|
||||
// entries, awaits the live pool, and drains the ordered commit lane —
|
||||
// including a commit already in progress when the program returned.
|
||||
await drive()
|
||||
}
|
||||
|
||||
// Read through a call, not a bare property: the abort state genuinely
|
||||
@@ -254,42 +369,91 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
throw new Error(`run_code run is over (${String(runController.signal.reason)}); ${name} not dispatched`)
|
||||
}
|
||||
const normalized = jsonNormalizeArgs(rawArgs)
|
||||
const outcome = await enqueue(async () => {
|
||||
const n = ++dispatches
|
||||
const subCallId = CallId(`${String(exec.callId)}:code:${n}`)
|
||||
const result = await registry.execute({
|
||||
callId: subCallId,
|
||||
name,
|
||||
arguments: normalized.dispatched,
|
||||
...exec.agent ? { agent: exec.agent } : {},
|
||||
parent: exec.token,
|
||||
signal: runController.signal,
|
||||
})
|
||||
for (const context of result.additionalContexts ?? []) {
|
||||
exec.deferContext(context)
|
||||
const n = ++dispatches
|
||||
const subCallId = CallId(`${String(exec.callId)}:code:${n}`)
|
||||
const input = {
|
||||
callId: subCallId,
|
||||
name,
|
||||
arguments: normalized.dispatched,
|
||||
...exec.agent ? { agent: exec.agent } : {},
|
||||
parent: exec.token,
|
||||
signal: runController.signal,
|
||||
}
|
||||
type DispatchOutcome = { isError: true; message: string } | { isError: false; value: JsonValue }
|
||||
const scheduler = registry[TOOL_REGISTRY_SCHEDULER]
|
||||
const outcome = await new Promise<DispatchOutcome>((resolve, reject) => {
|
||||
// Set by the dispatch stage (or start() for a pre-settled result): what commit() finalizes in submission order.
|
||||
let parked:
|
||||
| { kind: 'post-result' | 'final-result'; exec: ToolRunContext; result: ToolExecutionResult }
|
||||
| undefined
|
||||
const settle = (result: ToolExecutionResult): void => {
|
||||
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,
|
||||
// The registry deep-froze this projection at result finalization;
|
||||
// append snapshots it again, so the log copy stays detached.
|
||||
content: result.content,
|
||||
})
|
||||
resolve(result.isError
|
||||
? { isError: true, message: result.error.message }
|
||||
: { isError: false, value: result.value })
|
||||
}
|
||||
// Like the context forwarding above, cross-boundary facts travel on
|
||||
// the nested result and the composite forwards them: only a
|
||||
// successful nested result can carry the terminal marker
|
||||
// (ToolExecutionFailure types it never), so a policy-converted
|
||||
// failure cannot stop the turn through a recovering program.
|
||||
if (result.concludesTurn) exec.concludeTurn()
|
||||
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,
|
||||
// The registry deep-froze this projection at result finalization;
|
||||
// append snapshots it again, so the log copy stays detached.
|
||||
content: result.content,
|
||||
pendingQueue.push({
|
||||
flight: Promise.resolve(),
|
||||
settled: false,
|
||||
// Re-read per driver pass against the same agent view the SDK
|
||||
// declared; fail-closed exclusive when undeclared/invalid.
|
||||
classify: () => registry.executionMode(input).kind,
|
||||
abandon: () => {
|
||||
reject(new Error(`run_code run is over (${String(runController.signal.reason)}); ${name} tool call abandoned`))
|
||||
},
|
||||
async start(): Promise<void> {
|
||||
exec.agent?.session.append('tool/code-dispatch-start', {
|
||||
parentCallId: exec.callId,
|
||||
subCallId,
|
||||
name,
|
||||
arguments: normalized.logged,
|
||||
})
|
||||
// Ordered prepare runs INSIDE the driver lane: the next entry's
|
||||
// pre-execute waits for this resolution, as under the native
|
||||
// scheduler. Only the launched body below overlaps.
|
||||
const prepared = await scheduler.prepare(input)
|
||||
if (prepared.kind === 'dispatch') {
|
||||
this.flight = scheduler.dispatch(prepared.exec).then((dispatchOutcome) => {
|
||||
parked = { kind: dispatchOutcome.kind, exec: prepared.exec, result: dispatchOutcome.result }
|
||||
this.settled = true
|
||||
})
|
||||
return
|
||||
}
|
||||
parked = { kind: prepared.kind, exec: prepared.exec, result: prepared.result }
|
||||
this.settled = true
|
||||
},
|
||||
async commit(): Promise<void> {
|
||||
/* v8 ignore next -- commit() runs only after `settled` flipped, 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)
|
||||
}
|
||||
// Like the context forwarding above, cross-boundary facts travel
|
||||
// on the nested result and the composite forwards them: only a
|
||||
// successful nested result can carry the terminal marker
|
||||
// (ToolExecutionFailure types it never), so a policy-converted
|
||||
// failure cannot stop the turn through a recovering program.
|
||||
if (result.concludesTurn) exec.concludeTurn()
|
||||
settle(result)
|
||||
},
|
||||
})
|
||||
return result.isError
|
||||
? { isError: true as const, message: result.error.message }
|
||||
: { isError: false as const, value: result.value }
|
||||
wakeup()
|
||||
void drive()
|
||||
})
|
||||
// A budget expiry or outer cancel that lands while this call was in
|
||||
// flight already aborted the dispatch; stop the program now rather
|
||||
@@ -332,10 +496,11 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
signal: runController.signal,
|
||||
})
|
||||
} finally {
|
||||
// Abort sub-dispatches and drain the folded queue before closing the turn.
|
||||
// Abort sub-dispatches and drain every in-flight dispatch before
|
||||
// closing the turn (queued-unstarted ones are abandoned unlogged).
|
||||
// Binding failures remain observable through their individual promises.
|
||||
runController.abort('run_code settled')
|
||||
await queue
|
||||
await drainDispatches()
|
||||
}
|
||||
|
||||
if (result.error) {
|
||||
|
||||
Reference in New Issue
Block a user