/** * Projection carrier paths of the host ApiProxy: history tail pages snapshot * attached state or fold one cold inspected prefix, loadOlder omits the block, * and live unit changes push session/projection frames. */ import { describe, expect, it } from 'vitest' import { Context } from 'cordis' import { z } from 'zod' import AgentRegistry from '@deepseek-ai/dsh-agent' import { createUserMessage } from '@deepseek-ai/dsh-llm' import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session' import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection' import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection' import UserInteractionService from '@deepseek-ai/dsh-user-interaction' import type { MuxFrame, RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api' import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc' import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy' declare module '@deepseek-ai/dsh-session-projection/types' { interface SessionProjectionMap { 'test/last-user': { text: string } | null } } let nextRpc = 1 function request

(payload: P): RpcRequest

{ return { rpcId: RpcId(`proj-${String(nextRpc++)}`), payload } } /** Whole-value unit folding the latest user/message text; null before the first. */ type LastUserState = { text: string } | null const lastUserUnit = (): ProjectionDefinition<'test/last-user', LastUserState> => ({ key: 'test/last-user', schema: z.union([z.object({ text: z.string() }), z.null()]), init: () => null, apply: (state, event) => (event.type === 'user/message' ? { text: (event.data.content[0] as { text?: string }).text ?? '' } : state), view: state => state, stateVersion: 1, }) async function harness(withRegistry: boolean): Promise<{ ctx: Context; session: Session }> { const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(UserInteractionService) await ctx.plugin(AgentRegistry) if (withRegistry) await ctx.plugin(SessionProjectionRegistry) const session = ctx.sessions.create() return { ctx, session } } /** Append `count` user messages so the log has paginable message boundaries. */ function seedMessages(session: Session, count: number): void { for (let i = 0; i < count; i++) { session.append('user/message', createUserMessage({ content: [{ type: 'text', text: `m${i}` }], source: { kind: 'user' }, }), { surfaceOp: 'append' }) } } const api = (ctx: Context) => createApiProxy(ctx, { provider: 'p', model: 'm', cwd: '/tmp', workspaceRoot: '/tmp' }) describe('session.history projections block', () => { it('serves the unit value on the tail page with asOfSeq = last event seq', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register(lastUserUnit()) seedMessages(session, 3) const response = await api(ctx).sessions.history(request({ sessionId: session.id })) expect(response.result.ok).toBe(true) if (!response.result.ok) throw new Error('unreachable') const { events, projections } = response.result.value expect(projections).toBeDefined() expect(projections?.asOfSeq).toBe(session.seq - 1) expect(projections?.values['test/last-user']).toEqual({ text: 'm2' }) // asOfSeq IS the window tail: the last served event carries it. expect(events.at(-1)?.event.seq).toBe(projections?.asOfSeq) }) it('folds a cold inspected prefix without publishing an Agent', async () => { const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(UserInteractionService) await ctx.plugin(AgentRegistry) await ctx.plugin(SessionProjectionRegistry) ctx.sessionProjections.register(lastUserUnit()) const sessionId = SessionId('session-cold-history') const meta: SessionHeader = { version: 0, id: sessionId, createdAt: 1, cwd: '/tmp' } const events = [{ type: 'user/message', seq: 0, time: 2, data: createUserMessage({ content: [{ type: 'text', text: 'persisted' }], source: { kind: 'user' }, }), surfaceOp: 'append', }] as SessionEvent[] ctx.provide('sessionPersistence', { list: () => Promise.resolve([meta]), inspect: () => Promise.resolve({ meta, events }), } as never) const response = await api(ctx).sessions.history(request({ sessionId })) expect(response.result.ok).toBe(true) if (!response.result.ok) throw new Error('unreachable') expect(response.result.value.projections).toEqual({ asOfSeq: 0, values: { 'test/last-user': { text: 'persisted' } }, }) expect(ctx.agents.get(sessionId)).toBeUndefined() }) it('never carries the block on loadOlder pages (beforeSeq present)', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register(lastUserUnit()) seedMessages(session, 5) const older = await api(ctx).sessions.history(request({ sessionId: session.id, beforeSeq: 3, maxMessages: 2 })) expect(older.result.ok).toBe(true) if (!older.result.ok) throw new Error('unreachable') expect('projections' in older.result.value).toBe(false) }) it('serves no block when the composition has no projection registry', async () => { const { ctx, session } = await harness(false) seedMessages(session, 2) const response = await api(ctx).sessions.history(request({ sessionId: session.id })) expect(response.result.ok).toBe(true) if (!response.result.ok) throw new Error('unreachable') expect('projections' in response.result.value).toBe(false) }) it('drops a disposed registration from subsequent tail pages (empty block, key absent)', async () => { const { ctx, session } = await harness(true) const dispose = ctx.sessionProjections.register(lastUserUnit()) seedMessages(session, 1) const proxy = api(ctx) const before = await proxy.sessions.history(request({ sessionId: session.id })) if (!before.result.ok) throw new Error('unreachable') expect(before.result.value.projections?.values['test/last-user']).toEqual({ text: 'm0' }) dispose() const after = await proxy.sessions.history(request({ sessionId: session.id })) if (!after.result.ok) throw new Error('unreachable') // The registry is still mounted, so the block itself stays (asOfSeq cut // with zero keys); the disposed key reads as capability absence. expect(after.result.value.projections?.asOfSeq).toBe(session.seq - 1) expect(after.result.value.projections?.values).toEqual({}) }) }) describe('session.list projections column', () => { it('serves attached rows from the live registry cut, watermarked for client seeding', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register(lastUserUnit()) seedMessages(session, 1) const response = await api(ctx).sessions.list(request({})) if (!response.result.ok) throw new Error('unreachable') const row = response.result.value.items.find(item => item.sessionId === session.id) expect(row?.projections?.values['test/last-user']).toEqual({ text: 'm0' }) expect(row?.projections?.asOfSeq).toBe(session.seq - 1) }) it('omits the column entirely when no registry is mounted', async () => { const { ctx, session } = await harness(false) seedMessages(session, 1) const response = await api(ctx).sessions.list(request({})) if (!response.result.ok) throw new Error('unreachable') const row = response.result.value.items.find(item => item.sessionId === session.id) expect(row).toBeDefined() expect(row !== undefined && 'projections' in row).toBe(false) }) it('serves cold rows from the persisted projection cache with zero log loads', async () => { const { ctx } = await harness(true) const coldId = SessionId('session-cold-listing') const load = () => { throw new Error('list must not load event logs') } ctx.provide('sessionPersistence', { list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }], locate: () => undefined, load, inspect: load, readFrom: load, } as never) ctx.provide('sessionProjectionCache', { // The carrier hands the listed header through as the identity witness. cachedSnapshot: (meta: { id: unknown; createdAt: number }) => (meta.id === coldId && meta.createdAt === 5 ? { asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } } : undefined), } as never) const response = await api(ctx).sessions.list(request({})) if (!response.result.ok) throw new Error('unreachable') const row = response.result.value.items.find(item => item.sessionId === coldId) expect(row?.running).toBe(false) expect(row?.projections).toEqual({ asOfSeq: 7, values: { 'test/last-user': { text: 'cached' } } }) }) it('cold rows without a cache plugin (or without a stored row) just lack the column', async () => { const { ctx } = await harness(true) const coldId = SessionId('session-cold-uncached') ctx.provide('sessionPersistence', { list: async () => [{ version: 0, id: coldId, createdAt: 5, cwd: '/tmp' }], locate: () => undefined, } as never) const response = await api(ctx).sessions.list(request({})) if (!response.result.ok) throw new Error('unreachable') const row = response.result.value.items.find(item => item.sessionId === coldId) expect(row).toBeDefined() expect(row !== undefined && 'projections' in row).toBe(false) }) it('a throwing column read degrades that row, never the listing', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register({ ...lastUserUnit(), view: () => { throw new Error('unit exploded') }, }) seedMessages(session, 1) const response = await api(ctx).sessions.list(request({})) if (!response.result.ok) throw new Error('unreachable') const row = response.result.value.items.find(item => item.sessionId === session.id) expect(row).toBeDefined() expect(row !== undefined && 'projections' in row).toBe(false) }) }) describe('session/projection push frame', () => { /** Drain frames until `count` session/projection frames arrived. */ async function collect(iterable: AsyncIterable>, count: number, abort: AbortController): Promise { const frames: MuxFrame[] = [] for await (const envelope of iterable) { frames.push(envelope.payload) if (frames.filter(f => f.type === 'session/projection').length >= count) abort.abort() } return frames } it('broadcasts a frame per changed unit with the causing seq, and none for same-reference applies', async () => { const { ctx, session } = await harness(true) ctx.sessionProjections.register(lastUserUnit()) const proxy = api(ctx) // The gateway's onChanged subscription lives in an inject child whose // fiber activates asynchronously; yield until it lands before appending. await new Promise(resolve => setTimeout(resolve, 0)) const abort = new AbortController() const stream = proxy.events.mux({ rpcId: RpcId('t-proj-mux'), payload: {} }, abort.signal) const collected = collect(stream, 2, abort) seedMessages(session, 1) // Same-reference apply: turn/start does not concern the unit — no frame. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } }) seedMessages(session, 1) const frames = await collected const pushes = frames.filter( (f): f is Extract => f.type === 'session/projection', ) expect(pushes).toEqual([ { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 0 }, { type: 'session/projection', sessionId: session.id, key: 'test/last-user', value: { text: 'm0' }, seq: 2 }, ]) // Frame seq aligns with the tail block's asOfSeq vocabulary (higher-seq-wins compatible). const tail = await proxy.sessions.history(request({ sessionId: session.id })) if (!tail.result.ok) throw new Error('unreachable') expect(tail.result.value.projections?.asOfSeq).toBe(pushes.at(-1)?.seq) }) it('emits no projection frames when the composition has no registry', async () => { const { ctx, session } = await harness(false) const proxy = api(ctx) const abort = new AbortController() const stream = proxy.events.mux({ rpcId: RpcId('t-noproj-mux'), payload: {} }, abort.signal) const frames: MuxFrame[] = [] const drained = (async () => { for await (const envelope of stream) { frames.push(envelope.payload) if (frames.filter(f => f.type === 'session/event').length >= 2) abort.abort() } })() seedMessages(session, 2) await drained expect(frames.some(f => f.type === 'session/projection')).toBe(false) }) })