upload.spec.ts 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. import { afterEach, describe, expect, it } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import SessionStore, { Session, SessionId, type CreateSessionOptions, type SessionEvent } from '@deepseek-ai/dsh-session'
  4. import DeepSeekLlmApiExtensionRegistry from '@deepseek-ai/dsh-deepseek-llm-api-extensions'
  5. import * as SessionLogDeepSeek from '../src/index.ts'
  6. const contexts: Context[] = []
  7. const SIGNAL = new AbortController().signal
  8. afterEach(async () => {
  9. await Promise.all(contexts.splice(0).map(ctx => ctx.fiber.dispose()))
  10. })
  11. async function harness(id: string, seed?: readonly SessionEvent[], meta?: CreateSessionOptions['meta']): Promise<{
  12. ctx: Context
  13. session: Session
  14. disposeUpload: () => Promise<void>
  15. }> {
  16. const ctx = new Context()
  17. contexts.push(ctx)
  18. await ctx.plugin(SessionStore)
  19. await ctx.plugin(DeepSeekLlmApiExtensionRegistry)
  20. const upload = ctx.plugin(SessionLogDeepSeek, { enabled: true })
  21. await upload
  22. const options = seed === undefined
  23. ? undefined
  24. : { seed, ...meta === undefined ? {} : { meta } }
  25. const session = ctx.sessions.create(SessionId(id), options)
  26. return { ctx, session, disposeUpload: () => upload.dispose() }
  27. }
  28. function body(text = 'x'.repeat(300)) {
  29. return { messages: [{ role: 'user', content: text }] }
  30. }
  31. describe('incremental DeepSeek session-log upload', () => {
  32. it('does not contribute the session log under its default configuration', async () => {
  33. const ctx = new Context()
  34. contexts.push(ctx)
  35. await ctx.plugin(SessionStore)
  36. await ctx.plugin(DeepSeekLlmApiExtensionRegistry)
  37. await ctx.plugin(SessionLogDeepSeek)
  38. const session = ctx.sessions.create(SessionId('default-off'))
  39. session.append('turn/start', { turn: 1 })
  40. const prepared = await ctx.deepseekLlmApiExtensions.prepare({
  41. body: body(), signal: SIGNAL, sessionId: session.id,
  42. })
  43. expect(prepared.fields).not.toHaveProperty('dsh_session_log')
  44. })
  45. it('uploads the full first prefix, records acceptance, then sends only the appended suffix', async () => {
  46. const { ctx, session } = await harness('incremental')
  47. session.append('turn/start', { turn: 1 })
  48. session.append('step/start', { turn: 1, step: 1 })
  49. const first = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  50. const firstPayload = first.fields.dsh_session_log
  51. expect(firstPayload).toMatchObject({ afterSeq: -1, throughSeq: 1 })
  52. expect(firstPayload?.events).toHaveLength(2)
  53. await first.accept()
  54. expect(SessionLogDeepSeek.acceptedThrough(session)).toBe(1)
  55. session.append('step/end', { turn: 1, step: 1 })
  56. const second = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  57. expect(second.fields.dsh_session_log).toMatchObject({ afterSeq: 1, throughSeq: 3 })
  58. expect(second.fields.dsh_session_log?.events).toHaveLength(2)
  59. expect(second.fields.dsh_session_log?.events[0]).toMatchObject({
  60. type: 'session-log-deepseek/delivery-accepted',
  61. seq: 2,
  62. })
  63. })
  64. it('reconstructs a persisted cursor and ignores an inherited parent watermark in a fork', async () => {
  65. const first = await harness('parent')
  66. first.session.append('turn/start', { turn: 1 })
  67. const prepared = await first.ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: first.session.id })
  68. await prepared.accept()
  69. const seed = first.session.events
  70. const resumed = await harness('parent', seed)
  71. expect(SessionLogDeepSeek.acceptedThrough(resumed.session)).toBe(0)
  72. const resumedPayload = await resumed.ctx.deepseekLlmApiExtensions.prepare({
  73. body: body(), signal: SIGNAL, sessionId: resumed.session.id,
  74. })
  75. expect(resumedPayload.fields.dsh_session_log?.afterSeq).toBe(0)
  76. const fork = await harness('child', seed, { parentSession: first.session.id, seedLength: seed.length })
  77. expect(SessionLogDeepSeek.acceptedThrough(fork.session)).toBe(-1)
  78. const forkPayload = await fork.ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: fork.session.id })
  79. expect(forkPayload.fields.dsh_session_log).toMatchObject({ afterSeq: -1, throughSeq: fork.session.seq - 1 })
  80. })
  81. it('takes the maximum watermark when concurrent acceptances settle out of order', async () => {
  82. const { ctx, session } = await harness('concurrent')
  83. session.append('turn/start', { turn: 1 })
  84. const earlier = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  85. session.append('step/start', { turn: 1, step: 1 })
  86. const later = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  87. await later.accept()
  88. await earlier.accept()
  89. expect(SessionLogDeepSeek.acceptedThrough(session)).toBe(1)
  90. })
  91. it('folds only events appended after the cached acceptance scan', () => {
  92. const id = SessionId('incremental-fold')
  93. const events: SessionEvent[] = [
  94. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } },
  95. { type: 'session-log-deepseek/delivery-accepted', seq: 1, time: 2, data: { sessionId: id, throughSeq: 0 } },
  96. ]
  97. let reads = 0
  98. const observed = new Proxy(events, {
  99. get(target, property, receiver) {
  100. if (typeof property === 'string' && /^\d+$/.test(property)) reads++
  101. return Reflect.get(target, property, receiver) as unknown
  102. },
  103. })
  104. const session = { id, get events() { return observed } } as unknown as Session
  105. expect(SessionLogDeepSeek.acceptedThrough(session)).toBe(0)
  106. expect(reads).toBe(2)
  107. reads = 0
  108. expect(SessionLogDeepSeek.acceptedThrough(session)).toBe(0)
  109. expect(reads).toBe(0)
  110. events.push(
  111. { type: 'step/start', seq: 2, time: 3, data: { turn: 1, step: 1 } },
  112. { type: 'session-log-deepseek/delivery-accepted', seq: 3, time: 4, data: { sessionId: id, throughSeq: 2 } },
  113. )
  114. expect(SessionLogDeepSeek.acceptedThrough(session)).toBe(2)
  115. expect(reads).toBe(2)
  116. })
  117. it('omits the field for direct or stale requests and uploads the prior acceptance marker next', async () => {
  118. const { ctx, session } = await harness('edges')
  119. await expect(ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL }))
  120. .resolves.toMatchObject({ fields: {} })
  121. await expect(ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: 'missing' }))
  122. .resolves.toMatchObject({ fields: {} })
  123. await expect(ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id }))
  124. .resolves.toMatchObject({ fields: {} })
  125. session.append('turn/start', { turn: 1 })
  126. const first = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  127. await first.accept()
  128. const current = await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })
  129. expect(current.fields.dsh_session_log).toMatchObject({
  130. afterSeq: 0,
  131. throughSeq: 1,
  132. events: [{ type: 'session-log-deepseek/delivery-accepted' }],
  133. })
  134. })
  135. it('contributes complete events without reading request messages', async () => {
  136. const { ctx, session } = await harness('direct-events')
  137. session.append('turn/start', { turn: 1 })
  138. const prepared = await ctx.deepseekLlmApiExtensions.prepare({ body: {}, signal: SIGNAL, sessionId: session.id })
  139. expect(prepared.fields.dsh_session_log?.events).toEqual(session.events)
  140. })
  141. it('fails closed on a malformed persisted acceptance watermark', async () => {
  142. const malformed = [{
  143. type: 'session-log-deepseek/delivery-accepted',
  144. seq: 0,
  145. time: 1,
  146. data: { sessionId: 'malformed', throughSeq: 0 },
  147. }] as unknown as SessionEvent[]
  148. const session = Session.create(SessionId('malformed'), malformed)
  149. expect(() => SessionLogDeepSeek.acceptedThrough(session)).toThrow(/malformed acceptance watermark/)
  150. })
  151. it('withdraws its request field when the contributing plugin reloads', async () => {
  152. const { ctx, session, disposeUpload } = await harness('hmr')
  153. session.append('turn/start', { turn: 1 })
  154. expect((await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })).fields)
  155. .toHaveProperty('dsh_session_log')
  156. await disposeUpload()
  157. expect((await ctx.deepseekLlmApiExtensions.prepare({ body: body(), signal: SIGNAL, sessionId: session.id })).fields)
  158. .not.toHaveProperty('dsh_session_log')
  159. })
  160. })