1
/**2
* Lifecycle-edge publication for both subagent shapes: the contained emitter,3
* the one-shot run observer, and the continuable Activation observer.4
*5
* The public payload contracts ({@link SubagentRunInfo},6
* {@link SubagentRunEndInfo}) live in `./types.ts` with the rest of the seam's7
* consumer-facing types; this module owns only the implementation and the8
* package-private {@link ActivationObserver} the continuation manager consumes.9
* Keeping the internal control interface out of the published surface is10
* deliberate: the observer's `start`/`capture`/`settle` ordering is a contract11
* between this module and one in-package caller, not something a plugin may12
* depend on.13
*14
* @module @deepseek-ai/dsh-subagent/lifecycle15
*/17
import { randomUUID } from 'node:crypto'18
import type { Context } from '@deepseek-ai/cordis'19
import type { Agent } from '@deepseek-ai/dsh-agent'20
import type { ContentBlock } from '@deepseek-ai/dsh-llm'21
import { foldConsumedWork } from '@deepseek-ai/dsh-agent'22
import { SessionLogOffset } from '@deepseek-ai/dsh-session'23
import type { SessionEvent, SessionId, SessionLogOffset as SessionLogOffsetType } from '@deepseek-ai/dsh-session'24
import { finalAssistantOutput } from './assistant-output.ts'25
import { SubagentRunId } from './types.ts'26
import type { SubagentResult, SubagentRun, SubagentRunEndInfo, SubagentRunInfo } from './types.ts'28
/**29
* How one Activation's residency epoch ended, as both the terminal lifecycle30
* edge and the manager's own parent delivery report it.31
*/32
export interface ActivationTerminal {33
/** Why this epoch's last ordinary turn ended, or `error` when teardown failed. */34
readonly stopReason: SubagentResult['stopReason']35
/** The epoch's final assistant content, absent when it produced none or failed. */36
readonly output?: readonly ContentBlock[]37
}39
/**40
* Lifecycle observer for one Activation's residency epoch, so continuable41
* children emit the same start/end pair as one-shot runs. Package-private: the42
* continuation manager is the only consumer, and its call ordering is an43
* in-package contract rather than a published extension point.44
*/45
export interface ActivationObserver {46
/**47
* Publish the start edge once the epoch is resident.48
* @param child - the resident child agent, whose log suffix bounds this epoch.49
*/50
start(child: Agent): void51
/**52
* Snapshot the child-dependent terminal facts while the child is still53
* registered, because handle disposal unregisters it and consumers resolve it54
* to read the child's own log and scope.55
* @param child - the quiescent child agent about to be released.56
*/57
capture(child: Agent): void58
/**59
* Resolve the terminal facts {@link settle} will publish, without publishing60
* them. The manager's parent delivery must run before the ownership release61
* that lets the parent settle, which is earlier than the terminal edge; both62
* therefore read one computation instead of restating the failure rule.63
* @param failure - the teardown or durability failure, or `undefined` on success.64
* @returns this epoch's stop reason and final assistant content.65
*/66
terminal(failure: unknown): ActivationTerminal67
/**68
* Publish the terminal edge exactly once, pairing this epoch's {@link start},69
* after the disposal outcome is known. Called only for a resident epoch: a70
* failure before residency publishes no edge, because inventing one would71
* report a lifecycle the child never had.72
* @param failure - the teardown or durability failure, or `undefined` on success.73
*/74
settle(failure: unknown): void75
}77
/**78
* Publish one lifecycle edge with per-listener exception containment. Run edges79
* carry the delegating parent that keys scoped dispatch; provider removal has no80
* parent carrier and reaches listeners unscoped.81
*82
* The service owns this closure because scoped dispatch keys its carrier by the83
* exact service instance, whose own context filter composes into the carrier;84
* a narrowed stand-in would silently change scope filtering.85
*/86
export type LifecycleEmitter = {87
(name: 'subagent/start', info: SubagentRunInfo, parent: Agent): void88
(name: 'subagent/end', info: SubagentRunEndInfo, parent: Agent): void89
(name: 'subagent/provider-removed', info: string): void90
}92
/**93
* Build the contained lifecycle emitter this seam publishes every edge through.94
* Every listener is independently contained: a synchronous throw or a rejected95
* returned promise is logged without starving peer listeners, changing the run,96
* or — for provider removal, which fires from a disposer — breaking teardown.97
* @param ctx - the service's own context, owning dispatch and the logger.98
* @param carrier - resolve the scoped dispatch carrier for one delegating parent.99
* @returns the emitter both observers and the provider registry publish through.100
*/101
export function createLifecycleEmitter(102
ctx: Context,103
carrier: (parent: Agent) => object,104
): LifecycleEmitter {105
return (106
name: 'subagent/start' | 'subagent/end' | 'subagent/provider-removed',107
info: SubagentRunInfo | SubagentRunEndInfo | string,108
parent?: Agent,109
): void => {110
const dispatchArgs: unknown[] = parent === undefined111
? [name, info]112
: [carrier(parent), name, info]113
for (const callback of ctx.events.dispatch('emit', dispatchArgs)) {114
try {115
const returned: unknown = callback(info)116
void Promise.resolve(returned).catch((error: unknown) => {117
ctx.logger.warn(`subagent: ${name} listener rejected: ${renderThrown(error)}`)118
})119
} catch (error: unknown) {120
ctx.logger.warn(`subagent: ${name} listener threw: ${renderThrown(error)}`)121
}122
}123
}124
}126
/**127
* Emit the start/end lifecycle pair for one accepted one-shot run.128
* @param emit - the contained lifecycle emitter.129
* @param provider - the provider that established the run.130
* @param parent - the delegating parent keying scoped dispatch.131
* @param run - the published run whose settlement closes the pair.132
* @returns the same run, unchanged.133
*/134
export function observeRun(135
emit: LifecycleEmitter,136
provider: string,137
parent: Agent,138
run: SubagentRun,139
): SubagentRun {140
const identity = {141
runId: SubagentRunId(randomUUID()),142
provider,143
id: run.id,144
local: run.localAgent !== undefined,145
}146
// Attach the terminal observer before dispatching start. Promise reactions147
// still run after this synchronous start emission, preserving start → end.148
void run.result.then(149
(result) => {150
emit('subagent/end', {151
...identity,152
stopReason: result.stopReason,153
// Omit the field when no output exists, matching continuable epochs.154
...result.output.length === 0 ? {} : { lastAssistantMessage: result.output },155
}, parent)156
},157
() => {158
emit('subagent/end', { ...identity, stopReason: 'error' }, parent)159
},160
)161
emit('subagent/start', identity, parent)162
return run163
}165
/**166
* Build the observer for one continuable Activation's residency epoch. Observers167
* see the same vocabulary as a one-shot run, so a child's start and settlement168
* remain observable without exposing whether the manager materialized, woke, or169
* cold-resumed it. Creation failure before residency emits no lifecycle edge.170
* @param emit - the contained lifecycle emitter.171
* @param provider - the provider name recorded in the durable descriptor.172
* @param childId - the durable child session id.173
* @param parent - the exact live direct parent keying scoped dispatch.174
* @returns the observer whose edges this epoch publishes.175
*/176
export function createActivationObserver(177
emit: LifecycleEmitter,178
provider: string,179
childId: SessionId,180
parent: Agent,181
): ActivationObserver {182
const identity = { runId: SubagentRunId(randomUUID()), provider, id: childId, local: true }183
// A cold resume replays earlier turns, so this epoch's telemetry must come184
// from the suffix it actually produced — never the whole session, which185
// would report a previous epoch's answer when this one opened no turn.186
let boundary: SessionLogOffsetType = SessionLogOffset(0)187
// Assigned by `capture()`, which the disposal path always runs before188
// `settle()`; a resident epoch therefore always has its facts by then.189
let captured: ActivationTerminal = { stopReason: 'completed' }190
// Teardown failure overrides the epoch's own outcome and withholds its191
// output: an answer this harness could not durably release is not a result.192
const terminal = (failure: unknown): ActivationTerminal => failure === undefined193
? captured194
: { stopReason: 'error' }195
return {196
start: (child: Agent): void => {197
boundary = child.session.seq198
emit('subagent/start', identity, parent)199
},200
capture: (child: Agent): void => {201
// oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.202
const own = child.session.snapshotEvents(boundary)203
const output = finalAssistantOutput(own)204
captured = {205
stopReason: epochStopReason(own),206
...output === undefined ? {} : { output },207
}208
},209
terminal,210
settle: (failure: unknown): void => {211
const { stopReason, output } = terminal(failure)212
emit('subagent/end', {213
...identity,214
stopReason,215
...output === undefined ? {} : { lastAssistantMessage: output },216
}, parent)217
},218
}219
}221
/**222
* Why this child's epoch ended, for the terminal lifecycle edge and the223
* manager's own parent delivery. The child's own log is authoritative:224
* teardown succeeding says nothing about whether the model errored, hit its225
* token ceiling, or was cancelled, so deriving the reason from disposal would226
* report failed work as completed.227
*228
* {@link foldConsumedWork} supplies both halves the raw turn sequence cannot:229
* which turn accounts for the work this epoch consumed, and whether accepted230
* work was cancelled after it without any turn opening over it. A recorded231
* failure still wins over a cancellation — stopping a child that had already232
* failed does not turn its failure into a cancellation.233
* @param events - this epoch's own event suffix.234
* @returns its terminal stop reason; `completed` only for an epoch that both235
* closed cleanly and had nothing left to run.236
*/237
function epochStopReason(events: readonly SessionEvent[]): SubagentResult['stopReason'] {238
const { end, droppedUnrun } = foldConsumedWork(events)239
switch (end?.data.reason.kind) {240
case 'max-tokens':241
return 'max-tokens'242
case 'aborted':243
case 'interrupted':244
return 'aborted'245
case 'error':246
return 'error'247
// A pre-step rejection — a hook deny, a policy plugin — discarded input248
// this epoch had claimed: the work was declined, not done.249
case 'blocked':250
return 'refusal'251
// A clean ending and no accounting turn at all share one rule: the epoch252
// finished what it was given unless a cancelled queue says otherwise.253
case undefined:254
case 'completed':255
return droppedUnrun ? 'aborted' : 'completed'256
/* v8 ignore next 4 -- `forked` appears only in constructor seed history, while257
* this function reads an epoch-owned suffix. `TurnEndReason` is merge-extensible,258
* so a backend-added variant cannot be listed; treating an unnameable reason as259
* success would report failed work as completed. */260
default:261
return 'error'262
}263
}265
/** Render any listener-thrown value without letting coercion escape containment. */266
function renderThrown(value: unknown): string {267
try {268
return value instanceof Error ? `${value.name}: ${value.message}` : String(value)269
} catch {270
return '<unrenderable thrown value>'271
}272
}