helpers.ts 7.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. import { mkdtemp, rm } from 'node:fs/promises'
  2. import { tmpdir } from 'node:os'
  3. import { join } from 'node:path'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import { createAssistantMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  6. import type { MessageId } from '@deepseek-ai/dsh-llm/brand'
  7. import SessionStore, {
  8. SESSION_FORMAT_VERSION,
  9. Session,
  10. SessionId,
  11. type SessionEvent,
  12. type SessionHeader,
  13. } from '@deepseek-ai/dsh-session'
  14. import SessionPersistence, {
  15. SessionPersistenceRevision,
  16. type SessionInspection,
  17. type SessionLocation,
  18. type SessionPersistenceSnapshot,
  19. } from '@deepseek-ai/dsh-session-persistence'
  20. import Storage from '@deepseek-ai/dsh-storage'
  21. import * as StorageDomain from '@deepseek-ai/dsh-storage-domain'
  22. import * as StorageJson from '@deepseek-ai/dsh-storage-json'
  23. import MessageFeedbackService from '../src/index.ts'
  24. export interface MessageFixture {
  25. readonly session: Session
  26. readonly userMessageId: MessageId
  27. readonly assistantMessageIds: readonly [MessageId, MessageId]
  28. readonly emptyAssistantMessageId: MessageId
  29. readonly replacementAssistantMessageId: MessageId
  30. }
  31. /** Append one deterministic transcript used by target-validation tests. */
  32. export function appendMessageFixture(session: Session): Omit<MessageFixture, 'session'> {
  33. session.append('turn/start', { turn: 1 })
  34. session.append('step/start', { turn: 1, step: 1 })
  35. const user = createUserMessage({
  36. content: [{ type: 'text', text: 'Question' }],
  37. source: { kind: 'user' },
  38. })
  39. session.append('user/message', user, { surfaceOp: 'append' })
  40. const first = createAssistantMessage({
  41. content: [{ type: 'text', text: 'First answer' }],
  42. source: { provider: 'test', model: 'test' },
  43. })
  44. const firstEvent = session.append('assistant/message', {
  45. turn: 1,
  46. step: 1,
  47. message: first,
  48. }, { surfaceOp: 'append' })
  49. const second = createAssistantMessage({
  50. content: [{ type: 'text', text: 'Second answer' }],
  51. source: { provider: 'test', model: 'test' },
  52. })
  53. session.append('assistant/message', {
  54. turn: 1,
  55. step: 1,
  56. message: second,
  57. }, { surfaceOp: 'append' })
  58. const empty = createAssistantMessage({
  59. content: [],
  60. source: { provider: 'test', model: 'test' },
  61. })
  62. session.append('assistant/message', {
  63. turn: 1,
  64. step: 1,
  65. message: empty,
  66. }, { surfaceOp: 'append' })
  67. session.append('step/end', { turn: 1, step: 1 })
  68. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  69. const replacement = createAssistantMessage({
  70. content: [{ type: 'text', text: 'Model-only replacement' }],
  71. source: { provider: 'test', model: 'test' },
  72. })
  73. session.append('assistant/message', {
  74. turn: 1,
  75. step: 1,
  76. message: replacement,
  77. }, {
  78. surfaceOp: { op: 'replace', start: firstEvent.seq, end: firstEvent.seq },
  79. sourceEventSeqs: [firstEvent.seq],
  80. })
  81. return {
  82. userMessageId: user.id,
  83. assistantMessageIds: [first.id, second.id],
  84. emptyAssistantMessageId: empty.id,
  85. replacementAssistantMessageId: replacement.id,
  86. }
  87. }
  88. /** Construct one cold persistence fixture without publishing a live Session. */
  89. export function messageFixture(
  90. rawId: string,
  91. options: { readonly createdAt?: number; readonly cwd?: string } = {},
  92. ): MessageFixture {
  93. const id = SessionId(rawId)
  94. const header: SessionHeader = {
  95. version: SESSION_FORMAT_VERSION,
  96. id,
  97. createdAt: options.createdAt ?? 1_700_000_000_000,
  98. ...(options.cwd === undefined ? {} : { cwd: options.cwd }),
  99. }
  100. const session = Session.create(id, [], header)
  101. return { session, ...appendMessageFixture(session) }
  102. }
  103. /** Minimal controllable persistence provider for service-level tests. */
  104. class TestPersistence extends SessionPersistence {
  105. override readonly supportsRawArtifacts = false
  106. static inject = ['sessions']
  107. readonly durable = new Map<SessionId, SessionInspection>()
  108. readonly logical = new Map<SessionId, SessionInspection>()
  109. inspectFailure: Error | undefined
  110. inspectCalls = 0
  111. readFromCalls = 0
  112. onReadFrom: (() => void | Promise<void>) | undefined
  113. onListSnapshots: (() => void | Promise<void>) | undefined
  114. locate(_meta: SessionHeader): SessionLocation | undefined { return undefined }
  115. create(_meta: SessionHeader): Promise<void> { return Promise.resolve() }
  116. append(_id: SessionId, _events: readonly SessionEvent[]): Promise<void> { return Promise.resolve() }
  117. load(id: SessionId): Promise<SessionInspection> {
  118. return this.readFrom(id, 0)
  119. }
  120. inspect(id: SessionId): Promise<SessionInspection> {
  121. this.inspectCalls += 1
  122. if (this.inspectFailure !== undefined) return Promise.reject(this.inspectFailure)
  123. const explicit = this.logical.get(id)
  124. if (explicit !== undefined) return Promise.resolve(explicit)
  125. const live = this.ctx.sessions.get(id)
  126. if (live !== undefined) return Promise.resolve({ meta: live.header, events: live.events })
  127. const stored = this.durable.get(id)
  128. return stored === undefined
  129. ? Promise.reject(new Error(`test persistence: session '${id}' not found`))
  130. : Promise.resolve(stored)
  131. }
  132. borrowSession(_id: SessionId, _signal?: AbortSignal): ReturnType<SessionPersistence['borrowSession']> {
  133. return Promise.reject(new Error('not used'))
  134. }
  135. async readFrom(
  136. id: SessionId,
  137. fromSeq: number,
  138. ): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  139. this.readFromCalls += 1
  140. await this.onReadFrom?.()
  141. const stored = this.durable.get(id)
  142. return stored === undefined
  143. ? Promise.reject(new Error(`test persistence: session '${id}' not found`))
  144. : { meta: stored.meta, events: stored.events.filter(event => event.seq >= fromSeq) }
  145. }
  146. list(): Promise<SessionHeader[]> {
  147. return Promise.resolve([...this.durable.values()].map(value => value.meta))
  148. }
  149. async listSnapshots(): Promise<SessionPersistenceSnapshot[]> {
  150. await this.onListSnapshots?.()
  151. return [...this.durable.values()].map((value, index) => ({
  152. header: value.meta,
  153. revision: SessionPersistenceRevision(`test:${index}:${value.events.length}`),
  154. }))
  155. }
  156. persist(session: Session): void {
  157. this.durable.set(session.id, { meta: session.header, events: session.events })
  158. }
  159. setDurable(inspection: SessionInspection): void {
  160. this.durable.set(inspection.meta.id, inspection)
  161. }
  162. }
  163. export interface TestHarness {
  164. readonly ctx: Context
  165. readonly persistence: TestPersistence
  166. readonly root: string
  167. disposeFeedback(): Promise<void>
  168. dispose(): Promise<void>
  169. }
  170. /** Compose the service over the real storage hub/domain/JSON backend. */
  171. export async function setupHarness(maxNoteBytes = 64): Promise<TestHarness> {
  172. const root = await mkdtemp(join(tmpdir(), 'dsh-message-feedback-test-'))
  173. const ctx = new Context()
  174. let disposeFeedback: (() => Promise<void>) | undefined
  175. try {
  176. await ctx.plugin(SessionStore)
  177. await ctx.plugin(TestPersistence)
  178. await ctx.plugin(Storage)
  179. await ctx.plugin(StorageJson, { root })
  180. await ctx.plugin(StorageDomain, { backend: 'json' })
  181. const feedbackFiber = await ctx.plugin(MessageFeedbackService, { maxNoteBytes })
  182. disposeFeedback = feedbackFiber.dispose
  183. } catch (error) {
  184. await ctx.fiber.dispose()
  185. await rm(root, { recursive: true, force: true })
  186. throw error
  187. }
  188. if (disposeFeedback === undefined) throw new Error('message feedback test plugin did not load')
  189. return {
  190. ctx,
  191. persistence: ctx.sessionPersistence as unknown as TestPersistence,
  192. root,
  193. disposeFeedback,
  194. async dispose() {
  195. await ctx.fiber.dispose()
  196. await rm(root, { recursive: true, force: true })
  197. },
  198. }
  199. }