1
/**2
* Canonical Session-log feedback for finalized assistant messages.3
* @module @deepseek-ai/dsh-message-feedback4
*/6
import { Buffer } from 'node:buffer'7
import { randomUUID } from 'node:crypto'8
import { isDeepStrictEqual } from 'node:util'9
import { Context, Service } from '@deepseek-ai/cordis'10
import s from '@deepseek-ai/schemastery'11
import { z } from 'zod'12
import { FEEDBACK_CATEGORIES } from '@deepseek-ai/dsh-command-feedback'13
import { SessionSeq } from '@deepseek-ai/dsh-session/types'14
import { deriveEventMessage, isAppendSurfaceEvent } from '@deepseek-ai/dsh-session/surface'15
import type { SessionEvent, SessionId } from '@deepseek-ai/dsh-session/types'16
import type {} from '@deepseek-ai/dsh-session'17
import type { SessionInspection } from '@deepseek-ai/dsh-session-persistence'18
import { TypertRemoteService, Remote } from '@deepseek-ai/dsh-typert-protocol'19
import 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'37
export type * from './types.ts'39
/** Required deployment policy for optional notes. */40
export interface Config {41
/** Maximum UTF-8 byte length accepted for one note. */42
readonly maxNoteBytes: number43
}45
declare module '@deepseek-ai/cordis' {46
interface Context {47
messageFeedback: MessageFeedbackService48
}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 await53
* another message-feedback operation for this Session. The payload is borrowed54
* 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 parallel57
*/58
'feedback/committed'(inspection: SessionInspection): void59
}60
}62
const timestamp = z.number().int().nonnegative().max(Number.MAX_SAFE_INTEGER)63
const 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)72
const putSchema = z.object({ sessionId: z.string().min(1), item: itemSchema })73
const deleteSchema = z.object({ sessionId: z.string().min(1), messageId: z.string().min(1) })75
type FeedbackEvent = SessionEvent<'feedback/message-put' | 'feedback/message-delete'>76
type Mutation = Pick<SessionEvent<'feedback/message-put'>, 'type' | 'data'>77
| Pick<SessionEvent<'feedback/message-delete'>, 'type' | 'data'>78
type Append = (event?: Mutation) => Promise<void>79
type MissingSession = MessageFeedbackRejected<MessageFeedbackSessionNotFound>80
type ResolvedNote = MessageFeedbackSuccess<string | undefined>81
| MessageFeedbackRejected<MessageFeedbackNoteBlank | MessageFeedbackNoteTooLarge>83
/** Return a caller-owned immutable value, detached from the log. */84
function snapshotItem(item: MessageFeedbackItem): MessageFeedbackItem {85
return Object.freeze({ ...item })86
}88
function success<T>(value: T): MessageFeedbackSuccess<T> {89
return Object.freeze({ ok: true, value })90
}92
function rejected<E extends MessageFeedbackFailure>(error: E): MessageFeedbackRejected<E> {93
return Object.freeze({ ok: false, error: Object.freeze(error) })94
}96
/** Validate persisted payloads before deriving current, Session-owned feedback. */97
function 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
break105
case 'feedback/message-delete':106
deleteSchema.parse(event.data)107
if (event.data.sessionId === sessionId) items.delete(event.data.messageId)108
break109
default:110
// Other plugins' events do not change message feedback.111
break112
}113
}114
return [...items.values()]115
}117
/** Session-log service; cold operations never construct a Session or Agent. */118
export class MessageFeedbackService extends TypertRemoteService {119
static inject = ['sessionPersistence', 'sessions']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
})126
private readonly maxNoteBytes: number127
private readonly operationTails = new Map<SessionId, Promise<void>>()128
private mutationAdmissionOpen = true130
/**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.maxNoteBytes140
}142
protected [Service.init](): void {143
this.ctx.effect(() => async () => {144
this.mutationAdmissionOpen = false145
await Promise.all(this.operationTails.values())146
}, 'message-feedback.drain')147
}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
}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.value182
&& 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
}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
}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) === undefined227
&& await this.ctx.sessionPersistence.stat(sessionId) === undefined228
&& 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 ? undefined266
: { ...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
}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
}296
private versionConflict(current: MessageFeedbackItem | null): MessageFeedbackVersionConflict {297
return { code: 'version-conflict', current: current === null ? null : snapshotItem(current) }298
}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
}313
export default MessageFeedbackService