Merge latest origin/master into parallel-tool-call

# Conflicts:
#	docs/architecture.md
#	examples/acp-agent/tests/snapshots/bash-spill/session.jsonl
#	examples/acp-agent/tests/snapshots/escalation-approved/session.jsonl
#	examples/acp-agent/tests/snapshots/escalation-rejected/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-pretool-ask/session.jsonl
#	packages/core/agent-loop/src/loop.ts
This commit is contained in:
Tianyi Cui
2026-07-18 15:09:39 +08:00
198 changed files with 6086 additions and 4036 deletions
+58 -12
View File
@@ -8,7 +8,7 @@
import type { Context } from 'cordis'
import { agentEvents } from '@deepseek-ai/dsh-agent'
import type { AgentId, AgentOptions, AgentStatus, InjectOptions, SendOptions } from '@deepseek-ai/dsh-agent'
import type { AgentId, AgentOptions, AgentStatus, HookContext, InjectOptions, SendOptions } from '@deepseek-ai/dsh-agent'
import type { Agent } from '@deepseek-ai/dsh-agent'
import { deepFreeze } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
@@ -152,6 +152,10 @@ export class ReactLoopAgent implements Agent {
* this set before the lifecycle unregisters the agent or detaches its session.
*/
private pendingIdleFlushes = new Set<Promise<void>>()
/** Whether the current step is executing an assistant tool-call batch. */
private toolBatchActive = false
/** Open-turn injections waiting for the active assistant tool-call batch to close. */
private deferredInjections: HookContext[] = []
constructor(
private loopCtx: Context,
@@ -194,12 +198,11 @@ export class ReactLoopAgent implements Agent {
}
/**
* Accept one public send/steer payload as the exact detached record shared by
* the live notification and inbox. Lossless-JSON materialization reads every
* nested field once; deep freeze prevents an observer from rewriting queued
* work before the loop drains it.
* Accept one public message payload as a detached record. Lossless-JSON
* materialization reads every nested field once; deep freeze prevents later
* caller mutation before an inbox or deferred-injection queue drains it.
*/
private acceptInboxMessage(content: ContentBlock[], options?: SendOptions): InboxMessage {
private acceptMessage(content: ContentBlock[], options?: SendOptions): InboxMessage {
const source = this.resolveSource(options)
const accepted = snapshotJsonValue({ content, source })
if (accepted === undefined) {
@@ -208,6 +211,15 @@ export class ReactLoopAgent implements Agent {
return deepFreeze(accepted)
}
/** Detach one context before it can outlive its caller in the active-batch FIFO. */
private acceptContext(context: HookContext): HookContext {
const accepted = snapshotJsonValue(context)
if (accepted === undefined) {
throw new TypeError('agent context must be losslessly JSON-serializable')
}
return deepFreeze(accepted)
}
/** Reject a driving operation once teardown has synchronously closed the agent. */
private assertNotDisposed(): void {
if (this._status === 'disposed') throw new Error(`agent "${this.id}" is disposed`)
@@ -215,7 +227,7 @@ export class ReactLoopAgent implements Agent {
send(content: ContentBlock[], options?: SendOptions): void {
this.assertNotDisposed()
const accepted = this.acceptInboxMessage(content, options)
const accepted = this.acceptMessage(content, options)
this.#inbox.enqueue(accepted)
const info = { source: accepted.source, steering: false } as const
agentEvents(this.loopCtx, this).emit('agent/queued', accepted.content, info)
@@ -224,7 +236,7 @@ export class ReactLoopAgent implements Agent {
steer(content: ContentBlock[], options?: SendOptions): void {
this.assertNotDisposed()
if (this._status !== 'running') { this.send(content, options); return }
const accepted = this.acceptInboxMessage(content, options)
const accepted = this.acceptMessage(content, options)
this.#inbox.steer(accepted)
const info = { source: accepted.source, steering: true } as const
agentEvents(this.loopCtx, this).emit('agent/queued', accepted.content, info)
@@ -240,10 +252,15 @@ export class ReactLoopAgent implements Agent {
...options?.meta !== undefined ? { meta: options.meta } : {},
}
if (isTurnOpen(this.session)) {
// A turn is open in the LOG (decided from the log, not agent status —
// status can be `running` with no turn open): the context/message is
// turn-enclosed by that turn, so append it directly.
this.session.append('context/message', context, { surfaceOp: 'append' })
const accepted = this.acceptContext(context)
// Provider protocols require every assistant tool-call batch to be
// followed only by its tool results. Historical interrupted batches do
// not own new context; only the currently executing batch may defer it.
if (this.toolBatchActive) {
this.deferredInjections.push(accepted)
return
}
this.session.append('context/message', accepted, { surfaceOp: 'append' })
return
}
// No turn open: wrap the injection in a one-shot turn so every event stays
@@ -283,6 +300,34 @@ export class ReactLoopAgent implements Agent {
}
}
/** Append deferred open-turn injections after the loop closes a tool-result batch. */
private drainDeferredInjections(): void {
const pending = this.deferredInjections.splice(0)
for (const accepted of pending) {
this.session.append('context/message', accepted, { surfaceOp: 'append' })
}
}
/**
* Run one tool-call batch and drain its deferred context before settlement.
* The loop-owned acceptor remains valid after public disposal begins because
* the interrupted turn stays open until this batch settles.
*/
private async withToolBatch<T>(
run: (acceptContext: (context: HookContext) => void) => Promise<T>,
): Promise<T> {
this.toolBatchActive = true
const acceptContext = (context: HookContext): void => {
this.deferredInjections.push(this.acceptContext(context))
}
try {
return await run(acceptContext)
} finally {
this.toolBatchActive = false
this.drainDeferredInjections()
}
}
cancel(reason?: string): void {
// Arm only for current work; an idle marker would cancel the next prompt.
if (this._status === 'running' || this.currentAbort !== undefined || this.#inbox.hasQueued || this.#inbox.hasSteering) {
@@ -348,6 +393,7 @@ export class ReactLoopAgent implements Agent {
isCancelled: () => this.cancelRequested,
cancelReason: () => this.cancelReason,
clearCancel: () => { this.cancelRequested = false },
withToolBatch: run => this.withToolBatch(run),
// Pre-step cancellation re-parks without emitting a status transition.
settleIdle: () => { this.settleIdleWaiters() },
})
+55 -25
View File
@@ -10,7 +10,7 @@ import type { ContentBlock, FinishReason, GenerateOptions, LlmCallConfig, Messag
import { isDeepStrictEqual } from 'node:util'
import { BlockAssembler, HarnessError, deepFreeze } from '@deepseek-ai/dsh-llm'
import { agentEvents, assembleContextFor } from '@deepseek-ai/dsh-agent'
import type { AgentEventDispatch, ContinuationDecision, PromptDecision } from '@deepseek-ai/dsh-agent'
import type { AgentEventDispatch, ContinuationDecision, HookContext, PromptDecision } from '@deepseek-ai/dsh-agent'
import { canonicalHeader } from '@deepseek-ai/dsh-session'
import type { Session, TurnEndReason, TurnTrigger } from '@deepseek-ai/dsh-session'
import { createTransmissionLog, recordRequestHeader } from './request-log.ts'
@@ -89,6 +89,8 @@ export interface LoopHandle {
clearCancel(): void
/** Settle idle waiters when pre-running cancellation skips a turn, without emitting `agent/status`. */
settleIdle(): void
/** Run an active tool-call batch, accepting post-tool context into the FIFO drained before settlement. */
readonly withToolBatch: <T>(run: (acceptContext: (context: HookContext) => void) => Promise<T>) => Promise<T>
}
/**
@@ -334,8 +336,7 @@ async function runTurn(
let stepOutcome: { hadToolCalls: boolean; finish: FinishReason } | { error: Error }
try {
stepOutcome = await runStep(
ctx, events, agent, turn, step, assembly, fullSystemPrompt, boundaryMessages,
transmission, abort.signal, handle.maxParallelToolCalls)
ctx, events, agent, handle, turn, step, assembly, fullSystemPrompt, boundaryMessages, transmission, abort.signal)
} catch (error: unknown) {
stepOutcome = { error: toError(error) }
} finally {
@@ -472,6 +473,7 @@ async function runStep(
ctx: Context,
events: AgentEventDispatch,
agent: ReactLoopAgent,
handle: LoopHandle,
turn: number,
step: number,
assembly: PromptAssembly,
@@ -479,7 +481,6 @@ async function runStep(
boundaryMessages: Message[],
transmission: TransmissionLog,
signal: AbortSignal,
maxParallelToolCalls: number,
): Promise<{ hadToolCalls: boolean; finish: FinishReason }> {
const { session, options } = agent
@@ -541,7 +542,9 @@ async function runStep(
const assembled = assembler.message()
const assembledContent = structuredClone(assembled.content)
let message: Message = withoutToolCalls(assembled)
message = withoutToolCalls(await events.waterfall('agent/step-result', turn, step, message, () => Promise.resolve(message)))
message = withoutToolCalls(await processStepResult(
events, session, turn, step, header.config, assembledContent, message, assembler, chunkSeqs,
))
// Preserve usage even when max-token truncation produced no content.
recordAssistantMessage(session, turn, step, header.config, assembledContent, message, assembler, chunkSeqs)
return { hadToolCalls: false, finish: assembler.finish }
@@ -551,28 +554,55 @@ async function runStep(
const assembled = assembler.message()
const assembledContent = structuredClone(assembled.content)
let message: Message = assembled
message = await events.waterfall('agent/step-result', turn, step, message, () => Promise.resolve(message))
message = await processStepResult(
events, session, turn, step, header.config, assembledContent, message, assembler, chunkSeqs,
)
const toolCalls = message.content.filter(block => block.type === 'tool-call')
// Empty messages exist only to carry usage; the helper also omits empty chunk provenance.
// Every successful call records its completion anchor, including explicit
// empty chunk provenance for a contentless, usage-less provider response.
recordAssistantMessage(session, turn, step, header.config, assembledContent, message, assembler, chunkSeqs)
// Dispatch may overlap; policy, results, and context remain model-ordered.
const pendingContext = toolCalls.length > 0
? await executeToolCalls(ctx, agent, turn, step, toolCalls, signal, maxParallelToolCalls)
: []
// Dispatch may overlap; policy, durable results, and result context stay model-ordered.
const toolCalls = message.content.filter(block => block.type === 'tool-call')
if (toolCalls.length === 0) return { hadToolCalls: false, finish: assembler.finish }
return handle.withToolBatch(async (acceptContext) => {
await executeToolCalls(
ctx, agent, turn, step, toolCalls, signal, handle.maxParallelToolCalls, acceptContext,
)
return { hadToolCalls: true, finish: assembler.finish }
})
}
// Context follows the complete result batch to preserve call/result adjacency.
for (const context of pendingContext) {
agent.inject(context.content, {
source: context.source,
...context.envelope !== undefined ? { envelope: context.envelope } : {},
...context.meta !== undefined ? { meta: context.meta } : {},
})
/** Preserve successful-call accounting without retaining output that result processing rejected. */
async function processStepResult(
events: AgentEventDispatch,
session: Session,
turn: number,
step: number,
config: LlmCallConfig,
assembledContent: ContentBlock[],
message: Message,
assembler: BlockAssembler,
chunkSeqs: number[],
): Promise<Message> {
try {
return await events.waterfall(
'agent/step-result', turn, step, message, () => Promise.resolve(message),
)
} catch (error: unknown) {
recordAssistantMessage(
session,
turn,
step,
config,
assembledContent,
{ ...message, content: [] },
assembler,
chunkSeqs,
false,
)
throw error
}
return { hadToolCalls: toolCalls.length > 0, finish: assembler.finish }
}
/** Record one content-or-usage assistant message with replay-safe provenance. */
@@ -585,8 +615,8 @@ function recordAssistantMessage(
message: Message,
assembler: BlockAssembler,
chunkSeqs: number[],
preserveReplayState = true,
): void {
if (message.content.length === 0 && assembler.usage === undefined) return
session.append(
'assistant/message',
{
@@ -596,11 +626,11 @@ function recordAssistantMessage(
provenance: assistantProvenance(
config,
assembler.replayState,
isDeepStrictEqual(message.content, assembledContent),
preserveReplayState && isDeepStrictEqual(message.content, assembledContent),
),
...assembler.usage === undefined ? {} : { usage: assembler.usage },
},
{ surfaceOp: 'append', ...chunkSeqs.length > 0 ? { sourceEventSeqs: chunkSeqs } : {} },
{ surfaceOp: 'append', sourceEventSeqs: chunkSeqs },
)
}
+7 -15
View File
@@ -1,11 +1,12 @@
/**
* Per-loop-instance request-header bookkeeping for reconstructability. The
* comparison baseline is the header folded from the session log, so a fresh
* loop instance needs no special resume or fork state.
* comparison baseline is folded from the session log; a fresh instance anchors
* it with an initial/resume snapshot and later logs full changed snapshots.
*
* @module dsh-agent-loop/request-log
*/
import { diffHeader, headerEquals, applyHeaderDelta } from '@deepseek-ai/dsh-session'
import { headerEquals } from '@deepseek-ai/dsh-session'
import type { EpochHeader, Session } from '@deepseek-ai/dsh-session'
import type { Message } from '@deepseek-ai/dsh-llm'
@@ -32,10 +33,8 @@ export function createTransmissionLog(): TransmissionLog {
}
/**
* Append whatever header event makes the log reproduce this request's header.
* The first request from an instance always records a full `initial` or `resume`
* snapshot. Later requests record nothing when unchanged, a round-tripping
* delta when expressible, or a full `fallback` snapshot otherwise.
* Append the full header snapshot owed by this request: initial/resume for the
* instance's first request, nothing when unchanged, or change otherwise.
*
* @param session - the session whose log explains the request.
* @param state - this loop instance's bookkeeping (mutated on first log).
@@ -52,12 +51,5 @@ export function recordRequestHeader(session: Session, state: TransmissionLog, he
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const baseline = session.requestHeader()!
if (headerEquals(baseline, header)) return
const delta = diffHeader(baseline, header)
/* v8 ignore next -- headerEquals false ⟹ diffHeader defined: both compare the same four parts */
if (delta === undefined) return
if (headerEquals(applyHeaderDelta(baseline, delta), header)) {
session.append('request/header-delta', delta)
} else {
session.append('request/header', { header, reason: 'fallback' })
}
session.append('request/header', { header, reason: 'change' })
}
+13 -14
View File
@@ -1,8 +1,8 @@
/**
* 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 context remain model-ordered.
* Abort stops replenishment and drains started calls.
* Dispatch may overlap, while policy, results, and result context remain
* model-ordered. Abort stops replenishment and drains started calls.
*
* Each started call records `tool/call`; `tool/result` commits in model order,
* preserving derived history when audit events interleave with earlier results.
@@ -31,8 +31,8 @@ interface Slot {
/**
* Schedule one assistant step's tool calls by their live concurrency mode.
* Started calls receive ordered results; abort drains them, discards their
* buffered context, and rethrows so the turn owns final error handling.
* Started calls receive ordered results. Abort drains them and rethrows after
* accepting their context into the batch FIFO owned by the caller.
*
* @param ctx - loop context that owns the tool registry.
* @param agent - agent and session receiving the call lifecycle.
@@ -41,7 +41,7 @@ interface Slot {
* @param toolCalls - assistant calls in model order.
* @param signal - abort signal shared by the step.
* @param maxParallel - validated in-flight cap.
* @returns buffered contexts in model call order.
* @param acceptContext - accepts committed result context into the active batch.
*/
export async function executeToolCalls(
ctx: Context,
@@ -51,7 +51,8 @@ export async function executeToolCalls(
toolCalls: ToolCallBlock[],
signal: AbortSignal,
maxParallel: number,
): Promise<HookContext[]> {
acceptContext: (context: HookContext) => void,
): Promise<void> {
const { session } = agent
// Inputs are distinct because tools/execute wrappers may replace `exec.signal`.
@@ -66,7 +67,6 @@ export async function executeToolCalls(
},
}))
const pendingContext: HookContext[] = []
let next = 0
while (next < planned.length) {
// Commit before classifying again so registry changes affect unstarted calls.
@@ -74,9 +74,8 @@ export async function executeToolCalls(
const first = planned[next]!
const mode = ctx.tools.executionMode(first.exec).kind
const group = mode === 'parallel' ? planned.slice(next) : [first]
next += await runGroup(ctx, session, turn, step, group, mode, signal, maxParallel, pendingContext)
next += await runGroup(ctx, session, turn, step, group, mode, signal, maxParallel, acceptContext)
}
return pendingContext
}
/** Parse model arguments, preserving invalid JSON as text and mapping empty input to `{}`. */
@@ -92,8 +91,8 @@ function parseArguments(raw: string): unknown {
* Run one exclusive barrier or parallel pool. Later calls are reclassified
* before start; an exclusive reclassification waits for the current pool to
* drain and remains for the caller's next barrier. Results and contexts commit
* in model order. Abort stops starts, drains and commits started calls, discards
* their contexts, and throws.
* in model order. Abort stops starts, drains and commits started calls, accepts
* their contexts into the owning batch, and throws.
*/
async function runGroup(
ctx: Context,
@@ -104,7 +103,7 @@ async function runGroup(
mode: ToolExecutionMode['kind'],
signal: AbortSignal,
maxParallel: number,
pendingContext: HookContext[],
acceptContext: (context: HookContext) => void,
): Promise<number> {
/* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
if (signal.aborted) throw new Error(String(signal.reason ?? 'aborted'))
@@ -127,7 +126,7 @@ async function runGroup(
: ctx.tools[TOOL_REGISTRY_SCHEDULER].finish(slot.exec, slot.result)
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- bounded index
appendToolResult(session, turn, step, call!.block, result, callSeqs[committed]!)
pendingContext.push(...result.additionalContexts ?? [])
for (const context of result.additionalContexts ?? []) acceptContext(context)
committed++
}
}
@@ -189,7 +188,7 @@ async function runGroup(
}
if (aborted) {
// Started calls are committed; their context is discarded with the aborted step.
// Started calls and accepted context settle before the turn records the abort.
/* v8 ignore next -- signal.reason always set: cancel()/disposal provide a default */
throw new Error(String(signal.reason ?? 'aborted'))
}