1
/**2
* Semantic durability checkpoints for model requests, top-level tool dispatch,3
* and completed agent steps.4
* @module @deepseek-ai/dsh-session-checkpoint-policy5
*/7
import type { Context } from '@deepseek-ai/cordis'8
import type { Session } from '@deepseek-ai/dsh-session'9
import type { StreamChunk } from '@deepseek-ai/dsh-llm'10
import { TOOL_ABORTED_BEFORE_DISPATCH, type ToolExecutionResult } from '@deepseek-ai/dsh-tools'11
import type { PreStepDecision } from '@deepseek-ai/dsh-agent'12
import type {} from '@deepseek-ai/dsh-session-persistence'14
/** Cordis plugin name used by Loader diagnostics. */15
export const name = 'session-checkpoint-policy'17
/** Services whose request, tool, session, and persistence boundaries this policy joins. */18
export const inject = ['llm', 'sessionPersistence', 'sessions', 'tools']20
/**21
* Delay construction of the downstream model stream until the complete logged22
* 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
*/29
function 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
}40
/** Materialize the canonical result for a call cancelled before tool dispatch. */41
function 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
}52
/**53
* Install semantic checkpoint listeners. Loop-built model calls checkpoint the54
* logged request before adapter dispatch; top-level tool calls checkpoint their55
* recorded call before the tool body; the next request boundary checkpoints56
* 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-effect59
* boundaries: the downstream adapter or tool body is not invoked.60
*61
* @param ctx - plugin context that owns the listeners.62
*/63
export 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
})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
})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
}