message-projections.spec.ts 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117
  1. import { afterEach, describe, expect, it } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  4. import { deepFreeze } from '@deepseek-ai/dsh-util-values'
  5. import SessionStore, { Session, SessionId, SessionSeq, SessionLogOffset, foldSurface, deriveEventMessage } from '../src/index.ts'
  6. import type { SessionEvent, SessionMessageProjection } from '../src/index.ts'
  7. import { MESSAGE_PROJECTION_EVENT_TYPES } from '../src/known-event-types.ts'
  8. import { SurfaceManager } from '../src/surface.ts'
  9. declare module '@deepseek-ai/dsh-session/types' {
  10. interface SessionEventMap {
  11. 'test/project': { seq: SessionSeq; text: string }
  12. }
  13. }
  14. const projection: SessionMessageProjection<'test/project'> = {
  15. type: 'test/project',
  16. project(event, context) {
  17. const source = context.events[event.data.seq - context.baseSeq]!
  18. const original = deriveEventMessage(source, context.messages)!
  19. if (event.data.text === 'reject') throw new Error('rejected decision')
  20. return new Map([[source.seq, deepFreeze({ ...original, content: [{ type: 'text' as const, text: event.data.text }] })]])
  21. },
  22. }
  23. const contexts: Context[] = []
  24. afterEach(async () => {
  25. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  26. })
  27. function input(session: Session) {
  28. return session.append('user/message', createUserMessage({ content: [{ type: 'text', text: 'original' }], source: { kind: 'user' } }), { surfaceOp: 'append' })
  29. }
  30. describe('plugin-owned message projections', () => {
  31. it('refuses required events when no interpreter is supplied', () => {
  32. const type = [...MESSAGE_PROJECTION_EVENT_TYPES][0]!
  33. const event = { type, seq: SessionSeq(0), time: 0, data: {} } as SessionEvent
  34. const session = Session.create(SessionId('missing'))
  35. expect(() => session.append(event.type as 'test/project', event.data as never)).toThrow(/requires a message projection/)
  36. expect(() => foldSurface([event])).toThrow(/requires a message projection/)
  37. expect(() => Session.create(session.id, [event])).toThrow(/requires a message projection/)
  38. expect(session.seq).toBe(0)
  39. })
  40. it('applies generic decisions atomically and replays them through detached folds', () => {
  41. const definitions = [projection]
  42. const session = Session.create(SessionId('pure'), undefined, undefined, undefined, definitions)
  43. const source = input(session)
  44. const original = session.deriveMessages()
  45. expect(() => session.append('test/project', { seq: source.seq, text: 'reject' })).toThrow('rejected decision')
  46. expect(session.deriveMessages()).toEqual(original)
  47. session.append('test/project', { seq: source.seq, text: 'projected' })
  48. expect(session.deriveMessages()[0]?.content).toEqual([{ type: 'text', text: 'projected' }])
  49. expect(session.surface.replaceGeneration).toBe(0)
  50. expect(session.surface.contentGeneration).toBe(1)
  51. const folded = foldSurface(session.snapshotEvents(), definitions)
  52. expect(deriveEventMessage(source, folded.projectedMessages)).toEqual(session.deriveMessages()[0])
  53. expect(source.data.content).toEqual([{ type: 'text', text: 'original' }])
  54. })
  55. it('registers by fiber and refuses cached or pending decisions after disposal', async () => {
  56. const ctx = new Context()
  57. contexts.push(ctx)
  58. await ctx.plugin(SessionStore)
  59. const fiber = await ctx.plugin({
  60. inject: ['sessions'],
  61. apply(owner: Context) { owner.sessions.registerMessageProjection(projection) },
  62. })
  63. expect(() => ctx.sessions.registerMessageProjection(projection)).toThrow(/already registered/)
  64. const live = ctx.sessions.create(SessionId('live'))
  65. const source = input(live)
  66. live.append('test/project', { seq: source.seq, text: 'changed' })
  67. const pending = ctx.sessions.create(SessionId('pending'))
  68. input(pending)
  69. pending.append('test/project', { seq: source.seq, text: 'changed' })
  70. const before = live.deriveMessages()
  71. const child = ctx.sessions.fork(live)
  72. expect(child.deriveMessages()).toEqual(before)
  73. const restored = ctx.sessions.prepare(SessionId('restore'), {
  74. seed: [...live.snapshotEvents()], meta: { ...live.header, id: SessionId('restore') },
  75. inheritedEventCount: live.inheritedEventCount, eventState: 'shared-frozen',
  76. })
  77. expect(restored.deriveMessages()).toEqual(before)
  78. await fiber.dispose()
  79. expect(ctx.sessions.messageProjections).toEqual([])
  80. for (const session of [live, pending, child, restored]) {
  81. expect(() => session.deriveMessages()).toThrow(/was removed or replaced/)
  82. expect(() => session.deriveEventMessage(source)).toThrow(/was removed or replaced/)
  83. expect(() => session.surface.replaceGeneration).toThrow(/was removed or replaced/)
  84. expect(() => session.surface.contentGeneration).toThrow(/was removed or replaced/)
  85. expect(() => session.append('turn/start', { turn: 1 })).toThrow(/was removed or replaced/)
  86. }
  87. expect(before[0]?.content).toEqual([{ type: 'text', text: 'changed' }])
  88. })
  89. it('applies a supplied interpreter to a loaded window with absolute sequences', () => {
  90. const session = Session.create(SessionId('window'), undefined, undefined, undefined, [projection])
  91. input(session)
  92. const source = input(session)
  93. session.append('test/project', { seq: source.seq, text: 'window' })
  94. const events = session.snapshotEvents().slice(1)
  95. const surface = new SurfaceManager(events, SessionLogOffset(1), [projection])
  96. expect(surface.nodes).toEqual([source.seq])
  97. expect(surface.deriveEventMessage(source)?.content).toEqual([{ type: 'text', text: 'window' }])
  98. })
  99. it('does not retain an interpreter for a candidate that never committed', () => {
  100. const session = Session.create(SessionId('candidate'))
  101. const source = input(session)
  102. const definitions: SessionMessageProjection[] = [projection]
  103. const surface = new SurfaceManager(session.snapshotEvents(), undefined, definitions)
  104. surface.validateNext({ type: 'test/project', seq: SessionSeq(1), time: 0, data: { seq: source.seq, text: 'unused' } })
  105. definitions.length = 0
  106. expect(surface.deriveEventMessage(source)).toBe(source.data)
  107. })
  108. })