|
|
@@ -5,7 +5,8 @@
|
|
|
* root → 404, missing descendant → errored stream).
|
|
|
*/
|
|
|
|
|
|
-import { describe, expect, it } from 'vitest'
|
|
|
+import { randomBytes } from 'node:crypto'
|
|
|
+import { describe, expect, it, vi } from 'vitest'
|
|
|
import { Context } from '@deepseek-ai/cordis'
|
|
|
import { unzipSync, strFromU8 } from 'fflate'
|
|
|
import type { ImageAttachmentRef } from '@deepseek-ai/dsh-attachment'
|
|
|
@@ -13,8 +14,7 @@ import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
|
|
|
import type { SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
|
|
|
import type { SessionLineageNode } from '@deepseek-ai/dsh-session-query'
|
|
|
import type { SessionRawArtifact } from '@deepseek-ai/dsh-session-persistence'
|
|
|
-import { toFetchHandler } from '@deepseek-ai/dsh-host-apiproxy'
|
|
|
-import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
|
|
|
+import ApiProxyService, { createApiProxy, toFetchHandler } from '@deepseek-ai/dsh-host-apiproxy'
|
|
|
|
|
|
const sid = (id: string): SessionId => id as SessionId
|
|
|
|
|
|
@@ -59,8 +59,21 @@ async function buildApi(
|
|
|
descendants: SessionLineageNode[] = [],
|
|
|
services: {
|
|
|
query?: boolean
|
|
|
- persistence?: boolean | 'throw'
|
|
|
- attachments?: boolean | ((ref: ImageAttachmentRef) => Promise<ReturnType<typeof storedImage>>)
|
|
|
+ persistence?: boolean | 'throw' | 'unsupported'
|
|
|
+ attachments?: boolean | ((ref: ImageAttachmentRef, signal?: AbortSignal) => Promise<ReturnType<typeof storedImage>>)
|
|
|
+ sessions?: {
|
|
|
+ get(id: SessionId): { readonly id: SessionId } | undefined
|
|
|
+ flush(session: { readonly id: SessionId }): Promise<boolean>
|
|
|
+ }
|
|
|
+ readRaw?: (id: SessionId, signal?: AbortSignal) => Promise<SessionRawArtifact | undefined>
|
|
|
+ traceSession?: (id: SessionId, signal?: AbortSignal) => Promise<{
|
|
|
+ target: { header: SessionHeader; live: boolean; persisted: boolean }
|
|
|
+ ancestors: readonly SessionLineageNode[]
|
|
|
+ complete: boolean
|
|
|
+ root: { header: SessionHeader; live: boolean; persisted: boolean }
|
|
|
+ descendants: readonly SessionLineageNode[]
|
|
|
+ }>
|
|
|
+ compressionLevel?: 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9
|
|
|
} = {},
|
|
|
) {
|
|
|
const ctx = new Context()
|
|
|
@@ -69,21 +82,22 @@ async function buildApi(
|
|
|
const persistence = services.persistence ?? true
|
|
|
if (query) {
|
|
|
ctx.provide('sessionQuery', {
|
|
|
- traceSession: async () => ({
|
|
|
+ traceSession: services.traceSession ?? (async () => ({
|
|
|
target: { header: header('session-root'), live: false, persisted: true },
|
|
|
ancestors: [],
|
|
|
complete: true,
|
|
|
root: { header: header('session-root'), live: false, persisted: true },
|
|
|
descendants,
|
|
|
- }),
|
|
|
+ })),
|
|
|
} as never)
|
|
|
}
|
|
|
if (persistence) {
|
|
|
ctx.provide('sessionPersistence', {
|
|
|
- readRaw: async (id: SessionId) => {
|
|
|
+ supportsRawArtifacts: persistence !== 'unsupported',
|
|
|
+ readRaw: services.readRaw ?? (async (id: SessionId) => {
|
|
|
if (persistence === 'throw') throw new Error('/host/private/session.jsonl')
|
|
|
return artifacts[id]
|
|
|
- },
|
|
|
+ }),
|
|
|
} as never)
|
|
|
}
|
|
|
if (services.attachments !== false) {
|
|
|
@@ -97,9 +111,13 @@ async function buildApi(
|
|
|
readImage,
|
|
|
} as never)
|
|
|
}
|
|
|
+ if (services.sessions !== undefined) ctx.provide('sessions', services.sessions as never)
|
|
|
return createApiProxy(ctx, {
|
|
|
defaultModelSelection: () => ({ provider: 'p', model: 'm' }),
|
|
|
cwd: '/tmp',
|
|
|
+ ...services.compressionLevel === undefined
|
|
|
+ ? {}
|
|
|
+ : { sessionExportCompressionLevel: services.compressionLevel },
|
|
|
})
|
|
|
}
|
|
|
|
|
|
@@ -107,6 +125,19 @@ async function responseBytes(response: Response): Promise<Uint8Array> {
|
|
|
return new Uint8Array(await response.arrayBuffer())
|
|
|
}
|
|
|
|
|
|
+describe('session export compression config', () => {
|
|
|
+ it('defaults to level 6 and rejects values outside the integer 0-9 range', () => {
|
|
|
+ expect(ApiProxyService.Config({})).toEqual({ sessionExportCompressionLevel: 6 })
|
|
|
+ expect(ApiProxyService.Config({ sessionExportCompressionLevel: 0 }))
|
|
|
+ .toEqual({ sessionExportCompressionLevel: 0 })
|
|
|
+ expect(ApiProxyService.Config({ sessionExportCompressionLevel: 9 }))
|
|
|
+ .toEqual({ sessionExportCompressionLevel: 9 })
|
|
|
+ for (const value of [-1, 10, 1.5]) {
|
|
|
+ expect(() => ApiProxyService.Config({ sessionExportCompressionLevel: value } as never)).toThrow()
|
|
|
+ }
|
|
|
+ })
|
|
|
+})
|
|
|
+
|
|
|
describe('session.export download endpoint', () => {
|
|
|
it('streams a ZIP with the root artifact verbatim under its original filename', async () => {
|
|
|
const api = await buildApi({ 'session-root': artifact('session-root') })
|
|
|
@@ -121,6 +152,24 @@ describe('session.export download endpoint', () => {
|
|
|
expect(strFromU8(files['session.jsonl'] as Uint8Array)).toBe(artifact('session-root').content)
|
|
|
})
|
|
|
|
|
|
+ it('uses the resolved compression level for ZIP entries', async () => {
|
|
|
+ const root = artifact('session-root', undefined, 'compressible\n'.repeat(32 * 1024))
|
|
|
+ const storedApi = await buildApi({ 'session-root': root }, [], { compressionLevel: 0 })
|
|
|
+ const compressedApi = await buildApi({ 'session-root': root }, [], { compressionLevel: 9 })
|
|
|
+ const stored = await storedApi.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: false },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const compressed = await compressedApi.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: false },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const storedBytes = await responseBytes(stored)
|
|
|
+ const compressedBytes = await responseBytes(compressed)
|
|
|
+ expect(compressedBytes.byteLength).toBeLessThan(storedBytes.byteLength)
|
|
|
+ expect(strFromU8(unzipSync(compressedBytes)['session.jsonl'] as Uint8Array)).toBe(root.content)
|
|
|
+ })
|
|
|
+
|
|
|
it('includes descendant artifacts under subagents/<id>/ when requested', async () => {
|
|
|
const api = await buildApi({
|
|
|
'session-root': artifact('session-root'),
|
|
|
@@ -143,6 +192,55 @@ describe('session.export download endpoint', () => {
|
|
|
.toBe(artifact('child-a').content)
|
|
|
})
|
|
|
|
|
|
+ it('flushes each live root and descendant immediately before reading its artifact', async () => {
|
|
|
+ const stored: Record<string, SessionRawArtifact> = {
|
|
|
+ 'session-root': artifact('session-root', undefined, 'stale root'),
|
|
|
+ 'child-a': artifact('child-a', sid('session-root'), 'stale child'),
|
|
|
+ }
|
|
|
+ const durable: Record<string, SessionRawArtifact> = {
|
|
|
+ 'session-root': artifact('session-root', undefined, 'durable root'),
|
|
|
+ 'child-a': artifact('child-a', sid('session-root'), 'durable child'),
|
|
|
+ }
|
|
|
+ const flushed: SessionId[] = []
|
|
|
+ const api = await buildApi(stored, [node('child-a')], {
|
|
|
+ sessions: {
|
|
|
+ get: id => durable[id] === undefined ? undefined : { id },
|
|
|
+ flush: async (session) => {
|
|
|
+ const artifactAfterFlush = durable[session.id]
|
|
|
+ if (artifactAfterFlush === undefined) throw new Error('unexpected session')
|
|
|
+ flushed.push(session.id)
|
|
|
+ stored[session.id] = artifactAfterFlush
|
|
|
+ return true
|
|
|
+ },
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await toFetchHandler(api).fetch(
|
|
|
+ new Request('http://host/api/session.export?sessionId=session-root&includeDescendants=true'),
|
|
|
+ )
|
|
|
+ const files = unzipSync(await responseBytes(response))
|
|
|
+ expect(flushed).toEqual([sid('session-root'), sid('child-a')])
|
|
|
+ expect(strFromU8(files['session.jsonl'] as Uint8Array)).toBe('durable root')
|
|
|
+ expect(strFromU8(files['subagents/child-a/session.jsonl'] as Uint8Array)).toBe('durable child')
|
|
|
+ })
|
|
|
+
|
|
|
+ it('reads a cold artifact without asking the live-session store to flush', async () => {
|
|
|
+ const flush = vi.fn(async () => true)
|
|
|
+ const root = artifact('session-root')
|
|
|
+ const api = await buildApi({ 'session-root': root }, [], {
|
|
|
+ sessions: {
|
|
|
+ get: () => undefined,
|
|
|
+ flush,
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: false },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const files = unzipSync(await responseBytes(response))
|
|
|
+ expect(flush).not.toHaveBeenCalled()
|
|
|
+ expect(strFromU8(files['session.jsonl'] as Uint8Array)).toBe(root.content)
|
|
|
+ })
|
|
|
+
|
|
|
it('answers 404 for a missing root session', async () => {
|
|
|
const api = await buildApi({})
|
|
|
const response = await toFetchHandler(api).fetch(
|
|
|
@@ -151,6 +249,15 @@ describe('session.export download endpoint', () => {
|
|
|
expect(response.status).toBe(404)
|
|
|
})
|
|
|
|
|
|
+ it('answers 501 when the persistence backend has no per-session raw artifacts', async () => {
|
|
|
+ const api = await buildApi({}, [], { persistence: 'unsupported' })
|
|
|
+ const response = await toFetchHandler(api).fetch(
|
|
|
+ new Request('http://host/api/session.export?sessionId=session-root'),
|
|
|
+ )
|
|
|
+ expect(response.status).toBe(501)
|
|
|
+ expect(await response.text()).toContain('does not expose per-session raw artifacts')
|
|
|
+ })
|
|
|
+
|
|
|
it('answers 400 when the sessionId query parameter is absent', async () => {
|
|
|
const api = await buildApi({ 'session-root': artifact('session-root') })
|
|
|
const response = await toFetchHandler(api).fetch(
|
|
|
@@ -214,6 +321,37 @@ describe('session.export download endpoint', () => {
|
|
|
expect(strFromU8(files['session.jsonl'] as Uint8Array)).toBe(root.content)
|
|
|
})
|
|
|
|
|
|
+ it('waits for response pull capacity before reading the next archive entry', async () => {
|
|
|
+ const root = artifact('session-root', undefined, [
|
|
|
+ imageEventLine('after-root'),
|
|
|
+ randomBytes(512 * 1024).toString('base64'),
|
|
|
+ ].join('\n'))
|
|
|
+ let imageReads = 0
|
|
|
+ const api = await buildApi({ 'session-root': root }, [], {
|
|
|
+ attachments: async (ref) => {
|
|
|
+ imageReads += 1
|
|
|
+ return storedImage(String(ref.attachmentId), ref.mediaType)
|
|
|
+ },
|
|
|
+ })
|
|
|
+ vi.useFakeTimers()
|
|
|
+ let response: Response | undefined
|
|
|
+ try {
|
|
|
+ response = await toFetchHandler(api).fetch(
|
|
|
+ new Request('http://host/api/session.export?sessionId=session-root'),
|
|
|
+ )
|
|
|
+ // Exhausting timer turns must not advance a producer whose byte queue is
|
|
|
+ // full; only a consumer pull can release it.
|
|
|
+ await vi.runAllTimersAsync()
|
|
|
+ expect(imageReads).toBe(0)
|
|
|
+ } finally {
|
|
|
+ vi.useRealTimers()
|
|
|
+ }
|
|
|
+ if (response === undefined) throw new Error('missing export response')
|
|
|
+ const files = unzipSync(await responseBytes(response))
|
|
|
+ expect(imageReads).toBe(1)
|
|
|
+ expect(files['media/after-root.png']).toEqual(storedImage('after-root').data)
|
|
|
+ })
|
|
|
+
|
|
|
it('exports an empty artifact as an empty zip entry', async () => {
|
|
|
const root = { ...artifact('session-root'), content: '' }
|
|
|
const api = await buildApi({ 'session-root': root })
|
|
|
@@ -254,10 +392,179 @@ describe('session.export download endpoint', () => {
|
|
|
)
|
|
|
expect(response.status).toBe(500)
|
|
|
const body = await response.text()
|
|
|
- expect(body).toBe('session log export failed to read the stored artifact')
|
|
|
+ expect(body).toBe('session log export failed to prepare the stored artifact')
|
|
|
+ expect(body).not.toContain('/host/private/')
|
|
|
+ })
|
|
|
+
|
|
|
+ it('answers the private-error-safe 500 when the live root flush fails', async () => {
|
|
|
+ const api = await buildApi({ 'session-root': artifact('session-root') }, [], {
|
|
|
+ sessions: {
|
|
|
+ get: id => ({ id }),
|
|
|
+ flush: async () => { throw new Error('/host/private/flush-state') },
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await toFetchHandler(api).fetch(
|
|
|
+ new Request('http://host/api/session.export?sessionId=session-root'),
|
|
|
+ )
|
|
|
+ expect(response.status).toBe(500)
|
|
|
+ const body = await response.text()
|
|
|
+ expect(body).toBe('session log export failed to prepare the stored artifact')
|
|
|
expect(body).not.toContain('/host/private/')
|
|
|
})
|
|
|
|
|
|
+ it('forwards one request signal through root, lineage, and descendant reads', async () => {
|
|
|
+ const reads: Array<{ id: SessionId; signal: AbortSignal | undefined }> = []
|
|
|
+ const traces: AbortSignal[] = []
|
|
|
+ const api = await buildApi({}, [node('child-a')], {
|
|
|
+ readRaw: async (id, signal) => {
|
|
|
+ reads.push({ id, signal })
|
|
|
+ return id === sid('session-root')
|
|
|
+ ? artifact('session-root')
|
|
|
+ : artifact('child-a', sid('session-root'))
|
|
|
+ },
|
|
|
+ traceSession: async (_id, signal) => {
|
|
|
+ if (signal !== undefined) traces.push(signal)
|
|
|
+ return {
|
|
|
+ target: { header: header('session-root'), live: false, persisted: true },
|
|
|
+ ancestors: [],
|
|
|
+ complete: true,
|
|
|
+ root: { header: header('session-root'), live: false, persisted: true },
|
|
|
+ descendants: [node('child-a')],
|
|
|
+ }
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const controller = new AbortController()
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: true },
|
|
|
+ controller.signal,
|
|
|
+ )
|
|
|
+ await response.arrayBuffer()
|
|
|
+ const producerSignal = traces[0]
|
|
|
+ if (producerSignal === undefined) throw new Error('missing lineage signal')
|
|
|
+ expect(reads[0]).toEqual({ id: sid('session-root'), signal: controller.signal })
|
|
|
+ expect(reads[1]).toEqual({ id: sid('child-a'), signal: producerSignal })
|
|
|
+ const cancellation = new Error('request cancelled after response')
|
|
|
+ controller.abort(cancellation)
|
|
|
+ expect(producerSignal.aborted).toBe(true)
|
|
|
+ expect(producerSignal.reason).toBe(cancellation)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('preserves request cancellation instead of translating it to HTTP 500', async () => {
|
|
|
+ const api = await buildApi({ 'session-root': artifact('session-root') })
|
|
|
+ const controller = new AbortController()
|
|
|
+ const cancellation = new Error('request cancelled')
|
|
|
+ controller.abort(cancellation)
|
|
|
+ await expect(api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: false },
|
|
|
+ controller.signal,
|
|
|
+ )).rejects.toBe(cancellation)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('aborts descendant work and terminates ZIP production when its reader cancels', async () => {
|
|
|
+ let reportDescendantStarted!: (signal: AbortSignal) => void
|
|
|
+ const descendantStarted = new Promise<AbortSignal>((resolve) => {
|
|
|
+ reportDescendantStarted = resolve
|
|
|
+ })
|
|
|
+ const api = await buildApi({}, [node('child-a')], {
|
|
|
+ readRaw: async (id, signal) => {
|
|
|
+ if (id === sid('session-root')) return artifact('session-root')
|
|
|
+ if (signal === undefined) throw new Error('missing descendant signal')
|
|
|
+ reportDescendantStarted(signal)
|
|
|
+ return new Promise((_, reject) => {
|
|
|
+ signal.addEventListener('abort', () => {
|
|
|
+ reject(signal.reason as Error)
|
|
|
+ }, { once: true })
|
|
|
+ })
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: true },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const reader = response.body?.getReader()
|
|
|
+ if (reader === undefined) throw new Error('missing response body')
|
|
|
+ const descendantSignal = await descendantStarted
|
|
|
+ const cancellation = new Error('download consumer left')
|
|
|
+ await reader.cancel(cancellation)
|
|
|
+ expect(descendantSignal.aborted).toBe(true)
|
|
|
+ expect(descendantSignal.reason).toBe(cancellation)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('aborts attachment reads when its reader cancels', async () => {
|
|
|
+ let reportAttachmentStarted!: (signal: AbortSignal) => void
|
|
|
+ const attachmentStarted = new Promise<AbortSignal>((resolve) => {
|
|
|
+ reportAttachmentStarted = resolve
|
|
|
+ })
|
|
|
+ const root = artifact('session-root', undefined, [
|
|
|
+ '{"type":"session","version":0,"id":"session-root","createdAt":1000}',
|
|
|
+ imageEventLine('slow-img'),
|
|
|
+ ].join('\n') + '\n')
|
|
|
+ const api = await buildApi({ 'session-root': root }, [], {
|
|
|
+ attachments: async (_ref, signal) => {
|
|
|
+ if (signal === undefined) throw new Error('missing attachment signal')
|
|
|
+ reportAttachmentStarted(signal)
|
|
|
+ return new Promise((_, reject) => {
|
|
|
+ signal.addEventListener('abort', () => {
|
|
|
+ reject(signal.reason as Error)
|
|
|
+ }, { once: true })
|
|
|
+ })
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: false },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const reader = response.body?.getReader()
|
|
|
+ if (reader === undefined) throw new Error('missing response body')
|
|
|
+ const attachmentSignal = await attachmentStarted
|
|
|
+ const cancellation = new Error('download consumer left during attachment read')
|
|
|
+ await reader.cancel(cancellation)
|
|
|
+ expect(attachmentSignal.aborted).toBe(true)
|
|
|
+ expect(attachmentSignal.reason).toBe(cancellation)
|
|
|
+ })
|
|
|
+
|
|
|
+ it('uses a stable Error reason when its reader cancels without one', async () => {
|
|
|
+ let reportDescendantStarted!: (signal: AbortSignal) => void
|
|
|
+ const descendantStarted = new Promise<AbortSignal>((resolve) => {
|
|
|
+ reportDescendantStarted = resolve
|
|
|
+ })
|
|
|
+ const api = await buildApi({}, [node('child-a')], {
|
|
|
+ readRaw: async (id, signal) => {
|
|
|
+ if (id === sid('session-root')) return artifact('session-root')
|
|
|
+ if (signal === undefined) throw new Error('missing descendant signal')
|
|
|
+ reportDescendantStarted(signal)
|
|
|
+ return new Promise((_, reject) => {
|
|
|
+ signal.addEventListener('abort', () => {
|
|
|
+ reject(signal.reason as Error)
|
|
|
+ }, { once: true })
|
|
|
+ })
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: true },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ const reader = response.body?.getReader()
|
|
|
+ if (reader === undefined) throw new Error('missing response body')
|
|
|
+ const descendantSignal = await descendantStarted
|
|
|
+ await reader.cancel()
|
|
|
+ expect(descendantSignal.reason).toEqual(new Error('session log export stream cancelled'))
|
|
|
+ })
|
|
|
+
|
|
|
+ it('normalizes a non-Error descendant failure before erroring the stream', async () => {
|
|
|
+ const api = await buildApi({}, [node('child-a')], {
|
|
|
+ readRaw: async (id) => {
|
|
|
+ if (id === sid('session-root')) return artifact('session-root')
|
|
|
+ throw 'descendant read failed'
|
|
|
+ },
|
|
|
+ })
|
|
|
+ const response = await api.downloads.sessionLog(
|
|
|
+ { sessionId: sid('session-root'), includeDescendants: true },
|
|
|
+ new AbortController().signal,
|
|
|
+ )
|
|
|
+ await expect(response.arrayBuffer()).rejects.toEqual(new Error('descendant read failed'))
|
|
|
+ })
|
|
|
+
|
|
|
it('includes media objects referenced by the root log under media/<id>.<ext>', async () => {
|
|
|
const root = artifact('session-root', undefined, [
|
|
|
'{"type":"session","version":0,"id":"session-root","createdAt":1000}',
|