1
/**2
* Concrete agent-loop plugin: creates scoped ReactLoopAgents, publishes them3
* through the agent/session registries, and owns their ordered teardown.4
*5
* @module @deepseek-ai/dsh-agent-loop6
*/7
import type { Volatile } from '@deepseek-ai/cosmokit'9
import { Context, FiberState, Service } from '@deepseek-ai/cordis'10
import { randomUUID } from 'node:crypto'11
import z from '@deepseek-ai/schemastery'12
import { z as zod } from 'zod'13
import { brandString } from '@deepseek-ai/dsh-brand'14
import type {15
Agent,16
AgentFactory,17
AgentHandle,18
AgentOptions,19
AgentSetup,20
CreateAgentOptions,21
ResumeAgentOptions,22
SessionStartSource,23
TurnBoundaryProjection,24
} from '@deepseek-ai/dsh-agent'25
import { errorChain, ReasoningEffortId } from '@deepseek-ai/dsh-llm'26
import { interruptedTurnClosers, SessionLogOffset, SessionPreparation, SessionSeq } from '@deepseek-ai/dsh-session'27
import type { Session, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'28
import type {} from '@deepseek-ai/dsh-system-prompt'29
import type {} from '@deepseek-ai/dsh-tools'30
import type {} from '@deepseek-ai/dsh-session-projection'31
import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'32
import { SessionPersistenceNotFoundError } from '@deepseek-ai/dsh-session-persistence'33
import type { SessionHandle, SessionPersistence } from '@deepseek-ai/dsh-session-persistence'34
import { ReactLoopAgent } from './agent.ts'35
import { inboxProjectionDefinition } from './inbox.ts'36
import { DEFAULT_MAX_PARALLEL_TOOL_CALLS } from './constants.ts'37
import type {} from './runtime-context.ts'39
/** Fiber states that cannot own or serve a new lifecycle. */40
const INACTIVE_STATES: ReadonlySet<FiberState> = new Set([41
FiberState.UNLOADING,42
FiberState.DISPOSED,43
FiberState.FAILED,44
])46
const turnBoundaryProjectionSchema: zod.ZodType<TurnBoundaryProjection> = zod.object({47
openTurnStartSeq: zod.number().int().nonnegative().transform(SessionSeq).nullable(),48
lastStepStartSeq: zod.number().int().nonnegative().transform(SessionSeq).nullable(),49
lastStepBoundary: zod.object({50
kind: zod.union([zod.literal('start'), zod.literal('end')]),51
seq: zod.number().int().nonnegative().transform(SessionSeq),52
}).nullable(),53
lastTurn: zod.number().int().nonnegative(),54
})56
/** Host projection of agent turn and step boundaries. */57
export const turnBoundaryProjectionDefinition = {58
key: 'turnBoundary',59
stateVersion: 2,60
stateSchema: turnBoundaryProjectionSchema,61
init: () => ({62
openTurnStartSeq: null,63
lastStepStartSeq: null,64
lastStepBoundary: null,65
lastTurn: 0,66
}),67
apply: (state, event) => {68
switch (event.type) {69
case 'turn/start':70
return {71
...state,72
openTurnStartSeq: event.seq,73
lastTurn: event.data.turn,74
}75
case 'turn/end':76
return {77
...state,78
openTurnStartSeq: null,79
}80
case 'step/start':81
return {82
...state,83
lastStepStartSeq: event.seq,84
lastStepBoundary: { kind: 'start', seq: event.seq },85
}86
case 'step/end':87
return {88
...state,89
lastStepBoundary: { kind: 'end', seq: event.seq },90
}91
default:92
return state93
}94
},95
} satisfies ProjectionDefinition<'turnBoundary', TurnBoundaryProjection>97
/** Factory-level ownership: live agent teardowns plus config startup work. */98
class FactoryOwnership {99
private accepting = true100
private readonly teardown = new AbortController()101
private readonly inactive = Promise.withResolvers<void>()102
private readonly liveAgents = new Set<() => Promise<void>>()103
private startupTasks = new Set<Promise<void>>()105
constructor(private readonly fiber: Context['fiber']) {}107
/** Aborts (reason: `agent loop is not active` error) when factory teardown begins. */108
get signal(): AbortSignal {109
return this.teardown.signal110
}112
isActive(): boolean {113
return this.accepting && !INACTIVE_STATES.has(this.fiber.state)114
}116
/** Track one live agent's shared teardown until it has run. */117
track(dispose: () => Promise<void>): () => void {118
this.liveAgents.add(dispose)119
return () => { this.liveAgents.delete(dispose) }120
}122
/** Join config startup work that begins before an agent exists. */123
trackStartup(job: Promise<void>): void {124
this.startupTasks.add(job)125
const forget = () => { this.startupTasks.delete(job) }126
void job.then(forget, forget)127
}129
/** Join one public create/resume continuation; factory dispose awaits its settlement. */130
trackWrapper(job: Promise<unknown>): void {131
this.trackStartup(job.then(() => undefined, () => undefined))132
}134
/** Resolve `task`, or stop waiting when factory teardown begins. */135
async waitWhileActive(job: Promise<void>): Promise<void> {136
await Promise.race([job, this.inactive.promise])137
}139
async dispose(): Promise<void> {140
this.accepting = false141
this.teardown.abort(new Error('agent loop is not active'))142
this.inactive.resolve()143
await Promise.all([144
...[...this.liveAgents].map(dispose => dispose()),145
...this.startupTasks,146
])147
}148
}150
/** Await `operation`, or throw the signal's reason as soon as it aborts. */151
async function raceAbort<T>(operation: PromiseLike<T> | T, signal: AbortSignal, id: SessionId): Promise<T> {152
const toAbortError = (): Error => signal.reason instanceof Error153
? signal.reason154
: new Error(`agent "${id}" creation aborted`, { cause: signal.reason })155
if (signal.aborted) throw toAbortError()156
const aborted = Promise.withResolvers<never>()157
const listener = (): void => { aborted.reject(toAbortError()) }158
signal.addEventListener('abort', listener, { once: true })159
try {160
return await Promise.race([Promise.resolve(operation), aborted.promise])161
} finally {162
signal.removeEventListener('abort', listener)163
}164
}166
/** Start an abortable operation and release a value that arrives after cancellation. */167
async function raceAbortCall<T>(168
operation: () => PromiseLike<T> | T,169
signal: AbortSignal,170
id: SessionId,171
releaseAbandoned?: (value: T) => void,172
): Promise<T> {173
if (signal.aborted) {174
throw signal.reason instanceof Error175
? signal.reason176
: new Error(`agent "${id}" creation aborted`, { cause: signal.reason })177
}178
const pending = Promise.resolve().then(operation)179
try {180
return await raceAbort(pending, signal, id)181
} catch (error: unknown) {182
// oxlint-disable-next-line typescript/no-unnecessary-condition -- the signal can abort while the operation is awaited.183
if (signal.aborted && releaseAbandoned !== undefined) {184
void pending.then(releaseAbandoned, () => undefined)185
}186
throw error187
}188
}190
/** Reject an output-token cap that cannot be represented exactly on the request wire. */191
function assertAgentOptions(options: AgentOptions): void {192
if (options.maxTokens !== undefined193
&& (!Number.isSafeInteger(options.maxTokens) || options.maxTokens <= 0)) {194
throw new TypeError('agent maxTokens must be a positive safe integer')195
}196
}198
/** One session's owned write handle plus the count of events already stored through it. */199
interface StoredSession {200
readonly handle: SessionHandle201
storedCount: number202
}204
/** Prepared-but-unpublished agent resources sharing one memoized teardown. */205
interface PreparedAgent {206
agent: ReactLoopAgent207
/** Aborts when the factory unloads, the caller cancels, or teardown begins — ends any setup await. */208
signal: AbortSignal209
/** Enter both registries and await creation listeners. */210
publish(source: SessionStartSource): Promise<AgentHandle>211
/** Reverse teardown: stop the machine, unregister, unwind the scope. Memoized. */212
dispose(): Promise<void>213
}215
declare module '@deepseek-ai/cordis' {216
interface Context {217
agentLoop: AgentLoop218
/**219
* Launcher-owned exact session identities for configured agents, keyed by220
* the agent's config `id` and set with `ctx.provide()` before any Loader221
* entry mounts (see {@link CONFIGURED_AGENT_IDENTITIES_KEY}). A launcher222
* owns identity because only it knows whether the session already exists,223
* while the `cordis.yml` row keeps the model route as ordinary patchable224
* config. An entry with no matching key keeps its configured identity.225
*/226
configuredAgentIdentities?: ConfiguredAgentIdentities227
}228
interface Events {229
/**230
* A declarative agent entry failed before it could publish a live agent.231
* Consumers that buffer work for the configured identity use this232
* transient signal to reject that work instead of waiting forever. Normal233
* factory teardown suppresses failures from the cancelled startup attempt.234
* @param payload.sessionId - exact shared agent/session identity that failed startup.235
* @param payload.error - persistence, setup, or publication failure.236
* @mode emit237
*/238
'agent-loop/config-start-failed'(payload: { sessionId: SessionId; error: unknown }): void239
}240
}242
export { DEFAULT_MAX_PARALLEL_TOOL_CALLS }244
/**245
* One launcher-selected session identity for a configured agent. `resume`246
* distinguishes rehydrating existing persisted history from creating the247
* session fresh under that exact id, which the two config keys express as248
* `resumeSessionId` and `sessionId`.249
*/250
export interface LauncherAgentIdentity {251
/** Exact session id to create fresh or resume. */252
id: SessionId253
/** Resume existing persisted history instead of creating the session fresh. */254
resume: boolean255
}257
/** Launcher-selected identities keyed by the configured agent's `id`. */258
export interface ConfiguredAgentIdentities extends Readonly<Record<string, LauncherAgentIdentity>> {}260
/**261
* Context key a launcher sets before any Loader entry mounts262
* (`ctx.provide(CONFIGURED_AGENT_IDENTITIES_KEY, identities)`) to fix263
* configured agents' session identities without a config key, so an overlay264
* repointing the row's model route cannot drop them.265
*/266
export const CONFIGURED_AGENT_IDENTITIES_KEY = 'configuredAgentIdentities'268
/**269
* Apply launcher-owned identities over the configured agents, replacing both270
* identity keys for every entry the launcher named so a config-supplied271
* identity can never survive alongside a launcher-supplied one.272
* @param agents - the configured agent entries.273
* @param identities - launcher identities keyed by configured agent `id`, or `undefined`.274
* @returns the entries with launcher-owned identities applied.275
*/276
function applyLauncherIdentities(277
agents: Config['agents'],278
identities: ConfiguredAgentIdentities | undefined,279
): Config['agents'] {280
if (identities === undefined) return agents281
return agents.map((agent) => {282
const identity = identities[agent.id]283
if (identity === undefined) return agent284
const { sessionId: _sessionId, resumeSessionId: _resumeSessionId, ...rest } = agent285
return identity.resume286
? { ...rest, resumeSessionId: identity.id }287
: { ...rest, sessionId: identity.id }288
})289
}291
/** Agent-loop plugin configuration. */292
export interface Config {293
/**294
* Maximum parallel-safe calls in flight per agent step. `1` is serial;295
* omission defaults to {@link DEFAULT_MAX_PARALLEL_TOOL_CALLS}.296
*/297
maxParallelToolCalls: Volatile<number>298
/** Agents created or resumed at plugin startup. */299
agents: (AgentOptions & {300
/** Stable config label used in logs and as the fresh combined-id prefix. */301
id: string302
/** Optional stable identity; remounts resume its materialized history, while first use creates it fresh. */303
sessionId?: SessionId304
/** Optional workspace for a fresh session. */305
cwd?: string306
/** Persisted session to resume instead of creating a fresh session. */307
resumeSessionId?: SessionId308
})[]309
}311
/** Reject self-contained identity conflicts before any configured agent starts. */312
function validateConfiguredAgents(agents: Config['agents']): void {313
const exactIdentities = new Map<SessionId, string>()314
for (const { id, sessionId, resumeSessionId } of agents) {315
const hasResumeId = resumeSessionId !== undefined && resumeSessionId !== ''316
if (sessionId !== undefined && hasResumeId) {317
throw new Error(`agent "${id}": sessionId and resumeSessionId are mutually exclusive`)318
}319
const exactIdentity = hasResumeId ? resumeSessionId : sessionId320
if (exactIdentity === undefined) continue321
const firstId = exactIdentities.get(exactIdentity)322
if (firstId !== undefined) {323
throw new Error(`agents "${firstId}" and "${id}" use duplicate exact session identity "${exactIdentity}"`)324
}325
exactIdentities.set(exactIdentity, id)326
}327
}329
/** Concrete agent factory and driver service. */330
export class AgentLoop extends Service implements AgentFactory {331
static inject = ['agents', 'sessions', 'llm', 'tools', 'systemPrompt', 'sessionProjections']333
/** Runtime schema for declarative agents. */334
static Config: z<{ agents?: Config['agents']; maxParallelToolCalls?: number }, Config> = z.object({335
maxParallelToolCalls: z.number().step(1).min(1).default(DEFAULT_MAX_PARALLEL_TOOL_CALLS).volatile(),336
agents: z.array(z.object({337
id: z.string().required(),338
sessionId: z.string().min(1),339
provider: z.string(),340
model: z.string(),341
reasoningEffort: z.string().min(1) as z<ReturnType<typeof ReasoningEffortId>>,342
maxTokens: z.number().step(1).min(1).max(Number.MAX_SAFE_INTEGER),343
cwd: z.string(),344
resumeSessionId: z.string(),345
})).default([]),346
}) as z<{ agents?: Config['agents']; maxParallelToolCalls?: number }, Config>348
/** Validated configuration owned by the agent-loop service. */349
readonly config: Config350
private readonly ownership: FactoryOwnership351
/** Plain holder prevents Cordis from re-tracing the factory's dependency context through a caller shadow. */352
private readonly runtime: { ctx: Context }354
constructor(ctx: Context, config: Config) {355
super(ctx, 'agentLoop')357
this.config = {358
agents: applyLauncherIdentities(config.agents, ctx.get(CONFIGURED_AGENT_IDENTITIES_KEY)),359
maxParallelToolCalls: config.maxParallelToolCalls,360
}361
validateConfiguredAgents(this.config.agents)362
// Register only after every config validation above has passed, so a363
// rejected constructor leaves no projection unit behind.364
ctx.sessionProjections.register(turnBoundaryProjectionDefinition)365
ctx.sessionProjections.register(inboxProjectionDefinition)366
this.ownership = new FactoryOwnership(ctx.fiber)367
this.runtime = { ctx }368
ctx.effect(() => () => this.ownership.dispose(), 'agentLoop.transactions()')369
ctx.effect(() => ctx.agents.setFactory(this), 'agentLoop.setFactory()')370
ctx.systemPrompt.variable('provider', context => context.agent?.options.provider)371
ctx.systemPrompt.variable('model', context => context.agent?.options.model)372
ctx.systemPrompt.variable('cwd', context => context.agent?.session.header.cwd)374
for (const { id, sessionId, cwd, resumeSessionId, ...options } of this.config.agents) {375
const meta = cwd === undefined ? {} : { cwd }376
if (resumeSessionId === undefined || resumeSessionId === '') {377
const configuredId = sessionId ?? brandString<SessionId>(`${id}-session-${randomUUID()}`)378
const persistence = sessionId === undefined ? undefined : ctx.get('sessionPersistence')379
if (persistence === undefined) {380
const startup = this.create(configuredId, options, meta).then(() => undefined, (error: unknown) => {381
this.reportConfiguredStartupFailure(id, 'restore', configuredId, error)382
})383
this.ownership.trackStartup(startup)384
} else {385
const startup = this.restoreOrCreateConfigured(ctx, persistence, configuredId, options, meta).catch((error: unknown) => {386
this.reportConfiguredStartupFailure(id, 'restore', configuredId, error)387
})388
this.ownership.trackStartup(startup)389
}390
continue391
}392
ctx.effect(() => {393
const fiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {394
void this.resumeWith(ctx, childCtx.sessionPersistence, {395
resumeSessionId,396
agentOptions: options,397
}).catch((error: unknown) => {398
this.reportConfiguredStartupFailure(id, 'resume', resumeSessionId, error)399
})400
})401
return fiber.dispose402
}, `agentLoop.resume(${id})`)403
}404
}406
/** Report a contained declarative-start failure to identity-bound consumers. */407
private reportConfiguredStartupFailure(408
configId: string,409
action: 'restore' | 'resume',410
sessionId: SessionId,411
error: unknown,412
): void {413
if (!this.ownership.isActive()) return414
this.ctx.logger.warn(`agent "${configId}": config-driven ${action} of "${sessionId}" failed: ${errorChain(error)}`)415
const args: unknown[] = ['agent-loop/config-start-failed', { sessionId, error }]416
for (const callback of this.ctx.events.dispatch('emit', args)) {417
try {418
const returned: unknown = callback(...args)419
void Promise.resolve(returned).catch((listenerError: unknown) => {420
this.ctx.logger.warn(`agent "${configId}": config-start-failed listener rejected: ${errorChain(listenerError)}`)421
})422
} catch (listenerError: unknown) {423
this.ctx.logger.warn(`agent "${configId}": config-start-failed listener threw: ${errorChain(listenerError)}`)424
}425
}426
}428
/** Restore a materialized exact config identity on remount, or create it on first use. */429
private async restoreOrCreateConfigured(430
ownerCtx: Context,431
persistence: SessionPersistence,432
sessionId: SessionId,433
agentOptions: AgentOptions,434
meta: Pick<SessionHeader, 'cwd'>,435
): Promise<void> {436
await this.waitForDrainingConfiguredIdentity(ownerCtx, sessionId)437
if (!this.ownership.isActive()) return438
try {439
await this.resumeWith(ownerCtx, persistence, { resumeSessionId: sessionId, agentOptions })440
return441
} catch (error: unknown) {442
if (!this.ownership.isActive()) return443
// Only a genuinely absent stored session falls back to first creation;444
// corruption, ownership conflicts, and backend failures stay loud.445
if (!(error instanceof SessionPersistenceNotFoundError)) throw error446
}447
await this.create(sessionId, agentOptions, meta)448
}450
/** Wait for a draining same-id lifecycle to finish registry teardown. */451
private async waitForDrainingConfiguredIdentity(ownerCtx: Context, sessionId: SessionId): Promise<void> {452
// Only an id still occupying a registry needs waiting for; a live healthy453
// occupant is a collision the create/resume below will surface itself.454
if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) return456
const released = Promise.withResolvers<void>()457
const checkReleased = (): void => {458
if (ownerCtx.agents.get(sessionId) === undefined && ownerCtx.sessions.get(sessionId) === undefined) {459
released.resolve()460
}461
}462
const disposeAgentListener = ownerCtx.on('agent/disposed', () => { checkReleased() })463
const disposeSessionListener = ownerCtx.on('session/disposed', checkReleased)464
try {465
checkReleased()466
await this.ownership.waitWhileActive(released.promise)467
} finally {468
disposeAgentListener()469
disposeSessionListener()470
}471
}473
/**474
* Construct the driver, scope, and one memoized reverse teardown for a new475
* agent. The teardown is registered with the factory and the owner fiber476
* BEFORE publication, so a mid-setup unload rolls everything back; `signal`477
* fuses caller cancellation with lifecycle teardown for setup awaits.478
*/479
private prepare(480
ownerCtx: Context,481
id: SessionId,482
options: AgentOptions,483
session: Session,484
callerSignal?: AbortSignal,485
handle?: SessionHandle,486
parentAgent?: Agent,487
): PreparedAgent {488
assertAgentOptions(options)489
ownerCtx.fiber.assertActive()490
// Every caller reaches prepare() synchronously from a service method491
// whose Cordis dispatch already requires the live factory fiber, or492
// re-checks ownership itself after its awaits (resume's load barrier).493
/* v8 ignore next -- unreachable backstop, see above */494
if (!this.ownership.isActive()) throw new Error('agent loop is not active')495
if (callerSignal?.aborted) {496
throw callerSignal.reason instanceof Error497
? callerSignal.reason498
: new Error(`agent "${id}" creation aborted`, { cause: callerSignal.reason })499
}500
const loopCtx = this.runtime.ctx502
// Deactivation fuses three owners, each with its own reason: the caller's503
// cancellation signal, the owner fiber's unload, and factory teardown.504
// It is registered BEFORE any resource exists, over mutable slots, so an505
// unload arriving while the scope is still minting finds a working506
// disposer instead of a leak.507
const abort = new AbortController()508
const onCallerAbort = (): void => {509
abort.abort(callerSignal?.reason instanceof Error510
? callerSignal.reason511
: new Error(`agent "${id}" creation aborted`, { cause: callerSignal?.reason }))512
}513
const onFactoryTeardown = (): void => { abort.abort(this.ownership.signal.reason) }514
callerSignal?.addEventListener('abort', onCallerAbort, { once: true })515
this.ownership.signal.addEventListener('abort', onFactoryTeardown, { once: true })517
let machine: ReactLoopAgent | undefined518
let detachSession: (() => void) | undefined519
let detachAgent: (() => void) | undefined520
let disposing: Promise<void> | undefined521
let publication: ReturnType<typeof Promise.withResolvers<void>> | undefined522
const machineReady = Promise.withResolvers<void>()523
// Reverse teardown, memoized so every racing owner awaits one quiescence:524
// stop the machine, drain and close the session's write path, leave the525
// registries, unwind the scope, release bookkeeping.526
const dispose = (ownerTriggered = false): Promise<void> => (disposing ??= (async () => {527
abort.abort(new Error(`agent "${id}" lifecycle disposed`))528
callerSignal?.removeEventListener('abort', onCallerAbort)529
this.ownership.signal.removeEventListener('abort', onFactoryTeardown)530
// Teardown failures are collected, never swallowed: registry, scope,531
// and ownership cleanup always run to quiescence, then the memoized532
// disposal rejects with what failed so every racing owner observes it.533
const failures: unknown[] = []534
try {535
// Creation listeners retain the session and scope through their awaits.536
if (publication !== undefined) await publication.promise537
// Disposal IS a disposed-cause cancel followed by quiescence. New work538
// sent after this point is the sender's bug — the registries are about539
// to drop the agent, so nothing should still hold it.540
/* v8 ignore next -- Cordis effect teardown waits for synchronous setup before observing the machine slot. */541
if (machine === undefined) await machineReady.promise542
/* v8 ignore next -- setup failure untracks this disposer before resolving without a machine. */543
if (machine !== undefined) {544
machine.cancel({ kind: 'disposed' })545
await machine.whenIdle()546
await machine.scope.dispose()547
}548
} catch (error: unknown) {549
failures.push(error)550
}551
// The loop above committed its closing events synchronously into the552
// session; handle close drains them durably before releasing the write553
// path. The close drain can be the first operation that surfaces a554
// durability failure, so its error is retained, not logged away.555
try {556
await handle?.close()557
} catch (error: unknown) {558
failures.push(error)559
}560
try {561
detachAgent?.()562
detachSession?.()563
} finally {564
untrack()565
if (!ownerTriggered) await unfollowOwner()566
}567
if (failures.length === 1) throw failures[0]568
if (failures.length > 1) {569
throw new AggregateError(failures, `agent "${id}" disposal failed`)570
}571
})())572
const untrack = this.ownership.track(dispose)573
let unfollowOwner: () => Promise<void> | void574
try {575
unfollowOwner = ownerCtx.effect(function* () {576
machine = new ReactLoopAgent(loopCtx, id, options, session)577
machineReady.resolve()578
yield machine.scope.rawDispose579
yield () => {580
// Owner disposal owns the same quiescence boundary. Its teardown skips581
// unregistering this already-running owner effect from inside itself.582
if (disposing !== undefined) return583
abort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))584
return dispose(true)585
}586
}, `agentLoop.lifecycle(${id})`)587
/* v8 ignore start -- ctx.effect throws only on an inactive fiber, which assertActive() above already rejected */588
} catch (error: unknown) {589
machineReady.resolve()590
untrack()591
callerSignal?.removeEventListener('abort', onCallerAbort)592
this.ownership.signal.removeEventListener('abort', onFactoryTeardown)593
throw error594
}595
/* v8 ignore stop */597
const assertLive = (): void => {598
if (!abort.signal.aborted) return599
// Every fused abort source carries an Error reason: onCallerAbort and600
// raceAbort wrap non-Error caller reasons, and the factory/lifecycle601
// owners abort with constructed Errors.602
/* v8 ignore next -- unreachable String() arm, see above */603
throw abort.signal.reason instanceof Error ? abort.signal.reason : new Error(String(abort.signal.reason))604
}605
try {606
/* v8 ignore next -- a synchronous effect exhausts the generator before returning */607
if (machine === undefined) throw new Error(`agent "${id}" lifecycle did not construct its driver`)608
const agent = machine609
assertLive()611
return {612
agent,613
signal: abort.signal,614
publish: async (source) => {615
publication = Promise.withResolvers<void>()616
try {617
assertLive()618
detachSession = agent.ctx.sessions.enter(session)619
// The mounted backend routes announced live events into the active620
// write handle by session id; the loop only owns the handle itself.621
detachAgent = loopCtx.agents.enter(agent, parentAgent)622
agent.ctx.sessions.announce(session)623
assertLive()624
await loopCtx.agents.announce(agent, source, abort.signal)625
assertLive()626
return { agent, dispose }627
} finally {628
publication.resolve()629
publication = undefined630
}631
},632
dispose,633
}634
} catch (error: unknown) {635
machineReady.resolve()636
// Rollback swallows a disposal rejection: the setup failure is primary.637
void dispose().catch(() => {})638
throw error639
}640
}642
/**643
* Create an agent and session under one caller-supplied identity, owned by644
* the accessing fiber. Constructor-driven config calls mint a fresh combined645
* id before entering this boundary. When a persistence backend is mounted,646
* the session's durable identity and any seed are stored before publication.647
* @param id - shared agent/session identity.648
* @param options - concrete loop options.649
* @param meta - optional fresh-session workspace metadata.650
* @returns the published running agent.651
*/652
async create(id: SessionId, options: AgentOptions = {}, meta: Pick<SessionHeader, 'cwd'> = {}): Promise<Agent> {653
using preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(id, { meta }))654
const stored = await this.createStoredSession(preparation.session)655
let prepared: PreparedAgent656
try {657
prepared = this.prepare(this.ctx, id, options, preparation.session, undefined, stored?.handle)658
} catch (error: unknown) {659
await stored?.handle.close().catch(() => {})660
throw error661
}662
return (await this.initializeAgent(prepared, async () => {663
await this.appendUnstoredSuffix(stored, preparation.session)664
return await prepared.publish('startup')665
})).agent666
}668
/**669
* Take a fresh session's write ownership when persistence is mounted.670
* Nothing is appended here: the constructor seed (which never re-emits671
* through `session/event`) is stored by `appendUnstoredSuffix` at the672
* publication commit point, so a failed or cancelled validation or setup673
* closes an unmaterialized handle and leaves no stored residue — the same674
* id can be created again.675
* @param session - the unpublished session to store.676
* @param signal - optional cancellation forwarded to the backend create.677
* @returns the owned handle and stored cursor, or `undefined` without a backend.678
*/679
private async createStoredSession(session: Session, signal?: AbortSignal): Promise<StoredSession | undefined> {680
const persistence = this.runtime.ctx.get('sessionPersistence')681
if (persistence === undefined) return undefined682
const handle = await persistence.create(session.header, {683
inheritedEventCount: session.inheritedEventCount,684
...signal === undefined ? {} : { signal },685
})686
return { handle, storedCount: 0 }687
}689
/**690
* Durably store the session events appended since the last stored cursor.691
* Pre-publication appends (constructor seed markers, setup-window events692
* such as delegation policy records) never re-emit through `session/event`,693
* so publication must flush them through the handle before live events694
* start routing into it.695
* @param stored - the session's owned handle and stored cursor, if any.696
* @param session - the unpublished session whose suffix is stored.697
*/698
private async appendUnstoredSuffix(stored: StoredSession | undefined, session: Session): Promise<void> {699
if (stored === undefined) return700
// oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.701
const suffix = session.snapshotEvents(SessionLogOffset(stored.storedCount))702
if (suffix.length > 0) await stored.handle.append(suffix)703
// Advance by what was stored, not to `session.seq`: an event appended704
// during the await must stay unstored for the next flush.705
stored.storedCount += suffix.length706
}708
/**709
* Create an owned agent on a caller-supplied session id.710
* @param ownerCtx - caller context that structurally owns the lifecycle.711
* @param options - identities, optional live parent, session seed/metadata, loop options, setup, and cancellation.712
* @returns the published handle.713
*/714
async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {715
const preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(options.sessionId, {716
...options.seed === undefined ? {} : { seed: options.seed },717
...options.meta === undefined ? {} : { meta: options.meta },718
...options.inheritedEventCount === undefined ? {} : { inheritedEventCount: options.inheritedEventCount },719
}))720
const published = (async () => {721
let stored: StoredSession | undefined722
try {723
// raceAbortCall normalizes a pre-aborted or mid-create abort and724
// closes a handle that finishes creating after abandonment.725
stored = options.signal === undefined726
? await this.createStoredSession(preparation.session)727
: await raceAbortCall(728
() => this.createStoredSession(preparation.session, options.signal),729
options.signal,730
options.sessionId,731
(abandoned) => { void abandoned?.handle.close().catch(() => {}) },732
)733
} catch (error: unknown) {734
preparation[Symbol.dispose]()735
throw error736
}737
return this.setupAndPublish(738
ownerCtx,739
options.sessionId,740
preparation,741
options.agentOptions ?? {},742
options.setup,743
options.signal,744
'startup',745
stored,746
options.parentAgent,747
)748
})()749
this.ownership.trackWrapper(published)750
return published751
}753
/** Prepare one Agent around an acquired Session, run setup, and publish it. */754
private async setupAndPublish(755
ownerCtx: Context,756
id: SessionId,757
preparation: SessionPreparation,758
agentOptions: AgentOptions,759
setup: AgentSetup | undefined,760
signal: AbortSignal | undefined,761
source: SessionStartSource,762
stored?: StoredSession,763
parentAgent?: Agent,764
): Promise<AgentHandle> {765
using ownedPreparation = preparation766
const session = ownedPreparation.session767
let prepared: PreparedAgent768
try {769
prepared = this.prepare(ownerCtx, id, agentOptions, session, signal, stored?.handle, parentAgent)770
} catch (error: unknown) {771
await stored?.handle.close().catch(() => {})772
throw error773
}774
return await this.initializeAgent(prepared, async () => {775
const setupCommit = await raceAbort(setup?.(prepared.agent.ctx, prepared.agent), prepared.signal, id)776
setupCommit?.commit()777
await this.appendUnstoredSuffix(stored, session)778
return await prepared.publish(source)779
})780
}782
private async initializeAgent(prepared: PreparedAgent, initialize: () => Promise<AgentHandle>): Promise<AgentHandle> {783
try {784
return await prepared.agent.runMaintenance(async () => {785
try {786
return await initialize()787
} catch (error: unknown) {788
// Teardown owns inbox cleanup and may already have removed its projection.789
prepared.agent.cancel({ kind: 'disposed' }, { keepInbox: true })790
throw error791
}792
})793
} catch (error: unknown) {794
// Rollback swallows a disposal rejection (a failing final handle close):795
// the setup failure is the primary error the caller must see.796
await prepared.dispose().catch(() => {})797
throw error798
}799
}801
/**802
* Resume an owned agent from the configured persistence service.803
* @param ownerCtx - caller context that owns load, setup, and the live lifecycle.804
* @param options - persisted identity, optional live parent, loop options, setup, and cancellation.805
* @returns the published handle.806
*/807
async resume(ownerCtx: Context, options: ResumeAgentOptions): Promise<AgentHandle> {808
const persistence = this.runtime.ctx.get('sessionPersistence')809
if (persistence === undefined) {810
throw new Error('cannot resume: session persistence is not configured (load a dsh-session-persistence backend)')811
}812
return this.resumeWith(ownerCtx, persistence, options)813
}815
/** Resume through an explicit persistence handle used by the deferred config path. */816
private resumeWith(817
ownerCtx: Context,818
persistence: SessionPersistence,819
options: ResumeAgentOptions,820
): Promise<AgentHandle> {821
const id = options.resumeSessionId822
const published = (async () => {823
// The open and read may outlive their owner: race them against caller824
// cancellation, owner-fiber unload, and factory teardown so a825
// never-settling backend cannot pin the identity.826
const ownerAbort = new AbortController()827
const unfollowOwner = ownerCtx.effect(() => () => {828
ownerAbort.abort(new Error(`agent "${id}" setup aborted: owner disposed during setup`))829
}, `agentLoop.resume-load(${id})`)830
const fused = AbortSignal.any([831
...options.signal === undefined ? [] : [options.signal],832
ownerAbort.signal,833
this.ownership.signal,834
])835
let handle: SessionHandle | undefined836
let stored: StoredSession | undefined837
let preparation: SessionPreparation | undefined838
try {839
try {840
// Taking write ownership FIRST excludes a concurrent resume of the841
// same id (in this process, a live agent's handle holds the claim).842
handle = await raceAbortCall(843
() => persistence.open(id, 'write', { signal: fused }),844
fused,845
id,846
(abandoned) => { void abandoned.close() },847
)848
// Semantic crash repair is the agent layer's job: persistence hands849
// back the physically valid log; an interrupted final turn receives850
// synthetic closers (missing tool errors, step/end, turn/end) that851
// are appended through the same handle as an ordinary batch.852
const coldRead = await handle.read(0, undefined, { signal: fused })853
fused.throwIfAborted()854
const persisted = coldRead.events855
const closers = interruptedTurnClosers(persisted)856
if (closers.length > 0) await handle.append(closers)857
preparation = SessionPreparation.create(this.runtime.ctx.sessions.prepare(id, {858
seed: [...persisted, ...closers],859
meta: structuredClone(handle.header),860
inheritedEventCount: handle.inheritedEventCount,861
eventState: coldRead.eventState,862
}))863
stored = { handle, storedCount: persisted.length + closers.length }864
await this.appendUnstoredSuffix(stored, preparation.session)865
} finally {866
await unfollowOwner()867
}868
ownerCtx.fiber.assertActive()869
if (!this.ownership.isActive()) throw new Error('agent loop is not active')870
const owned = stored871
handle = undefined // ownership passes to setupAndPublish/prepare872
return await this.setupAndPublish(873
ownerCtx,874
id,875
preparation,876
options.agentOptions ?? {},877
options.setup,878
options.signal,879
'resume',880
owned,881
options.parentAgent,882
)883
} finally {884
preparation?.[Symbol.dispose]()885
await handle?.close().catch(() => {})886
}887
})()888
this.ownership.trackWrapper(published)889
return published890
}891
}893
export default AgentLoop