返回源码地图

packages/feedback/message-feedback/src/index.ts

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

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

1/**
2 * Canonical Session-log feedback for finalized assistant messages.
3 * @module @deepseek-ai/dsh-message-feedback
4 */
5
6import { Buffer } from 'node:buffer'
7import { randomUUID } from 'node:crypto'
8import { isDeepStrictEqual } from 'node:util'
9import { Context, Service } from '@deepseek-ai/cordis'
10import s from '@deepseek-ai/schemastery'
11import { z } from 'zod'
12import { FEEDBACK_CATEGORIES } from '@deepseek-ai/dsh-command-feedback'
13import { SessionSeq } from '@deepseek-ai/dsh-session/types'
14import { deriveEventMessage, isAppendSurfaceEvent } from '@deepseek-ai/dsh-session/surface'
15import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session/types'
16import type {} from '@deepseek-ai/dsh-session'
17import type { SessionInspection } from '@deepseek-ai/dsh-session-persistence'
18import { TypertRemoteService, Remote } from '@deepseek-ai/dsh-typert-protocol'
19import type {
20 MessageFeedbackDeleteRequest,
21 MessageFeedbackDeleteResult,
22 MessageFeedbackFailure,
23 MessageFeedbackItem,
24 MessageFeedbackListRequest,
25 MessageFeedbackListResult,
26 MessageFeedbackNoteBlank,
27 MessageFeedbackNoteTooLarge,
28 MessageFeedbackPutRequest,
29 MessageFeedbackPutResult,
30 MessageFeedbackRejected,
31 MessageFeedbackSessionNotFound,
32 MessageFeedbackSuccess,
33 MessageFeedbackVersion,
34 MessageFeedbackVersionConflict,
35} from './types.ts'
36
37export type * from './types.ts'
38
39/** Required deployment policy for optional notes. */
40export interface Config {
41 /** Maximum UTF-8 byte length accepted for one note. */
42 readonly maxNoteBytes: number
43}
44
45declare module '@deepseek-ai/cordis' {
46 interface Context {
47 messageFeedback: MessageFeedbackService
48 }
49 interface Events {
50 /**
51 * Observe a durable cold feedback mutation without publishing a live Session.
52 * Observers run before write ownership is released and must not await
53 * another message-feedback operation for this Session. The payload is borrowed
54 * read-only; deep-clone it before transferring ownership (for example, to Session.fromRestore).
55 * @param inspection - committed canonical prefix, including the feedback as its last event.
56 * @mode parallel
57 */
58 'feedback/committed'(inspection: SessionInspection): void
59 }
60}
61
62const timestamp = z.number().int().nonnegative().max(Number.MAX_SAFE_INTEGER)
63const itemSchema = z.object({
64 messageId: z.string().min(1),
65 rating: z.enum(['positive', 'negative']),
66 note: z.string().refine(note => note.trim().length > 0).optional(),
67 category: z.enum(FEEDBACK_CATEGORIES).optional(),
68 version: z.uuid(),
69 createdAt: timestamp,
70 updatedAt: timestamp,
71}).refine(item => item.updatedAt >= item.createdAt)
72const putSchema = z.object({ sessionId: z.string().min(1), item: itemSchema })
73const deleteSchema = z.object({ sessionId: z.string().min(1), messageId: z.string().min(1) })
74
75type FeedbackEvent = SessionEvent<'feedback/message-put' | 'feedback/message-delete'>
76type Mutation = Pick<SessionEvent<'feedback/message-put'>, 'type' | 'data'>
77 | Pick<SessionEvent<'feedback/message-delete'>, 'type' | 'data'>
78type Append = (event?: Mutation) => Promise<void>
79type MissingSession = MessageFeedbackRejected<MessageFeedbackSessionNotFound>
80type ResolvedNote = MessageFeedbackSuccess<string | undefined>
81 | MessageFeedbackRejected<MessageFeedbackNoteBlank | MessageFeedbackNoteTooLarge>
82
83/** Return a caller-owned immutable value, detached from the log. */
84function snapshotItem(item: MessageFeedbackItem): MessageFeedbackItem {
85 return Object.freeze({ ...item })
86}
87
88function success<T>(value: T): MessageFeedbackSuccess<T> {
89 return Object.freeze({ ok: true, value })
90}
91
92function rejected<E extends MessageFeedbackFailure>(error: E): MessageFeedbackRejected<E> {
93 return Object.freeze({ ok: false, error: Object.freeze(error) })
94}
95
96/** Validate persisted payloads before deriving current, Session-owned feedback. */
97function currentItems(sessionId: SessionId, events: readonly SessionEvent[]): MessageFeedbackItem[] {
98 const items = new Map<MessageFeedbackItem['messageId'], MessageFeedbackItem>()
99 for (const event of events) {
100 switch (event.type) {
101 case 'feedback/message-put':
102 putSchema.parse(event.data)
103 if (event.data.sessionId === sessionId) items.set(event.data.item.messageId, event.data.item)
104 break
105 case 'feedback/message-delete':
106 deleteSchema.parse(event.data)
107 if (event.data.sessionId === sessionId) items.delete(event.data.messageId)
108 break
109 default:
110 // Other plugins' events do not change message feedback.
111 break
112 }
113 }
114 return [...items.values()]
115}
116
117/** Session-log service; cold operations never construct a Session or Agent. */
118export class MessageFeedbackService extends TypertRemoteService {
119 static inject = ['sessionPersistence', 'sessions']
120
121 /** Loader validation for the required note-size policy. */
122 static Config: s<Config> = s.object({
123 maxNoteBytes: s.number().step(1).min(1).required(),
124 })
125
126 private readonly maxNoteBytes: number
127 private readonly operationTails = new Map<SessionId, Promise<void>>()
128 private mutationAdmissionOpen = true
129
130 /**
131 * @param ctx - Host context carrying Session persistence and live owners.
132 * @param config - Required note-size policy.
133 */
134 constructor(ctx: Context, config: Config) {
135 super(ctx, 'messageFeedback')
136 if (!Number.isSafeInteger(config.maxNoteBytes) || config.maxNoteBytes < 1) {
137 throw new TypeError('message-feedback: maxNoteBytes must be a positive safe integer')
138 }
139 this.maxNoteBytes = config.maxNoteBytes
140 }
141
142 protected [Service.init](): void {
143 this.ctx.effect(() => async () => {
144 this.mutationAdmissionOpen = false
145 await Promise.all(this.operationTails.values())
146 }, 'message-feedback.drain')
147 }
148
149 /**
150 * Read current feedback from the canonical log.
151 * @param request - Session to inspect.
152 * @returns immutable items or a definite persistence miss.
153 */
154 @Remote('list')
155 list(request: MessageFeedbackListRequest): Promise<MessageFeedbackListResult> {
156 return this.enqueue(request.sessionId, () => this.withSession(request.sessionId, false, events =>
157 success(Object.freeze({ items: Object.freeze(currentItems(request.sessionId, events).map(snapshotItem)) }))))
158 }
159
160 /**
161 * Create or replace feedback after checking its current version.
162 * Matching no-ops retain the version and append no event.
163 * @param request - Target, desired value, and observed item version.
164 * @returns the durable item or an explicit business failure.
165 */
166 @Remote('put')
167 put(request: MessageFeedbackPutRequest): Promise<MessageFeedbackPutResult> {
168 const note = this.resolveNote(request.note)
169 if (!note.ok) return Promise.resolve(note)
170 return this.enqueue(request.sessionId, () => this.withSession(request.sessionId, true, async (events, append) => {
171 const items = currentItems(request.sessionId, events)
172 if (!events.some(event => event.type === 'assistant/message'
173 && isAppendSurfaceEvent(event)
174 && deriveEventMessage(event)?.id === request.messageId)) {
175 return rejected({ code: 'target-not-found', sessionId: request.sessionId, messageId: request.messageId })
176 }
177 const existing = items.find(item => item.messageId === request.messageId)
178 if (request.ifVersion !== (existing?.version ?? null)) {
179 return rejected(this.versionConflict(existing ?? null))
180 }
181 if (existing !== undefined && existing.rating === request.rating && existing.note === note.value
182 && existing.category === request.category) {
183 await append()
184 return success(snapshotItem(existing))
185 }
186 const now = Date.now()
187 const item: MessageFeedbackItem = {
188 messageId: request.messageId,
189 rating: request.rating,
190 ...(note.value === undefined ? {} : { note: note.value }),
191 ...(request.category === undefined ? {} : { category: request.category }),
192 version: randomUUID() as MessageFeedbackVersion,
193 createdAt: existing?.createdAt ?? now,
194 updatedAt: existing === undefined ? now : Math.max(now, existing.updatedAt),
195 }
196 await append({ type: 'feedback/message-put', data: { sessionId: request.sessionId, item } })
197 return success(snapshotItem(item))
198 }))
199 }
200
201 /**
202 * Delete one item after checking its version; absence succeeds without an event.
203 * @param request - Session, message, and observed item version.
204 * @returns the stable absent postcondition or an explicit failure.
205 */
206 @Remote('delete')
207 delete(request: MessageFeedbackDeleteRequest): Promise<MessageFeedbackDeleteResult> {
208 return this.enqueue(request.sessionId, () => this.withSession(request.sessionId, true, async (events, append) => {
209 const existing = currentItems(request.sessionId, events).find(item => item.messageId === request.messageId)
210 if (existing !== undefined) {
211 if (request.ifVersion !== existing.version) return rejected(this.versionConflict(existing))
212 await append({ type: 'feedback/message-delete', data: { sessionId: request.sessionId, messageId: request.messageId } })
213 } else {
214 await append()
215 }
216 return success(Object.freeze({ absent: true as const }))
217 }))
218 }
219
220 /** Hold cold write ownership across read/compare/append; use live owners directly. */
221 private async withSession<T>(
222 sessionId: SessionId,
223 write: boolean,
224 operation: (events: readonly SessionEvent[], append: Append) => T | Promise<T>,
225 ): Promise<T | MissingSession> {
226 if (this.ctx.sessions.get(sessionId) === undefined
227 && await this.ctx.sessionPersistence.stat(sessionId) === undefined
228 && this.ctx.sessions.get(sessionId) === undefined) {
229 return rejected({ code: 'session-not-found', sessionId })
230 }
231 const live = this.ctx.sessions.get(sessionId)
232 if (live !== undefined) {
233 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
234 return operation(live.snapshotEvents(), async (event) => {
235 if (event !== undefined) {
236 live.append(event.type, event.data)
237 }
238 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
239 const last = live.snapshotEvents().at(-1)
240 if (!(await this.ctx.sessions.flush(live))) {
241 throw new Error(
242 `message-feedback: no durability listener participated for live session '${sessionId}'`,
243 )
244 }
245 // Listener participation alone does not prove this Session has a persistence writer.
246 const handle = await this.ctx.sessionPersistence.open(sessionId, 'read')
247 try {
248 const { events: stored } = await handle.read(last?.seq ?? 0, 1)
249 if (!isDeepStrictEqual(
250 [handle.header.id, handle.header.createdAt, handle.header.cwd],
251 [live.header.id, live.header.createdAt, live.header.cwd],
252 )
253 || (last !== undefined && !isDeepStrictEqual(stored[0], last))) {
254 throw new Error(`message-feedback: feedback prefix is not durable for live session '${sessionId}'`)
255 }
256 } finally {
257 await handle.close()
258 }
259 })
260 }
261 const handle = await this.ctx.sessionPersistence.open(sessionId, write ? 'write' : 'read')
262 try {
263 const { events } = await handle.read()
264 return await operation(events, async (event) => {
265 const entry: FeedbackEvent | undefined = event === undefined ? undefined
266 : { ...event, seq: SessionSeq(events.length), time: Date.now() }
267 if (entry !== undefined) await handle.append([entry])
268 await handle.flush()
269 if (entry !== undefined) {
270 try {
271 await this.ctx.parallel('feedback/committed', {
272 meta: handle.header,
273 inheritedEventCount: handle.inheritedEventCount,
274 events: [...events, entry],
275 })
276 } catch (error) {
277 this.ctx.logger.warn('message-feedback: committed feedback observer failed', error)
278 }
279 }
280 })
281 } finally {
282 await handle.close()
283 }
284 }
285
286 private resolveNote(note: string | undefined): ResolvedNote {
287 if (note === undefined) return success(undefined)
288 if (note.trim().length === 0) return rejected({ code: 'note-blank' })
289 const actualBytes = Buffer.byteLength(note, 'utf8')
290 if (actualBytes > this.maxNoteBytes) {
291 return rejected({ code: 'note-too-large', maxBytes: this.maxNoteBytes, actualBytes })
292 }
293 return success(note)
294 }
295
296 private versionConflict(current: MessageFeedbackItem | null): MessageFeedbackVersionConflict {
297 return { code: 'version-conflict', current: current === null ? null : snapshotItem(current) }
298 }
299
300 /** Serialize complete operations and drain their handles before disposal. */
301 private enqueue<T>(sessionId: SessionId, operation: () => Promise<T>): Promise<T> {
302 if (!this.mutationAdmissionOpen) return Promise.reject(new Error('message-feedback: service is disposing'))
303 const previous = this.operationTails.get(sessionId) ?? Promise.resolve()
304 const result = previous.then(operation)
305 const tail = result.then(() => undefined, () => undefined)
306 this.operationTails.set(sessionId, tail)
307 return result.finally(() => {
308 if (this.operationTails.get(sessionId) === tail) this.operationTails.delete(sessionId)
309 })
310 }
311}
312
313export default MessageFeedbackService