Merge branch 'code-mode-ui/shiki' into code-mode-ui/trajectory-spans
This commit is contained in:
@@ -14,7 +14,7 @@ SlotsService gives the renderer separate bare observables for `useSessions` and
|
||||
|
||||
## Code Mode sub-dispatch index
|
||||
|
||||
`ConversationSnapshot.codeDispatches` groups a `run_code` call's sub-dispatches under their parent callId, in start order, using the native call-block shapes: a started-but-unsettled sub-call is a `RunningToolCall` (rows derive the running ring from the shape) and its `tool/code-dispatch` settlement replaces it in place with the `ToolResultNode` form, `callTime` carrying the paired start's time. Live mux frames and history replay build the identical index; sub-calls never join the surface `nodes` flow; per-parent array and map references are memo-stable across unrelated snapshot swaps.
|
||||
`ConversationSnapshot.codeDispatches` groups a `run_code` call's sub-dispatches under their parent callId, in start order, using the native call-block shapes: a `tool/code-dispatch-start` event lands the `RunningToolCall` form (rows derive the running ring from the shape) and its `tool/code-dispatch` settlement replaces it in place with the `ToolResultNode` form, `callTime` carrying the paired start's time. A settle whose start fell outside the replay window appends directly with `callTime: null` (duration unknown — never a fabricated zero). Live mux frames and history replay build the identical index; sub-calls never join the surface `nodes` flow; per-parent array and map references are memo-stable across unrelated snapshot swaps.
|
||||
|
||||
## Session title projection
|
||||
|
||||
|
||||
@@ -685,6 +685,9 @@ describe('run_code sub-dispatch indexing', () => {
|
||||
expect(subs?.[0]).toMatchObject({
|
||||
kind: 'tool-result', callId: 'p1:code:1',
|
||||
call: { name: 'bash', argsRaw: '{"command":"ls","description":"列目录"}' },
|
||||
// The settle event carries no start time: callTime stays null (never a
|
||||
// fabricated zero-duration).
|
||||
callTime: null,
|
||||
isError: false, content: [{ type: 'text', text: 'demo.txt' }],
|
||||
})
|
||||
expect(subs?.[1]).toMatchObject({ callId: 'p1:code:2', isError: true })
|
||||
|
||||
@@ -18,21 +18,26 @@ import langBash from '@shikijs/langs/shellscript'
|
||||
import langJson from '@shikijs/langs/json'
|
||||
import type { HighlighterCore } from 'shiki/core'
|
||||
|
||||
/** Language ids (and aliases) the singleton registers; everything else renders plain. */
|
||||
const LANG_ALIASES: Record<string, string> = {
|
||||
typescript: 'typescript',
|
||||
ts: 'typescript',
|
||||
tsx: 'typescript',
|
||||
javascript: 'typescript',
|
||||
js: 'typescript',
|
||||
shellscript: 'shellscript',
|
||||
bash: 'shellscript',
|
||||
sh: 'shellscript',
|
||||
shell: 'shellscript',
|
||||
zsh: 'shellscript',
|
||||
json: 'json',
|
||||
jsonc: 'json',
|
||||
}
|
||||
/**
|
||||
* Language ids (and aliases) the singleton registers; everything else renders
|
||||
* plain. A Map, not an object: fence info strings are assistant-authored, so
|
||||
* a label like `constructor` or `__proto__` must miss instead of resolving an
|
||||
* inherited property and crashing the renderer inside shiki.
|
||||
*/
|
||||
const LANG_ALIASES = new Map<string, string>([
|
||||
['typescript', 'typescript'],
|
||||
['ts', 'typescript'],
|
||||
['tsx', 'typescript'],
|
||||
['javascript', 'typescript'],
|
||||
['js', 'typescript'],
|
||||
['shellscript', 'shellscript'],
|
||||
['bash', 'shellscript'],
|
||||
['sh', 'shellscript'],
|
||||
['shell', 'shellscript'],
|
||||
['zsh', 'shellscript'],
|
||||
['json', 'json'],
|
||||
['jsonc', 'json'],
|
||||
])
|
||||
|
||||
/** All token colors resolve through `--shiki-*` custom properties (theme package sheets). */
|
||||
const cssVariablesTheme = createCssVariablesTheme({
|
||||
@@ -43,7 +48,7 @@ const cssVariablesTheme = createCssVariablesTheme({
|
||||
|
||||
let singleton: HighlighterCore | undefined
|
||||
|
||||
/** The lazily-created synchronous highlighter (one instance per document). */
|
||||
/** The synchronous highlighter (one instance per document); pre-warmed below, lazy as the fallback. */
|
||||
function highlighter(): HighlighterCore {
|
||||
singleton ??= createHighlighterCoreSync({
|
||||
themes: [cssVariablesTheme],
|
||||
@@ -53,6 +58,15 @@ function highlighter(): HighlighterCore {
|
||||
return singleton
|
||||
}
|
||||
|
||||
// Engine + grammar construction costs a long task (~120-175ms); building it
|
||||
// during the first finalized fence's render would jank exactly when a stream
|
||||
// completes. Warm the singleton in a deferred task at module load (= plugin
|
||||
// boot) instead; the lazy path above stays as the correctness fallback for a
|
||||
// fence that renders before the timer fires. `unref` (Node-only) keeps a
|
||||
// non-browser import from pinning the event loop.
|
||||
const warmupTimer = setTimeout(() => { highlighter() }, 0)
|
||||
;(warmupTimer as { unref?: () => void }).unref?.()
|
||||
|
||||
/**
|
||||
* Highlight `code` into shiki's HTML (a single `<pre class="shiki">` tree)
|
||||
* when `lang` maps to a registered grammar; `undefined` means the caller
|
||||
@@ -62,7 +76,7 @@ function highlighter(): HighlighterCore {
|
||||
* @returns the highlighted HTML, or `undefined` for unknown languages.
|
||||
*/
|
||||
export function highlightToHtml(code: string, lang: string | undefined): string | undefined {
|
||||
const resolved = lang === undefined ? undefined : LANG_ALIASES[lang.toLowerCase()]
|
||||
const resolved = lang === undefined ? undefined : LANG_ALIASES.get(lang.toLowerCase())
|
||||
if (resolved === undefined) return undefined
|
||||
return highlighter().codeToHtml(code, { lang: resolved, theme: 'css-variables' })
|
||||
}
|
||||
|
||||
@@ -64,6 +64,15 @@ describe('MarkdownText', () => {
|
||||
expect(screen.getByRole('link', { name: 'https://deepseek.com' })).toBeTruthy()
|
||||
})
|
||||
|
||||
it('a fence labeled with an inherited object key renders plain, never crashing shiki', () => {
|
||||
for (const label of ['constructor', '__proto__', 'toString', 'hasOwnProperty']) {
|
||||
const { container, unmount } = render(<MarkdownText text={'```' + label + '\ncode body\n```'} />)
|
||||
expect(container.querySelector('pre.shiki')).toBeNull()
|
||||
expect(container.querySelector('pre code')?.textContent).toContain('code body')
|
||||
unmount()
|
||||
}
|
||||
})
|
||||
|
||||
it('an empty fence keeps the stock pre; a language-less fence renders the plain CodeBlock arm', () => {
|
||||
const empty = render(<MarkdownText text={'```\n```'} />)
|
||||
expect(empty.container.querySelector('pre')?.outerHTML).toBe('<pre><code></code></pre>')
|
||||
|
||||
@@ -249,108 +249,119 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
|
||||
let dispatches = 0
|
||||
// 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.
|
||||
// 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, prepare, dispatch (body overlaps), park for commit. */
|
||||
/** 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>
|
||||
/** Set once the dispatch stage settles; commit() runs after this resolves. */
|
||||
dispatched?: 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>>()
|
||||
/** 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 pumping = false
|
||||
/** Ordered commit cursor: drain the head-of-line settled dispatches one at a time. */
|
||||
const commitReady = async (): Promise<void> => {
|
||||
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
|
||||
/* v8 ignore next -- entries join commitQueue only after start() set dispatched (see pump). */
|
||||
if (head.dispatched === undefined) break
|
||||
await head.dispatched
|
||||
commitQueue.shift()
|
||||
await head.commit()
|
||||
}
|
||||
} finally {
|
||||
committing = false
|
||||
}
|
||||
let driving = false
|
||||
let driverRun: Promise<void> = Promise.resolve()
|
||||
let wake: (() => void) | undefined
|
||||
const wakeup = (): void => {
|
||||
const release = wake
|
||||
wake = undefined
|
||||
release?.()
|
||||
}
|
||||
const pump = (): void => {
|
||||
// Defensive re-entry guard: today every caller (binding submission,
|
||||
// flight.finally, drain) runs off promise callbacks, never while pump
|
||||
// is on the stack, so this cannot fire — kept against a future
|
||||
// synchronous caller.
|
||||
/* v8 ignore next -- see the re-entry note above. */
|
||||
if (pumping) return
|
||||
pumping = true
|
||||
try {
|
||||
for (;;) {
|
||||
const head = pendingQueue[0]
|
||||
if (head === undefined) return
|
||||
if (runController.signal.aborted) {
|
||||
pendingQueue.shift()
|
||||
head.abandon()
|
||||
continue
|
||||
/**
|
||||
* 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
|
||||
}
|
||||
// Reclassify at start time (fail-closed on registry changes).
|
||||
const mode = head.classify()
|
||||
if (exclusiveActive || inFlight.size >= (mode === 'exclusive' ? 1 : maxParallel)) return
|
||||
// The guard above already returned for an exclusive head with any
|
||||
// in-flight sibling, so claiming the barrier here is race-free.
|
||||
if (mode === 'exclusive') exclusiveActive = true
|
||||
pendingQueue.shift()
|
||||
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()
|
||||
})
|
||||
// Joined AFTER start() ran synchronously, so every commitQueue
|
||||
// entry already carries its `dispatched` promise.
|
||||
commitQueue.push(head)
|
||||
inFlight.add(flight)
|
||||
} finally {
|
||||
driving = false
|
||||
wake = undefined
|
||||
}
|
||||
} finally {
|
||||
pumping = false
|
||||
}
|
||||
})()
|
||||
return driverRun
|
||||
}
|
||||
/** Every in-flight dispatch settled AND committed; nothing can start (the run is aborted at call time). */
|
||||
/** Every dispatch settled AND committed; nothing can start (the run is aborted at call time). */
|
||||
const drainDispatches = async (): Promise<void> => {
|
||||
// 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)
|
||||
}
|
||||
// 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()
|
||||
// Every settle's shaped append lands inside the open run_code turn
|
||||
// (tasks self-remove on settlement).
|
||||
while (logWork.size > 0) await Promise.allSettled([...logWork])
|
||||
}
|
||||
|
||||
// Read through a call, not a bare property: the abort state genuinely
|
||||
@@ -376,7 +387,7 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
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 start(): what commit() finalizes in submission order.
|
||||
// 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
|
||||
@@ -392,7 +403,7 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
: { isError: false, value: result.value })
|
||||
const agent = exec.agent
|
||||
if (agent === undefined) return
|
||||
logWork.add((async () => {
|
||||
const task: Promise<void> = (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.
|
||||
@@ -414,37 +425,41 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
isError: result.isError,
|
||||
content: logged,
|
||||
})
|
||||
})())
|
||||
})().finally(() => { logWork.delete(task) })
|
||||
logWork.add(task)
|
||||
}
|
||||
pendingQueue.push({
|
||||
// Re-read per pump pass against the same agent view the SDK
|
||||
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`))
|
||||
},
|
||||
start(): Promise<void> {
|
||||
async start(): Promise<void> {
|
||||
exec.agent?.session.append('tool/code-dispatch-start', {
|
||||
parentCallId: exec.callId,
|
||||
subCallId,
|
||||
name,
|
||||
arguments: normalized.logged,
|
||||
})
|
||||
// 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)
|
||||
// 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 }
|
||||
return
|
||||
}
|
||||
parked = { kind: prepared.kind, exec: prepared.exec, result: prepared.result }
|
||||
})()
|
||||
return this.dispatched
|
||||
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 this.dispatched resolved, which set parked. */
|
||||
/* 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)
|
||||
@@ -453,9 +468,16 @@ export function createRunCodeTool(registry: ToolRegistry, requireRuntime: () =>
|
||||
exec.deferContext(context)
|
||||
}
|
||||
settle(result)
|
||||
// Backpressure on the shaped-append side channel: pending log
|
||||
// tasks (each retaining a full result while a slow backend
|
||||
// stores it) are bounded by the pool cap — beyond it the
|
||||
// ordered lane waits, so later sub-calls cannot start and
|
||||
// pending I/O/memory cannot grow without bound.
|
||||
while (logWork.size > maxParallel) await Promise.race(logWork)
|
||||
},
|
||||
})
|
||||
pump()
|
||||
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
|
||||
|
||||
@@ -289,8 +289,10 @@ export type ToolExecutionMode =
|
||||
* One settled `run_code` sub-dispatch about to be logged, as seen by the
|
||||
* `tools/code-dispatch-log` waterfall: the parent execution (session owner,
|
||||
* outer call identity), the sub-call identity, and the outcome whose durable
|
||||
* copy a listener may reshape. The complete `content` is what the program
|
||||
* already received; only the `tool/code-dispatch` event's copy changes.
|
||||
* copy a listener may reshape. `content` is the RENDERED result projection
|
||||
* (what a native `tool/result` would carry) — the program itself received
|
||||
* the structured `value` (or just the error message on failure); only the
|
||||
* `tool/code-dispatch` event's copy changes.
|
||||
*/
|
||||
export interface CodeDispatchLog {
|
||||
/** The outer `run_code` execution. */
|
||||
@@ -670,6 +672,15 @@ interface FusedToolSignal {
|
||||
dispose(): void
|
||||
}
|
||||
|
||||
/** Resolve the run_code overlap cap at the owning config boundary (direct construction bypasses the Loader schema). */
|
||||
function resolveMaxParallelSubCalls(value: number | undefined): number {
|
||||
const maxParallelSubCalls = value ?? 10
|
||||
if (!Number.isInteger(maxParallelSubCalls) || maxParallelSubCalls < 1) {
|
||||
throw new Error('maxParallelSubCalls must be a positive integer')
|
||||
}
|
||||
return maxParallelSubCalls
|
||||
}
|
||||
|
||||
/**
|
||||
* Tool registry and execution pipeline. Scoped registrations shadow globals;
|
||||
* one visibility resolver feeds presentation, lookup, and dispatch.
|
||||
@@ -716,7 +727,7 @@ export class ToolRegistry extends Service {
|
||||
// the filterable global/scoped capability layers.
|
||||
this.codeTransport = this.mode === 'native'
|
||||
? undefined
|
||||
: createRunCodeTool(this, () => this.requireCodeRuntime(), config.maxParallelSubCalls ?? 10)
|
||||
: createRunCodeTool(this, () => this.requireCodeRuntime(), resolveMaxParallelSubCalls(config.maxParallelSubCalls))
|
||||
ctx.systemPrompt.tools(context => this.wireSchemas(context.scope))
|
||||
if (this.mode !== 'native') {
|
||||
ctx.systemPrompt.section({
|
||||
@@ -982,7 +993,7 @@ export class ToolRegistry extends Service {
|
||||
() => Promise.resolve(dispatch.content),
|
||||
)
|
||||
} catch (error: unknown) {
|
||||
this.ctx.logger.warn(`tools: code-dispatch-log listener failed for ${dispatch.name}: ${String(error)}; logging the unshaped content`)
|
||||
this.ctx.logger.warn(`tools: code-dispatch-log listener failed for ${dispatch.name}: ${errorMessage(error)}; logging the unshaped content`)
|
||||
return dispatch.content
|
||||
}
|
||||
}
|
||||
|
||||
@@ -510,6 +510,113 @@ describe('the sub-dispatch scheduler (native concurrency contract)', () => {
|
||||
expect(calls).toEqual([])
|
||||
})
|
||||
|
||||
it('ordered pre-execute never overlaps: a slow policy on one call delays the next start', async () => {
|
||||
const { ctx, runtime } = await setup({ mode: 'code' })
|
||||
const gated = registerGated(ctx, 'safe_read', true)
|
||||
const stages: string[] = []
|
||||
let releaseGate: (() => void) | undefined
|
||||
ctx.on('tools/pre-execute', async (preExec, next) => {
|
||||
if (preExec.name !== 'safe_read') return next()
|
||||
stages.push(`pre-enter:${String(preExec.callId)}`)
|
||||
if (releaseGate === undefined) {
|
||||
// The FIRST call's policy awaits an asynchronous decision.
|
||||
await new Promise<void>((resolve) => { releaseGate = resolve })
|
||||
}
|
||||
stages.push(`pre-exit:${String(preExec.callId)}`)
|
||||
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' })])
|
||||
// Both submissions are in; the second pre-execute must NOT have entered
|
||||
// while the first is still awaiting its policy decision.
|
||||
await expect.poll(() => stages.length).toBeGreaterThanOrEqual(1)
|
||||
expect(stages).toEqual(['pre-enter:call-1:code:1'])
|
||||
releaseGate!()
|
||||
await expect.poll(() => gated.pending()).toBe(2)
|
||||
gated.releaseAll()
|
||||
await all
|
||||
return { logs: [], value: 'ordered-prepare' }
|
||||
}
|
||||
const result = await runCode(ctx, 'program')
|
||||
expect(result.isError).toBe(false)
|
||||
expect(stages).toEqual([
|
||||
'pre-enter:call-1:code:1', 'pre-exit:call-1:code:1',
|
||||
'pre-enter:call-1:code:2', 'pre-exit:call-1:code:2',
|
||||
])
|
||||
})
|
||||
|
||||
it('an exclusive call holds its barrier through post-execute: the next start waits for the commit', async () => {
|
||||
const { ctx, runtime } = await setup({ mode: 'code' })
|
||||
const writer = registerGated(ctx, 'writer', false)
|
||||
const reader = registerGated(ctx, 'safe_read', true)
|
||||
const stages: string[] = []
|
||||
let releasePost: (() => void) | undefined
|
||||
ctx.on('tools/post-execute', async (postExec, _result, next): Promise<PostToolDecision> => {
|
||||
if (postExec.name === 'writer') {
|
||||
stages.push('post-enter:writer')
|
||||
await new Promise<void>((resolve) => { releasePost = resolve })
|
||||
stages.push('post-exit:writer')
|
||||
}
|
||||
return next()
|
||||
})
|
||||
runtime.behavior = async (request) => {
|
||||
const tools = request.bindings[0]!.functions
|
||||
const w = tools.writer!({ id: 'w' })
|
||||
const r = tools.safe_read!({ id: 'r' })
|
||||
await expect.poll(() => writer.pending()).toBe(1)
|
||||
writer.release()
|
||||
// The writer's body is done and its async post-execute is running; the
|
||||
// parallel read must not have STARTED (no pre/body) while the exclusive
|
||||
// call's pipeline is still open.
|
||||
await expect.poll(() => stages).toContain('post-enter:writer')
|
||||
expect(reader.pending()).toBe(0)
|
||||
releasePost!()
|
||||
await w
|
||||
await expect.poll(() => reader.pending()).toBe(1)
|
||||
reader.releaseAll()
|
||||
await r
|
||||
return { logs: [], value: 'barrier-through-commit' }
|
||||
}
|
||||
const result = await runCode(ctx, 'program')
|
||||
expect(result.isError).toBe(false)
|
||||
expect(stages).toEqual(['post-enter:writer', 'post-exit:writer'])
|
||||
})
|
||||
|
||||
it('run settlement drains a commit already in progress: the settle event lands inside the turn', async () => {
|
||||
const { ctx, runtime } = await setup({ mode: 'code' })
|
||||
const gated = registerGated(ctx, 'safe_read', true)
|
||||
const { agent, events } = fakeAgent()
|
||||
let releasePost: (() => void) | undefined
|
||||
ctx.on('tools/post-execute', async (postExec, _result, next): Promise<PostToolDecision> => {
|
||||
if (postExec.name === 'safe_read') {
|
||||
await new Promise<void>((resolve) => { releasePost = resolve })
|
||||
}
|
||||
return next()
|
||||
})
|
||||
runtime.behavior = async (request) => {
|
||||
// Fire-and-forget: the program returns while the sub-call's async
|
||||
// post-execute commit is mid-flight.
|
||||
request.bindings[0]!.functions.safe_read!({ id: 'a' }).catch(() => 'run-over')
|
||||
await expect.poll(() => gated.pending()).toBe(1)
|
||||
gated.release()
|
||||
await expect.poll(() => releasePost !== undefined).toBe(true)
|
||||
queueMicrotask(() => { releasePost!() })
|
||||
return { logs: [], value: 'returned-early' }
|
||||
}
|
||||
const result = await runCode(ctx, 'program', { agent })
|
||||
expect(result.isError).toBe(false)
|
||||
// The drain awaited the in-progress commit: the settle event exists and
|
||||
// preceded the run_code turn closing (all appends happen inside
|
||||
// execute()). The run's settlement aborted the sub-call's signal while
|
||||
// its post-execute was mid-flight, so the native cancellation contract
|
||||
// replaces the successful outcome with the aborted result — the event is
|
||||
// still durable and in-turn, which is the invariant under test.
|
||||
const settles = events.filter(event => event.type === 'tool/code-dispatch')
|
||||
expect(settles).toHaveLength(1)
|
||||
expect(settles[0]?.data).toMatchObject({ name: 'safe_read', isError: true })
|
||||
})
|
||||
|
||||
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)
|
||||
@@ -1308,6 +1415,13 @@ describe('the run_code dispatch bridge', () => {
|
||||
expect(derived[0]?.role).toBe('user')
|
||||
})
|
||||
|
||||
it('direct construction rejects a non-positive parallel sub-call cap at load', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt, {})
|
||||
expect(() => new ToolRegistry(ctx, { mode: 'code', maxParallelSubCalls: 0 }))
|
||||
.toThrow('maxParallelSubCalls must be a positive integer')
|
||||
})
|
||||
|
||||
it('direct construction in code mode defaults the parallel sub-call cap', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt, {})
|
||||
|
||||
@@ -29,9 +29,12 @@ const testToolSignal = new AbortController().signal
|
||||
class StubStore extends SpillStore {
|
||||
saves: SaveTextSpill[] = []
|
||||
fail = false
|
||||
/** Per-save hang hook: each call awaits the returned promise before completing. */
|
||||
gate: (() => Promise<void>) | undefined
|
||||
|
||||
async saveText(input: SaveTextSpill): Promise<SpillRef> {
|
||||
if (this.fail) throw new Error('disk full')
|
||||
await this.gate?.()
|
||||
this.saves.push(input)
|
||||
return {
|
||||
locator: SpillLocator(`/spill/${input.suggestedName}`),
|
||||
@@ -367,6 +370,63 @@ describe('the durable dispatch-log arm', () => {
|
||||
expect(smallAfterHuge).toBe(true)
|
||||
})
|
||||
|
||||
it('a sustained slow backend backpressures the run instead of accumulating unbounded log tasks', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt)
|
||||
// Cap 1: once the hung shaped-append backlog exceeds the cap, the ordered
|
||||
// lane holds inside the second commit, so the THIRD dispatch cannot start
|
||||
// until a pending save drains — the bound is observable as its missing
|
||||
// start event.
|
||||
await ctx.plugin(ToolRegistry, { mode: 'code', maxParallelSubCalls: 1 })
|
||||
await ctx.plugin(StubStore)
|
||||
await ctx.plugin(SpillPolicy, { maxInlineBytes: 100 })
|
||||
await ctx.plugin(WorkerCodeRuntime, {})
|
||||
const store = ctx.spillStore as StubStore
|
||||
const releases: (() => void)[] = []
|
||||
store.gate = () => new Promise<void>((resolve) => { releases.push(resolve) })
|
||||
const events: { type: string; data: unknown }[] = []
|
||||
const agent = {
|
||||
session: {
|
||||
header: { id: SessionId('dispatch-spill-bound'), cwd: '/workspace' },
|
||||
append: (type: string, data: unknown) => { events.push({ type, data }) },
|
||||
},
|
||||
}
|
||||
ctx.tools.register(textTool('huge_read', 'H'.repeat(2_000)))
|
||||
const started = (n: number): boolean => events.some(event => event.type === 'tool/code-dispatch-start'
|
||||
&& (event.data as { subCallId: string }).subCallId.endsWith(`:code:${n}`))
|
||||
const runPromise = ctx.tools.execute({
|
||||
signal: testToolSignal,
|
||||
callId: CallId('parent-bound'),
|
||||
name: 'run_code',
|
||||
arguments: {
|
||||
code: 'await tools.huge_read({}); await tools.huge_read({}); await tools.huge_read({}); return "done"',
|
||||
description: 'Three oversized reads against a hung backend',
|
||||
},
|
||||
agent: agent as never,
|
||||
})
|
||||
// Two hung saves = backlog above the cap: the lane must hold before
|
||||
// starting dispatch 3.
|
||||
await vi.waitFor(() => {
|
||||
if (releases.length < 2) throw new Error('second hung save not reached yet')
|
||||
})
|
||||
expect(started(2)).toBe(true)
|
||||
expect(started(3)).toBe(false)
|
||||
releases.shift()!()
|
||||
// Draining one pending save releases the lane; dispatch 3 starts.
|
||||
await vi.waitFor(() => {
|
||||
if (!started(3)) throw new Error('third dispatch not started yet')
|
||||
})
|
||||
while (releases.length > 0) releases.shift()!()
|
||||
const result = await runPromise
|
||||
expect(result.isError).toBe(false)
|
||||
await vi.waitFor(() => {
|
||||
if (releases.length > 0) { while (releases.length > 0) releases.shift()!() }
|
||||
if (events.filter(event => event.type === 'tool/code-dispatch').length !== 3) {
|
||||
throw new Error('settle events still pending')
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
it('a saveText failure keeps the complete content in the durable log (best-effort)', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(SystemPrompt)
|
||||
|
||||
Reference in New Issue
Block a user