agent.ts 24 KB

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