fix(session-export): propagate download cancellation

Only the root raw-artifact read received the request signal. Lineage discovery and descendant reads could continue after disconnect, response-body cancellation did not stop the producer, and the root error boundary converted an abort rejection into an ordinary HTTP 500.

Combine request and response-consumer cancellation into the ZIP producer signal, forward it through every cancellable read, check it around the attachment seam, and terminate fflate exactly once when production stops. The pre-stream boundary now rethrows the original abort instead of translating it. Regression tests cover signal propagation, exact cancellation identity at the HTTP boundary, and a reader cancellation interrupting an in-flight descendant read; the bilingual host contract records these lifecycle semantics.
This commit is contained in:
Tianyi Cui
2026-08-11 15:14:10 +08:00
parent 52f0b09e76
commit 192840e198
9 changed files with 176 additions and 24 deletions

View File

@@ -2,5 +2,5 @@
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write packages/host/apiproxy/README.md
README.md: 1d9790beba0b72bec15c3e6f3b35a4d1f0f67d61
README.zh.md: 478088f1167edd4e3c2b55c33e914e85349d2bba
README.md: c7b655816099d786a08ecfae6d794f35f2a9c8e3
README.zh.md: 7577eb025fb84ac40de206e3ed780d92c2ab237a

View File

@@ -28,7 +28,7 @@ Question responses are validated against their pending request before the first
`session.history`'s tail page (`beforeSeq` absent) additionally carries an optional `projections` block — the watermark snapshot of every unit registered on `ctx.sessionProjections` (`@deepseek-ai/dsh-session-projection`), with `asOfSeq` = the last event seq the values reflect (`-1` on an empty log). The gateway also subscribes to the registry's change feed and mints a `session/projection` mux frame per changed unit (`{sessionId, key, value, seq}` — live push state, never logged; clients hold one generic per-session value store under higher-seq-wins). The carrier holds zero domain knowledge (each value passed its unit's own schema inside the registry; the wire schemas keep `values`/`value` wide); loadOlder pages never carry the block, and a composition without the registry serves histories without either surface.
Session-log export is a host-only download surface, not an RPC: `GET /api/session.export?sessionId=…&includeDescendants=true` streams a ZIP whose files are each session's stored artifact text verbatim (the persistence backend's `readRaw` — exact durable bytes decoded from the physical encoding, never a reconstruction from parsed events), root under its original base name plus each subagent descendant under `subagents/<id>/`, and every image any included log references under `media/<attachmentId>.<ext>` (read and verified from the attachment store; a shared image appears once). Each live root or descendant crosses the authoritative `SessionStore.flush` durability barrier immediately before its raw artifact read; cold sessions have no in-memory work to flush. Compression runs on the host with fflate's streaming Zip API, so the response is chunked as it is produced and the host never holds the whole archive in one buffer, and production yields whenever the response queue fills, so a slow consumer bounds the accumulation (fflate's callback is synchronous — the drain point is the only backpressure). It requires the persistence, session-query, and attachment services: a deployment without any answers 500, a persistence backend without per-session raw artifacts answers 501, a missing root session answers 404, and a descendant without a stored artifact or a referenced image that cannot be read fails the stream (fail-loud, never silent under-export). The carrier mounts the endpoint; `ApiProxy.downloads.sessionLog` implements it.
Session-log export is a host-only download surface, not an RPC: `GET /api/session.export?sessionId=…&includeDescendants=true` streams a ZIP whose files are each session's stored artifact text verbatim (the persistence backend's `readRaw` — exact durable bytes decoded from the physical encoding, never a reconstruction from parsed events), root under its original base name plus each subagent descendant under `subagents/<id>/`, and every image any included log references under `media/<attachmentId>.<ext>` (read and verified from the attachment store; a shared image appears once). Each live root or descendant crosses the authoritative `SessionStore.flush` durability barrier immediately before its raw artifact read; cold sessions have no in-memory work to flush. Compression runs on the host with fflate's streaming Zip API, so the response is chunked as it is produced and the host never holds the whole archive in one buffer, and production yields whenever the response queue fills, so a slow consumer bounds the accumulation (fflate's callback is synchronous — the drain point is the only backpressure). Request abort and response-body cancellation stop lineage and artifact work, terminate the active compressor, and propagate as cancellation rather than an HTTP 500. It requires the persistence, session-query, and attachment services: a deployment without any answers 500, a persistence backend without per-session raw artifacts answers 501, a missing root session answers 404, and a descendant without a stored artifact or a referenced image that cannot be read fails the stream (fail-loud, never silent under-export). The carrier mounts the endpoint; `ApiProxy.downloads.sessionLog` implements it.
Session titles ride the generic projection pair like every other domain — the history-tail `projections` block plus `session/projection` frames under the `title` key. Titles do not join `session.list`; cold sessions remain metadata-only there until opening or resuming attaches their logs. `session.rename` accepts an explicit user title (resuming a cold session first), delegating to `ctx.sessionTitle.rename` — the accepted `session/title` event pins the title against automatic regeneration — and returns the normalized title plus its event seq so a client settles its `title` projection cell ahead of the push frame; a title that normalizes to empty returns `title-invalid`.

