| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304 |
- /** Raw Session journal transport and message-aligned pagination coverage. */
- import { describe, expect, it, vi } from 'vitest'
- import { Context } from '@deepseek-ai/cordis'
- import AgentRegistry from '@deepseek-ai/dsh-agent'
- import SessionStore from '@deepseek-ai/dsh-session'
- import { CallId, createMessage, createToolResultMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
- import type { Session, SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
- import { SessionHistoryController } from '@deepseek-ai/dsh-api-session-controller/src/history.ts'
- import type { SessionFollowFrame } from '@deepseek-ai/dsh-api-session-controller/types'
- import { createSessionTestRemote, installSessionReadTestServices } from './test-remote.ts'
- /** Append a production-shaped human prompt to the session surface. */
- function appendUserText(session: Session, text: string): SessionEvent {
- return session.append('user/message', createUserMessage({
- content: [{ type: 'text', text }], source: { kind: 'user' },
- }), { surfaceOp: 'append' })
- }
- /** Append a production-shaped assistant message to the session surface. */
- function appendAssistantText(session: Session, text: string, step: number): SessionEvent {
- return session.append('assistant/message', {
- turn: 1,
- step,
- message: createMessage({
- role: 'assistant',
- content: [{ type: 'text', text }],
- source: { kind: 'model', provider: 'p', model: 'm' },
- }),
- }, { surfaceOp: 'append' })
- }
- /**
- * Append a plugin-owned log-only event. The host proxy is projection-only, so it
- * declares no compaction vocabulary; the cast writes the real event shape without
- * depending on the owning package.
- */
- function appendExtension(session: Session, type: string, data: unknown): SessionEvent {
- return (session.append as unknown as (type: string, data: unknown) => SessionEvent)(type, data)
- }
- async function harness(): Promise<{ ctx: Context }> {
- const ctx = new Context()
- await ctx.plugin(SessionStore)
- await ctx.plugin(AgentRegistry)
- installSessionReadTestServices(ctx)
- return { ctx }
- }
- /** Drain one Session follow until `count` event frames arrive. */
- async function collect(
- iterable: AsyncIterable<SessionFollowFrame>,
- count: number,
- abort: AbortController,
- ): Promise<SessionFollowFrame[]> {
- const frames: SessionFollowFrame[] = []
- for await (const frame of iterable) {
- frames.push(frame)
- if (frames.filter(candidate => candidate.type === 'event').length >= count) abort.abort()
- }
- return frames
- }
- /** Open follow and wait until its cursor is fixed before appending fixtures. */
- async function openFollow(
- history: SessionHistoryController,
- sessionId: SessionId,
- signal: AbortSignal,
- ): Promise<AsyncIterable<SessionFollowFrame>> {
- const iterator = history.follow({
- address: { kind: 'session', sessionId },
- }, signal)[Symbol.asyncIterator]()
- await expect(iterator.next()).resolves.toMatchObject({
- done: false,
- value: { type: 'snapshot' },
- })
- return { [Symbol.asyncIterator]: () => iterator }
- }
- describe('Session history raw journal', () => {
- it('follows raw tool events and preserves result metadata without a Tools service', async () => {
- const { ctx } = await harness()
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
- const abort = new AbortController()
- const stream = await openFollow(history, session.id, abort.signal)
- const collected = collect(stream, 2, abort)
- const call = session.append('tool/call', {
- turn: 1, step: 1, callId: CallId('raw-call'), name: 'custom', arguments: '{malformed',
- })
- const result = session.append('tool/result', {
- turn: 1, step: 1,
- message: createToolResultMessage({
- callId: CallId('raw-call'),
- content: [{ type: 'text', text: 'raw output' }],
- isError: false,
- }),
- meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] },
- }, { surfaceOp: 'append' })
- const frames = await collected
- expect(frames).toEqual([
- { type: 'event', event: call },
- { type: 'event', event: result },
- ])
- expect((frames[1] as Extract<SessionFollowFrame, { type: 'event' }>).event.data)
- .toMatchObject({ meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] } })
- })
- it('follows live results without rescanning Session history', async () => {
- const { ctx } = await harness()
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
- const abort = new AbortController()
- const stream = await openFollow(history, session.id, abort.signal)
- const iterator = stream[Symbol.asyncIterator]()
- session.append('tool/call', {
- turn: 1, step: 1, callId: CallId('live-fast'), name: 'term', arguments: '{"cmd":"pwd"}',
- })
- await expect(iterator.next()).resolves.toMatchObject({
- value: { type: 'event', event: { type: 'tool/call', data: { callId: 'live-fast' } } },
- })
- const events = vi.spyOn(session, 'events', 'get').mockImplementation(() => {
- throw new Error('live result rescanned Session history')
- })
- try {
- session.append('tool/result', {
- turn: 1, step: 1,
- message: createToolResultMessage({
- callId: CallId('live-fast'),
- content: [{ type: 'text', text: 'ok' }],
- isError: false,
- }),
- }, { surfaceOp: 'append' })
- await expect(iterator.next()).resolves.toMatchObject({
- value: { type: 'event', event: { type: 'tool/result', data: { message: { source: { callId: 'live-fast' } } } } },
- })
- } finally {
- events.mockRestore()
- abort.abort()
- await iterator.next()
- await ctx.fiber.dispose()
- }
- })
- it('serves raw call and result entries without parsing tool arguments', async () => {
- const { ctx } = await harness()
- const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- const start = session.append('turn/start', { turn: 1 })
- const call = session.append('tool/call', {
- turn: 1, step: 1, callId: CallId('history-call'), name: 'custom', arguments: '{broken',
- })
- const result = session.append('tool/result', {
- turn: 1, step: 1,
- message: createToolResultMessage({
- callId: CallId('history-call'),
- content: [{ type: 'text', text: 'failed raw output' }],
- isError: true,
- }),
- meta: { persisted: true, count: 3 },
- }, { surfaceOp: 'append' })
- const response = await remote.page({
- address: { kind: 'session', sessionId: session.id },
- throughSeq: session.seq - 1,
- })
- expect(response.ok).toBe(true)
- if (!response.ok) throw new Error('unreachable')
- expect(response.value.events).toEqual([
- { event: start },
- { event: call },
- { event: result },
- ])
- })
- it('counts only append-origin messages toward maxMessages and keeps each compaction summary with its replacement', async () => {
- const { ctx } = await harness()
- const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- session.append('turn/start', { turn: 1 })
- const first = appendUserText(session, 'first prompt')
- appendAssistantText(session, 'first reply', 1)
- const third = appendUserText(session, 'second prompt')
- appendAssistantText(session, 'second reply', 2)
- const shadowed = [...session.surface.nodes]
- // A compaction transaction: a log-only summary record immediately followed by the
- // replacement that shadows the range.
- const summary = appendExtension(session, 'compaction/summary', {
- summary: [{ type: 'text', text: 'summary' }],
- shadowedRange: { start: shadowed[0], end: shadowed.at(-1) },
- shadowedSeqs: shadowed,
- shadowedTokenCount: 0,
- provider: 'p',
- model: 'm',
- })
- session.append('user/message', createUserMessage({
- content: [{ type: 'text', text: '<context_checkpoint>summary</context_checkpoint>' }],
- source: { kind: 'plugin', plugin: 'compact' },
- }), {
- surfaceOp: { op: 'replace', start: shadowed[0] as number, end: shadowed.at(-1) as number },
- sourceEventSeqs: [...shadowed, summary.seq],
- })
- const response = await remote.page({
- address: { kind: 'session', sessionId: session.id },
- throughSeq: session.seq - 1,
- maxMessages: 2,
- })
- if (!response.ok) throw new Error('unreachable')
- const page = response.value.events.map(entry => entry.event)
- // Two append-origin messages fill the page even though a replacement copy of
- // the same event type sits in the window: the copy is model-only.
- const messages = page.filter(event => event.type === 'user/message' || event.type === 'assistant/message')
- expect(messages.map(event => event.seq)).toEqual([third.seq, third.seq + 1, third.seq + 3])
- expect(page.some(event => event.seq === first.seq)).toBe(false)
- expect(response.value.hasMore).toBe(true)
- // The range stays contiguous, so the checkpoint's summary record is readable on
- // the same page as the checkpoint itself.
- const summaryIndex = page.findIndex(event => event.seq === summary.seq)
- expect(summaryIndex).toBeGreaterThan(-1)
- expect(page[summaryIndex + 1]?.seq).toBe(summary.seq + 1)
- expect(page.map(event => event.seq)).toEqual(page.map((_event, index) => third.seq + index))
- })
- it('paginates a message with many provenance sources without variadic argument expansion', async () => {
- const { ctx } = await harness()
- const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- session.append('turn/start', { turn: 1 })
- const sources = Array.from({ length: 128 }, (_unused, index) => session.append('assistant/chunk', {
- turn: 1,
- step: 1,
- chunk: { type: 'text-delta', index, text: 'x' },
- }).seq)
- const message = session.append('assistant/message', {
- turn: 1,
- step: 1,
- message: createMessage({
- role: 'assistant',
- content: [{ type: 'text', text: 'x'.repeat(sources.length) }],
- source: { kind: 'model', provider: 'p', model: 'm' },
- }),
- }, { surfaceOp: 'append', sourceEventSeqs: sources })
- const scalarMin = Math.min
- const min = vi.spyOn(Math, 'min').mockImplementation((...values) => {
- if (values.length > 2) throw new RangeError('variadic minimum rejected by regression harness')
- return scalarMin(...values)
- })
- try {
- const response = await remote.page({
- address: { kind: 'session', sessionId: session.id },
- throughSeq: message.seq,
- maxMessages: 1,
- })
- if (!response.ok) throw new Error('unreachable')
- expect(response.value.events.map(entry => entry.event.seq)).toEqual([...sources, message.seq])
- expect(response.value.hasMore).toBe(true)
- } finally {
- min.mockRestore()
- }
- })
- it('follows a result after turn/end without reading the addressed Session log', async () => {
- const { ctx } = await harness()
- const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
- const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
- const abort = new AbortController()
- const stream = await openFollow(history, session.id, abort.signal)
- const iterator = stream[Symbol.asyncIterator]()
- session.append('turn/start', { turn: 1 })
- await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'turn/start' } } })
- session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-late'), name: 'term', arguments: '{"cmd":"tail"}' })
- await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'tool/call' } } })
- session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
- await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'turn/end' } } })
- const events = vi.spyOn(session, 'events', 'get').mockImplementation(() => {
- throw new Error('live result rescanned Session history')
- })
- try {
- const result = session.append('tool/result', {
- turn: 1, step: 1,
- message: createToolResultMessage({
- callId: CallId('c-late'),
- content: [{ type: 'text', text: 'ok' }],
- isError: false,
- }),
- }, { surfaceOp: 'append' })
- await expect(iterator.next()).resolves.toEqual({
- done: false,
- value: { type: 'event', event: result },
- })
- } finally {
- events.mockRestore()
- abort.abort()
- await iterator.next()
- await ctx.fiber.dispose()
- }
- })
- })
|