1
/**2
* Same-session goal-round driver over public agent, session, and goal services.3
* @module @deepseek-ai/dsh-goal-round-driver4
*/6
import { isDeepStrictEqual } from 'node:util'7
import { FiberState } from '@deepseek-ai/cordis'8
import type { Context } from '@deepseek-ai/cordis'9
import type { Agent, PreStepDecision } from '@deepseek-ai/dsh-agent'10
import type { GoalMessageSource, GoalRef, GoalView } from '@deepseek-ai/dsh-goal'11
import { createUserMessage } from '@deepseek-ai/dsh-llm'12
import type { ContentBlock, MessageId, MessageSource } from '@deepseek-ai/dsh-llm'13
import type { Session, SessionEvent, UserMessage } from '@deepseek-ai/dsh-session'14
import { renderGoalRoundPrompt } from './prompt.ts'16
export { renderGoalRoundPrompt } from './prompt.ts'18
export const name = 'goal-round-driver'19
export const inject = ['agents', 'goals', 'sessions']21
/** Identity reserved before a goal continuation enters the agent inbox. */22
interface RoundIdentity {23
readonly goalId: GoalRef['id']24
readonly revision: number25
readonly round: number26
}28
/** One queued, claimed, or admitted goal message retained until whole-agent quiescence. */29
interface RoundAttempt extends RoundIdentity {30
readonly messageId: MessageId31
readonly content: ContentBlock[]32
phase: 'queued' | 'claimed' | 'admitted'33
cancelled: boolean34
stale: boolean35
}37
/** Serialized process-local scheduling state for one exact Agent lifecycle. */38
interface DriverState {39
readonly agent: Agent40
attempt: RoundAttempt | undefined41
competingQueued: boolean42
needsCheckpoint: boolean43
requested: boolean44
run: Promise<void> | undefined45
stopping: boolean46
}48
/** Whether a source identifies an automatic, positive-numbered goal round. */49
function isGoalRoundSource(source: MessageSource): source is GoalMessageSource {50
return source.kind === 'goal' && source.round > 051
}53
/** Compare a source to one reserved identity. */54
function sameRound(source: GoalMessageSource, round: RoundIdentity): boolean {55
return source.goalId === round.goalId56
&& source.revision === round.revision57
&& source.round === round.round58
}60
/** Compare the complete queued record to the driver's reservation. */61
function sameQueued(content: readonly ContentBlock[], source: MessageSource, attempt: RoundAttempt): boolean {62
return isGoalRoundSource(source) && sameRound(source, attempt) && isDeepStrictEqual(content, attempt.content)63
}65
/** Exact current ref for a view. */66
function goalRef(goal: GoalView): GoalRef {67
return { id: goal.id, revision: goal.revision }68
}70
/** Human-readable unexpected values for logs. */71
function renderThrown(value: unknown): string {72
return value instanceof Error ? value.message : String(value)73
}75
/** Install automatic same-session continuation and its race fences. */76
export function apply(ctx: Context): void {77
const states = new Map<Agent, DriverState>()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 existing83
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 state94
}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 undefined99
return ctx.goals.get(state.agent)100
}102
/** Whether this exact lifecycle is quiescent with no competing prompt. */103
function readyToDrive(state: DriverState): boolean {104
return ctx.fiber.state === FiberState.ACTIVE105
&& !state.stopping106
&& ctx.agents.get(state.agent.id) === state.agent107
&& state.agent.status === 'idle'108
&& !state.competingQueued109
}111
/** Recheck every condition that an awaited checkpoint may have changed. */112
function readyAfterCheckpoint(state: DriverState): boolean {113
return readyToDrive(state) && !state.needsCheckpoint114
}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
}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 !== messageId129
&& !(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)) continue133
agent.inbox.prepend('next-step', message)134
}135
}137
/** Process admitted work at quiescence, then reserve at most one next round. */138
async function drive(state: DriverState): Promise<void> {139
const { agent } = state140
if (!readyToDrive(state)) return142
if (state.needsCheckpoint) {143
state.needsCheckpoint = false144
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
return150
}151
// A mutation or ordinary prompt may have arrived while the checkpoint152
// was settling. Give it its own checkpoint / turn before reserving.153
if (!readyAfterCheckpoint(state)) return154
}156
const attempt = state.attempt157
if (attempt !== undefined) {158
state.attempt = undefined159
state.needsCheckpoint = true160
state.requested = true161
return162
}164
const goal = currentGoal(state)165
if (goal === undefined || goal.phase !== 'active' || goal.activation !== 'armed') return166
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
return172
}174
const round = goal.roundsStarted + 1175
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 = reservation191
try {192
agent.followup(message)193
} catch (error: unknown) {194
state.attempt = undefined195
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.revision198
&& 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
}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) return211
state.requested = true212
if (state.run !== undefined) return213
let run: Promise<void>214
try {215
run = ctx.agents.withoutInitiator(async () => {216
while (state.requested && !state.stopping) {217
state.requested = false218
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
return230
}231
state.run = run232
const retire = (): void => {233
state.run = undefined234
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
}243
// One composite effect keeps the step fence installed until this244
// 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
})251
ctx.on('agent/disposed', ({ agent }) => { states.delete(agent) })252
ctx.on('agent/created', ({ agent }) => {253
const state = stateFor(agent)254
state.attempt = undefined255
state.competingQueued = false256
state.needsCheckpoint = false257
})258
ctx.on('agent/status', ({ agent, status }) => {259
const state = stateFor(agent)260
if (status === 'idle') {261
state.competingQueued = false262
const attempt = state.attempt263
const goal = currentGoal(state)264
// Fence the pause to the exact dropped attempt's ref. A resume bumps265
// the revision, so a host pause followed by an immediate resume (before266
// the aborted turn converges to idle) must not re-pause the resumed goal.267
const pause = attempt !== undefined268
&& (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.revision271
// 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 = undefined275
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 = true289
// A host-initiated pause stops goal execution: abort the live turn so the290
// model cannot keep acting or resume in the same turn. A model-initiated291
// 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
})299
ctx.on('agent/inbox/inserted', ({ agent, message }) => {300
if (!agent.inbox.nextTurn.some(candidate => candidate.id === message.id)) return301
const state = stateFor(agent)302
const attempt = state.attempt303
if (attempt !== undefined && sameQueued(message.content, message.source, attempt)) return304
state.competingQueued = true305
if (attempt?.phase === 'queued') attempt.stale = true306
})307
ctx.on('agent/inbox/claimed', ({ agent, message }) => {308
const state = stateFor(agent)309
const attempt = state.attempt310
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.attempt317
if (attempt !== undefined && sameQueued(message.content, message.source, attempt)) {318
attempt.cancelled = true319
}320
})322
ctx.on('session/event', (session: Session, event: SessionEvent) => {323
const agent = ctx.agents.get(session.id)324
if (agent === undefined || agent.session !== session) return325
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
return332
case 'turn/end':333
if (event.data.reason.kind === 'max-tokens') {334
disarm(state)335
return336
}337
if (event.data.reason.kind !== 'aborted') return338
if (state.attempt?.phase === 'claimed' || state.attempt?.phase === 'admitted') {339
state.attempt.cancelled = true340
}341
else disarm(state)342
return343
default:344
return345
}346
})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.attempt355
const goal = currentGoal(state)356
return ctx.fiber.state === FiberState.ACTIVE357
&& !state.stopping && attempt !== undefined && attempt.phase === 'claimed'358
&& !attempt.stale && sameQueued(content, source, attempt)359
&& goal !== undefined && goal.id === source.goalId && goal.revision === source.revision360
&& goal.phase === 'active' && goal.activation === 'armed'361
&& source.round === goal.roundsStarted + 1362
}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 } = submitted369
const state = stateFor(agent)370
let valid = false371
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.attempt379
if (attempt !== undefined && sameRound(source, attempt)) {380
attempt.stale = true381
state.attempt = undefined382
}383
restoreOtherClaimed(agent, messages, submitted.id)384
requestDrive(state)385
return { kind: 'reject' }386
}387
let decision: PreStepDecision388
try {389
decision = await next()390
} catch (error: unknown) {391
if (signal.aborted) throw error392
// A throwing downstream hook drops the whole step proposal. Clear the393
// reservation before the balanced no-step turn returns to idle so the394
// next drive pass can reschedule the round.395
state.attempt = undefined396
requestDrive(state)397
throw error398
}399
if (signal.aborted) {400
if (decision.kind === 'enter') restoreOtherClaimed(agent, decision.messages, submitted.id)401
return decision402
}403
if (decision.kind === 'reject') {404
state.attempt = undefined405
const goal = currentGoal(state)406
if (goal !== undefined && goal.id === source.goalId && goal.revision === source.revision407
&& 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 decision414
}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 = false421
}422
if (!valid) {423
state.attempt = undefined424
restoreOtherClaimed(agent, decision.messages, submitted.id)425
requestDrive(state)426
return { kind: 'reject' }427
}428
return { ...decision, startsRequestSeries: true }429
})431
// Loading a lifecycle driver over existing agents never inherits hidden432
// automatic authority from an earlier producer instance.433
for (const agent of ctx.agents.list()) {434
const state = stateFor(agent)435
disarm(state)436
}438
// Yielded after listener registration, so this close runs first and the439
// 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 = true444
disarm(state)445
const attempt = state.attempt446
if (attempt !== undefined) {447
attempt.stale = true448
/* 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
}