api-proxy-projections.spec.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  1. /**
  2. * Projection carrier paths of the host ApiProxy: history tail pages snapshot
  3. * attached state or fold one cold inspected prefix, loadOlder omits the block,
  4. * and live unit changes push session/projection frames.
  5. */
  6. import { describe, expect, it } from 'vitest'
  7. import { Context } from 'cordis'
  8. import { z } from 'zod'
  9. import AgentRegistry from '@deepseek-ai/dsh-agent'
  10. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  11. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  12. import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
  13. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  14. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  15. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  16. import type { MuxFrame, RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api'
  17. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  18. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  19. declare module '@deepseek-ai/dsh-session-projection/types' {
  20. interface SessionProjectionMap {
  21. 'test/last-user': { text: string } | null
  22. }
  23. }
  24. let nextRpc = 1
  25. function request<P>(payload: P): RpcRequest<P> {
  26. return { rpcId: RpcId(`proj-${String(nextRpc++)}`), payload }
  27. }
  28. /** Whole-value unit folding the latest user/message text; null before the first. */
  29. type LastUserState = { text: string } | null
  30. const lastUserUnit = (): ProjectionDefinition<'test/last-user', LastUserState> => ({
  31. key: 'test/last-user',
  32. schema: z.union([z.object({ text: z.string() }), z.null()]),
  33. init: () => null,
  34. apply: (state, event) => (event.type === 'user/message'
  35. ? { text: (event.data.content[0] as { text?: string }).text ?? '' }
  36. : state),
  37. view: state => state,
  38. stateVersion: 1,
  39. })
  40. async function harness(withRegistry: boolean): Promise<{ ctx: Context; session: Session }> {
  41. const ctx = new Context()
  42. await ctx.plugin(SessionStore)
  43. await ctx.plugin(UserInteractionService)
  44. await ctx.plugin(AgentRegistry)
  45. if (withRegistry) await ctx.plugin(SessionProjectionRegistry)
  46. const session = ctx.sessions.create()
  47. return { ctx, session }
  48. }
  49. /** Append `count` user messages so the log has paginable message boundaries. */
  50. function seedMessages(session: Session, count: number): void {
  51. for (let i = 0; i < count; i++) {
  52. session.append('user/message', createUserMessage({
  53. content: [{ type: 'text', text: `m${i}` }],
  54. source: { kind: 'user' },
  55. }), { surfaceOp: 'append' })
  56. }
  57. }
  58. const api = (ctx: Context) => createApiProxy(ctx, { provider: 'p', model: 'm', cwd: '/tmp', workspaceRoot: '/tmp' })
  59. describe('session.history projections block', () => {
  60. it('serves the unit value on the tail page with asOfSeq = last event seq', async () => {
  61. const { ctx, session } = await harness(true)
  62. ctx.sessionProjections.register(lastUserUnit())
  63. seedMessages(session, 3)
  64. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  65. expect(response.result.ok).toBe(true)
  66. if (!response.result.ok) throw new Error('unreachable')
  67. const { events, projections } = response.result.value
  68. expect(projections).toBeDefined()
  69. expect(projections?.asOfSeq).toBe(session.seq - 1)
  70. expect(projections?.values['test/last-user']).toEqual({ text: 'm2' })
  71. // asOfSeq IS the window tail: the last served event carries it.
  72. expect(events.at(-1)?.event.seq).toBe(projections?.asOfSeq)
  73. })
  74. it('folds a cold inspected prefix without publishing an Agent', async () => {
  75. const ctx = new Context()
  76. await ctx.plugin(SessionStore)
  77. await ctx.plugin(UserInteractionService)
  78. await ctx.plugin(AgentRegistry)
  79. await ctx.plugin(SessionProjectionRegistry)
  80. ctx.sessionProjections.register(lastUserUnit())
  81. const sessionId = SessionId('session-cold-history')
  82. const meta: SessionHeader = { version: 0, id: sessionId, createdAt: 1, cwd: '/tmp' }
  83. const events = [{
  84. type: 'user/message',
  85. seq: 0,
  86. time: 2,
  87. data: createUserMessage({
  88. content: [{ type: 'text', text: 'persisted' }],
  89. source: { kind: 'user' },
  90. }),
  91. surfaceOp: 'append',
  92. }] as SessionEvent[]
  93. ctx.provide('sessionPersistence', {
  94. list: () => Promise.resolve([meta]),
  95. inspect: () => Promise.resolve({ meta, events }),
  96. } as never)
  97. const response = await api(ctx).sessions.history(request({ sessionId }))
  98. expect(response.result.ok).toBe(true)
  99. if (!response.result.ok) throw new Error('unreachable')
  100. expect(response.result.value.projections).toEqual({
  101. asOfSeq: 0,
  102. values: { 'test/last-user': { text: 'persisted' } },
  103. })
  104. expect(ctx.agents.get(sessionId)).toBeUndefined()
  105. })
  106. it('never carries the block on loadOlder pages (beforeSeq present)', async () => {
  107. const { ctx, session } = await harness(true)
  108. ctx.sessionProjections.register(lastUserUnit())
  109. seedMessages(session, 5)
  110. const older = await api(ctx).sessions.history(request({ sessionId: session.id, beforeSeq: 3, maxMessages: 2 }))
  111. expect(older.result.ok).toBe(true)
  112. if (!older.result.ok) throw new Error('unreachable')
  113. expect('projections' in older.result.value).toBe(false)
  114. })
  115. it('serves no block when the composition has no projection registry', async () => {
  116. const { ctx, session } = await harness(false)
  117. seedMessages(session, 2)
  118. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  119. expect(response.result.ok).toBe(true)
  120. if (!response.result.ok) throw new Error('unreachable')
  121. expect('projections' in response.result.value).toBe(false)
  122. })
  123. it('drops a disposed registration from subsequent tail pages (empty block, key absent)', async () => {
  124. const { ctx, session } = await harness(true)
  125. const dispose = ctx.sessionProjections.register(lastUserUnit())
  126. seedMessages(session, 1)
  127. const proxy = api(ctx)
  128. const before = await proxy.sessions.history(request({ sessionId: session.id }))
  129. if (!before.result.ok) throw new Error('unreachable')
  130. expect(before.result.value.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  131. dispose()
  132. const after = await proxy.sessions.history(request({ sessionId: session.id }))
  133. if (!after.result.ok) throw new Error('unreachable')
  134. // The registry is still mounted, so the block itself stays (asOfSeq cut
  135. // with zero keys); the disposed key reads as capability absence.
  136. expect(after.result.value.projections?.asOfSeq).toBe(session.seq - 1)
  137. expect(after.result.value.projections?.values).toEqual({})
  138. })
  139. })
  140. describe('session.list projections column', () => {
  141. it('serves attached rows from the live registry cut, watermarked for client seeding', async () => {
  142. const { ctx, session } = await harness(true)
  143. ctx.sessionProjections.register(lastUserUnit())
  144. seedMessages(session, 1)
  145. const response = await api(ctx).sessions.list(request({}))
  146. if (!response.result.ok) throw new Error('unreachable')
  147. const row = response.result.value.items.find(item => item.sessionId === session.id)
  148. expect(row?.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  149. expect(row?.projections?.asOfSeq).toBe(session.seq - 1)
  150. })
  151. it('omits the column entirely when no registry is mounted', async () => {
  152. const { ctx, session } = await harness(false)
  153. seedMessages(session, 1)
  154. const response = await api(ctx).sessions.list(request({}))
  155. if (!response.result.ok) throw new Error('unreachable')
  156. const row = response.result.value.items.find(item => item.sessionId === session.id)
  157. expect(row).toBeDefined()
  158. expect(row !== undefined && 'projections' in row).toBe(false)
  159. })
  160. it('serves cold rows from the persisted projection cache with zero log loads', async () => {
  161. const { ctx } = await harness(true)
  162. const coldId = SessionId('session-cold-listing')
  163. const load = () => { throw new Error('list must not load event logs') }
  164. ctx.provide('sessionPersistence', {
  165. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  166. locate: () => undefined,
  167. load,
  168. inspect: load,
  169. readFrom: load,
  170. } as never)
  171. ctx.provide('sessionProjectionCache', {
  172. // The carrier hands the listed header through as the identity witness.
  173. cachedSnapshot: (meta: { id: unknown; createdAt: number }) =>
  174. (meta.id === coldId && meta.createdAt === 5
  175. ? { asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } }
  176. : undefined),
  177. } as never)
  178. const response = await api(ctx).sessions.list(request({}))
  179. if (!response.result.ok) throw new Error('unreachable')
  180. const row = response.result.value.items.find(item => item.sessionId === coldId)
  181. expect(row?.running).toBe(false)
  182. expect(row?.projections).toEqual({ asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } })
  183. })
  184. it('cold rows without a cache plugin (or without a stored row) just lack the column', async () => {
  185. const { ctx } = await harness(true)
  186. const coldId = SessionId('session-cold-uncached')
  187. ctx.provide('sessionPersistence', {
  188. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  189. locate: () => undefined,
  190. } as never)
  191. const response = await api(ctx).sessions.list(request({}))
  192. if (!response.result.ok) throw new Error('unreachable')
  193. const row = response.result.value.items.find(item => item.sessionId === coldId)
  194. expect(row).toBeDefined()
  195. expect(row !== undefined && 'projections' in row).toBe(false)
  196. })
  197. it('a throwing column read degrades that row, never the listing', async () => {
  198. const { ctx, session } = await harness(true)
  199. ctx.sessionProjections.register({
  200. ...lastUserUnit(),
  201. view: () => { throw new Error('unit exploded') },
  202. })
  203. seedMessages(session, 1)
  204. const response = await api(ctx).sessions.list(request({}))
  205. if (!response.result.ok) throw new Error('unreachable')
  206. const row = response.result.value.items.find(item => item.sessionId === session.id)
  207. expect(row).toBeDefined()
  208. expect(row !== undefined && 'projections' in row).toBe(false)
  209. })
  210. })
  211. describe('session/projection push frame', () => {
  212. /** Drain frames until `count` session/projection frames arrived. */
  213. async function collect(iterable: AsyncIterable<RpcRequest<MuxFrame>>, count: number, abort: AbortController): Promise<MuxFrame[]> {
  214. const frames: MuxFrame[] = []
  215. for await (const envelope of iterable) {
  216. frames.push(envelope.payload)
  217. if (frames.filter(f => f.type === 'session/projection').length >= count) abort.abort()
  218. }
  219. return frames
  220. }
  221. it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => {
  222. const { ctx, session } = await harness(true)
  223. ctx.sessionProjections.register(lastUserUnit())
  224. const proxy = api(ctx)
  225. // The gateway's onChanged subscription lives in an inject child whose
  226. // fiber activates asynchronously; yield until it lands before appending.
  227. await new Promise(resolve => setTimeout(resolve, 0))
  228. const abort = new AbortController()
  229. const stream = proxy.events.mux({ rpcId: RpcId('t-proj-mux'), payload: {} }, abort.signal)
  230. const collected = collect(stream, 2, abort)
  231. seedMessages(session, 1)
  232. // Same-reference apply: turn/start does not concern the unit — no frame.
  233. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  234. seedMessages(session, 1)
  235. const frames = await collected
  236. const pushes = frames.filter(
  237. (f): f is Extract<MuxFrame, { type: 'session/projection' }> => f.type === 'session/projection',
  238. )
  239. expect(pushes).toEqual([
  240. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 0 },
  241. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 2 },
  242. ])
  243. // Frame seq aligns with the tail block's asOfSeq vocabulary (higher-seq-wins compatible).
  244. const tail = await proxy.sessions.history(request({ sessionId: session.id }))
  245. if (!tail.result.ok) throw new Error('unreachable')
  246. expect(tail.result.value.projections?.asOfSeq).toBe(pushes.at(-1)?.seq)
  247. })
  248. it('emits no projection frames when the composition has no registry', async () => {
  249. const { ctx, session } = await harness(false)
  250. const proxy = api(ctx)
  251. const abort = new AbortController()
  252. const stream = proxy.events.mux({ rpcId: RpcId('t-noproj-mux'), payload: {} }, abort.signal)
  253. const frames: MuxFrame[] = []
  254. const drained = (async () => {
  255. for await (const envelope of stream) {
  256. frames.push(envelope.payload)
  257. if (frames.filter(f => f.type === 'session/event').length >= 2) abort.abort()
  258. }
  259. })()
  260. seedMessages(session, 2)
  261. await drained
  262. expect(frames.some(f => f.type === 'session/projection')).toBe(false)
  263. })
  264. })