session-projections.host.spec.ts 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600
  1. /**
  2. * Session Controller projection paths: 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 through the control stream.
  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 { agentPresetProjectionDefinition } from '@deepseek-ai/dsh-agent-presets'
  15. import type { Agent } from '@deepseek-ai/dsh-agent'
  16. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  17. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  18. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  19. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  20. import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
  21. import { SessionControlController } from '@deepseek-ai/dsh-api-session-controller/src/control.ts'
  22. import type { SessionControlFrame, SessionFollowFrame } from '@deepseek-ai/dsh-api-session-controller/types'
  23. import { createSessionTestRemote, testSessionPersistence, type TestSessionRemote } from './test-remote.ts'
  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. function request<P>(payload: P): P {
  34. return payload
  35. }
  36. function page(
  37. remote: TestSessionRemote,
  38. request: { sessionId: SessionId; throughSeq: number; beforeSeq?: number; maxMessages?: number },
  39. ) {
  40. return remote.page({
  41. address: { kind: 'session', sessionId: request.sessionId },
  42. throughSeq: request.throughSeq,
  43. ...(request.beforeSeq === undefined ? {} : { beforeSeq: request.beforeSeq }),
  44. ...(request.maxMessages === undefined ? {} : { maxMessages: request.maxMessages }),
  45. })
  46. }
  47. /** Read and close one snapshot-first follow generation. */
  48. async function opening(
  49. remote: TestSessionRemote,
  50. sessionId: SessionId,
  51. maxMessages?: number,
  52. ): Promise<Extract<SessionFollowFrame, { type: 'snapshot' }>> {
  53. const abort = new AbortController()
  54. const iterator = remote.follow({
  55. address: { kind: 'session', sessionId },
  56. ...(maxMessages === undefined ? {} : { maxMessages }),
  57. }, abort.signal)[Symbol.asyncIterator]()
  58. const first = await iterator.next()
  59. abort.abort()
  60. await iterator.return?.()
  61. if (first.done || first.value.type !== 'snapshot') throw new Error('follow did not open with a snapshot')
  62. return first.value
  63. }
  64. /** Whole-value unit folding the latest user/message text; null before the first. */
  65. type LastUserState = { text: string } | null
  66. const lastUserUnit = () => ({
  67. key: 'test/last-user',
  68. stateSchema: z.union([z.object({ text: z.string() }), z.null()]),
  69. init: () => null,
  70. apply: (state, event) => (event.type === 'user/message'
  71. ? { text: (event.data.content[0] as { text?: string }).text ?? '' }
  72. : state),
  73. wire: {
  74. viewSchema: z.union([z.object({ text: z.string() }), z.null()]),
  75. view: state => state,
  76. },
  77. stateVersion: 1,
  78. }) satisfies ProjectionDefinition<'test/last-user', LastUserState>
  79. const internalCountUnit = () => ({
  80. key: 'test/internal-count',
  81. stateSchema: z.number().int().nonnegative(),
  82. init: () => 0,
  83. apply: (state: number) => state + 1,
  84. stateVersion: 1,
  85. }) satisfies ProjectionDefinition<'test/internal-count', number>
  86. async function harness(withRegistry: boolean): Promise<{ ctx: Context; session: Session }> {
  87. const ctx = new Context()
  88. await ctx.plugin(SessionStore)
  89. await ctx.plugin(AgentRegistry)
  90. if (withRegistry) await ctx.plugin(SessionProjectionRegistry)
  91. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  92. const agent = {
  93. id: session.id,
  94. session,
  95. inbox: { nextTurn: [], nextStep: [], hasPending: false } as never,
  96. status: 'idle',
  97. ctx,
  98. } as unknown as Agent
  99. if (withRegistry) Object.assign(agent, { inbox: new Inbox(ctx, agent.session, agentEvents(ctx, agent)) })
  100. ctx.agents.register(agent)
  101. return { ctx, session }
  102. }
  103. /** Append `count` user messages so the log has paginable message boundaries. */
  104. function seedMessages(session: Session, count: number): void {
  105. for (let i = 0; i < count; i++) {
  106. session.append('user/message', createUserMessage({
  107. content: [{ type: 'text', text: `m${i}` }],
  108. source: { kind: 'user' },
  109. }), { surfaceOp: 'append' })
  110. }
  111. }
  112. const remote = (ctx: Context) => createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  113. describe('session.history projections block', () => {
  114. it('tracks pending and used model selections across repeated request headers', async () => {
  115. const { ctx, session } = await harness(true)
  116. remote(ctx)
  117. await new Promise(resolve => setTimeout(resolve, 0))
  118. const selected = { provider: 'p', model: 'next' }
  119. session.append('model/selection', selected)
  120. session.append('model/selection', selected)
  121. session.append('request/header', {
  122. header: { config: { provider: 'p', model: 'used' } }, reason: 'initial',
  123. })
  124. session.append('request/header', {
  125. header: { config: { provider: 'p', model: 'used' } }, reason: 'initial',
  126. })
  127. expect(ctx.sessionProjections.snapshot(session).values.modelSelection).toEqual({
  128. lastUsed: { provider: 'p', model: 'used' },
  129. next: selected,
  130. })
  131. session.append('request/header', {
  132. header: { config: selected }, reason: 'initial',
  133. })
  134. expect(ctx.sessionProjections.snapshot(session).values.modelSelection).toEqual({
  135. lastUsed: selected,
  136. next: selected,
  137. })
  138. })
  139. it('serves the unit value on the tail page with asOfSeq = last event seq', async () => {
  140. const { ctx, session } = await harness(true)
  141. ctx.sessionProjections.register(lastUserUnit())
  142. seedMessages(session, 3)
  143. const snapshot = await opening(remote(ctx), session.id)
  144. const { records, projections } = snapshot
  145. expect(projections.asOfSeq).toBe(session.seq - 1)
  146. expect(projections.values['test/last-user']).toEqual({ text: 'm2' })
  147. // asOfSeq IS the window tail: the last served event carries it.
  148. const last = records.at(-1)
  149. expect(last?.event.seq).toBe(projections.asOfSeq)
  150. })
  151. it('reconstructs a cold persisted queue without publishing or resuming an Agent', async () => {
  152. const { ctx } = await harness(true)
  153. const coldId = SessionId('cold-persisted-queue')
  154. const meta = { version: 0 as const, id: coldId, createdAt: 1, cwd: '/tmp' }
  155. const message = createUserMessage({
  156. content: [{ type: 'text', text: 'survive process restart' }],
  157. source: { kind: 'user' },
  158. })
  159. const events: SessionEvent[] = [{
  160. type: 'agent/inbox/spliced',
  161. seq: 0,
  162. time: 2,
  163. data: { target: 'next-turn', start: 0, inserted: [message] },
  164. }]
  165. ctx.provide('sessionPersistence', testSessionPersistence(ctx, {
  166. list: () => Promise.resolve([meta]),
  167. inspect: () => Promise.resolve({ meta, events }),
  168. }) as never)
  169. const snapshot = await opening(remote(ctx), coldId)
  170. expect(snapshot.projections.values.inbox).toEqual({
  171. 'next-turn': [message],
  172. 'next-step': [],
  173. })
  174. expect(ctx.agents.get(coldId)).toBeUndefined()
  175. expect(ctx.sessions.get(coldId)).toBeUndefined()
  176. })
  177. it('removes claimed steering from the pending Inbox projection immediately', async () => {
  178. const { ctx, session } = await harness(true)
  179. const proxy = remote(ctx)
  180. const message = createUserMessage({
  181. content: [{ type: 'text', text: 'apply this now' }],
  182. source: { kind: 'user' },
  183. })
  184. const agent = ctx.agents.get(session.id)
  185. if (agent === undefined) throw new Error('missing Agent')
  186. agent.inbox.append('next-step', message)
  187. agent.inbox.claim('next-step', 1)
  188. const during = await opening(proxy, session.id)
  189. expect(during.projections.values.inbox).toEqual({
  190. 'next-turn': [],
  191. 'next-step': [],
  192. })
  193. session.append('user/message', message, { surfaceOp: 'append' })
  194. const settled = await opening(proxy, session.id)
  195. expect(settled.projections.values.inbox).toEqual({
  196. 'next-turn': [],
  197. 'next-step': [],
  198. })
  199. const rejected = createUserMessage({
  200. content: [{ type: 'text', text: 'reject this pre-step' }],
  201. source: { kind: 'user' },
  202. })
  203. session.append('turn/start', { turn: 1 })
  204. agent.inbox.append('next-step', rejected)
  205. agent.inbox.claim('next-step', 1)
  206. session.append('turn/end', { turn: 1, reason: { kind: 'blocked' } })
  207. const closed = await opening(proxy, session.id)
  208. expect(closed.projections.values.inbox).toEqual({
  209. 'next-turn': [],
  210. 'next-step': [],
  211. })
  212. })
  213. it('returns a complete current replacement cut on each follow generation', async () => {
  214. const { ctx, session } = await harness(true)
  215. ctx.sessionProjections.register(lastUserUnit())
  216. seedMessages(session, 2)
  217. const snapshot = await opening(remote(ctx), session.id)
  218. expect(snapshot.records.map(record => record.event.seq)).toEqual([0, 1])
  219. expect(snapshot.projections.asOfSeq).toBe(1)
  220. expect(snapshot.projections.values).toEqual(
  221. expect.objectContaining({ 'test/last-user': { text: 'm1' } }),
  222. )
  223. })
  224. it('projects an empty log at cursor -1', async () => {
  225. const { ctx, session } = await harness(true)
  226. ctx.sessionProjections.register(lastUserUnit())
  227. const snapshot = await opening(remote(ctx), session.id)
  228. expect(snapshot.records).toEqual([])
  229. expect(snapshot.projections.asOfSeq).toBe(-1)
  230. expect(snapshot.projections.values).toEqual(
  231. expect.objectContaining({ 'test/last-user': null }),
  232. )
  233. })
  234. it('publishes the attachments imageLimits as a constant unit while both seams are composed', async () => {
  235. const { ctx, session } = await harness(true)
  236. const limits = {
  237. maxImageBytes: 5 * 1024 * 1024,
  238. maxImagesPerMessage: 20,
  239. maxMessageImageBytes: 100 * 1024 * 1024,
  240. maxImagePixels: 40_000_000,
  241. maxImageDimension: 2000,
  242. mediaTypes: ['image/png'] as const,
  243. }
  244. await ctx.plugin(class extends AttachmentStore {
  245. readonly imageLimits = limits
  246. validateImage(): Promise<void> { return Promise.resolve() }
  247. saveImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  248. readImage(): Promise<never> { return Promise.reject(new Error('unused')) }
  249. })
  250. const gateway = remote(ctx)
  251. await new Promise(resolve => setTimeout(resolve, 0))
  252. seedMessages(session, 2)
  253. const snapshot = await opening(gateway, session.id)
  254. expect(snapshot.projections.values['imageLimits']).toEqual(limits)
  255. // Constant unit: appending events must never broadcast an imageLimits projection.
  256. await new Promise(resolve => setTimeout(resolve, 0))
  257. const abort = new AbortController()
  258. const iterator = gateway.control(abort.signal)[Symbol.asyncIterator]()
  259. await iterator.next()
  260. const next = iterator.next()
  261. seedMessages(session, 1)
  262. await new Promise(resolve => setTimeout(resolve, 0))
  263. await expect(next).resolves.toMatchObject({
  264. done: false,
  265. value: { type: 'projection', key: 'sessionListMetadata' },
  266. })
  267. const extra = iterator.next()
  268. const quiet = Symbol('quiet')
  269. expect(await Promise.race([
  270. extra,
  271. new Promise<typeof quiet>(resolve => setTimeout(() => { resolve(quiet) }, 0)),
  272. ])).toBe(quiet)
  273. abort.abort()
  274. await expect(extra).resolves.toEqual({ done: true, value: undefined })
  275. })
  276. it('leaves the imageLimits key absent while no attachment service is composed', async () => {
  277. const { ctx, session } = await harness(true)
  278. seedMessages(session, 1)
  279. const snapshot = await opening(remote(ctx), session.id)
  280. expect('imageLimits' in snapshot.projections.values).toBe(false)
  281. })
  282. it('never carries the block on loadOlder pages (beforeSeq present)', async () => {
  283. const { ctx, session } = await harness(true)
  284. ctx.sessionProjections.register(lastUserUnit())
  285. seedMessages(session, 5)
  286. const older = await page(remote(ctx), request({
  287. sessionId: session.id, throughSeq: session.seq - 1, beforeSeq: 3, maxMessages: 2,
  288. }))
  289. expect(older.ok).toBe(true)
  290. if (!older.ok) throw new Error('unreachable')
  291. expect('projections' in older.value).toBe(false)
  292. })
  293. it('serves no block when the composition has no projection registry', async () => {
  294. const { ctx, session } = await harness(false)
  295. seedMessages(session, 2)
  296. const response = await page(remote(ctx), request({ sessionId: session.id, throughSeq: session.seq - 1 }))
  297. expect(response.ok).toBe(true)
  298. if (!response.ok) throw new Error('unreachable')
  299. expect('projections' in response.value).toBe(false)
  300. })
  301. it('never exposes a host-only unit through history, listing, or push frames', async () => {
  302. const { ctx, session } = await harness(true)
  303. ctx.sessionProjections.register(internalCountUnit())
  304. const proxy = remote(ctx)
  305. await new Promise(resolve => setTimeout(resolve, 0))
  306. const abort = new AbortController()
  307. const iterator = proxy.control(abort.signal)[Symbol.asyncIterator]()
  308. const baseline = await iterator.next()
  309. if (baseline.done || baseline.value.type !== 'baseline') {
  310. throw new Error('control stream ended before its baseline')
  311. }
  312. expect('test/internal-count' in (baseline.value.value.projections[session.id]?.values ?? {}))
  313. .toBe(false)
  314. seedMessages(session, 1)
  315. const changed = await iterator.next()
  316. expect(changed).toMatchObject({
  317. done: false,
  318. value: { type: 'projection', key: 'sessionListMetadata' },
  319. })
  320. abort.abort()
  321. await iterator.return?.()
  322. const history = await opening(proxy, session.id)
  323. expect('test/internal-count' in history.projections.values).toBe(false)
  324. const listing = await proxy.list(request({}))
  325. if (!listing.ok) throw new Error('listing failed')
  326. const row = listing.value.items.find(item => item.sessionId === session.id)
  327. expect('test/internal-count' in (row?.projections?.values ?? {})).toBe(false)
  328. })
  329. it('drops a disposed registration from subsequent tail pages (empty block, key absent)', async () => {
  330. const { ctx, session } = await harness(true)
  331. const dispose = ctx.sessionProjections.register(lastUserUnit())
  332. seedMessages(session, 1)
  333. const proxy = remote(ctx)
  334. const before = await opening(proxy, session.id)
  335. expect(before.projections.values['test/last-user']).toEqual({ text: 'm0' })
  336. dispose()
  337. const after = await opening(proxy, session.id)
  338. // The registry stays mounted; only the disposed key leaves while the
  339. // gateway-owned Session-list unit remains.
  340. expect(after.projections.asOfSeq).toBe(session.seq - 1)
  341. expect('test/last-user' in after.projections.values).toBe(false)
  342. expect(after.projections.values.sessionListMetadata).toEqual({
  343. blank: true,
  344. lastPromptAt: session.events.at(-1)?.time,
  345. })
  346. })
  347. it('removes the gateway-owned Session-list unit when the gateway fiber unloads', async () => {
  348. const { ctx, session } = await harness(true)
  349. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  350. const fiber = ctx.plugin(Object.assign((gatewayCtx: Context) => {
  351. createSessionTestRemote(gatewayCtx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  352. }, { inject: ['sessions', 'agents', 'sessionProjections'] }))
  353. await fiber.await()
  354. await vi.waitFor(() => {
  355. expect(ctx.sessionProjections.snapshot(session).values.sessionListMetadata)
  356. .toEqual({ blank: true, lastPromptAt: null })
  357. })
  358. await fiber.dispose()
  359. expect('sessionListMetadata' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  360. })
  361. })
  362. describe('session.list projections column', () => {
  363. it('serves every already-materialized wire value from the live registry without folding', async () => {
  364. const { ctx, session } = await harness(true)
  365. ctx.sessionProjections.register(lastUserUnit())
  366. const gateway = remote(ctx)
  367. await new Promise(resolve => setTimeout(resolve, 0))
  368. session.append('turn/start', { turn: 1 })
  369. seedMessages(session, 1)
  370. const response = await gateway.list(request({}))
  371. if (!response.ok) throw new Error('unreachable')
  372. const row = response.value.items.find(item => item.sessionId === session.id)
  373. expect(row?.projections?.values['test/last-user']).toEqual({ text: 'm0' })
  374. expect(row?.projections?.values.sessionListMetadata).toEqual({
  375. blank: false,
  376. lastPromptAt: session.events.at(-1)?.time,
  377. })
  378. expect(row?.projections?.asOfSeq).toBe(session.seq - 1)
  379. })
  380. it('lists the latest preset selected by a blank Session instead of its creation preset', async () => {
  381. const { ctx } = await harness(true)
  382. const session = ctx.sessions.create(SessionId('preset-list'), {
  383. meta: { cwd: '/workspace', agentPreset: 'standard' },
  384. })
  385. ctx.sessionProjections.register(agentPresetProjectionDefinition)
  386. const gateway = remote(ctx)
  387. await new Promise(resolve => setTimeout(resolve, 0))
  388. session.append('agent-preset/selected', { agentPreset: 'minimal' })
  389. const response = await gateway.list(request({}))
  390. if (!response.ok) throw new Error('unreachable')
  391. const row = response.value.items.find(item => item.sessionId === session.id)
  392. expect(row?.projections?.values.agentPreset).toBe('minimal')
  393. })
  394. it('omits an unmaterialized live projection instead of folding history for listing', async () => {
  395. const { ctx, session } = await harness(true)
  396. seedMessages(session, 1)
  397. const unit = lastUserUnit()
  398. const apply = vi.fn(unit.apply)
  399. ctx.sessionProjections.register({ ...unit, apply })
  400. const response = await remote(ctx).list(request({}))
  401. if (!response.ok) throw new Error('unreachable')
  402. const row = response.value.items.find(item => item.sessionId === session.id)
  403. expect(row).toBeDefined()
  404. expect('test/last-user' in (row?.projections?.values ?? {})).toBe(false)
  405. expect(apply).not.toHaveBeenCalled()
  406. })
  407. it('omits the column entirely when no registry is mounted', async () => {
  408. const { ctx, session } = await harness(false)
  409. seedMessages(session, 1)
  410. const response = await remote(ctx).list(request({}))
  411. if (!response.ok) throw new Error('unreachable')
  412. const row = response.value.items.find(item => item.sessionId === session.id)
  413. expect(row).toBeDefined()
  414. expect(row !== undefined && 'projections' in row).toBe(false)
  415. })
  416. it('serves every available cold projection hint from the cache with zero log loads', async () => {
  417. const { ctx } = await harness(true)
  418. const coldId = SessionId('session-cold-listing')
  419. const load = () => { throw new Error('list must not load event logs') }
  420. ctx.provide('sessionPersistence', {
  421. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  422. locate: () => undefined,
  423. load,
  424. inspect: load,
  425. readFrom: load,
  426. } as never)
  427. ctx.provide('sessionProjectionCache', {
  428. // The carrier hands the listed header through as the identity witness.
  429. cachedSnapshot: (meta: { id: unknown; createdAt: number }) =>
  430. (meta.id === coldId && meta.createdAt === 5
  431. ? {
  432. asOfSeq: 7,
  433. values: {
  434. 'test/last-user': { text: 'cached' },
  435. sessionListMetadata: { blank: false, lastPromptAt: 6 },
  436. title: 'Cached title',
  437. },
  438. }
  439. : undefined),
  440. } as never)
  441. const response = await remote(ctx).list(request({}))
  442. if (!response.ok) throw new Error('unreachable')
  443. const row = response.value.items.find(item => item.sessionId === coldId)
  444. expect(row?.running).toBe(false)
  445. expect(row?.projections).toEqual({
  446. asOfSeq: 7,
  447. values: {
  448. 'test/last-user': { text: 'cached' },
  449. sessionListMetadata: { blank: false, lastPromptAt: 6 },
  450. title: 'Cached title',
  451. },
  452. })
  453. })
  454. it('cold rows without a cache plugin (or without a stored row) just lack the column', async () => {
  455. const { ctx } = await harness(true)
  456. const coldId = SessionId('session-cold-uncached')
  457. ctx.provide('sessionPersistence', {
  458. list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }],
  459. locate: () => undefined,
  460. } as never)
  461. const response = await remote(ctx).list(request({}))
  462. if (!response.ok) throw new Error('unreachable')
  463. const row = response.value.items.find(item => item.sessionId === coldId)
  464. expect(row).toBeDefined()
  465. expect(row !== undefined && 'projections' in row).toBe(false)
  466. })
  467. it('a throwing column read degrades that row, never the listing', async () => {
  468. const { ctx, session } = await harness(true)
  469. ctx.sessionProjections.register({
  470. ...lastUserUnit(),
  471. wire: {
  472. viewSchema: z.union([z.object({ text: z.string() }), z.null()]),
  473. view: () => { throw new Error('unit exploded') },
  474. },
  475. })
  476. seedMessages(session, 1)
  477. const response = await remote(ctx).list(request({}))
  478. if (!response.ok) throw new Error('unreachable')
  479. const row = response.value.items.find(item => item.sessionId === session.id)
  480. expect(row).toBeDefined()
  481. expect(row !== undefined && 'projections' in row).toBe(false)
  482. })
  483. })
  484. describe('Session control projection frames', () => {
  485. /** Drain frames until `count` projection replacements arrive. */
  486. async function collect(
  487. iterable: AsyncIterable<SessionControlFrame>,
  488. count: number,
  489. abort: AbortController,
  490. ): Promise<SessionControlFrame[]> {
  491. const frames: SessionControlFrame[] = []
  492. for await (const frame of iterable) {
  493. frames.push(frame)
  494. if (frames.filter(candidate => candidate.type === 'projection').length >= count) abort.abort()
  495. }
  496. return frames
  497. }
  498. it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => {
  499. const { ctx, session } = await harness(true)
  500. ctx.sessionProjections.register(lastUserUnit())
  501. const proxy = remote(ctx)
  502. // The controller's onChanged subscription lives in an inject child whose
  503. // fiber activates asynchronously; yield until it lands before appending.
  504. await new Promise(resolve => setTimeout(resolve, 0))
  505. const abort = new AbortController()
  506. const stream = proxy.control(abort.signal)
  507. const collected = collect(stream, 5, abort)
  508. const now = vi.spyOn(Date, 'now').mockReturnValue(100)
  509. seedMessages(session, 1)
  510. now.mockReturnValue(200)
  511. session.append('turn/start', { turn: 1 })
  512. now.mockReturnValue(300)
  513. seedMessages(session, 1)
  514. now.mockRestore()
  515. const frames = await collected
  516. const pushes = frames.filter(
  517. (f): f is Extract<SessionControlFrame, { type: 'projection' }> =>
  518. f.type === 'projection' && f.key === 'test/last-user',
  519. )
  520. expect(pushes).toEqual([
  521. { type: 'projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 0 },
  522. { type: 'projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 2 },
  523. ])
  524. expect(frames.filter(
  525. (f): f is Extract<SessionControlFrame, { type: 'projection' }> =>
  526. f.type === 'projection' && f.key === 'sessionListMetadata',
  527. )).toEqual([
  528. { type: 'projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: true, lastPromptAt: 100 }, seq: 0 },
  529. { type: 'projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 100 }, seq: 1 },
  530. { type: 'projection', sessionId: session.id, key: 'sessionListMetadata', value: { blank: false, lastPromptAt: 300 }, seq: 2 },
  531. ])
  532. // Frame seq aligns with the tail block's asOfSeq vocabulary (higher-seq-wins compatible).
  533. const tail = await opening(proxy, session.id)
  534. expect(tail.projections.asOfSeq).toBe(pushes.at(-1)?.seq)
  535. })
  536. it('emits no projection frames when the composition has no registry', async () => {
  537. const { ctx, session } = await harness(false)
  538. const control = new SessionControlController(ctx)
  539. const abort = new AbortController()
  540. const iterator = control.control(abort.signal)[Symbol.asyncIterator]()
  541. const baseline = await iterator.next()
  542. const next = iterator.next()
  543. seedMessages(session, 2)
  544. await new Promise(resolve => setTimeout(resolve, 0))
  545. abort.abort()
  546. if (baseline.done) throw new Error('Control stream ended before its baseline')
  547. expect(baseline.value.type).toBe('baseline')
  548. await expect(next).resolves.toEqual({ done: true, value: undefined })
  549. })
  550. })