返回源码地图

packages/session/session-checkpoint-policy/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 16;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Semantic durability checkpoints for model requests, top-level tool dispatch,
3 * and completed agent steps.
4 * @module @deepseek-ai/dsh-session-checkpoint-policy
5 */
6
7import type { Context } from '@deepseek-ai/cordis'
8import type { Session } from '@deepseek-ai/dsh-session'
9import type { StreamChunk } from '@deepseek-ai/dsh-llm'
10import { TOOL_ABORTED_BEFORE_DISPATCH, type ToolExecutionResult } from '@deepseek-ai/dsh-tools'
11import type { PreStepDecision } from '@deepseek-ai/dsh-agent'
12import type {} from '@deepseek-ai/dsh-session-persistence'
13
14/** Cordis plugin name used by Loader diagnostics. */
15export const name = 'session-checkpoint-policy'
16
17/** Services whose request, tool, session, and persistence boundaries this policy joins. */
18export const inject = ['llm', 'sessionPersistence', 'sessions', 'tools']
19
20/**
21 * Delay construction of the downstream model stream until the complete logged
22 * request prefix is durable. A checkpoint rejection prevents adapter dispatch.
23 *
24 * @param ctx - plugin context that owns the session store.
25 * @param session - live session named by the model request.
26 * @param next - downstream `llm/stream` chain.
27 * @returns a stream that checkpoints before requesting its first chunk.
28 */
29function afterCheckpoint(
30 ctx: Context,
31 session: Session,
32 next: () => AsyncIterable<StreamChunk>,
33): AsyncIterable<StreamChunk> {
34 return (async function* (): AsyncIterable<StreamChunk> {
35 await ctx.sessions.flush(session)
36 yield* next()
37 })()
38}
39
40/** Materialize the canonical result for a call cancelled before tool dispatch. */
41function abortedBeforeDispatchResult(): ToolExecutionResult {
42 return {
43 content: [{ type: 'text', text: 'Error: tool call aborted before dispatch' }],
44 isError: true,
45 error: {
46 message: 'tool call aborted before dispatch',
47 info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
48 },
49 }
50}
51
52/**
53 * Install semantic checkpoint listeners. Loop-built model calls checkpoint the
54 * logged request before adapter dispatch; top-level tool calls checkpoint their
55 * recorded call before the tool body; the next request boundary checkpoints
56 * the preceding response/result batch. Nested tool dispatches reuse the durable outer call.
57 *
58 * Checkpoint failures are fail-closed at the model and tool side-effect
59 * boundaries: the downstream adapter or tool body is not invoked.
60 *
61 * @param ctx - plugin context that owns the listeners.
62 */
63export function apply(ctx: Context): void {
64 ctx.on('llm/stream', (options, next): AsyncIterable<StreamChunk> => {
65 if (options.sessionId === undefined) return next()
66 const session = ctx.sessions.get(options.sessionId)
67 return session === undefined ? next() : afterCheckpoint(ctx, session, next)
68 })
69
70 ctx.on('tools/execute', async (exec, next): Promise<ToolExecutionResult> => {
71 if (exec.agent === undefined || exec.parent !== undefined) return next()
72 await ctx.sessions.flush(exec.agent.session)
73 if (exec.signal.aborted) return abortedBeforeDispatchResult()
74 return next()
75 })
76
77 // Before each request, persist everything committed by the preceding step;
78 // the first step's call is an intentional no-op beyond any prompt intake.
79 ctx.on('agent/pre-step', async ({ agent }, next): Promise<PreStepDecision> => {
80 await ctx.sessions.flush(agent.session)
81 return next()
82 })
83}