返回源码地图

packages/goal/goal-round-driver/src/index.ts

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

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

1/**
2 * Same-session goal-round driver over public agent, session, and goal services.
3 * @module @deepseek-ai/dsh-goal-round-driver
4 */
5
6import { isDeepStrictEqual } from 'node:util'
7import { FiberState } from '@deepseek-ai/cordis'
8import type { Context } from '@deepseek-ai/cordis'
9import type { Agent, PreStepDecision } from '@deepseek-ai/dsh-agent'
10import type { GoalMessageSource, GoalRef, GoalView } from '@deepseek-ai/dsh-goal'
11import { createUserMessage } from '@deepseek-ai/dsh-llm'
12import type { ContentBlock, MessageId, MessageSource } from '@deepseek-ai/dsh-llm'
13import type { Session, SessionEvent, UserMessage } from '@deepseek-ai/dsh-session'
14import { renderGoalRoundPrompt } from './prompt.ts'
15
16export { renderGoalRoundPrompt } from './prompt.ts'
17
18export const name = 'goal-round-driver'
19export const inject = ['agents', 'goals', 'sessions']
20
21/** Identity reserved before a goal continuation enters the agent inbox. */
22interface RoundIdentity {
23 readonly goalId: GoalRef['id']
24 readonly revision: number
25 readonly round: number
26}
27
28/** One queued, claimed, or admitted goal message retained until whole-agent quiescence. */
29interface RoundAttempt extends RoundIdentity {
30 readonly messageId: MessageId
31 readonly content: ContentBlock[]
32 phase: 'queued' | 'claimed' | 'admitted'
33 cancelled: boolean
34 stale: boolean
35}
36
37/** Serialized process-local scheduling state for one exact Agent lifecycle. */
38interface DriverState {
39 readonly agent: Agent
40 attempt: RoundAttempt | undefined
41 competingQueued: boolean
42 needsCheckpoint: boolean
43 requested: boolean
44 run: Promise<void> | undefined
45 stopping: boolean
46}
47
48/** Whether a source identifies an automatic, positive-numbered goal round. */
49function isGoalRoundSource(source: MessageSource): source is GoalMessageSource {
50 return source.kind === 'goal' && source.round > 0
51}
52
53/** Compare a source to one reserved identity. */
54function sameRound(source: GoalMessageSource, round: RoundIdentity): boolean {
55 return source.goalId === round.goalId
56 && source.revision === round.revision
57 && source.round === round.round
58}
59
60/** Compare the complete queued record to the driver's reservation. */
61function sameQueued(content: readonly ContentBlock[], source: MessageSource, attempt: RoundAttempt): boolean {
62 return isGoalRoundSource(source) && sameRound(source, attempt) && isDeepStrictEqual(content, attempt.content)
63}
64
65/** Exact current ref for a view. */
66function goalRef(goal: GoalView): GoalRef {
67 return { id: goal.id, revision: goal.revision }
68}
69
70/** Human-readable unexpected values for logs. */
71function renderThrown(value: unknown): string {
72 return value instanceof Error ? value.message : String(value)
73}
74
75/** Install automatic same-session continuation and its race fences. */
76export function apply(ctx: Context): void {
77 const states = new Map<Agent, DriverState>()
78
79 /** Create state for an exact currently live agent. */
80 function stateFor(agent: Agent): DriverState {
81 const existing = states.get(agent)
82 if (existing !== undefined) return existing
83 const state: DriverState = {
84 agent,
85 attempt: undefined,
86 competingQueued: false,
87 needsCheckpoint: false,
88 requested: false,
89 run: undefined,
90 stopping: false,
91 }
92 states.set(agent, state)
93 return state
94 }
95
96 /** Read only when the exact Agent remains live. */
97 function currentGoal(state: DriverState): GoalView | undefined {
98 if (ctx.agents.get(state.agent.id) !== state.agent) return undefined
99 return ctx.goals.get(state.agent)
100 }
101
102 /** Whether this exact lifecycle is quiescent with no competing prompt. */
103 function readyToDrive(state: DriverState): boolean {
104 return ctx.fiber.state === FiberState.ACTIVE
105 && !state.stopping
106 && ctx.agents.get(state.agent.id) === state.agent
107 && state.agent.status === 'idle'
108 && !state.competingQueued
109 }
110
111 /** Recheck every condition that an awaited checkpoint may have changed. */
112 function readyAfterCheckpoint(state: DriverState): boolean {
113 return readyToDrive(state) && !state.needsCheckpoint
114 }
115
116 /** Remove automatic authority while preserving the durable phase. */
117 function disarm(state: DriverState): void {
118 try {
119 const goal = currentGoal(state)
120 if (goal?.activation === 'armed') ctx.goals.disarm(state.agent)
121 } catch (error: unknown) {
122 ctx.logger.warn(`goal-round-driver: could not disarm agent "${state.agent.id}": ${renderThrown(error)}`)
123 }
124 }
125
126 /** Preserve claimed step context when this driver drops only its own round. */
127 function restoreOtherClaimed(agent: Agent, messages: UserMessage[], messageId: MessageId): void {
128 const retained = messages.filter(message => message.id !== messageId
129 && !(message.source.kind === 'goal' && message.source.round === 0))
130 for (const message of retained.toReversed()) {
131 if (agent.inbox.nextStep.some(candidate => candidate.id === message.id)
132 || agent.inbox.nextTurn.some(candidate => candidate.id === message.id)) continue
133 agent.inbox.prepend('next-step', message)
134 }
135 }
136
137 /** Process admitted work at quiescence, then reserve at most one next round. */
138 async function drive(state: DriverState): Promise<void> {
139 const { agent } = state
140 if (!readyToDrive(state)) return
141
142 if (state.needsCheckpoint) {
143 state.needsCheckpoint = false
144 try {
145 await ctx.sessions.flush(agent.session)
146 } catch (error: unknown) {
147 ctx.logger.warn(`goal-round-driver: durability checkpoint failed for agent "${agent.id}": ${renderThrown(error)}`)
148 disarm(state)
149 return
150 }
151 // A mutation or ordinary prompt may have arrived while the checkpoint
152 // was settling. Give it its own checkpoint / turn before reserving.
153 if (!readyAfterCheckpoint(state)) return
154 }
155
156 const attempt = state.attempt
157 if (attempt !== undefined) {
158 state.attempt = undefined
159 state.needsCheckpoint = true
160 state.requested = true
161 return
162 }
163
164 const goal = currentGoal(state)
165 if (goal === undefined || goal.phase !== 'active' || goal.activation !== 'armed') return
166 if (goal.roundsStarted >= goal.maxGoalRounds) {
167 ctx.goals.block(agent, goalRef(goal), {
168 code: 'round-limit',
169 message: `Goal reached its configured limit of ${goal.maxGoalRounds} rounds.`,
170 })
171 return
172 }
173
174 const round = goal.roundsStarted + 1
175 const content = renderGoalRoundPrompt(goal, round)
176 const message = createUserMessage({
177 content,
178 source: { kind: 'goal', goalId: goal.id, revision: goal.revision, round },
179 })
180 const reservation: RoundAttempt = {
181 goalId: goal.id,
182 revision: goal.revision,
183 round,
184 messageId: message.id,
185 content,
186 phase: 'queued',
187 cancelled: false,
188 stale: false,
189 }
190 state.attempt = reservation
191 try {
192 agent.followup(message)
193 } catch (error: unknown) {
194 state.attempt = undefined
195 ctx.logger.warn(`goal-round-driver: could not queue round ${round} for agent "${agent.id}": ${renderThrown(error)}`)
196 const latest = currentGoal(state)
197 if (latest !== undefined && latest.id === goal.id && latest.revision === goal.revision
198 && latest.phase === 'active' && latest.activation === 'armed') {
199 ctx.goals.block(agent, goalRef(latest), {
200 code: 'queue-failed',
201 message: `Could not queue goal round ${round}: ${renderThrown(error)}`,
202 })
203 }
204 }
205 }
206
207 /** Coalesce triggers onto one agent-local serialized driver. */
208 function requestDrive(state: DriverState): void {
209 /* v8 ignore next -- teardown may race a final trigger after synchronously closing the step fence */
210 if (state.stopping) return
211 state.requested = true
212 if (state.run !== undefined) return
213 let run: Promise<void>
214 try {
215 run = ctx.agents.withoutInitiator(async () => {
216 while (state.requested && !state.stopping) {
217 state.requested = false
218 try {
219 await drive(state)
220 } catch (error: unknown) {
221 ctx.logger.warn(`goal-round-driver: driver failed for agent "${state.agent.id}": ${renderThrown(error)}`)
222 disarm(state)
223 }
224 }
225 })
226 } catch (error: unknown) {
227 ctx.logger.warn(`goal-round-driver: could not start driver for agent "${state.agent.id}": ${renderThrown(error)}`)
228 disarm(state)
229 return
230 }
231 state.run = run
232 const retire = (): void => {
233 state.run = undefined
234 if (state.requested && !state.stopping) requestDrive(state)
235 }
236 void run.then(retire, (error: unknown) => {
237 ctx.logger.warn(`goal-round-driver: driver task rejected for agent "${state.agent.id}": ${renderThrown(error)}`)
238 disarm(state)
239 retire()
240 })
241 }
242
243 // One composite effect keeps the step fence installed until this
244 // plugin's own scheduling tasks settle.
245 ctx.effect(function* () {
246 ctx.on('agent/error', ({ agent }) => {
247 const state = stateFor(agent)
248 disarm(state)
249 })
250
251 ctx.on('agent/disposed', ({ agent }) => { states.delete(agent) })
252 ctx.on('agent/created', ({ agent }) => {
253 const state = stateFor(agent)
254 state.attempt = undefined
255 state.competingQueued = false
256 state.needsCheckpoint = false
257 })
258 ctx.on('agent/status', ({ agent, status }) => {
259 const state = stateFor(agent)
260 if (status === 'idle') {
261 state.competingQueued = false
262 const attempt = state.attempt
263 const goal = currentGoal(state)
264 // Fence the pause to the exact dropped attempt's ref. A resume bumps
265 // the revision, so a host pause followed by an immediate resume (before
266 // the aborted turn converges to idle) must not re-pause the resumed goal.
267 const pause = attempt !== undefined
268 && (attempt.phase === 'queued' || attempt.phase === 'claimed' || attempt.cancelled)
269 && goal !== undefined && goal.phase === 'active' && goal.activation === 'armed'
270 && attempt.goalId === goal.id && attempt.revision === goal.revision
271 // A reservation still queued when the agent reaches idle cannot run:
272 // withdraw it so human input queued behind it is not stranded.
273 if (pause || attempt?.phase === 'queued') {
274 state.attempt = undefined
275 try {
276 if (attempt.phase === 'queued') agent.inbox.remove(attempt.messageId)
277 if (pause) ctx.goals.pause(agent, goalRef(goal))
278 } catch (error: unknown) {
279 ctx.logger.warn(`goal-round-driver: could not settle cancelled goal round for agent "${agent.id}": ${renderThrown(error)}`)
280 disarm(state)
281 }
282 }
283 requestDrive(state)
284 }
285 })
286 ctx.on('goal/changed', ({ agent, change }) => {
287 const state = stateFor(agent)
288 state.needsCheckpoint = true
289 // A host-initiated pause stops goal execution: abort the live turn so the
290 // model cannot keep acting or resume in the same turn. A model-initiated
291 // pause (update_goal inside its own turn) finishes normally.
292 if (change.operation === 'pause' && agent.status === 'running'
293 && ctx.agents.currentInitiator() !== agent) {
294 agent.cancel({ kind: 'user' }, { keepInbox: true })
295 }
296 requestDrive(state)
297 })
298
299 ctx.on('agent/inbox/inserted', ({ agent, message }) => {
300 if (!agent.inbox.nextTurn.some(candidate => candidate.id === message.id)) return
301 const state = stateFor(agent)
302 const attempt = state.attempt
303 if (attempt !== undefined && sameQueued(message.content, message.source, attempt)) return
304 state.competingQueued = true
305 if (attempt?.phase === 'queued') attempt.stale = true
306 })
307 ctx.on('agent/inbox/claimed', ({ agent, message }) => {
308 const state = stateFor(agent)
309 const attempt = state.attempt
310 if (attempt !== undefined && sameQueued(message.content, message.source, attempt)) {
311 attempt.phase = 'claimed'
312 }
313 })
314 ctx.on('agent/inbox/discarded', ({ agent, message }) => {
315 const state = stateFor(agent)
316 const attempt = state.attempt
317 if (attempt !== undefined && sameQueued(message.content, message.source, attempt)) {
318 attempt.cancelled = true
319 }
320 })
321
322 ctx.on('session/event', (session: Session, event: SessionEvent) => {
323 const agent = ctx.agents.get(session.id)
324 if (agent === undefined || agent.session !== session) return
325 const state = stateFor(agent)
326 switch (event.type) {
327 case 'user/message':
328 if (state.attempt !== undefined && event.data.id === state.attempt.messageId) {
329 state.attempt.phase = 'admitted'
330 }
331 return
332 case 'turn/end':
333 if (event.data.reason.kind === 'max-tokens') {
334 disarm(state)
335 return
336 }
337 if (event.data.reason.kind !== 'aborted') return
338 if (state.attempt?.phase === 'claimed' || state.attempt?.phase === 'admitted') {
339 state.attempt.cancelled = true
340 }
341 else disarm(state)
342 return
343 default:
344 return
345 }
346 })
347
348 /** Fail closed unless the queued prompt still owns the exact live revision. */
349 function validReservation(
350 state: DriverState,
351 content: readonly ContentBlock[],
352 source: GoalMessageSource,
353 ): boolean {
354 const attempt = state.attempt
355 const goal = currentGoal(state)
356 return ctx.fiber.state === FiberState.ACTIVE
357 && !state.stopping && attempt !== undefined && attempt.phase === 'claimed'
358 && !attempt.stale && sameQueued(content, source, attempt)
359 && goal !== undefined && goal.id === source.goalId && goal.revision === source.revision
360 && goal.phase === 'active' && goal.activation === 'armed'
361 && source.round === goal.roundsStarted + 1
362 }
363
364 ctx.on('agent/pre-step', async ({ agent, messages, signal }, next): Promise<PreStepDecision> => {
365 const submitted = messages.find((message): message is UserMessage & { source: GoalMessageSource } =>
366 isGoalRoundSource(message.source))
367 if (submitted === undefined) return next()
368 const { content, source } = submitted
369 const state = stateFor(agent)
370 let valid = false
371 try {
372 valid = validReservation(state, content, source)
373 } catch (error: unknown) {
374 ctx.logger.warn(`goal-round-driver: pre-step check failed for agent "${agent.id}": ${renderThrown(error)}`)
375 disarm(state)
376 }
377 if (!valid) {
378 const attempt = state.attempt
379 if (attempt !== undefined && sameRound(source, attempt)) {
380 attempt.stale = true
381 state.attempt = undefined
382 }
383 restoreOtherClaimed(agent, messages, submitted.id)
384 requestDrive(state)
385 return { kind: 'reject' }
386 }
387 let decision: PreStepDecision
388 try {
389 decision = await next()
390 } catch (error: unknown) {
391 if (signal.aborted) throw error
392 // A throwing downstream hook drops the whole step proposal. Clear the
393 // reservation before the balanced no-step turn returns to idle so the
394 // next drive pass can reschedule the round.
395 state.attempt = undefined
396 requestDrive(state)
397 throw error
398 }
399 if (signal.aborted) {
400 if (decision.kind === 'enter') restoreOtherClaimed(agent, decision.messages, submitted.id)
401 return decision
402 }
403 if (decision.kind === 'reject') {
404 state.attempt = undefined
405 const goal = currentGoal(state)
406 if (goal !== undefined && goal.id === source.goalId && goal.revision === source.revision
407 && goal.phase === 'active' && goal.activation === 'armed') {
408 ctx.goals.block(agent, goalRef(goal), {
409 code: 'prompt-rejected',
410 message: 'Goal round was rejected before entering its step.',
411 })
412 }
413 return decision
414 }
415 try {
416 valid = validReservation(state, content, source)
417 } catch (error: unknown) {
418 ctx.logger.warn(`goal-round-driver: post-decision check failed for agent "${agent.id}": ${renderThrown(error)}`)
419 disarm(state)
420 valid = false
421 }
422 if (!valid) {
423 state.attempt = undefined
424 restoreOtherClaimed(agent, decision.messages, submitted.id)
425 requestDrive(state)
426 return { kind: 'reject' }
427 }
428 return { ...decision, startsRequestSeries: true }
429 })
430
431 // Loading a lifecycle driver over existing agents never inherits hidden
432 // automatic authority from an earlier producer instance.
433 for (const agent of ctx.agents.list()) {
434 const state = stateFor(agent)
435 disarm(state)
436 }
437
438 // Yielded after listener registration, so this close runs first and the
439 // composite effect removes listeners only after its promise settles.
440 yield async () => {
441 const waits: Promise<void>[] = []
442 for (const state of states.values()) {
443 state.stopping = true
444 disarm(state)
445 const attempt = state.attempt
446 if (attempt !== undefined) {
447 attempt.stale = true
448 /* v8 ignore next -- followup reserves the live agent before publishing a queued attempt */
449 if (state.agent.status === 'running') {
450 state.agent.cancel({ kind: 'parent' })
451 waits.push(state.agent.whenIdle())
452 }
453 }
454 if (state.run !== undefined) waits.push(state.run)
455 }
456 await Promise.allSettled(waits)
457 states.clear()
458 }
459 }, 'goal-round-driver lifecycle')
460}