|
|
@@ -2,10 +2,12 @@
|
|
|
* Schedules one assistant step's tool calls. Exclusive calls form barriers;
|
|
|
* parallel calls use a bounded rolling pool and are reclassified before start.
|
|
|
* Dispatch may overlap, while policy, results, and result context remain
|
|
|
- * model-ordered. Abort stops replenishment and drains started calls.
|
|
|
+ * model-ordered. Abort or an internal scheduler failure stops replenishment
|
|
|
+ * and drains started calls.
|
|
|
*
|
|
|
- * Each advertised call records a balanced `tool/call`/`tool/result` pair. Calls
|
|
|
- * skipped after abort receive synthetic error results so replay stays valid.
|
|
|
+ * Abort records synthetic error results for skipped calls so replay stays
|
|
|
+ * valid. A terminal scheduler failure preserves already-recorded `tool/call`
|
|
|
+ * events without fabricating results.
|
|
|
* @module dsh-agent-loop/tool-calls
|
|
|
*/
|
|
|
|
|
|
@@ -37,10 +39,13 @@ interface GroupOutcome {
|
|
|
|
|
|
/**
|
|
|
* Schedule one assistant step's tool calls by their live concurrency mode.
|
|
|
- * Started calls receive ordered results. Abort drains them, records synthetic
|
|
|
- * results for unstarted calls, and returns with the signal still aborted after
|
|
|
- * accepting started-call context through the caller-supplied acceptor (the
|
|
|
- * machine stages it on its outbox for the next step boundary).
|
|
|
+ * Ordinary completion and abort commit started-call results in order. Abort
|
|
|
+ * drains them, records synthetic results for unstarted calls, and returns with
|
|
|
+ * the signal still aborted after accepting started-call context through the
|
|
|
+ * caller-supplied acceptor (the machine stages it on its outbox for the next
|
|
|
+ * step boundary). An internal scheduler failure stops new dispatches, drains
|
|
|
+ * already-started dispatches, and rejects with the first failure without
|
|
|
+ * fabricating tool results.
|
|
|
* The committed step's AgentLoop driver boundary supplies the initiating Agent
|
|
|
* that becomes each explicit {@link ToolExecutionInput.agent}.
|
|
|
*
|
|
|
@@ -110,7 +115,8 @@ function parseArguments(raw: string): unknown {
|
|
|
* drain and remains for the caller's next barrier. Results and contexts commit
|
|
|
* in model order. Abort stops starts, drains and commits started calls, accepts
|
|
|
* their contexts into the owning batch, records results for skipped calls, and
|
|
|
- * returns an aborted outcome.
|
|
|
+ * returns an aborted outcome. Scheduler failure drains dispatches without
|
|
|
+ * committing synthetic recovery results.
|
|
|
*/
|
|
|
async function runGroup(
|
|
|
ctx: Context,
|
|
|
@@ -131,6 +137,10 @@ async function runGroup(
|
|
|
let started = 0
|
|
|
let aborted: boolean = signal.aborted
|
|
|
let concluded = false
|
|
|
+ let schedulerFailure: { error: unknown } | undefined
|
|
|
+ const throwSchedulerFailure = (): void => {
|
|
|
+ if (schedulerFailure !== undefined) throw schedulerFailure.error
|
|
|
+ }
|
|
|
|
|
|
// `committed` advances only across contiguous model-order slots.
|
|
|
const commitReady = async (): Promise<void> => {
|
|
|
@@ -157,12 +167,19 @@ async function runGroup(
|
|
|
callSeqs[index] = appendToolCall(session, turn, step, call.block)
|
|
|
started++
|
|
|
const prepared = await ctx.tools[TOOL_REGISTRY_SCHEDULER].prepare(call.exec)
|
|
|
+ throwSchedulerFailure()
|
|
|
switch (prepared.kind) {
|
|
|
case 'dispatch': {
|
|
|
- const promise = ctx.tools[TOOL_REGISTRY_SCHEDULER].dispatch(prepared.exec).then((outcome) => {
|
|
|
- slots[index] = { exec: prepared.exec, result: outcome.result, needsPost: outcome.kind === 'post-result' }
|
|
|
- return index
|
|
|
- })
|
|
|
+ const promise = ctx.tools[TOOL_REGISTRY_SCHEDULER].dispatch(prepared.exec).then(
|
|
|
+ (outcome) => {
|
|
|
+ slots[index] = { exec: prepared.exec, result: outcome.result, needsPost: outcome.kind === 'post-result' }
|
|
|
+ return index
|
|
|
+ },
|
|
|
+ (error: unknown) => {
|
|
|
+ schedulerFailure ??= { error }
|
|
|
+ return index
|
|
|
+ },
|
|
|
+ )
|
|
|
inFlight.set(index, promise)
|
|
|
break
|
|
|
}
|
|
|
@@ -187,24 +204,34 @@ async function runGroup(
|
|
|
&& ctx.tools.executionMode(nextCall.exec).kind !== 'parallel') break
|
|
|
await startCall(nextToStart)
|
|
|
nextToStart++
|
|
|
+ throwSchedulerFailure()
|
|
|
await commitReady()
|
|
|
+ throwSchedulerFailure()
|
|
|
// Abort may arrive while pre-execute awaits.
|
|
|
if (signal.aborted) aborted = true
|
|
|
}
|
|
|
}
|
|
|
|
|
|
- // Ordered pre-execute may await; only dispatch/body overlaps.
|
|
|
- // TODO: Drain every started call before rethrowing a scheduler error; tool
|
|
|
- // bodies must not outlive the failed turn.
|
|
|
- await fillPool()
|
|
|
- while (inFlight.size > 0) {
|
|
|
- const settledIndex = await Promise.race(inFlight.values())
|
|
|
- inFlight.delete(settledIndex)
|
|
|
- await commitReady()
|
|
|
- // Abort may arrive while a tool or ordered commit awaits.
|
|
|
-
|
|
|
- if (signal.aborted) aborted = true
|
|
|
+ // Ordered pre-execute may await; only dispatch/body overlaps. A scheduler
|
|
|
+ // failure stops new dispatches and reaches the turn boundary after every
|
|
|
+ // already-started dispatch settles.
|
|
|
+ try {
|
|
|
await fillPool()
|
|
|
+ while (inFlight.size > 0) {
|
|
|
+ const settledIndex = await Promise.race(inFlight.values())
|
|
|
+ inFlight.delete(settledIndex)
|
|
|
+ throwSchedulerFailure()
|
|
|
+ await commitReady()
|
|
|
+ throwSchedulerFailure()
|
|
|
+ // Abort may arrive while a tool or ordered commit awaits.
|
|
|
+
|
|
|
+ if (signal.aborted) aborted = true
|
|
|
+ await fillPool()
|
|
|
+ }
|
|
|
+ } catch (error: unknown) {
|
|
|
+ schedulerFailure ??= { error }
|
|
|
+ await Promise.allSettled(inFlight.values())
|
|
|
+ throw schedulerFailure.error
|
|
|
}
|
|
|
|
|
|
if (aborted) {
|