返回源码地图

packages/interaction/user-questions/src/index.ts

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

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

1/**
2 * Service Definition for the user-questions capability seam (`ctx.userQuestions`): a UI-backed service for
3 * pausing an agent tool call until the human answers a question. The model-
4 * facing tool lives in `@deepseek-ai/dsh-tool-ask-user`; UI packages compose
5 * answerers on the Agent-scoped Cordis waterfall.
6 *
7 * @module @deepseek-ai/dsh-user-questions
8 */
9
10import { Context } from '@deepseek-ai/cordis'
11import type {} from '@deepseek-ai/dsh-agent'
12import { createUserMessage, HarnessError, type MessageId, type ToolCallId, type UserMessage } from '@deepseek-ai/dsh-llm'
13import type { Session } from '@deepseek-ai/dsh-session'
14import { scopeTarget } from '@deepseek-ai/dsh-scope'
15import { Remote, TypertRemoteService } from '@deepseek-ai/dsh-typert-protocol'
16import type { Agent } from '@deepseek-ai/dsh-agent'
17import type {} from '@deepseek-ai/dsh-session-projection'
18import z from '@deepseek-ai/schemastery'
19import { userQuestionProjectionDefinition } from './projection.ts'
20import { TimedQuestionWait } from './timed-wait.ts'
21
22declare module '@deepseek-ai/cordis' {
23 interface Context {
24 userQuestions: UserQuestionService
25 }
26}
27
28import type {
29 AskUserQuestionAnswer, AskUserQuestionRequestEvent, PendingUserQuestion,
30} from './types.ts'
31
32export type {
33 AskUserQuestionAnswer, AskUserQuestionAnswerItem, AskUserQuestionIntent, AskUserQuestionItem,
34 AskUserQuestionOption, PendingUserQuestion, SettledUserQuestion, UserQuestionProjectionView,
35 UserQuestionState,
36} from './types.ts'
37export { isTimedAskUserQuestionSchema, TIMED_WAIT_PARAMETER } from './projection.ts'
38
39/** Request for a human answer. */
40export interface AskUserQuestionRequest extends AskUserQuestionRequestEvent {}
41
42/** Timed ask result returned when the foreground answer window closes. */
43export type TimedUserQuestionResult = AskUserQuestionAnswer | { pending: true; callId: ToolCallId }
44
45/** Stable error taxonomy for user-questions failures. */
46export class UserQuestionError extends HarnessError {
47 constructor(message: string, code: string, options?: ErrorOptions) {
48 super(message, code, options)
49 this.name = 'UserQuestionError'
50 }
51}
52
53function abortedQuestion(cause?: unknown): UserQuestionError {
54 return new UserQuestionError(
55 'ask_user_question was aborted before the user answered',
56 'ASK_ABORTED',
57 cause === undefined ? undefined : { cause },
58 )
59}
60
61function isRecord(value: unknown): value is Record<string, unknown> {
62 return typeof value === 'object' && value !== null && !Array.isArray(value)
63}
64
65function restoreUserQuestionError(reason: unknown): unknown {
66 if (reason instanceof UserQuestionError) return reason
67 if (isRecord(reason)
68 && reason.name === 'UserQuestionError'
69 && typeof reason.message === 'string'
70 && typeof reason.code === 'string') {
71 return new UserQuestionError(reason.message, reason.code, { cause: reason })
72 }
73 return reason
74}
75
76type QueuedReply = { messageId: MessageId; claimedTurn?: number }
77
78/** `ctx.userQuestions`: validation plus the scoped answerer waterfall. */
79export class UserQuestionService extends TypertRemoteService {
80 static Config = z.object({})
81 private readonly waits = new Map<Agent, Map<ToolCallId, TimedQuestionWait>>()
82 private readonly queuedReplies = new WeakMap<Session, Map<ToolCallId, QueuedReply>>()
83
84 constructor(ctx: Context) {
85 super(ctx, 'userQuestions')
86 ctx.inject(['sessionProjections'], (projectionCtx) => {
87 projectionCtx.sessionProjections.register(userQuestionProjectionDefinition)
88 })
89 ctx.effect(() => () => {
90 for (const calls of this.waits.values()) {
91 for (const wait of calls.values()) wait.close(abortedQuestion())
92 }
93 this.waits.clear()
94 }, 'userQuestions: foreground waits')
95 ctx.on('agent/inbox/claimed', ({ agent, message, turn }) => {
96 const source = message.source
97 if (source.kind !== 'user-question-reply') return
98 const calls = this.queuedReplies.get(agent.session) ?? new Map<ToolCallId, QueuedReply>()
99 const reply = calls.get(source.callId)
100 if (reply === undefined) {
101 calls.set(source.callId, { messageId: message.id, claimedTurn: turn })
102 this.queuedReplies.set(agent.session, calls)
103 } else if (reply.messageId === message.id) {
104 reply.claimedTurn = turn
105 }
106 }, { global: true })
107 ctx.on('agent/inbox/discarded', ({ agent, message }) => {
108 const source = message.source
109 if (source.kind === 'user-question-reply') this.releaseReply(agent.session, source.callId, message.id)
110 }, { global: true })
111 ctx.on('session/event', (session, event) => {
112 if (event.type === 'user/message') {
113 const source = event.data.source
114 if (source.kind === 'user-question-reply') this.releaseReply(session, source.callId, event.data.id)
115 } else if (event.type === 'turn/end') {
116 const calls = this.queuedReplies.get(session)
117 if (calls === undefined) return
118 for (const [callId, reply] of calls) {
119 if (reply.claimedTurn === event.data.turn) calls.delete(callId)
120 }
121 }
122 }, { global: true })
123 }
124
125 private releaseReply(session: Session, callId: ToolCallId, messageId: MessageId): void {
126 const calls = this.queuedReplies.get(session)
127 if (calls?.get(callId)?.messageId !== messageId) return
128 calls.delete(callId)
129 }
130
131 private assertLiveRoot(agent: Agent): void {
132 const agents = this.ctx.get('agents')
133 if (agents === undefined || agents.get(agent.id) !== agent) {
134 throw new UserQuestionError(
135 'human interaction requires the exact live calling agent when an agent is supplied',
136 'CALLER_NOT_LIVE')
137 }
138 if (!agents.roots().includes(agent)) {
139 throw new UserQuestionError(
140 'human interaction is unavailable while the calling agent is owned by another live agent; '
141 + "include the unresolved question or decision in the child agent's final result",
142 'DELEGATED_CALLER')
143 }
144 }
145
146 private continued(agent: Agent): readonly PendingUserQuestion[] {
147 const state = this.ctx.get('sessionProjections')?.stateOf(agent.session, 'userQuestions')
148 return (state?.questions.active ?? []).filter(question => question.state === 'continued')
149 }
150
151 /**
152 * Answer a continued question. The reply is steered into the agent as a
153 * user message whose source names the call; that message is also the
154 * record that closes the question in the projection.
155 * @param agent - Live root agent for the owning Session.
156 * @param callId - Continued question identity.
157 * @param answer - Complete structured answer batch, one item per question of the call.
158 * @returns Whether the question is still continued; an accepted reply stays
159 * queued until the agent admits its user message.
160 * @throws {UserQuestionError} `BAD_ANSWER` when the batch does not name each
161 * question of the call exactly once, or `REPLY_QUEUED` when a reply is
162 * already waiting for admission.
163 */
164 @Remote
165 answer(agent: Agent, callId: ToolCallId, answer: AskUserQuestionAnswer): boolean {
166 this.assertLiveRoot(agent)
167 const question = this.continued(agent).find(item => item.callId === callId)
168 if (question === undefined) return false
169 const queued = this.queuedReplies.get(agent.session)
170 const matches = (message: UserMessage): boolean =>
171 message.source.kind === 'user-question-reply' && message.source.callId === callId
172 if (queued?.has(callId) || agent.inbox.nextTurn.some(matches) || agent.inbox.nextStep.some(matches)) {
173 throw new UserQuestionError('a reply is already queued for this question', 'REPLY_QUEUED')
174 }
175 // The gateway validated the batch's shape from the type; the model-facing
176 // contract also promises one item per question, which only this owner of
177 // the asked questions can check before the batch reaches the model.
178 const answered = new Set(answer.answers.map(item => item.id))
179 if (answered.size !== answer.answers.length
180 || question.questions.length !== answer.answers.length
181 || !question.questions.every(item => answered.has(item.id))) {
182 throw new UserQuestionError(
183 `the answer batch for ${callId} must name each of its ${String(question.questions.length)} questions exactly once`,
184 'BAD_ANSWER')
185 }
186 const message = createUserMessage({
187 source: { kind: 'user-question-reply', callId, outcome: 'answered' },
188 content: [{
189 type: 'text',
190 text: JSON.stringify({
191 kind: 'answer_to_pending_question', tool: 'ask_user_question', callId,
192 questions: question.questions, answers: answer.answers,
193 }),
194 }],
195 })
196 const calls = queued ?? new Map<ToolCallId, QueuedReply>()
197 calls.set(callId, { messageId: message.id })
198 this.queuedReplies.set(agent.session, calls)
199 try {
200 agent.steer(message)
201 } catch (error: unknown) {
202 this.releaseReply(agent.session, callId, message.id)
203 throw error
204 }
205 return true
206 }
207
208 /**
209 * Let one answer UI hold a live timed wait. Closing the stream releases its claim.
210 * @param agent - Live root agent owning the question.
211 * @param callId - Foreground tool call to attach to.
212 * @param signal - Remote stream cancellation, including Client disconnect.
213 * @returns One Host-computed remaining duration, or no frames once the wait ended.
214 */
215 @Remote({ mode: 'stream' })
216 async *attachWait(agent: Agent, callId: ToolCallId, signal: AbortSignal): AsyncIterable<{ remainingMs: number }> {
217 this.assertLiveRoot(agent)
218 const wait = this.waits.get(agent)?.get(callId)
219 if (wait !== undefined) yield* wait.attach(signal)
220 }
221
222 /**
223 * Foreground wait whose first settlement the Client decides: the Client
224 * rejects with `ASK_TIMED_OUT` when its countdown ends, and this method maps
225 * that code to the pending result.
226 * @param request - Questions, live owner agent, and abort signal.
227 * @param callId - Tool call identity the Client card is keyed by.
228 * @param timeoutMs - Positive foreground wait in milliseconds.
229 * @returns The answer when it arrives inside the window, otherwise a pending
230 * result, also when no connected Client claimed the request by the deadline.
231 * @throws {UserQuestionError} `BAD_TIMEOUT` for a non-integer, non-positive,
232 * or oversized wait.
233 */
234 async askTimed(
235 request: AskUserQuestionRequest & { agent: Agent },
236 callId: ToolCallId,
237 timeoutMs: number,
238 ): Promise<TimedUserQuestionResult> {
239 if (!Number.isSafeInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 2_147_483_647) {
240 throw new UserQuestionError('timeout must fit a positive platform timer', 'BAD_TIMEOUT')
241 }
242 this.assertLiveRoot(request.agent)
243 const calls = this.waits.get(request.agent) ?? new Map<ToolCallId, TimedQuestionWait>()
244 if (calls.has(callId)) throw new UserQuestionError('the question call already has a foreground wait', 'DUPLICATE_WAIT')
245 const wait = new TimedQuestionWait(Date.now() + timeoutMs, request.signal,
246 new UserQuestionError('ask_user_question timed out before the user answered', 'ASK_TIMED_OUT'))
247 calls.set(callId, wait)
248 this.waits.set(request.agent, calls)
249 try {
250 try {
251 return await this.ask({ ...request, signal: wait.signal, wait: { callId, timed: true } })
252 } catch (error) {
253 if (wait.signal.aborted) throw wait.signal.reason
254 if (error instanceof UserQuestionError && error.code === 'NO_PROVIDER') {
255 await wait.done
256 throw wait.signal.reason
257 }
258 throw error
259 }
260 } catch (error) {
261 if (error instanceof UserQuestionError && error.code === 'ASK_TIMED_OUT') return { pending: true, callId }
262 if (wait.signal.aborted) throw abortedQuestion(error)
263 throw error
264 } finally {
265 wait.close(abortedQuestion())
266 calls.delete(callId)
267 if (calls.size === 0) this.waits.delete(request.agent)
268 }
269 }
270
271 /**
272 * Ask the scoped answerer waterfall and wait for the user's answer.
273 *
274 * When a caller supplies an agent, human interaction is valid only for the
275 * exact live runtime root. Runtime ownership, not durable session lineage,
276 * decides this boundary: an owned child has no human answerer and would
277 * block forever, while a lineage-bearing session resumed as a new runtime
278 * root may ask normally.
279 *
280 * @param request Questions, owner agent, and abort signal.
281 * @returns The answer chosen or typed by the human.
282 * @throws {UserQuestionError} code `ASK_ABORTED` when the supplied signal
283 * is already or becomes aborted, `CALLER_NOT_LIVE` when a supplied agent
284 * is not the registry's exact live instance, or `DELEGATED_CALLER` when
285 * that live agent is owned by another agent.
286 */
287 async ask(request: AskUserQuestionRequest): Promise<AskUserQuestionAnswer> {
288 if (request.signal?.aborted) {
289 throw abortedQuestion()
290 }
291 if (request.questions.length === 0) {
292 throw new UserQuestionError('ask_user_question requires at least one question', 'EMPTY_QUESTIONS')
293 }
294 const agent = request.agent
295 if (agent !== undefined) this.assertLiveRoot(agent)
296 // A presentation intent asserts two things the types cannot: that the
297 // named approve label is one of this question's own options, and that a
298 // plan-review carries the plan it is a review of. A UI honouring the
299 // intent answers with that label, and shows that detail as the plan, so
300 // either gap would put a choice the asker never offered — or an approval of
301 // something invisible — in front of the user. Caught at the asker, where
302 // the mistake is, rather than in each UI.
303 for (const question of request.questions) {
304 const intent = question.intent
305 if (intent === undefined) continue
306 if (!(question.options ?? []).some(option => option.label === intent.approve)) {
307 throw new UserQuestionError(
308 `question ${question.id} declares intent ${intent.kind} whose approve label `
309 + `${JSON.stringify(intent.approve)} names none of its options`,
310 'BAD_INTENT')
311 }
312 if (question.detail === undefined) {
313 throw new UserQuestionError(
314 `question ${question.id} declares intent ${intent.kind} without the detail it reviews`,
315 'BAD_INTENT')
316 }
317 }
318 const noAnswerer = () => Promise.reject(new UserQuestionError(
319 'no user-questions answerer accepted the request',
320 'NO_PROVIDER',
321 ))
322 try {
323 return await (agent === undefined
324 ? this.ctx.waterfall('user-questions/request', request, noAnswerer)
325 : this.ctx.waterfall(
326 scopeTarget(agent, agent),
327 'user-questions/request',
328 { ...request, agent },
329 noAnswerer,
330 ))
331 } catch (error) {
332 const restored = restoreUserQuestionError(error)
333 if (restored instanceof UserQuestionError) throw restored
334 if (request.signal?.aborted) {
335 throw abortedQuestion(error)
336 }
337 throw restored
338 }
339 }
340}
341
342export default UserQuestionService