View File

@@ -28,7 +28,7 @@ Settings 分节中的 `reasoningEffort` 在 agent-default-model 插件配置中
`session.history` 的尾页(不带 `beforeSeq`)额外携带一个可选的 `projections` 块——`ctx.sessionProjections``@deepseek-ai/dsh-session-projection`)上每个已注册单元的水位线快照,`asOfSeq` = 这些值共同反映到的最后一个事件 seq空日志为 `-1`)。网关还订阅注册表的变更流,为每个状态发生变化的单元生成一个 `session/projection` mux 帧(`{sessionId, key, value, seq}`——实时推送状态,绝不入日志;客户端按 seq 高者胜维护一个按会话的通用值仓)。载体不持有任何领域知识(每个值在注册表内部已过其单元自己的 schema协议 schema 对 `values`/`value` 保持宽松loadOlder 页永不携带该块,未装注册表的组合则两个面都不提供。
会话日志导出是宿主侧的下载面,不是 RPC`GET /api/session.export?sessionId=…&includeDescendants=true` 流式返回一个 ZIP其中每个文件都是会话存储工件的逐字原文持久化后端的 `readRaw`——按物理编码解码的确切持久化字节,绝非从解析后事件重建),根会话放在其原始基础文件名下,每个子代理后代放在 `subagents/<id>/` 下,每个被任何包含的日志引用的图片放在 `media/<attachmentId>.<ext>` 下(从附件存储读取并校验;共享图片只出现一次)。每个实时根会话或后代都会在读取原始工件前立即通过权威的 `SessionStore.flush` 持久性屏障;冷会话没有需要 flush 的内存工作。压缩在宿主侧用 fflate 的流式 Zip API 完成响应边生成边分块写出宿主从不把整个归档放进单个缓冲区且每当响应队列填满时生产会让出慢消费者因此只产生有界的积压fflate 的回调是同步的——让出点是唯一的背压手段。它要求同时挂载持久化、session-query 与附件服务:任一缺失应答 500持久化后端不提供每会话原始工件时应答 501根会话缺失时应答 404后代缺少存储工件或引用的图片无法读取则整个流失败fail-loud绝不静默少导出。端点由传输层挂载`ApiProxy.downloads.sessionLog` 实现它。
会话日志导出是宿主侧的下载面,不是 RPC`GET /api/session.export?sessionId=…&includeDescendants=true` 流式返回一个 ZIP其中每个文件都是会话存储工件的逐字原文持久化后端的 `readRaw`——按物理编码解码的确切持久化字节,绝非从解析后事件重建),根会话放在其原始基础文件名下,每个子代理后代放在 `subagents/<id>/` 下,每个被任何包含的日志引用的图片放在 `media/<attachmentId>.<ext>` 下(从附件存储读取并校验;共享图片只出现一次)。每个实时根会话或后代都会在读取原始工件前立即通过权威的 `SessionStore.flush` 持久性屏障;冷会话没有需要 flush 的内存工作。压缩在宿主侧用 fflate 的流式 Zip API 完成响应边生成边分块写出宿主从不把整个归档放进单个缓冲区且每当响应队列填满时生产会让出慢消费者因此只产生有界的积压fflate 的回调是同步的——让出点是唯一的背压手段)。请求中止或响应 body 取消会停止血缘与工件工作、终止活跃压缩器,并继续按取消传播,而不会变成 HTTP 500。它要求同时挂载持久化、session-query 与附件服务:任一缺失应答 500持久化后端不提供每会话原始工件时应答 501根会话缺失时应答 404后代缺少存储工件或引用的图片无法读取则整个流失败fail-loud绝不静默少导出。端点由传输层挂载`ApiProxy.downloads.sessionLog` 实现它。
会话标题与其他所有领域一样搭乘这对通用投影机制——历史尾页的 `projections` 块外加 `title` 键下的 `session/projection` 帧。标题不会加入 `session.list`;冷会话在其中仍只有元数据,直到打开或恢复操作附加其日志。`session.rename` 接受用户显式标题(冷会话先恢复),委托给 `ctx.sessionTitle.rename`——被接受的 `session/title` 事件将标题钉住、不再被自动生成覆盖——并返回规范化后的标题及其事件 seq让 client 在推送帧到达前就结算自己的 `title` 投影格;规范化后为空的标题返回 `title-invalid`

View File

