返回源码地图

packages/core/agent-loop/src/agent.ts

main snapshot · da00f7f5358f · 正文引用章节 02 / 05 / 06 / 06 / 07 / 07 / 07 / 07 / 07 / 07 / 07 / 07 / 07;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Default Agent driver over queued turns and step-boundary input. Every request
3 * is derived from the session log.
4 * @module dsh-agent-loop/agent
5 */
6
7import type {
8 Agent,
9 AgentCancelCause,
10 AgentEventDispatch,
11 AgentOptions,
12 AgentStatus,
13 CancelOptions,
14 InboxTarget,
15 PreStepDecision,
16 RequestErrorAction,
17} from '@deepseek-ai/dsh-agent'
18import { agentEvents, assembleContextFor } from '@deepseek-ai/dsh-agent'
19import type { GenerateOptions, LlmCallConfig, Message, PreparedLlmCall } from '@deepseek-ai/dsh-llm'
20import {
21 LlmError,
22 createAssistantMessage,
23 createDeveloperMessage,
24 errorChain,
25 markAgentLoopRequest,
26} from '@deepseek-ai/dsh-llm'
27import { assertNever, deepFreeze } from '@deepseek-ai/dsh-util-values'
28import type { Scope } from '@deepseek-ai/dsh-scope'
29import { createScope } from '@deepseek-ai/dsh-scope'
30import type { EpochHeader, RequestContext, Session, SessionId, SessionSeq, TurnEndReason, UserMessage } from '@deepseek-ai/dsh-session'
31import { canonicalHeader, headerEquals, ToolCallRecovery } from '@deepseek-ai/dsh-session'
32import { joinContextSections, renderContextSections, renderPrompt } from '@deepseek-ai/dsh-system-prompt'
33import type { PromptAssembly } from '@deepseek-ai/dsh-system-prompt'
34import type {} from '@deepseek-ai/dsh-session-projection'
35import type { Context } from '@deepseek-ai/cordis'
36import { ReactLoopInbox } from './inbox.ts'
37import { RuntimeContextProjection } from './runtime-context.ts'
38import { AssistantStreamAttempt } from './assistant-stream.ts'
39import { SystemPromptProjection } from './runtime-context.ts'
40import { executeToolCalls } from './tool-calls.ts'
41
42type Phase =
43 | { kind: 'idle'; lastTurn: number }
44 | {
45 kind: 'maintenance'
46 abort: AbortController
47 lastTurn: number
48 wakeRequested: boolean
49 }
50 | { kind: 'running'; abort: AbortController; turn: number; step: number; wakeRequested: boolean }
51
52type StepEndReason = Extract<TurnEndReason, { kind: 'completed' | 'max-tokens' }>
53
54type PreparedStep =
55 | { kind: 'reject' }
56 | {
57 kind: 'enter'
58 messages: UserMessage[]
59 startsRequestSeries?: true
60 assembly: PromptAssembly
61 }
62
63/** Remove adapter-derived values before plugins propose the next request config. */
64function requestProposal(header: EpochHeader): LlmCallConfig {
65 if (header.adapterDefaults === undefined) return header.config
66 const proposal = { ...header.config }
67 if (header.adapterDefaults.reasoningEffort === true) delete proposal.reasoningEffort
68 if (header.adapterDefaults.maxTokens === true) delete proposal.maxTokens
69 return proposal
70}
71
72/**
73 * Read the cause `cancel()` passed when aborting a loop-owned signal, copying
74 * only the fields `turn/end` records. The live reason stays the caller's
75 * object, and Node's fetch assigns a `stack` onto it that `Session.append`
76 * would either log or reject as data JSON cannot hold.
77 * @param signal - a turn or maintenance signal this loop owns.
78 * @returns the copied cause, or undefined while the signal is still live.
79 */
80function abortedCancelCause(signal: AbortSignal): AgentCancelCause | undefined {
81 if (!signal.aborted) return undefined
82 // `cancel()` is the only aborter of the signals this loop owns.
83 const cause = signal.reason as AgentCancelCause
84 switch (cause.kind) {
85 case 'user':
86 case 'parent':
87 case 'disposed':
88 return { kind: cause.kind }
89 case 'hook':
90 return { kind: 'hook', reason: cause.reason }
91 /* v8 ignore next -- cancel accepts the closed AgentCancelCause union */
92 default:
93 return assertNever(cause)
94 }
95}
96
97/** Drives one session through turn and step boundaries. */
98export class ReactLoopAgent implements Agent {
99 readonly inbox: ReactLoopInbox
100 private phase: Phase
101 private activityDone: Promise<void> = Promise.resolve()
102
103 /** The agent-scoped registration boundary; the lifecycle owner unwinds it after the driver exits. */
104 readonly scope: Scope
105 readonly ctx: Context
106
107 /** Fused dispatcher, built once in the constructor so hot-path dispatches never allocate. */
108 private readonly dispatch: AgentEventDispatch
109
110 /** Whether this loop instance has appended its initial/resume request anchor. */
111 private requestHeaderLogged = false
112 /** Surface generation at attachment or the preceding built request. */
113 private requestSurfaceGeneration: number
114 private readonly runtimeContext: RuntimeContextProjection
115 /** Process-local revision of assistant frames for this attached Session. */
116 private assistantStreamRevision = 0
117 private assistantAttemptCounter = 0
118 private readonly systemPrompt: SystemPromptProjection
119 /** Identities fully frozen by this loop; weak references do not retain replaced history. */
120 private readonly frozenMessages = new WeakSet<Message>()
121
122 constructor(
123 private loopCtx: Context,
124 public readonly id: SessionId,
125 public readonly options: AgentOptions,
126 public readonly session: Session,
127 ) {
128 this.requestSurfaceGeneration = session.surface.contentGeneration
129 this.dispatch = agentEvents(loopCtx, this)
130 this.scope = createScope(loopCtx, this)
131 this.ctx = this.scope.ctx
132 this.inbox = new ReactLoopInbox(this.ctx.sessionProjections, session, this.dispatch)
133 /* v8 ignore next -- the loop registers its own turnBoundary unit, so the key is always present */
134 const lastTurn = this.loopCtx.sessionProjections.stateOf(session, 'turnBoundary')?.lastTurn ?? 0
135 this.phase = { kind: 'idle', lastTurn }
136 this.runtimeContext = new RuntimeContextProjection(this.ctx, session)
137 this.systemPrompt = new SystemPromptProjection(session)
138 }
139
140 get status(): AgentStatus {
141 return this.phase.kind === 'idle' || this.phase.kind === 'maintenance' ? 'idle' : 'running'
142 }
143
144 /** Commit a phase and publish its externally visible status transition. */
145 private setPhase(next: Phase): void {
146 const previousStatus = this.status
147 this.phase = next
148 const status = this.status
149 if (status !== previousStatus) {
150 this.dispatch.emit('agent/status', { status })
151 }
152 }
153
154 send(message: UserMessage, target: InboxTarget, wakeup: boolean): void {
155 // Waking input cannot join an aborted activity, so it starts the next turn.
156 // Captured before the insertion so a reentrant cancel from a splice observer cannot reclassify it.
157 const wakingAfterAbort = wakeup && this.phase.kind !== 'idle' && this.phase.abort.signal.aborted
158 const resolvedTarget = wakingAfterAbort ? 'next-turn' : target
159 this.inbox.splice(resolvedTarget, Infinity, 0, [message])
160 if (wakeup) this.wakeDriver(wakingAfterAbort)
161 }
162
163 followup(input: UserMessage): void {
164 this.send(input, 'next-turn', true)
165 }
166
167 steer(input: UserMessage): void {
168 this.send(input, 'next-step', true)
169 }
170
171 inject(input: UserMessage): void {
172 this.send(input, 'next-step', false)
173 }
174
175 cancel(cause: AgentCancelCause, options: CancelOptions = {}): void {
176 if (!options.keepInbox) {
177 this.inbox.clear()
178 if (this.phase.kind !== 'idle') this.phase.wakeRequested = false
179 }
180 if (this.phase.kind !== 'idle') this.phase.abort.abort(cause)
181 }
182
183 runMaintenance<T>(job: (signal: AbortSignal) => Promise<T>): Promise<T> {
184 if (this.phase.kind !== 'idle') throw new Error(`agent "${this.id}" already has active work`)
185 const done = Promise.withResolvers<void>()
186 const maintenance: Phase = {
187 kind: 'maintenance',
188 abort: new AbortController(),
189 lastTurn: this.phase.lastTurn,
190 wakeRequested: false,
191 }
192 this.setPhase(maintenance)
193 this.activityDone = done.promise
194 return (async () => {
195 try {
196 return await job(maintenance.abort.signal)
197 } finally {
198 this.setPhase({ kind: 'idle', lastTurn: maintenance.lastTurn })
199 const cause = abortedCancelCause(maintenance.abort.signal)
200 if (cause?.kind !== 'disposed' && maintenance.wakeRequested && this.inbox.hasPending) this.wakeDriver()
201 done.resolve()
202 }
203 })()
204 }
205
206 /**
207 * Start one driver, or latch its wake behind maintenance or an aborted
208 * activity. A wake sent while idle always opens its turn boundary, even
209 * when its message was cleared; only a latched replay is suppressed when
210 * the queue no longer holds the wake.
211 * @param wakeAfterAbort - the {@link send} classification, captured before
212 * the inbox insertion so a reentrant cancel cannot reclassify it.
213 */
214 private wakeDriver(wakeAfterAbort = false): void {
215 if (this.phase.kind !== 'idle') {
216 // Maintenance and aborted drivers cannot deliver the wake: latch it for
217 // replay at convergence. Live drivers claim queued work themselves;
218 // disposal never latches, so teardown waits on no model turn.
219 const reason = abortedCancelCause(this.phase.abort.signal)
220 if (reason?.kind !== 'disposed' && (this.phase.kind === 'maintenance' || wakeAfterAbort)) {
221 this.phase.wakeRequested = true
222 }
223 return
224 }
225 const driver = Promise.withResolvers<void>()
226 this.activityDone = driver.promise
227 this.setPhase({
228 kind: 'running',
229 abort: new AbortController(),
230 turn: this.phase.lastTurn,
231 step: 0,
232 wakeRequested: false,
233 })
234 this.loopCtx.agents.withInitiator(this, () => this.kick()).then(driver.resolve, driver.reject)
235 }
236
237 async whenIdle(): Promise<void> {
238 let activity: Promise<void>
239 do {
240 await (activity = this.activityDone)
241 } while (activity !== this.activityDone)
242 }
243
244 /** Report one failure at its live boundary, then preserve it for driver containment. */
245 private throwError(error: unknown): never {
246 const turn = this.phase.kind === 'running' ? this.phase.turn : this.phase.lastTurn
247 const step = this.phase.kind === 'running' ? this.phase.step : 0
248 this.dispatch.emit('agent/error', { turn, step, error })
249 throw error
250 }
251
252 private async kick(): Promise<void> {
253 try {
254 while (await this.turn()) {}
255 } catch (_error) {
256 // Reported failures and cancellation are contained at the driver boundary.
257 } finally {
258 /* v8 ignore next -- kick owns a running phase until this driver boundary */
259 if (this.phase.kind === 'running') {
260 const { turn, wakeRequested } = this.phase
261 this.setPhase({ kind: 'idle', lastTurn: turn })
262 if (wakeRequested && this.inbox.hasPending) this.wakeDriver()
263 }
264 }
265 }
266
267 private async preStep(target: InboxTarget, position: { turn: number; step: number }): Promise<PreparedStep> {
268 /* v8 ignore next -- private callers establish the running phase before proposing a step */
269 if (this.phase.kind !== 'running') throw new Error(`agent "${this.id}": pre-step outside running phase`)
270 const signal = this.phase.abort.signal
271 const claimed = this.inbox.claim(target, position.turn)
272 const assembly = await this.loopCtx.systemPrompt.assemble(assembleContextFor(this, signal))
273 signal.throwIfAborted()
274 const sections = renderContextSections(assembly)
275 const context = this.runtimeContext.project(joinContextSections(sections), sections)
276 const decision = await this.dispatch.waterfall(
277 'agent/pre-step', { messages: claimed, ...position, signal },
278 (): Promise<PreStepDecision> => Promise.resolve<PreStepDecision>({
279 kind: 'enter',
280 messages: context === undefined ? claimed : [...claimed, context],
281 }),
282 )
283 signal.throwIfAborted()
284 if (decision.kind === 'reject') return decision
285 return { ...decision, assembly }
286 }
287
288 /** Whether the assembled tool schemas differ from the logged request header's. */
289 private toolsChanged(tools: PromptAssembly['tools']): boolean {
290 const baseline = this.session.requestHeader()
291 if (baseline === undefined) return false
292 return !headerEquals(baseline, canonicalHeader({ ...baseline, tools: [...tools] }))
293 }
294
295 /** Open one turn before claiming its first proposed step. */
296 private async turn(): Promise<boolean> {
297 if (this.phase.kind !== 'running') {
298 this.throwError(new Error(`agent "${this.id}": turn without driver reservation`))
299 }
300 const phase = this.phase
301 const { signal } = phase.abort
302 signal.throwIfAborted()
303 const turn = phase.turn + 1
304 try {
305 this.session.append('turn/start', { turn })
306 } catch (error: unknown) {
307 this.throwError(error)
308 }
309 phase.turn = turn
310 let turnEnds: TurnEndReason | null = null
311 let target: InboxTarget = 'next-turn'
312 try {
313 while (true) {
314 signal.throwIfAborted()
315 const step = phase.step + 1
316 const decision = await this.preStep(target, { turn, step })
317 if (decision.kind === 'reject') {
318 turnEnds = { kind: 'blocked' }
319 return false
320 }
321 if (turnEnds && decision.messages.length === 0) break
322 // A removed waking message or an enter decision rewritten to empty
323 // still owns the initial turn boundary, but it spends no model call.
324 if (phase.step === 0 && decision.messages.length === 0) {
325 turnEnds = { kind: 'completed' }
326 return false
327 }
328 signal.throwIfAborted()
329 this.session.append('step/start', { turn, step })
330 phase.step = step
331 const toolRecovery = new ToolCallRecovery()
332 const stopRecovery = this.ctx.on('session/event', (session, event) => {
333 if (session === this.session) toolRecovery.observe(event)
334 })
335 try {
336 // max-tokens is sticky: once any step hits the ceiling, later steps
337 // that complete normally must not downgrade the turn outcome.
338 const stepEnd = await this.step(decision)
339 // max-tokens stays sticky: a later completed step must not
340 // downgrade the turn outcome.
341 if (turnEnds === null || turnEnds.kind !== 'max-tokens') turnEnds = stepEnd
342 } catch (error: unknown) {
343 try {
344 for (const event of toolRecovery.results()) {
345 this.session.append('tool/result', event.data, {
346 surfaceOp: 'append',
347 ...event.sourceEventSeqs === undefined ? {} : { sourceEventSeqs: event.sourceEventSeqs },
348 })
349 }
350 } catch (recoveryError: unknown) {
351 throw new AggregateError([error, recoveryError], 'Step failed and its pending tool results could not be recorded', { cause: error })
352 }
353 throw error
354 } finally {
355 stopRecovery()
356 this.session.append('step/end', { turn, step })
357 }
358 signal.throwIfAborted()
359 if (turnEnds && this.inbox.nextStep.length === 0) {
360 await this.dispatch.serial('agent/turn-stopping', { turn, signal })
361 signal.throwIfAborted()
362 }
363 if (turnEnds && this.inbox.nextStep.length === 0) break
364 target = 'next-step'
365 }
366 } catch (error: unknown) {
367 // A cause is present exactly while the signal is aborted.
368 const cause = abortedCancelCause(signal)
369 if (cause !== undefined) {
370 turnEnds = { kind: 'aborted', reason: cause }
371 throw error
372 }
373 // Every failure is structured: an `LlmError` keeps its facts, anything
374 // else flattens to `errorChain` text under the `UNKNOWN` code.
375 turnEnds = {
376 kind: 'error',
377 error: error instanceof LlmError
378 ? error.failure
379 : { message: errorChain(error), code: 'UNKNOWN' },
380 }
381 this.throwError(error)
382 } finally {
383 try {
384 // oxlint-disable-next-line typescript/no-non-null-assertion -- every exit assigns a turn ending
385 this.session.append('turn/end', { turn, reason: turnEnds! })
386 } catch (error: unknown) {
387 this.throwError(error)
388 }
389 }
390 if (!this.inbox.hasPending) return false
391 phase.abort = new AbortController()
392 // A fresh controller makes a latch set on the old one stale: the live driver claims the queue itself.
393 phase.wakeRequested = false
394 phase.step = 0
395 return true
396 }
397
398 private async step(decision: Extract<PreparedStep, { kind: 'enter' }>): Promise<StepEndReason | null> {
399 /* v8 ignore next -- private callers establish the running phase before executing a step */
400 if (this.phase.kind !== 'running') throw new Error(`agent "${this.id}": step outside running phase`)
401 const { turn, step, abort: { signal } } = this.phase
402 signal.throwIfAborted()
403
404 const { assembly } = decision
405 const renderedPrompt = renderPrompt(assembly)
406 let firstAttempt = true
407 while (true) {
408 const { config, preparedCall } = await this.prepareRequest(turn, step, signal)
409 const startsRequestSeries = firstAttempt && decision.startsRequestSeries === true
410 const commits = this.systemPrompt.project(renderedPrompt, {
411 inHistory: preparedCall?.systemPromptUpdate === 'in-history',
412 startsSeries: startsRequestSeries
413 || this.requestSurfaceGeneration !== this.session.surface.contentGeneration
414 || (preparedCall?.toolUpdate === undefined && this.toolsChanged(assembly.tools)),
415 })
416 for (const { message, intent } of commits) {
417 this.session.append('system/message', { turn, step, message }, intent)
418 }
419 if (firstAttempt) {
420 for (const message of decision.messages) {
421 this.session.append('user/message', message, { surfaceOp: 'append' })
422 }
423 }
424 firstAttempt = false
425 const request = this.buildRequest(config, preparedCall, assembly.tools, { turn, step }, startsRequestSeries, signal)
426 const live = new AssistantStreamAttempt(
427 this.session.id,
428 ++this.assistantAttemptCounter,
429 () => ++this.assistantStreamRevision,
430 turn,
431 step,
432 (frame) => { this.dispatch.emit('agent/assistant-stream', { frame }) },
433 )
434 let started = false
435 try {
436 const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request)
437 signal.throwIfAborted()
438 live.start()
439 started = true
440 for await (const chunk of stream) {
441 signal.throwIfAborted()
442 live.push(chunk)
443 }
444 signal.throwIfAborted()
445 } catch (error: unknown) {
446 if (!started) throw error
447 try {
448 if (signal.aborted) {
449 const content = live.interruptedBlocks()
450 if (content.length > 0) {
451 live.settle('assistant/message', () => this.session.append('assistant/message', {
452 turn,
453 step,
454 message: createAssistantMessage({
455 content,
456 source: {
457 provider: request.provider,
458 model: request.model,
459 ...live.replayState === undefined ? {} : { replayState: live.replayState },
460 },
461 }),
462 interrupted: true,
463 ...live.usage === undefined ? {} : { usage: live.usage },
464 stream: live.stream,
465 }, { surfaceOp: 'append' }).seq)
466 } else {
467 live.settle(
468 'assistant/attempt',
469 () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq,
470 )
471 }
472 } else {
473 live.settle(
474 'assistant/attempt',
475 () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq,
476 )
477 }
478 } catch (settlementError: unknown) {
479 throw new AggregateError(
480 [error, settlementError],
481 'Assistant stream failed and its durable settlement was rejected',
482 { cause: error },
483 )
484 }
485 throw error
486 }
487 try {
488 const finish = live.finish
489 if (finish.kind === 'error' || finish.kind === 'aborted') {
490 live.settle(
491 'assistant/attempt',
492 () => this.session.append('assistant/attempt', { turn, step, stream: live.stream }).seq,
493 )
494 const action = await this.dispatch.waterfall(
495 'agent/request-error', {
496 turn,
497 step,
498 provider: request.provider,
499 failure: finish.failure,
500 retryPolicy: preparedCall?.retryPolicy,
501 signal,
502 },
503 () => Promise.resolve<RequestErrorAction>(undefined),
504 )
505 signal.throwIfAborted()
506 if (action?.kind !== 'retry') {
507 throw new LlmError(finish.failure.message, finish.failure.code, finish.failure)
508 }
509 continue
510 }
511
512 const message = createAssistantMessage({
513 content: live.blocks(),
514 source: {
515 provider: request.provider,
516 model: request.model,
517 ...live.replayState !== undefined ? { replayState: live.replayState } : {},
518 },
519 })
520 live.settle(
521 'assistant/message',
522 () => this.session.append('assistant/message', {
523 turn,
524 step,
525 message,
526 ...live.usage === undefined ? {} : { usage: live.usage },
527 stream: live.stream,
528 }, { surfaceOp: 'append' }).seq,
529 )
530 if (finish.kind === 'max-tokens') return { kind: 'max-tokens' }
531
532 const toolCalls = message.content.filter(block => block.type === 'tool-call')
533 if (toolCalls.length === 0) return { kind: 'completed' }
534 const { concluded } = await executeToolCalls(
535 this.loopCtx, turn, step, toolCalls, signal,
536 context => this.inbox.splice('next-step', this.inbox.nextStep.length, 0, [context]),
537 )
538 return concluded ? { kind: 'completed' } : null
539 } catch (error: unknown) {
540 if (!live.ended) live.abandon()
541 throw error
542 }
543 }
544 }
545
546 /** Resolve request config and bind its adapter before admitting model-visible input. */
547 private async prepareRequest(
548 turn: number,
549 step: number,
550 signal: AbortSignal,
551 ): Promise<{ config: LlmCallConfig; preparedCall?: PreparedLlmCall }> {
552 const { session } = this
553
554 // A loop instance starts from its declared route, restoring only an explicit
555 // effort owned by that exact model. Later steps re-resolve marked defaults.
556 const persistedHeader = session.requestHeader()
557 const persistedConfig = persistedHeader?.config
558 const route = { provider: this.options.provider ?? '', model: this.options.model ?? '' }
559 const persistedReasoningEffort = persistedConfig?.provider === route.provider
560 && persistedConfig.model === route.model
561 && persistedHeader?.adapterDefaults?.reasoningEffort !== true
562 ? persistedConfig.reasoningEffort
563 : undefined
564 const reasoningEffort = this.options.reasoningEffort ?? persistedReasoningEffort
565 const maxTokens = this.options.maxTokens
566 const seedConfig = deepFreeze(structuredClone(
567 this.requestHeaderLogged
568 // oxlint-disable-next-line typescript/no-non-null-assertion -- the instance logged the header it now folds
569 ? requestProposal(persistedHeader!)
570 : {
571 ...route,
572 ...reasoningEffort === undefined ? {} : { reasoningEffort },
573 ...maxTokens === undefined ? {} : { maxTokens },
574 },
575 ))
576 const proposedConfig = await this.dispatch.waterfall(
577 'agent/request', { turn, step, signal },
578 () => Promise.resolve(seedConfig),
579 )
580 signal.throwIfAborted()
581 if (!proposedConfig.provider || !proposedConfig.model) {
582 throw new Error(`agent "${this.id}" has no provider/model: set AgentOptions.provider and AgentOptions.model or supply both via the agent/request waterfall`)
583 }
584 let config: LlmCallConfig
585 let preparedCall: PreparedLlmCall | undefined
586 try {
587 preparedCall = await this.loopCtx.llm.prepareCall(proposedConfig, signal)
588 config = preparedCall.config
589 } catch (error: unknown) {
590 // Middleware may serve an unregistered route; terminal dispatch still requires an adapter.
591 if (!(error instanceof LlmError) || error.code !== 'NO_ADAPTER') throw error
592 config = proposedConfig
593 }
594 signal.throwIfAborted()
595 return { config, ...preparedCall === undefined ? {} : { preparedCall } }
596 }
597
598 /** Log the resolved envelope and derive a frozen request from the admitted surface. */
599 private buildRequest(
600 config: LlmCallConfig,
601 preparedCall: PreparedLlmCall | undefined,
602 tools: GenerateOptions['tools'] & object,
603 position: { turn: number; step: number },
604 startsRequestSeries: boolean,
605 signal: AbortSignal,
606 ): GenerateOptions {
607 const { session } = this
608 const surfaceGeneration = session.surface.contentGeneration
609 const header = canonicalHeader({
610 config,
611 ...preparedCall === undefined ? {} : { adapterDefaults: preparedCall.adapterDefaults },
612 ...tools.length > 0 ? { tools } : {},
613 })
614 const baseline = this.session.requestHeader()
615 const startsSeries = startsRequestSeries
616 || this.requestSurfaceGeneration !== surfaceGeneration
617 let headerSeq: SessionSeq | undefined
618 if (!this.requestHeaderLogged) {
619 // Compaction during the first resumed pre-step must still mark a new series.
620 headerSeq = this.session.append('request/header', {
621 header,
622 reason: baseline === undefined ? 'initial' : 'resume',
623 ...startsSeries ? { startsSeries: true } : {},
624 }).seq
625 this.requestHeaderLogged = true
626 } else if (baseline === undefined || !headerEquals(baseline, header)) {
627 headerSeq = this.session.append('request/header', {
628 header,
629 reason: 'change',
630 ...startsSeries ? { startsSeries: true } : {},
631 }).seq
632 } else if (startsSeries) {
633 this.session.append('request/header', { header, reason: 'series' })
634 }
635 if (baseline !== undefined && headerSeq !== undefined) {
636 const previousNames = new Set(baseline.tools?.map(tool => tool.name))
637 const currentNames = new Set(tools.map(tool => tool.name))
638 const additions = tools.filter(tool => !previousNames.has(tool.name))
639 .map(tool => ({ type: 'tool-addition' as const, toolName: tool.name }))
640 const removals = (baseline.tools ?? []).filter(tool => !currentNames.has(tool.name))
641 .map(tool => ({ type: 'tool-removal' as const, toolName: tool.name }))
642 if (additions.length > 0 || removals.length > 0) {
643 session.append('developer/message', {
644 ...position,
645 message: createDeveloperMessage({ source: { kind: 'tool-registry' }, content: [...additions, ...removals] }),
646 ...additions.length > 0 ? { headerSeq } : {},
647 }, { surfaceOp: 'append' })
648 }
649 }
650 this.requestSurfaceGeneration = surfaceGeneration
651
652 const contextWindow = preparedCall?.context?.contextWindow
653 const systemPromptUpdate = preparedCall?.systemPromptUpdate
654 const requestContext: RequestContext = {
655 provider: config.provider,
656 model: config.model,
657 ...contextWindow === undefined ? {} : { contextWindow },
658 ...systemPromptUpdate === undefined ? {} : { systemPromptUpdate },
659 }
660 const previousContext = session.requestContext()
661 if (previousContext?.provider !== requestContext.provider
662 || previousContext.model !== requestContext.model
663 || previousContext.contextWindow !== requestContext.contextWindow
664 || previousContext.systemPromptUpdate !== requestContext.systemPromptUpdate) {
665 session.append('request/context', requestContext)
666 }
667 signal.throwIfAborted()
668
669 // canonicalHeader is shallow; append logs a detached snapshot, not these local values.
670 deepFreeze(header)
671 const boundaryMessages = session.deriveMessages()
672 for (const message of boundaryMessages) {
673 if (this.frozenMessages.has(message)) continue
674 deepFreeze(message)
675 this.frozenMessages.add(message)
676 }
677 Object.freeze(boundaryMessages)
678 const request = markAgentLoopRequest(Object.freeze({
679 ...header.config,
680 messages: boundaryMessages,
681 toolHistory: session.toolHistory(),
682 ...header.tools !== undefined ? { tools: header.tools } : {},
683 sessionId: this.session.id,
684 signal,
685 }))
686 return request
687 }
688}