api-proxy-projections.spec.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. /**
  2. * Projection carrier paths of the host ApiProxy: the history tail page's
  3. * projections block reads the registry's watermark snapshot (asOfSeq = last
  4. * event seq, one consistent cut); loadOlder pages never carry the block; a
  5. * composition without the registry serves histories without it; a disposed
  6. * registration's key leaves subsequent responses; and every unit change is
  7. * pushed to mux consumers as a session/projection frame minted here.
  8. */
  9. import { describe, expect, it, vi } from 'vitest'
  10. import { Context } from '@deepseek-ai/cordis'
  11. import { z } from 'zod'
  12. import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent'
  13. import { AttachmentStore } from '@deepseek-ai/dsh-attachment'
  14. import type { Agent } from '@deepseek-ai/dsh-agent'
  15. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  16. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  17. import type { Session } from '@deepseek-ai/dsh-session'
  18. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  19. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  20. import UserQuestionService from '@deepseek-ai/dsh-user-questions'
  21. import type { MuxFrame, RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api'
  22. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  23. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  24. declare module '@deepseek-ai/dsh-session-projection/types' {
  25. interface SessionProjectionMap {
  26. 'test/last-user': { text: string } | null
  27. }
  28. }
  29. let nextRpc = 1
  30. function request<P>(payload: P): RpcRequest<P> {
  31. return { rpcId: RpcId(`proj-${String(nextRpc++)}`), payload }
  32. }
  33. /** Whole-value unit folding the latest user/message text; null before the first. */
  34. type LastUserState = { text: string } | null
  35. const lastUserUnit = (): ProjectionDefinition<'test/last-user', LastUserState> => ({
  36. key: 'test/last-user',
  37. schema: z.union([z.object({ text: z.string() }), z.null()]),
  38. init: () => null,
  39. apply: (state, event) => (event.type === 'user/message'
  40. ? { text: (event.data.content[0] as { text?: string }).text ?? '' }
  41. : state),
  42. view: state => state,
  43. stateVersion: 1,
  44. })
  45. async function harness(withRegistry: boolean): Promise<{ ctx: Context; session: Session }> {
  46. const ctx = new Context()
  47. await ctx.plugin(SessionStore)
  48. await ctx.plugin(UserQuestionService)
  49. await ctx.plugin(AgentRegistry)
  50. if (withRegistry) await ctx.plugin(SessionProjectionRegistry)
  51. const session = ctx.sessions.create()
  52. // The gateway reads both the session and durable inbox baseline.
  53. ctx.agents.register({ id: session.id, session, inbox: new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }), status: 'idle', ctx } as Agent)
  54. return { ctx, session }
  55. }
  56. /** Append `count` user messages so the log has paginable message boundaries. */
  57. function seedMessages(session: Session, count: number): void {
  58. for (let i = 0; i < count; i++) {
  59. session.append('user/message', createUserMessage({
  60. content: [{ type: 'text', text: `m${i}` }],
  61. source: { kind: 'user' },
  62. }), { surfaceOp: 'append' })
  63. }
  64. }
  65. const api = (ctx: Context) => createApiProxy(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  66. describe('session.history projections block', () => {
  67. it('serves the unit value on the tail page with asOfSeq = last event seq', async () => {
  68. const { ctx, session } = await harness(true)
  69. ctx.sessionProjections.register(lastUserUnit())
  70. seedMessages(session, 3)
  71. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  72. expect(response.result.ok).toBe(true)
  73. if (!response.result.ok) throw new Error('unreachable')
  74. const { events, projections } = response.result.value
  75. expect(projections).toBeDefined()
  76. expect(projections?.asOfSeq).toBe(session.seq - 1)
  77. expect(projections?.values['test/last-user']).toEqual({ text: 'm2' })
  78. // asOfSeq IS the window tail: the last served event carries it.
  79. expect(events.at(-1)?.event.seq).toBe(projections?.asOfSeq)
  80. })
  81. it('publishes the attachments imageLimits as a constant unit while both seams are composed', async () => {
  82. const { ctx, session } = await harness(true)
  83. const limits = {
  84. maxImageBytes: 5 * 1024 * 1024,
  85. maxImagesPerMessage: 20,
  86. maxMessageImageBytes: 100 * 1024 * 1024,
  87. maxImagePixels: 40_000_000,
  88. mediaTypes: ['image/png'] as const,
  89. }
  90. await ctx.plugin(class extends AttachmentStore {
  91. readonly imageLimits = limits
  92. validateImage(): Promise<void> { return Promise.resolve() }
  93. saveImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  94. readImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  95. })
  96. const gateway = api(ctx)
  97. seedMessages(session, 2)
  98. const response = await gateway.sessions.history(request({ sessionId: session.id }))
  99. if (!response.result.ok) throw new Error('history failed')
  100. expect(response.result.value.projections?.values['imageLimits']).toEqual(limits)
  101. // Constant unit: appending events must never broadcast an imageLimits frame.
  102. await new Promise(resolve => setTimeout(resolve, 0))
  103. const abort = new AbortController()
  104. const stream = gateway.events.mux({ rpcId: RpcId('t-limits-mux'), payload: {} }, abort.signal)
  105. const frames: MuxFrame[] = []
  106. const drained = (async () => {
  107. for await (const envelope of stream) {
  108. frames.push(envelope.payload)
  109. if (frames.some(f => f.type === 'session/event')) abort.abort()
  110. }
  111. })().catch(() => {})
  112. seedMessages(session, 1)
  113. await drained
  114. expect(frames.some(f => f.type === 'session/projection' && f.key === 'imageLimits')).toBe(false)
  115. })
  116. it('leaves the imageLimits key absent while no attachment service is composed', async () => {
  117. const { ctx, session } = await harness(true)
  118. seedMessages(session, 1)
  119. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  120. if (!response.result.ok) throw new Error('history failed')
  121. expect(response.result.value.projections).toBeDefined()
  122. expect('imageLimits' in (response.result.value.projections?.values ?? {})).toBe(false)
  123. })
  124. it('never carries the block on loadOlder pages (beforeSeq present)', async () => {
  125. const { ctx, session } = await harness(true)
  126. ctx.sessionProjections.register(lastUserUnit())
  127. seedMessages(session, 5)
  128. const older = await api(ctx).sessions.history(request({ sessionId: session.id, beforeSeq: 3, maxMessages: 2 }))
  129. expect(older.result.ok).toBe(true)
  130. if (!older.result.ok) throw new Error('unreachable')
  131. expect('projections' in older.result.value).toBe(false)
  132. })
  133. it('serves no block when the composition has no projection registry', async () => {
  134. const { ctx, session } = await harness(false)
  135. seedMessages(session, 2)
  136. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  137. expect(response.result.ok).toBe(true)
  138. if (!response.result.ok) throw new Error('unreachable')
  139. expect('projections' in response.result.value).toBe(false)
  140. })
  141. it('drops a disposed registration from subsequent tail pages (empty block, key absent)', async () => {
  142. const { ctx, session } = await harness(true)
  143. const dispose = ctx.sessionProjections.register(lastUserUnit())
  144. seedMessages(session, 1)
  145. const proxy = api(ctx)
  146. const before = await proxy.sessions.history(request({ sessionId: session.id }))
  147. if (!before.result.ok) throw new Error('unreachable')
  148. expect(before.result.value.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  149. dispose()
  150. const after = await proxy.sessions.history(request({ sessionId: session.id }))
  151. if (!after.result.ok) throw new Error('unreachable')
  152. // The registry stays mounted; only the disposed key leaves while the
  153. // gateway-owned Session-list unit remains.
  154. expect(after.result.value.projections?.asOfSeq).toBe(session.seq - 1)
  155. expect('test/last-user' in (after.result.value.projections?.values ?? {})).toBe(false)
  156. expect(after.result.value.projections?.values.sessionListMetadata).toEqual({
  157. blank: true,
  158. lastPromptAt: session.events.at(-1)?.time,
  159. })
  160. })
  161. it('removes the gateway-owned Session-list unit when the gateway fiber unloads', async () => {
  162. const { ctx, session } = await harness(true)
  163. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  164. const fiber = ctx.plugin(Object.assign((gatewayCtx: Context) => {
  165. createApiProxy(gatewayCtx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  166. }, { inject: ['sessions', 'agents', 'userQuestions', 'sessionProjections'] }))
  167. await fiber.await()
  168. await vi.waitFor(() => {
  169. expect(ctx.sessionProjections.snapshot(session).values.sessionListMetadata)
  170. .toEqual({ blank: true, lastPromptAt: null })
  171. })
  172. await fiber.dispose()
  173. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  174. })
  175. })
  176. describe('session.list projections column', () => {
  177. it('serves attached rows from the live registry cut, watermarked for client seeding', async () => {
  178. const { ctx, session } = await harness(true)
  179. ctx.sessionProjections.register(lastUserUnit())
  180. const gateway = api(ctx)
  181. await new Promise(resolve => setTimeout(resolve, 0))
  182. session.append('turn/start', { turn: 1 })
  183. seedMessages(session, 1)
  184. const response = await gateway.sessions.list(request({}))
  185. if (!response.result.ok) throw new Error('unreachable')
  186. const row = response.result.value.items.find(item => item.sessionId === session.id)
  187. expect(row?.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  188. expect(row?.projections?.values.sessionListMetadata).toEqual({
  189. blank: false,
  190. lastPromptAt: session.events.at(-1)?.time,
  191. })
  192. expect(row?.projections?.asOfSeq).toBe(session.seq - 1)
  193. })
  194. it('omits the column entirely when no registry is mounted', async () => {
  195. const { ctx, session } = await harness(false)
  196. seedMessages(session, 1)
  197. const response = await api(ctx).sessions.list(request({}))
  198. if (!response.result.ok) throw new Error('unreachable')
  199. const row = response.result.value.items.find(item => item.sessionId === session.id)
  200. expect(row).toBeDefined()
  201. expect(row !== undefined && 'projections' in row).toBe(false)
  202. })
  203. it('serves cold rows from the persisted projection cache with zero log loads', async () => {
  204. const { ctx } = await harness(true)
  205. const coldId = SessionId('session-cold-listing')
  206. const load = () => { throw new Error('list must not load event logs') }
  207. ctx.provide('sessionPersistence', {
  208. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  209. locate: () => undefined,
  210. load,
  211. inspect: load,
  212. readFrom: load,
  213. } as never)
  214. ctx.provide('sessionProjectionCache', {
  215. // The carrier hands the listed header through as the identity witness.
  216. cachedSnapshot: (meta: { id: unknown; createdAt: number }) =>
  217. (meta.id === coldId && meta.createdAt === 5
  218. ? { asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } }
  219. : undefined),
  220. } as never)
  221. const response = await api(ctx).sessions.list(request({}))
  222. if (!response.result.ok) throw new Error('unreachable')
  223. const row = response.result.value.items.find(item => item.sessionId === coldId)
  224. expect(row?.running).toBe(false)
  225. expect(row?.projections).toEqual({ asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } })
  226. })
  227. it('cold rows without a cache plugin (or without a stored row) just lack the column', async () => {
  228. const { ctx } = await harness(true)
  229. const coldId = SessionId('session-cold-uncached')
  230. ctx.provide('sessionPersistence', {
  231. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  232. locate: () => undefined,
  233. } as never)
  234. const response = await api(ctx).sessions.list(request({}))
  235. if (!response.result.ok) throw new Error('unreachable')
  236. const row = response.result.value.items.find(item => item.sessionId === coldId)
  237. expect(row).toBeDefined()
  238. expect(row !== undefined && 'projections' in row).toBe(false)
  239. })
  240. it('a throwing column read degrades that row, never the listing', async () => {
  241. const { ctx, session } = await harness(true)
  242. ctx.sessionProjections.register({
  243. ...lastUserUnit(),
  244. view: () => { throw new Error('unit exploded') },
  245. })
  246. seedMessages(session, 1)
  247. const response = await api(ctx).sessions.list(request({}))
  248. if (!response.result.ok) throw new Error('unreachable')
  249. const row = response.result.value.items.find(item => item.sessionId === session.id)
  250. expect(row).toBeDefined()
  251. expect(row !== undefined && 'projections' in row).toBe(false)
  252. })
  253. })
  254. describe('session/projection push frame', () => {
  255. /** Drain frames until `count` session/projection frames arrived. */
  256. async function collect(iterable: AsyncIterable<RpcRequest<MuxFrame>>, count: number, abort: AbortController): Promise<MuxFrame[]> {
  257. const frames: MuxFrame[] = []
  258. for await (const envelope of iterable) {
  259. frames.push(envelope.payload)
  260. if (frames.filter(f => f.type === 'session/projection').length >= count) abort.abort()
  261. }
  262. return frames
  263. }
  264. it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => {
  265. const { ctx, session } = await harness(true)
  266. ctx.sessionProjections.register(lastUserUnit())
  267. const proxy = api(ctx)
  268. // The gateway's onChanged subscription lives in an inject child whose
  269. // fiber activates asynchronously; yield until it lands before appending.
  270. await new Promise(resolve => setTimeout(resolve, 0))
  271. const abort = new AbortController()
  272. const stream = proxy.events.mux({ rpcId: RpcId('t-proj-mux'), payload: {} }, abort.signal)
  273. const collected = collect(stream, 5, abort)
  274. const now = vi.spyOn(Date, 'now').mockReturnValue(100)
  275. seedMessages(session, 1)
  276. now.mockReturnValue(200)
  277. session.append('turn/start', { turn: 1 })
  278. now.mockReturnValue(300)
  279. seedMessages(session, 1)
  280. now.mockRestore()
  281. const frames = await collected
  282. const pushes = frames.filter(
  283. (f): f is Extract<MuxFrame, { type: 'session/projection' }> =>
  284. f.type === 'session/projection' && f.key === 'test/last-user',
  285. )
  286. expect(pushes).toEqual([
  287. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 0 },
  288. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 2 },
  289. ])
  290. expect(frames.filter(
  291. (f): f is Extract<MuxFrame, { type: 'session/projection' }> =>
  292. f.type === 'session/projection' && f.key === 'sessionListMetadata',
  293. )).toEqual([
  294. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: true, lastPromptAt: 100 }, seq: 0 },
  295. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 100 }, seq: 1 },
  296. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 300 }, seq: 2 },
  297. ])
  298. // Frame seq aligns with the tail block's asOfSeq vocabulary (higher-seq-wins compatible).
  299. const tail = await proxy.sessions.history(request({ sessionId: session.id }))
  300. if (!tail.result.ok) throw new Error('unreachable')
  301. expect(tail.result.value.projections?.asOfSeq).toBe(pushes.at(-1)?.seq)
  302. })
  303. it('emits no projection frames when the composition has no registry', async () => {
  304. const { ctx, session } = await harness(false)
  305. const proxy = api(ctx)
  306. const abort = new AbortController()
  307. const stream = proxy.events.mux({ rpcId: RpcId('t-noproj-mux'), payload: {} }, abort.signal)
  308. const frames: MuxFrame[] = []
  309. const drained = (async () => {
  310. for await (const envelope of stream) {
  311. frames.push(envelope.payload)
  312. if (frames.filter(f => f.type === 'session/event').length >= 2) abort.abort()
  313. }
  314. })()
  315. seedMessages(session, 2)
  316. await drained
  317. expect(frames.some(f => f.type === 'session/projection')).toBe(false)
  318. })
  319. })