helpers.ts 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231
  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, { SessionLogOffset,
  8. SESSION_FORMAT_VERSION,
  9. Session,
  10. SessionId,
  11. type SessionEvent,
  12. type SessionHeader,
  13. } from '@deepseek-ai/dsh-session'
  14. import SessionPersistence, {
  15. SessionAlreadyExistsError,
  16. SessionHandleClosedError,
  17. SessionPersistenceNotFoundError,
  18. SessionPersistenceRevision,
  19. SessionReadOnlyError,
  20. type SessionAccess,
  21. type SessionHandle,
  22. type SessionPersistenceSnapshot,
  23. } from '@deepseek-ai/dsh-session-persistence'
  24. import MessageFeedbackService from '../src/index.ts'
  25. export interface MessageFixture {
  26. readonly session: Session
  27. readonly userMessageId: MessageId
  28. readonly assistantMessageIds: readonly [MessageId, MessageId]
  29. readonly emptyAssistantMessageId: 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. session.append('assistant/message', {
  45. stream: [],
  46. turn: 1,
  47. step: 1,
  48. message: first,
  49. }, { surfaceOp: 'append' })
  50. const second = createAssistantMessage({
  51. content: [{ type: 'text', text: 'Second answer' }],
  52. source: { provider: 'test', model: 'test' },
  53. })
  54. session.append('assistant/message', {
  55. stream: [],
  56. turn: 1,
  57. step: 1,
  58. message: second,
  59. }, { surfaceOp: 'append' })
  60. const empty = createAssistantMessage({
  61. content: [],
  62. source: { provider: 'test', model: 'test' },
  63. })
  64. session.append('assistant/message', {
  65. stream: [],
  66. turn: 1,
  67. step: 1,
  68. message: empty,
  69. }, { surfaceOp: 'append' })
  70. session.append('step/end', { turn: 1, step: 1 })
  71. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  72. return {
  73. userMessageId: user.id,
  74. assistantMessageIds: [first.id, second.id],
  75. emptyAssistantMessageId: empty.id,
  76. }
  77. }
  78. /** Construct one cold persistence fixture without publishing a live Session. */
  79. export function messageFixture(
  80. rawId: string,
  81. options: { readonly createdAt?: number; readonly cwd?: string } = {},
  82. ): MessageFixture {
  83. const id = SessionId(rawId)
  84. const header: SessionHeader = {
  85. version: SESSION_FORMAT_VERSION,
  86. id,
  87. createdAt: options.createdAt ?? 1_700_000_000_000,
  88. isSeeded: false,
  89. ...(options.cwd === undefined ? {} : { cwd: options.cwd }),
  90. }
  91. const session = Session.create(id, [], header)
  92. return { session, ...appendMessageFixture(session) }
  93. }
  94. /** One stored session in the in-memory test backend. */
  95. interface StoredSession {
  96. readonly meta: SessionHeader
  97. events: readonly SessionEvent[]
  98. }
  99. /** Minimal controllable persistence provider for service-level tests. */
  100. class TestPersistence extends SessionPersistence {
  101. readonly durable = new Map<SessionId, StoredSession>()
  102. readFailure: Error | undefined
  103. appendFailure: Error | undefined
  104. flushFailure: Error | undefined
  105. openCalls: SessionAccess[] = []
  106. closeCalls = 0
  107. appendCalls = 0
  108. statCalls = 0
  109. readCalls = 0
  110. onRead: (() => void | Promise<void>) | undefined
  111. onStat: (() => void | Promise<void>) | undefined
  112. async create(header: SessionHeader): Promise<SessionHandle> {
  113. if (this.durable.has(header.id)) throw new SessionAlreadyExistsError(header.id)
  114. const stored: StoredSession = { meta: header, events: [] }
  115. this.durable.set(header.id, stored)
  116. return this.handle(stored, 'write')
  117. }
  118. // Appends are durable on resolution here; nothing buffers, so the service-wide flush is a no-op.
  119. async flush(): Promise<void> {}
  120. async open(id: SessionId, access: SessionAccess): Promise<SessionHandle> {
  121. this.openCalls.push(access)
  122. const stored = this.durable.get(id)
  123. if (stored === undefined) throw new SessionPersistenceNotFoundError(id)
  124. return this.handle(stored, access)
  125. }
  126. async stat(id: SessionId): Promise<SessionPersistenceSnapshot | undefined> {
  127. this.statCalls += 1
  128. await this.onStat?.()
  129. const stored = this.durable.get(id)
  130. if (stored === undefined) return undefined
  131. return { header: stored.meta, revision: SessionPersistenceRevision(`test:${id}:${stored.events.length}`) }
  132. }
  133. async list(): Promise<SessionPersistenceSnapshot[]> {
  134. return [...this.durable.entries()].map(([id, stored]) => ({
  135. header: stored.meta,
  136. revision: SessionPersistenceRevision(`test:${id}:${stored.events.length}`),
  137. }))
  138. }
  139. private handle(stored: StoredSession, access: SessionAccess): SessionHandle {
  140. let closed = false
  141. const handle: SessionHandle = {
  142. id: stored.meta.id,
  143. header: stored.meta,
  144. inheritedEventCount: SessionLogOffset(0),
  145. access,
  146. read: async (offset = 0, length?: number) => {
  147. if (closed) throw new SessionHandleClosedError(stored.meta.id, 'read')
  148. this.readCalls += 1
  149. if (this.readFailure !== undefined) throw this.readFailure
  150. await this.onRead?.()
  151. const events = stored.events.filter(event => event.seq >= offset)
  152. return {
  153. eventState: 'detached',
  154. events: structuredClone(length === undefined ? events : events.slice(0, length)),
  155. }
  156. },
  157. append: async (events) => {
  158. if (closed) throw new SessionHandleClosedError(stored.meta.id, 'append')
  159. if (access !== 'write') throw new SessionReadOnlyError(stored.meta.id, 'append')
  160. if (this.appendFailure !== undefined) throw this.appendFailure
  161. this.appendCalls += 1
  162. stored.events = [...stored.events, ...events]
  163. },
  164. flush: async () => {
  165. if (closed) throw new SessionHandleClosedError(stored.meta.id, 'flush')
  166. if (access !== 'write') throw new SessionReadOnlyError(stored.meta.id, 'flush')
  167. if (this.flushFailure !== undefined) throw this.flushFailure
  168. },
  169. close: async () => { closed = true; this.closeCalls += 1 },
  170. [Symbol.asyncDispose]() { return handle.close() },
  171. }
  172. return handle
  173. }
  174. persist(session: Session): void {
  175. this.durable.set(session.id, { meta: session.header, events: [...session.snapshotEvents()] })
  176. }
  177. setDurable(stored: StoredSession): void {
  178. this.durable.set(stored.meta.id, stored)
  179. }
  180. }
  181. export interface TestHarness {
  182. readonly ctx: Context
  183. readonly persistence: TestPersistence
  184. readonly root: string
  185. disposeFeedback(): Promise<void>
  186. dispose(): Promise<void>
  187. }
  188. /** Compose feedback over a controllable Session persistence backend. */
  189. export async function setupHarness(maxNoteBytes = 64): Promise<TestHarness> {
  190. const root = await mkdtemp(join(tmpdir(), 'dsh-message-feedback-test-'))
  191. const ctx = new Context()
  192. let disposeFeedback: (() => Promise<void>) | undefined
  193. try {
  194. await ctx.plugin(SessionStore)
  195. await ctx.plugin(TestPersistence)
  196. const feedbackFiber = await ctx.plugin(MessageFeedbackService, { maxNoteBytes })
  197. disposeFeedback = feedbackFiber.dispose
  198. } catch (error) {
  199. await ctx.fiber.dispose()
  200. await rm(root, { recursive: true, force: true })
  201. throw error
  202. }
  203. if (disposeFeedback === undefined) throw new Error('message feedback test plugin did not load')
  204. return {
  205. ctx,
  206. persistence: ctx.sessionPersistence as unknown as TestPersistence,
  207. root,
  208. disposeFeedback,
  209. async dispose() {
  210. await ctx.fiber.dispose()
  211. await rm(root, { recursive: true, force: true })
  212. },
  213. }
  214. }