api-proxy-projections.spec.ts 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479
  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, { agentEvents, 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 SessionProjectionStateMap {
  26. 'test/last-user': LastUserState
  27. 'test/internal-count': number
  28. }
  29. interface SessionProjectionMap {
  30. 'test/last-user': { text: string } | null
  31. }
  32. }
  33. let nextRpc = 1
  34. function request<P>(payload: P): RpcRequest<P> {
  35. return { rpcId: RpcId(`proj-${String(nextRpc++)}`), payload }
  36. }
  37. /** Whole-value unit folding the latest user/message text; null before the first. */
  38. type LastUserState = { text: string } | null
  39. const lastUserUnit = () => ({
  40. key: 'test/last-user',
  41. stateSchema: z.union([z.object({ text: z.string() }), z.null()]),
  42. init: () => null,
  43. apply: (state, event) => (event.type === 'user/message'
  44. ? { text: (event.data.content[0] as { text?: string }).text ?? '' }
  45. : state),
  46. wire: {
  47. viewSchema: z.union([z.object({ text: z.string() }), z.null()]),
  48. view: state => state,
  49. },
  50. stateVersion: 1,
  51. }) satisfies ProjectionDefinition<'test/last-user', LastUserState>
  52. const internalCountUnit = () => ({
  53. key: 'test/internal-count',
  54. stateSchema: z.number().int().nonnegative(),
  55. init: () => 0,
  56. apply: (state: number) => state + 1,
  57. stateVersion: 1,
  58. }) satisfies ProjectionDefinition<'test/internal-count', number>
  59. async function harness(withRegistry: boolean): Promise<{ ctx: Context; session: Session }> {
  60. const ctx = new Context()
  61. await ctx.plugin(SessionStore)
  62. await ctx.plugin(UserQuestionService)
  63. await ctx.plugin(AgentRegistry)
  64. if (withRegistry) {
  65. await ctx.plugin(SessionProjectionRegistry)
  66. }
  67. const session = ctx.sessions.create()
  68. const agent = {
  69. id: session.id,
  70. options: {},
  71. session,
  72. inbox: { nextTurn: [], nextStep: [], hasPending: false } as never,
  73. status: 'idle',
  74. ctx,
  75. send() {},
  76. followup() {},
  77. steer() {},
  78. inject() {},
  79. cancel() {},
  80. runMaintenance: task => task(new AbortController().signal),
  81. whenIdle: () => Promise.resolve(),
  82. } satisfies Agent
  83. if (withRegistry) Object.assign(agent, { inbox: new Inbox(ctx, agent.session, agentEvents(ctx, agent)) })
  84. ctx.agents.register(agent)
  85. return { ctx, session }
  86. }
  87. /** Append `count` user messages so the log has paginable message boundaries. */
  88. function seedMessages(session: Session, count: number): void {
  89. for (let i = 0; i < count; i++) {
  90. session.append('user/message', createUserMessage({
  91. content: [{ type: 'text', text: `m${i}` }],
  92. source: { kind: 'user' },
  93. }), { surfaceOp: 'append' })
  94. }
  95. }
  96. const api = (ctx: Context) => createApiProxy(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  97. describe('session.history projections block', () => {
  98. it('serves the unit value on the tail page with asOfSeq = last event seq', async () => {
  99. const { ctx, session } = await harness(true)
  100. ctx.sessionProjections.register(lastUserUnit())
  101. seedMessages(session, 3)
  102. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  103. expect(response.result.ok).toBe(true)
  104. if (!response.result.ok) throw new Error('unreachable')
  105. const { events, projections } = response.result.value
  106. expect(projections).toBeDefined()
  107. expect(projections?.asOfSeq).toBe(session.seq - 1)
  108. expect(projections?.values['test/last-user']).toEqual({ text: 'm2' })
  109. // asOfSeq IS the window tail: the last served event carries it.
  110. expect(events.at(-1)?.event.seq).toBe(projections?.asOfSeq)
  111. })
  112. it('reconstructs a cold persisted queue without publishing or resuming an Agent', async () => {
  113. const { ctx } = await harness(true)
  114. const coldId = SessionId('cold-persisted-queue')
  115. const meta = { version: 0 as const, id: coldId, createdAt: 1, cwd: '/tmp' }
  116. const message = createUserMessage({
  117. content: [{ type: 'text', text: 'survive process restart' }],
  118. source: { kind: 'user' },
  119. })
  120. const events = [{
  121. type: 'agent/inbox/spliced',
  122. seq: 0,
  123. time: 2,
  124. data: { target: 'next-turn', start: 0, inserted: [message] },
  125. }] as const
  126. ctx.provide('sessionPersistence', {
  127. list: () => Promise.resolve([meta]),
  128. inspect: () => Promise.resolve({ meta, events }),
  129. } as never)
  130. const response = await api(ctx).sessions.history(request({ sessionId: coldId }))
  131. if (!response.result.ok) throw new Error('history failed')
  132. expect(response.result.value.projections?.values.inbox).toEqual({
  133. 'next-turn': [message],
  134. 'next-step': [],
  135. })
  136. expect(ctx.agents.get(coldId)).toBeUndefined()
  137. expect(ctx.sessions.get(coldId)).toBeUndefined()
  138. })
  139. it('removes claimed steering from the pending Inbox projection immediately', async () => {
  140. const { ctx, session } = await harness(true)
  141. const proxy = api(ctx)
  142. const message = createUserMessage({
  143. content: [{ type: 'text', text: 'apply this now' }],
  144. source: { kind: 'user' },
  145. })
  146. const agent = ctx.agents.get(session.id)
  147. if (agent === undefined) throw new Error('missing Agent')
  148. agent.inbox.append('next-step', message)
  149. agent.inbox.claim('next-step', 1)
  150. const during = await proxy.sessions.history(request({ sessionId: session.id }))
  151. if (!during.result.ok) throw new Error('history failed')
  152. expect(during.result.value.projections?.values.inbox).toEqual({
  153. 'next-turn': [],
  154. 'next-step': [],
  155. })
  156. session.append('user/message', message, { surfaceOp: 'append' })
  157. const settled = await proxy.sessions.history(request({ sessionId: session.id }))
  158. if (!settled.result.ok) throw new Error('history failed')
  159. expect(settled.result.value.projections?.values.inbox).toEqual({
  160. 'next-turn': [],
  161. 'next-step': [],
  162. })
  163. const rejected = createUserMessage({
  164. content: [{ type: 'text', text: 'reject this pre-step' }],
  165. source: { kind: 'user' },
  166. })
  167. session.append('turn/start', { turn: 1 })
  168. agent.inbox.append('next-step', rejected)
  169. agent.inbox.claim('next-step', 1)
  170. session.append('turn/end', { turn: 1, reason: { kind: 'blocked' } })
  171. const closed = await proxy.sessions.history(request({ sessionId: session.id }))
  172. if (!closed.result.ok) throw new Error('history failed')
  173. expect(closed.result.value.projections?.values.inbox).toEqual({
  174. 'next-turn': [],
  175. 'next-step': [],
  176. })
  177. })
  178. it('publishes the attachments imageLimits as a constant unit while both seams are composed', async () => {
  179. const { ctx, session } = await harness(true)
  180. const limits = {
  181. maxImageBytes: 5 * 1024 * 1024,
  182. maxImagesPerMessage: 20,
  183. maxMessageImageBytes: 100 * 1024 * 1024,
  184. maxImagePixels: 40_000_000,
  185. maxImageDimension: 2000,
  186. mediaTypes: ['image/png'] as const,
  187. }
  188. await ctx.plugin(class extends AttachmentStore {
  189. readonly imageLimits = limits
  190. validateImage(): Promise<void> { return Promise.resolve() }
  191. saveImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  192. readImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  193. })
  194. const gateway = api(ctx)
  195. seedMessages(session, 2)
  196. const response = await gateway.sessions.history(request({ sessionId: session.id }))
  197. if (!response.result.ok) throw new Error('history failed')
  198. expect(response.result.value.projections?.values['imageLimits']).toEqual(limits)
  199. // Constant unit: appending events must never broadcast an imageLimits frame.
  200. await new Promise(resolve => setTimeout(resolve, 0))
  201. const abort = new AbortController()
  202. const stream = gateway.events.mux({ rpcId: RpcId('t-limits-mux'), payload: {} }, abort.signal)
  203. const frames: MuxFrame[] = []
  204. const drained = (async () => {
  205. for await (const envelope of stream) {
  206. frames.push(envelope.payload)
  207. if (frames.some(f => f.type === 'session/event')) abort.abort()
  208. }
  209. })().catch(() => {})
  210. seedMessages(session, 1)
  211. await drained
  212. expect(frames.some(f => f.type === 'session/projection' && f.key === 'imageLimits')).toBe(false)
  213. })
  214. it('leaves the imageLimits key absent while no attachment service is composed', async () => {
  215. const { ctx, session } = await harness(true)
  216. seedMessages(session, 1)
  217. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  218. if (!response.result.ok) throw new Error('history failed')
  219. expect(response.result.value.projections).toBeDefined()
  220. expect('imageLimits' in (response.result.value.projections?.values ?? {})).toBe(false)
  221. })
  222. it('never carries the block on loadOlder pages (beforeSeq present)', async () => {
  223. const { ctx, session } = await harness(true)
  224. ctx.sessionProjections.register(lastUserUnit())
  225. seedMessages(session, 5)
  226. const older = await api(ctx).sessions.history(request({ sessionId: session.id, beforeSeq: 3, maxMessages: 2 }))
  227. expect(older.result.ok).toBe(true)
  228. if (!older.result.ok) throw new Error('unreachable')
  229. expect('projections' in older.result.value).toBe(false)
  230. })
  231. it('serves no block when the composition has no projection registry', async () => {
  232. const { ctx, session } = await harness(false)
  233. seedMessages(session, 2)
  234. const response = await api(ctx).sessions.history(request({ sessionId: session.id }))
  235. expect(response.result.ok).toBe(true)
  236. if (!response.result.ok) throw new Error('unreachable')
  237. expect('projections' in response.result.value).toBe(false)
  238. })
  239. it('never exposes a host-only unit through history, listing, or push frames', async () => {
  240. const { ctx, session } = await harness(true)
  241. ctx.sessionProjections.register(internalCountUnit())
  242. const proxy = api(ctx)
  243. await new Promise(resolve => setTimeout(resolve, 0))
  244. const abort = new AbortController()
  245. const frames: MuxFrame[] = []
  246. const drained = (async () => {
  247. for await (const envelope of proxy.events.mux({ rpcId: RpcId('t-host-only-mux'), payload: {} }, abort.signal)) {
  248. frames.push(envelope.payload)
  249. if (envelope.payload.type === 'session/event') abort.abort()
  250. }
  251. })().catch(() => {})
  252. seedMessages(session, 1)
  253. await drained
  254. const history = await proxy.sessions.history(request({ sessionId: session.id }))
  255. if (!history.result.ok) throw new Error('history failed')
  256. expect('test/internal-count' in (history.result.value.projections?.values ?? {})).toBe(false)
  257. const listing = await proxy.sessions.list(request({}))
  258. if (!listing.result.ok) throw new Error('listing failed')
  259. const row = listing.result.value.items.find(item => item.sessionId === session.id)
  260. expect('test/internal-count' in (row?.projections?.values ?? {})).toBe(false)
  261. expect(frames.some(frame => frame.type === 'session/projection' && frame.key === 'test/internal-count')).toBe(false)
  262. })
  263. it('drops a disposed registration from subsequent tail pages (empty block, key absent)', async () => {
  264. const { ctx, session } = await harness(true)
  265. const dispose = ctx.sessionProjections.register(lastUserUnit())
  266. seedMessages(session, 1)
  267. const proxy = api(ctx)
  268. const before = await proxy.sessions.history(request({ sessionId: session.id }))
  269. if (!before.result.ok) throw new Error('unreachable')
  270. expect(before.result.value.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  271. dispose()
  272. const after = await proxy.sessions.history(request({ sessionId: session.id }))
  273. if (!after.result.ok) throw new Error('unreachable')
  274. // The registry stays mounted; only the disposed key leaves while the
  275. // gateway-owned Session-list unit remains.
  276. expect(after.result.value.projections?.asOfSeq).toBe(session.seq - 1)
  277. expect('test/last-user' in (after.result.value.projections?.values ?? {})).toBe(false)
  278. expect(after.result.value.projections?.values.sessionListMetadata).toEqual({
  279. blank: true,
  280. lastPromptAt: session.events.at(-1)?.time,
  281. })
  282. })
  283. it('removes the gateway-owned Session-list unit when the gateway fiber unloads', async () => {
  284. const { ctx, session } = await harness(true)
  285. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  286. const fiber = ctx.plugin(Object.assign((gatewayCtx: Context) => {
  287. createApiProxy(gatewayCtx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  288. }, { inject: ['sessions', 'agents', 'userQuestions', 'sessionProjections'] }))
  289. await fiber.await()
  290. await vi.waitFor(() => {
  291. expect(ctx.sessionProjections.snapshot(session).values.sessionListMetadata)
  292. .toEqual({ blank: true, lastPromptAt: null })
  293. })
  294. await fiber.dispose()
  295. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  296. })
  297. })
  298. describe('session.list projections column', () => {
  299. it('serves attached rows from the live registry cut, watermarked for client seeding', async () => {
  300. const { ctx, session } = await harness(true)
  301. ctx.sessionProjections.register(lastUserUnit())
  302. const gateway = api(ctx)
  303. await new Promise(resolve => setTimeout(resolve, 0))
  304. session.append('turn/start', { turn: 1 })
  305. seedMessages(session, 1)
  306. const response = await gateway.sessions.list(request({}))
  307. if (!response.result.ok) throw new Error('unreachable')
  308. const row = response.result.value.items.find(item => item.sessionId === session.id)
  309. expect(row?.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  310. expect(row?.projections?.values.sessionListMetadata).toEqual({
  311. blank: false,
  312. lastPromptAt: session.events.at(-1)?.time,
  313. })
  314. expect(row?.projections?.asOfSeq).toBe(session.seq - 1)
  315. })
  316. it('omits the column entirely when no registry is mounted', async () => {
  317. const { ctx, session } = await harness(false)
  318. seedMessages(session, 1)
  319. const response = await api(ctx).sessions.list(request({}))
  320. if (!response.result.ok) throw new Error('unreachable')
  321. const row = response.result.value.items.find(item => item.sessionId === session.id)
  322. expect(row).toBeDefined()
  323. expect(row !== undefined && 'projections' in row).toBe(false)
  324. })
  325. it('serves cold rows from the persisted projection cache with zero log loads', async () => {
  326. const { ctx } = await harness(true)
  327. const coldId = SessionId('session-cold-listing')
  328. const load = () => { throw new Error('list must not load event logs') }
  329. ctx.provide('sessionPersistence', {
  330. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  331. locate: () => undefined,
  332. load,
  333. inspect: load,
  334. readFrom: load,
  335. } as never)
  336. ctx.provide('sessionProjectionCache', {
  337. // The carrier hands the listed header through as the identity witness.
  338. cachedSnapshot: (meta: { id: unknown; createdAt: number }) =>
  339. (meta.id === coldId && meta.createdAt === 5
  340. ? { asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } }
  341. : undefined),
  342. } as never)
  343. const response = await api(ctx).sessions.list(request({}))
  344. if (!response.result.ok) throw new Error('unreachable')
  345. const row = response.result.value.items.find(item => item.sessionId === coldId)
  346. expect(row?.running).toBe(false)
  347. expect(row?.projections).toEqual({ asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } })
  348. })
  349. it('cold rows without a cache plugin (or without a stored row) just lack the column', async () => {
  350. const { ctx } = await harness(true)
  351. const coldId = SessionId('session-cold-uncached')
  352. ctx.provide('sessionPersistence', {
  353. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  354. locate: () => undefined,
  355. } as never)
  356. const response = await api(ctx).sessions.list(request({}))
  357. if (!response.result.ok) throw new Error('unreachable')
  358. const row = response.result.value.items.find(item => item.sessionId === coldId)
  359. expect(row).toBeDefined()
  360. expect(row !== undefined && 'projections' in row).toBe(false)
  361. })
  362. it('a throwing column read degrades that row, never the listing', async () => {
  363. const { ctx, session } = await harness(true)
  364. ctx.sessionProjections.register({
  365. ...lastUserUnit(),
  366. wire: {
  367. viewSchema: z.union([z.object({ text: z.string() }), z.null()]),
  368. view: () => { throw new Error('unit exploded') },
  369. },
  370. })
  371. seedMessages(session, 1)
  372. const response = await api(ctx).sessions.list(request({}))
  373. if (!response.result.ok) throw new Error('unreachable')
  374. const row = response.result.value.items.find(item => item.sessionId === session.id)
  375. expect(row).toBeDefined()
  376. expect(row !== undefined && 'projections' in row).toBe(false)
  377. })
  378. })
  379. describe('session/projection push frame', () => {
  380. /** Drain frames until `count` session/projection frames arrived. */
  381. async function collect(iterable: AsyncIterable<RpcRequest<MuxFrame>>, count: number, abort: AbortController): Promise<MuxFrame[]> {
  382. const frames: MuxFrame[] = []
  383. for await (const envelope of iterable) {
  384. frames.push(envelope.payload)
  385. if (frames.filter(f => f.type === 'session/projection').length >= count) abort.abort()
  386. }
  387. return frames
  388. }
  389. it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => {
  390. const { ctx, session } = await harness(true)
  391. ctx.sessionProjections.register(lastUserUnit())
  392. const proxy = api(ctx)
  393. // The gateway's onChanged subscription lives in an inject child whose
  394. // fiber activates asynchronously; yield until it lands before appending.
  395. await new Promise(resolve => setTimeout(resolve, 0))
  396. const abort = new AbortController()
  397. const stream = proxy.events.mux({ rpcId: RpcId('t-proj-mux'), payload: {} }, abort.signal)
  398. const collected = collect(stream, 5, abort)
  399. const now = vi.spyOn(Date, 'now').mockReturnValue(100)
  400. seedMessages(session, 1)
  401. now.mockReturnValue(200)
  402. session.append('turn/start', { turn: 1 })
  403. now.mockReturnValue(300)
  404. seedMessages(session, 1)
  405. now.mockRestore()
  406. const frames = await collected
  407. const pushes = frames.filter(
  408. (f): f is Extract<MuxFrame, { type: 'session/projection' }> =>
  409. f.type === 'session/projection' && f.key === 'test/last-user',
  410. )
  411. expect(pushes).toEqual([
  412. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 0 },
  413. { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 2 },
  414. ])
  415. expect(frames.filter(
  416. (f): f is Extract<MuxFrame, { type: 'session/projection' }> =>
  417. f.type === 'session/projection' && f.key === 'sessionListMetadata',
  418. )).toEqual([
  419. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: true, lastPromptAt: 100 }, seq: 0 },
  420. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 100 }, seq: 1 },
  421. { type: 'session/projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 300 }, seq: 2 },
  422. ])
  423. // Frame seq aligns with the tail block's asOfSeq vocabulary (higher-seq-wins compatible).
  424. const tail = await proxy.sessions.history(request({ sessionId: session.id }))
  425. if (!tail.result.ok) throw new Error('unreachable')
  426. expect(tail.result.value.projections?.asOfSeq).toBe(pushes.at(-1)?.seq)
  427. })
  428. it('emits no projection frames when the composition has no registry', async () => {
  429. const { ctx, session } = await harness(false)
  430. const proxy = api(ctx)
  431. const abort = new AbortController()
  432. const stream = proxy.events.mux({ rpcId: RpcId('t-noproj-mux'), payload: {} }, abort.signal)
  433. const frames: MuxFrame[] = []
  434. const drained = (async () => {
  435. for await (const envelope of stream) {
  436. frames.push(envelope.payload)
  437. if (frames.filter(f => f.type === 'session/event').length >= 2) abort.abort()
  438. }
  439. })()
  440. seedMessages(session, 2)
  441. await drained
  442. expect(frames.some(f => f.type === 'session/projection')).toBe(false)
  443. })
  444. })