@@ -3506,7 +3506,9 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro
try {
await flushLiveSessionLog(deps, request.sessionId, signal)
root = await deps.sessionPersistence.readRaw(request.sessionId, signal)
signal.throwIfAborted()
} catch {
signal.throwIfAborted()
// Backend read failure: answer 500 without echoing the error, which
// may carry absolute host paths into the browser error bar.
return new Response('session log export failed to read the stored artifact', { status: 500 })

View File

@@ -8,7 +8,9 @@
* every file is byte-identical to the backend's durable artifact or attachment
* store and self-describing through its own header line or media type. Before
* each live session's artifact read, the SessionStore flush barrier makes the
* current in-memory log durable; cold sessions need no barrier.
* current in-memory log durable; cold sessions need no barrier. Request abort
* and response-consumer cancellation share one producer signal and terminate
* the active compressor.
* Compression runs on the host with fflate's streaming Zip API, so the archive
* bytes are produced incrementally and the host never holds the whole archive
* in one buffer; production yields to the consumer whenever the response queue
@@ -206,7 +208,7 @@ export function sessionLogZipFilename(sessionId: string): string {
* missing-session path can answer cleanly before streaming starts).
* @param sessionId - the root session id.
* @param includeDescendants - whether to include every subagent descendant.
* @param signal - optional cancellation for read work.
* @param signal - optional cancellation forwarded to lineage and persistence reads.
* @returns the export entries in zip order.
*/
export async function* sessionLogZipEntries(
@@ -233,7 +235,8 @@ export async function* sessionLogZipEntries(
if (seen.has(id)) continue
seen.add(id)
await flushLiveSessionLog(deps, id, signal)
const raw = await deps.sessionPersistence.readRaw(id)
const raw = await deps.sessionPersistence.readRaw(id, signal)
signal?.throwIfAborted()
if (raw === undefined) {
throw new Error(`subagent "${id}" has no stored log artifact`)
}
@@ -245,12 +248,14 @@ export async function* sessionLogZipEntries(
yield* collect(node.descendants)
}
}
const lineage = await deps.sessionQuery.traceSession(sessionId)
const lineage = await deps.sessionQuery.traceSession(sessionId, signal)
signal?.throwIfAborted()
yield* collect(lineage.descendants)
}
for (const ref of media.values()) {
signal?.throwIfAborted()
const stored = await deps.attachments.readImage(ref)
signal?.throwIfAborted()
yield { path: mediaEntryPath(ref), data: stored.data }
}
}
@@ -335,7 +340,7 @@ async function pushArtifactChunks(
* @param root - the already-read root artifact (first zip entry).
* @param sessionId - the root session id.
* @param includeDescendants - whether to include every subagent descendant.
* @param signal - optional cancellation for read work.
* @param signal - request cancellation combined with response-consumer cancellation.
* @returns the zip byte stream.
*/
export function streamSessionLogZip(
@@ -343,15 +348,24 @@ export function streamSessionLogZip(
root: SessionRawArtifact,
sessionId: SessionId,
includeDescendants: boolean,
signal?: AbortSignal,
signal: AbortSignal,
): ReadableStream<Uint8Array> {
const consumerAbort = new AbortController()
const producerSignal = AbortSignal.any([signal, consumerAbort.signal])
let zip: Zip | undefined
let zipTerminated = false
const terminateZip = (): void => {
if (zip === undefined || zipTerminated) return
zipTerminated = true
zip.terminate()
}
return new ReadableStream<Uint8Array>({
start(controller) {
// fflate invokes the callback synchronously per compressed chunk, so a
// single push can enqueue ahead of a slow consumer; pushArtifactChunks
// yields between chunks once the queue is over-full, bounding the
// accumulation to the queue high-water mark plus one push.
const zip = new Zip((error, data, final) => {
const archive = new Zip((error, data, final) => {
/* v8 ignore next 3 -- fflate reports only internal zip failures, unreachable for valid inputs */
if (error) {
controller.error(error)
@@ -361,25 +375,33 @@ export function streamSessionLogZip(
if (data.byteLength > 0) controller.enqueue(data)
if (final) controller.close()
})
zip = archive
void (async () => {
try {
for await (const entry of sessionLogZipEntries(deps, root, sessionId, includeDescendants, signal)) {
for await (const entry of sessionLogZipEntries(deps, root, sessionId, includeDescendants, producerSignal)) {
const deflate = new ZipDeflate(entry.path, { level: 6 })
zip.add(deflate)
archive.add(deflate)
if ('content' in entry) {
await pushArtifactChunks(deflate, entry.content, controller, signal)
await pushArtifactChunks(deflate, entry.content, controller, producerSignal)
} else {
await pushBinaryChunks(deflate, entry.data, controller, signal)
await pushBinaryChunks(deflate, entry.data, controller, producerSignal)
}
}
zip.end()
archive.end()
} catch (error) {
// A mid-stream failure (missing descendant, cancellation, read
// error) must fail the download rather than ship a truncated archive.
/* v8 ignore next -- typed backends reject with Error, and DOMException is one in Node */
terminateZip()
controller.error(error instanceof Error ? error : new Error(String(error)))
}
})()
},
cancel(reason) {
consumerAbort.abort(
reason instanceof Error ? reason : new Error('session log export stream cancelled'),
)
terminateZip()
},
})
}

View File

@@ -65,6 +65,14 @@ async function buildApi(
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[]
}>
} = {},
) {
const ctx = new Context()
@@ -73,22 +81,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', {
supportsRawArtifacts: persistence !== 'unsupported',
readRaw: async (id: SessionId) => {
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) {
@@ -322,6 +330,126 @@ describe('session.export download endpoint', () => {
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('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}',