返回源码地图

packages/subagent/subagent/src/lifecycle.ts

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

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

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's
7 * consumer-facing types; this module owns only the implementation and the
8 * package-private {@link ActivationObserver} the continuation manager consumes.
9 * Keeping the internal control interface out of the published surface is
10 * deliberate: the observer's `start`/`capture`/`settle` ordering is a contract
11 * between this module and one in-package caller, not something a plugin may
12 * depend on.
13 *
14 * @module @deepseek-ai/dsh-subagent/lifecycle
15 */
16
17import { randomUUID } from 'node:crypto'
18import type { Context } from '@deepseek-ai/cordis'
19import type { Agent } from '@deepseek-ai/dsh-agent'
20import type { ContentBlock } from '@deepseek-ai/dsh-llm'
21import { foldConsumedWork } from '@deepseek-ai/dsh-agent'
22import { SessionLogOffset } from '@deepseek-ai/dsh-session'
23import type { SessionEvent, SessionId, SessionLogOffset as SessionLogOffsetType } from '@deepseek-ai/dsh-session'
24import { finalAssistantOutput } from './assistant-output.ts'
25import { SubagentRunId } from './types.ts'
26import type { SubagentResult, SubagentRun, SubagentRunEndInfo, SubagentRunInfo } from './types.ts'
27
28/**
29 * How one Activation's residency epoch ended, as both the terminal lifecycle
30 * edge and the manager's own parent delivery report it.
31 */
32export 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}
38
39/**
40 * Lifecycle observer for one Activation's residency epoch, so continuable
41 * children emit the same start/end pair as one-shot runs. Package-private: the
42 * continuation manager is the only consumer, and its call ordering is an
43 * in-package contract rather than a published extension point.
44 */
45export 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): void
51 /**
52 * Snapshot the child-dependent terminal facts while the child is still
53 * registered, because handle disposal unregisters it and consumers resolve it
54 * 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): void
58 /**
59 * Resolve the terminal facts {@link settle} will publish, without publishing
60 * them. The manager's parent delivery must run before the ownership release
61 * that lets the parent settle, which is earlier than the terminal edge; both
62 * 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): ActivationTerminal
67 /**
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: a
70 * failure before residency publishes no edge, because inventing one would
71 * report a lifecycle the child never had.
72 * @param failure - the teardown or durability failure, or `undefined` on success.
73 */
74 settle(failure: unknown): void
75}
76
77/**
78 * Publish one lifecycle edge with per-listener exception containment. Run edges
79 * carry the delegating parent that keys scoped dispatch; provider removal has no
80 * parent carrier and reaches listeners unscoped.
81 *
82 * The service owns this closure because scoped dispatch keys its carrier by the
83 * exact service instance, whose own context filter composes into the carrier;
84 * a narrowed stand-in would silently change scope filtering.
85 */
86export type LifecycleEmitter = {
87 (name: 'subagent/start', info: SubagentRunInfo, parent: Agent): void
88 (name: 'subagent/end', info: SubagentRunEndInfo, parent: Agent): void
89 (name: 'subagent/provider-removed', info: string): void
90}
91
92/**
93 * Build the contained lifecycle emitter this seam publishes every edge through.
94 * Every listener is independently contained: a synchronous throw or a rejected
95 * 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 */
101export 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 === undefined
111 ? [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}
125
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 */
134export 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 reactions
147 // 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 run
163}
164
165/**
166 * Build the observer for one continuable Activation's residency epoch. Observers
167 * see the same vocabulary as a one-shot run, so a child's start and settlement
168 * remain observable without exposing whether the manager materialized, woke, or
169 * 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 */
176export 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 come
184 // from the suffix it actually produced — never the whole session, which
185 // 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 before
188 // `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 its
191 // output: an answer this harness could not durably release is not a result.
192 const terminal = (failure: unknown): ActivationTerminal => failure === undefined
193 ? captured
194 : { stopReason: 'error' }
195 return {
196 start: (child: Agent): void => {
197 boundary = child.session.seq
198 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}
220
221/**
222 * Why this child's epoch ended, for the terminal lifecycle edge and the
223 * manager's own parent delivery. The child's own log is authoritative:
224 * teardown succeeding says nothing about whether the model errored, hit its
225 * token ceiling, or was cancelled, so deriving the reason from disposal would
226 * 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 accepted
230 * work was cancelled after it without any turn opening over it. A recorded
231 * failure still wins over a cancellation — stopping a child that had already
232 * 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 both
235 * closed cleanly and had nothing left to run.
236 */
237function 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 input
248 // 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 epoch
252 // 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, while
257 * this function reads an epoch-owned suffix. `TurnEndReason` is merge-extensible,
258 * so a backend-added variant cannot be listed; treating an unnameable reason as
259 * success would report failed work as completed. */
260 default:
261 return 'error'
262 }
263}
264
265/** Render any listener-thrown value without letting coercion escape containment. */
266function 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}