Skip to content

Durable Event Store ​

Durable Event Store 是可恢复执行内核的基础设施。它提供稳定的事件 envelope、 Session 内单调序列、compare-and-append、cursor 分页读取,以及确定性的 Session 生命周期投影。

当前集成阶段

Session 只有在显式设置 SessionOptions.durableEventStore 时才写入 durable 事件;消息历史独立保存在原子 SessionState 投影中。resumeSession() 会自动恢复已接受但尚未跨过 request_started 边界的 Request。已开始但尚无 Turn 的 Request,以及活动 Turn,必须先通过 Recovery Coordinator 原子 rollover;待决权限、未知工具结果、 未知模型结果和已完成 Turn 的 Request 仍需显式消解。non_idempotent 工具和结果未知的模型调用绝不会被自动重放。

安装与导入 ​

协议类型和解析器可从根入口或浏览器安全的 /browser 导入。Node.js JSONL adapter 从 /advanced 导入:

ts
import {
  CommandId,
  DurableEventSubscription,
  type DurableEventDataMap,
  DurableSessionRecoveryCoordinator,
  DurableSessionProjector,
  DurableSessionJournal,
  DurableEventType,
  EventSequence,
  InputId,
  ModelAttemptId,
  PermissionRequestId,
  RequestId,
  SessionId,
  ToolAttemptId,
  TurnId,
} from '@blade-ai/agent-sdk';
import { JsonlDurableEventStore } from '@blade-ai/agent-sdk/advanced';

Event envelope ​

ts
interface DurableEventEnvelope<TType extends DurableEventType> {
  schemaVersion: 4;
  eventId: EventId;
  sequence: EventSequence;
  sessionId: SessionId;
  type: TType;
  data: DurableEventDataMap[TType];
  recordedAt: string;
  occurredAt: string;
  commandId?: CommandId;
  requestId?: RequestId;
  turnId?: TurnId;
  modelAttemptId?: ModelAttemptId;
  toolAttemptId?: ToolAttemptId;
  causationEventId?: EventId;
}
  • eventId 标识单个事件。
  • sequence 在一个 Session 内从 1 开始严格递增。
  • recordedAt 是 Store 提交时间。
  • occurredAt 是业务事件发生时间,未提供时等于 recordedAt。
  • 关联 ID 使用 branded types,防止不同 ID 被误用。

每种事件都有独立的严格 payload,并与 envelope 中必需的关联 ID 组成 DurableEventDraft / DurableEventEnvelope 判别联合。未知字段、缺失 scope ID、 无效枚举和非有限数值都会在追加前被拒绝。

事件必需 scope关键 payload
session_createdSessionsource?、parentSessionId?
session_closedSessionreason
request_acceptedrequestId、commandIdinputId、input、priority、maxTurns?、model?、context?、recovery?
request_startedrequestId空对象
request_completedrequestIdoutput?、usage?
request_failedrequestIderror
request_interruptedrequestIdreason、byInputId?
turn_startedrequestId、turnIdturn、model?
turn_completedrequestId、turnIdturn、hasToolCalls
turn_abortedrequestId、turnIdturn、reason
model_request_startedRequest、Turn、modelAttemptIdmodel、可选 modelIdentity、streaming
model_request_completedRequest、Turn、modelAttemptId完整模型 response
model_request_failedRequest、Turn、modelAttemptIderror
model_request_abortedRequest、Turn、modelAttemptIdreason
tool_scheduledRequest、Turn、modelAttemptId、toolAttemptIdtoolCallId、toolName、modelInput、input、sideEffect、interruptBehavior
tool_startedRequest、Turn、toolAttemptId工具标识、最终 input、解析后的 sideEffect
tool_completedRequest、Turn、toolAttemptId工具标识、result
tool_failedRequest、Turn、toolAttemptId工具标识、error
tool_cancelledRequest、Turn、toolAttemptId工具标识、reason
tool_outcome_unknownRequest、Turn、toolAttemptId工具标识、reason
permission_requestedRequest、Turn、toolAttemptIdpermissionRequestId、工具标识、input、可选 message
permission_resolvedRequest、Turn、toolAttemptIdpermissionRequestId、decision、可选 message
input_appliedrequestId、可选 turnIdinputId、priority

request_accepted.recovery 使用 { requestId, turnId, turn } wire shape。首个 Turn 前的 Request rollover 会 在同一 command 中写入一个 synthetic Turn 作为 provenance;projector 通过 recoveryKind: 'pre_turn_request' 暴露该语义,而不扩展持久化的 recovery 字段。

追加事件 ​

ts
const store = new JsonlDurableEventStore('/var/lib/my-agent');
// 可选:分别限制锁获取和完整 Store 调用。
const boundedWaitStore = new JsonlDurableEventStore('/var/lib/my-agent', {
  lockTimeoutMs: 15_000,
  operationTimeoutMs: 30_000,
});
const sessionId = SessionId('session-123');
const requestId = RequestId('request-123');
const commandId = CommandId('command-123');
const inputId = InputId('input-123');

const result = await store.append(
  sessionId,
  [
    {
      type: DurableEventType.SESSION_CREATED,
      data: { source: 'create' },
    },
    {
      type: DurableEventType.REQUEST_ACCEPTED,
      requestId,
      commandId,
      data: {
        inputId,
        input: 'Run the deployment checks',
        priority: 'next',
      },
    },
    {
      type: DurableEventType.REQUEST_STARTED,
      requestId,
      data: {},
    },
  ],
  {
    expectedLastSequence: null,
  },
);

console.log(result.previousSequence); // null
console.log(result.lastSequence);     // 3

一次 append() 的事件会写入同一个 batch。Store 在写入前完成 schema 校验并 分配连续 sequence。

Compare-and-append ​

expectedLastSequence 提供乐观并发控制:

值语义
undefined在当前 head 后追加
null仅允许空事件流
EventSequence(n)当前 head 必须严格等于 n

前置条件失败时抛出 DurableEventSequenceConflictError,并包含 expected 与 actual sequence。

ts
await store.append(sessionId, events, {
  expectedLastSequence: EventSequence(12),
});

Command journal ​

生产代码应优先通过 DurableSessionJournal 提交生命周期事件。Journal 在 Store 之上增加:

  • 每个实例内的串行提交。
  • 写入前的 lifecycle transition 预演。
  • 所有事件统一写入调用方提供的 commandId。
  • CAS 冲突后的有界刷新与重试。
  • 相同 command 的幂等 replay。
  • DURABLE_EVENT_WRITE_FAILED 后的 read-after-failure 对账。
ts
const journal = await DurableSessionJournal.open(store, sessionId);

const committed = await journal.commit({
  commandId: CommandId('create-session-123'),
  events: [
    {
      type: DurableEventType.SESSION_CREATED,
      data: { source: 'create' },
    },
  ],
});

console.log(committed.status); // committed | replayed | reconciled
console.log(journal.getProjection().status); // open

相同 commandId 和相同事件再次提交时返回 replayed,不会追加第二份事件。 若 command 已存在但内容不同,抛出 DurableCommandConflictError。 一个 command 的事件必须在日志中连续;同一 ID 分散在多个区间也按冲突处理。 依赖当前 projection 生成的 command 应通过 expectedHeadSequence 固定其观察 到的 head,避免同进程或其他 writer 更新状态后仍提交旧决策。Recovery Coordinator 对首次提交的恢复 command 强制设置该前置条件;相同 command 已由 竞争者提交时走幂等 replay,不再校验 head,并仍返回 reconciled。

当底层写入报错后,Journal 会重新读取 canonical log:

  • 能找到完整且匹配的 command:返回 reconciled。
  • 找不到 command 或无法重新读取:抛出 DurableCommandOutcomeUnknownError,且不会自动重试。

后一种情况必须由调用方或更高层 recovery coordinator 对账。自动重试可能重复 执行已经生效但暂时不可见的写入。Journal 会在内存中记录 getUncertainCommandId() 并拒绝其他 command;相同 command 只能再次触发读取 对账。maxConflictRetries 只控制明确 compare-and-append 冲突的重试,不适用于 unknown outcome。

Cursor 读取 ​

ts
const page = await store.read(sessionId, {
  after: EventSequence(20),
  limit: 100,
});

for (const event of page.events) {
  consume(event);
}

console.log(page.nextCursor);
console.log(page.headSequence);
console.log(page.hasMore);

after 是 exclusive cursor。单页限制为 1 到 1000 条;cursor 超过当前 head 会被拒绝,而不是静默返回空结果。

可重连事件订阅 ​

DurableEventSubscription 将 cursor 分页读取封装为 pull-based AsyncIterableIterator。订阅打开时固定一个 replay head,先回放到该位置(含 head 事件本身),再发送一次 caught_up barrier,之后到达的事件标记为 live:

ts
const subscription = await DurableEventSubscription.open(store, sessionId, {
  after: savedCursor,
  pageSize: 100,
  pollIntervalMs: 250,
});

for await (const message of subscription) {
  if (message.type === 'caught_up') {
    markClientReady(message.headSequence);
    continue;
  }

  await deliver(message.event);
  await saveCursor(message.cursor);
}

也可以从已初始化的 Session 创建相同订阅:

ts
const subscription = await session.subscribeDurableEvents({
  after: savedCursor,
});

cursor 是严格版本化的 JSON 值,包含 version、sessionId、sequence 和 eventId,请用 durableEventCursor() 生成、用 parseDurableEventCursor() 解析:解析要求恰好这 4 个键,缺少 version 会 fail closed。 重连时会验证 cursor 指向的事件仍是 canonical log 中的同一事件;跨 Session、 超前、被替换或产生 sequence gap 的 cursor 会 fail closed,而不是跳过数据。

订阅只在消费者请求下一项时读取下一页,内存和读取压力由 pageSize 有界控制。 收到 session_closed 后流自动结束;close() 正常结束等待中的读取, AbortSignal 则以 AbortError 终止。follow: false 只回放订阅创建时已经存在 的快照。应用应在成功处理事件后再持久化该 delivery 的 cursor,从而在断线重连 时获得 at-least-once 交付。

状态投影与恢复分类 ​

DurableSessionProjector 可以逐页消费事件;projectDurableSession() 是一次性 便利函数。投影器会重新执行所有生命周期约束,而不是信任调用方构造的类型:

ts
const projector = new DurableSessionProjector();
let after: EventSequence | undefined;

do {
  const page = await store.read(sessionId, { after, limit: 100 });
  projector.apply(page.events);
  after = page.nextCursor ?? undefined;
  if (!page.hasMore) break;
} while (true);

const projection = projector.snapshot();
const recovery = projector.recoveryPlan();

投影器 fail closed,至少验证:

  • 首条事件必须创建 Session,关闭后不能继续追加生命周期事件。
  • 一个 Session 同时只能有一个活动 Request,一个 Request 同时只能有一个活动 Turn。
  • Request、Turn、Tool Attempt、Permission Request、Command 和已应用 Input ID 不可复用。
  • Turn 编号必须连续;关联的 Session、Request、Turn 与 Tool 标识必须一致。
  • input_applied.turnId 若存在,必须匹配当前 active Turn。
  • Request 终态若携带 causationEventId,必须指向该 Request 最后一次持久化边界 (例如 request_accepted、request_started、input_applied 或 Turn 终止事件)。
  • tool_started 前必须完成权限决策;未完成的工具会阻止 Turn 结束。
  • causationEventId 只能引用同一日志中已出现的事件。

任一事件校验失败后 projector 实例会保持 failed 状态;调用方必须丢弃该实例, 修复 canonical journal 后从头重新投影,不能跳过坏事件继续运行。 Session writer 会把所有独立 Request 终态绑定到最后一次 Request 边界, reconcileRequestOutcome() 则绑定调用方确认过的最后 Turn 终止事件。Journal preview 会拒绝新的无锚点或 stale-boundary 写入。同一 command 中紧邻的 Turn/Request 原子终止不需要引用尚未分配的 event ID;直接绕过 Journal 使用 Store 的写入方必须自行遵守同一契约。

recoveryPlan() 返回以下动作之一:

动作含义
none没有未完成工作
resume_requestRequest 已接受且未开始,可恢复同一 Request
rollover_requestRequest 已开始但首个 Turn 尚未开始,可安全转成新 Request
resume_turnTurn 尚未调用模型,或模型结果已知且工具可安全继续
resolve_permissions必须重新呈现或按策略处理未决权限
reconcile_tool_outcomes工具已开始但没有可靠终态,禁止自动重试
reconcile_model_outcome模型请求已开始但没有可靠终态,必须先向 provider 或业务记录对账
reconcile_request_inputs首个 Turn 前的已应用输入集合不明确,必须显式对账
reconcile_request_outcomeTurn 已结束但 Request 没有终态,必须先确认最终结果

Recovery plan 还返回 activeModelAttempt,并分别返回 retryableToolAttempts、cancelableToolAttempts、unknownToolAttempts 和 pendingPermissions。started 或 tool_outcome_unknown 状态的 pure / idempotent 工具进入 retryable 集合, 前提是它所属的 Model Attempt 已经 completed——模型结果尚未确认的工具不进入该集合, 而是进入 cancelable 集合,其 continuation 标记为 discarded_unconfirmed_model_response; non_idempotent 工具进入 unknown 集合,必须在外部对账后由 tool_completed、tool_failed 或 tool_cancelled 解析。在此之前投影器不会 允许 Turn 结束。

模型调用使用独立的 ModelAttemptId。Session 在调用 provider 前提交 model_request_started,并在调用返回后提交 model_request_completed、 model_request_failed 或 model_request_aborted。活动 model attempt 会阻止 Turn 结束;进程崩溃留下 started attempt 时,plan 返回 reconcile_model_outcome,不会把可能已经计费或完成的调用当作从未发生。 一次 model attempt 表示包含内部 HTTP 重试的一次逻辑模型调用;反应式压缩后的 重新调用会创建新的 attempt。高频 token delta 仍是临时流,完整响应在任何后续 Turn 终态前持久化。 新的 model attempt event 会保留逻辑 Provider 和 API adapter,避免恢复时丢失 模型来源;缺少这些可选字段的旧事件仍然有效。 活动 model attempt 的对账优先于同一 Turn 的权限和工具结果;模型终态提交后, plan 会继续暴露尚未消解的下一层恢复动作。

Recovery Coordinator ​

DurableSessionRecoveryCoordinator 在 projector 的恢复分类之上提供受约束的 状态变更 API:

ts
const coordinator = await DurableSessionRecoveryCoordinator.open(
  store,
  sessionId,
);

const decision = coordinator.planResume();
if (decision.action === 'resume_accepted_request') {
  console.log(decision.request.input);
}

await coordinator.reconcileToolOutcome({
  commandId: CommandId('reconcile-deploy-42'),
  toolAttemptId: ToolAttemptId('attempt-42'),
  outcome: {
    status: 'completed',
    result: { deploymentId: 'dep-42' },
  },
});

对结果未知的模型调用,必须查询 provider request 日志或上层业务记录后显式 对账。命令同时绑定 Request、Turn 和 Model Attempt,且通过 Journal CAS 拒绝 stale decision:

ts
await coordinator.reconcileModelOutcome({
  commandId: CommandId('reconcile-model-42'),
  requestId: RequestId('request-42'),
  turnId: TurnId('turn-42'),
  modelAttemptId: ModelAttemptId('model-attempt-42'),
  outcome: {
    status: 'completed',
    response: {
      content: 'Inspected result',
      usage: {
        promptTokens: 120,
        completionTokens: 30,
        totalTokens: 150,
      },
    },
  },
});

也可以确认 failed 或 aborted。对账后的 Turn 才能进入 prepareTurnRecovery();continuation 会携带有界的已确认模型响应和所有工具 结果,因此不会把未知 provider 调用静默重发。重试同一操作必须复用原 commandId。

reconcileToolOutcome() 只允许 projector 当前仍可解析的 Tool Attempt,并复用 Journal 的 command 幂等和 CAS 语义。completed、failed 与 cancelled 三种结果都可显式提交;调用方重试时必须复用同一个 commandId。

待决权限通过 resolvePermission() 处理。deny 或 cancel 会把 permission_resolved 和对应的 tool_cancelled 放在同一 command/batch 中, 不会暴露“权限已拒绝但工具仍处于 scheduled”的中间状态:

ts
await coordinator.resolvePermission({
  commandId: CommandId('deny-deploy-42'),
  permissionRequestId: PermissionRequestId('permission-42'),
  decision: 'deny',
  message: 'Deployment window closed',
});

权限恢复为 allow 后,应用必须先调用并等待 startToolAttempt({ commandId, toolAttemptId }),再使用返回 projection 中完全 相同的输入执行工具。该方法只使用已持久化的操作事实并把 tool_started 作为 独立幂等 command 提交,从而保持 persist-before-side-effect。经过权限恢复的 工具会保守标记为 non_idempotent,调用方不能在恢复时降低副作用等级。已经处于 started 的 pure / idempotent 工具可直接安全重放,再通过 reconcileToolOutcome() 写入终态;已经处于 started 的 non_idempotent 工具只能查询外部系统后对账,禁止再次执行。

planResume() 仅将没有 input_applied / request_started 且具备完整执行快照 的 accepted Request 标记为 resume_accepted_request。Session 会使用 durable 输入以及持久化的 maxTurns、模型和 Runtime Context 恢复同一个 requestId。 旧日志中缺少执行快照的 Request 继续返回 recovery_required。

request_started 会在 UserPromptSubmit、附件展开、首轮压缩以及 AgentLoop 产生 turn_start 之前持久化。因此“尚无 Turn”只证明主模型调用尚未开始, 不能证明 pre-Turn Hook 或其他准备步骤没有产生副作用。仅当 durable appliedInputIds 恰好为初始输入时,plan 才返回 rollover_request;缺失初始 应用记录或存在额外 steering 输入时返回 reconcile_request_inputs。 初始输入和 steering 输入都在各自 Hook/附件准备前提交 input_applied;提交 失败时不会运行准备副作用,提交成功后即使准备失败也不会自动重放该输入。 若输入在一个已完成 Turn 与下一个 Turn 之间进入准备阶段,plan 同样返回 reconcile_request_inputs;sourceLastTurn 将 rollover 固定到调用方观察到的 Turn 编号。 该保证依赖 writer 遵守相同顺序;若自定义或旧版 writer 曾在 input_applied 前执行输入副作用,journal 无法证明该副作用,调用方必须 fail closed,而不能仅依据 rollover_request 自动继续。

两种动作都通过 prepareRequestRecovery() 收敛,但调用方必须先对账整个 pre-Turn 准备阶段,提供最终准备好的输入,并原样回传观察到的 activeRequest.appliedInputIds:

ts
const request = coordinator.getProjection().activeRequest;
if (!request) throw new Error('No active Request');
const preparedInput = 'prepared input after external reconciliation';

await coordinator.prepareRequestRecovery({
  commandId: CommandId('recover-request-42'),
  requestId: RequestId('request-source-42'),
  inputId: InputId('input-source-42'),
  sourceLastTurn: request.lastTurn,
  recoveryTurnId: TurnId('synthetic-recovery-turn-42'),
  recoveryRequestId: RequestId('request-recovery-42'),
  recoveryInputId: InputId('input-recovery-42'),
  preparation: {
    status: 'reconciled',
    appliedInputIds: request.appliedInputIds ?? [],
    input: preparedInput,
  },
});

该 API 在一个 CAS command 中写入可选的 request_started(源 Request 尚为 accepted 时),随后写入 turn_started (synthetic) → turn_aborted → request_interrupted → request_accepted。source Request/Input ID、appliedInputIds 和 Journal head 都是前置条件;并发出现真实 turn_started 或新的输入应用时,旧决策会失败。 恢复后的 Session 会过滤已经应用或对账的旧 pending 输入,并跳过 UserPromptSubmit、附件展开和首轮 beforeTurn 准备,确保调用方提供的 preparedInput 只执行一次。原始多模态 content parts 保持原结构。

如果 Request 已经完成过至少一个 Turn、当前却没有 active Turn, recoveryPlan() 返回 reconcile_request_outcome。此时最终模型输出可能已经 发生而 Request 终态尚未持久化,SDK 不会自动重试。查询 provider 或上层业务 记录后,使用稳定 commandId 显式提交 completed、failed 或 interrupted:

ts
const terminalPending = coordinator.getProjection().activeRequest;
if (!terminalPending?.lastTurnEventId) {
  throw new Error('No terminal-pending Turn');
}

await coordinator.reconcileRequestOutcome({
  commandId: CommandId('reconcile-request-42'),
  requestId: terminalPending.requestId,
  lastTurnEventId: terminalPending.lastTurnEventId,
  outcome: {
    status: 'completed',
    output: 'already completed',
    usage: { inputTokens: 120, outputTokens: 30, totalTokens: 150 },
  },
});

lastTurnEventId 把结论绑定到调用方实际检查过的 Turn 终止事件;若期间出现更新的 Turn,或重试时改用其他锚点,提交会 fail closed。

对于 resume_turn,调用 prepareTurnRecovery() 可在同一个 CAS command 中:

  1. 取消尚未获得可信终态的 pure / idempotent 或未开始工具;
  2. 以 process_restart 终止旧 Turn 和 Request;
  3. 接受一个包含原始输入、工具终态和 source Request/Turn provenance 的新 continuation Request。

continuation 会把从未执行的工具标记为 not_started,把已开始但可安全重试的 工具标记为 interrupted_before_trusted_completion,并保留已完成工具的权威 结果。恢复后的消息历史会丢弃没有配对结果的旧 tool call;上述 durable 状态随 新的 user continuation 一起送入模型,因此不会构造跨 Store 的伪造 tool result, 也不会留下 provider 不接受的悬空 tool call。多模态原始输入仍以原始 content parts 传递,不会降级成 JSON 文本。

ts
await coordinator.prepareTurnRecovery({
  commandId: CommandId('recover-turn-42'),
  requestId: RequestId('request-source-42'),
  turnId: TurnId('turn-source-42'),
  recoveryRequestId: RequestId('request-recovery-42'),
  recoveryInputId: InputId('input-recovery-42'),
});

const session = await resumeSession({
  ...options,
  sessionId,
});
for await (const event of session.stream()) {
  // 新 accepted Request 通过现有恢复路径执行。
}

整个 rollover 可按相同 commandId 重试,竞争进程只能提交一次。active-Turn provenance 只有在同一 command 中紧邻 turn_aborted → request_interrupted → request_accepted 时才有效;pre-Turn provenance 还要求前置 synthetic turn_started。requestId 和 turnId 是 调用方已观察状态的前置条件,不能隐式切换到更新后的 active Turn。任何已 completed / failed,或在开始执行后变为 cancelled 的 non_idempotent 工具都会触发 DURABLE_RECOVERY_UNSAFE_ROLLOVER,必须保持 fail-closed;该 API 不使用提示词 绕过未知副作用。

当前 runtime 只读写 schema v4。tool_scheduled.modelAttemptId 显式绑定产生该 工具调用的 Model Attempt,modelInput 保存 provider 原始参数,input 保存参数 修复后的执行值;model_request_started.modelIdentity 可记录 provider 身份。 projector 以 canonical JSON 校验工具 ID、 名称和原始参数与已确认模型响应完全一致。即使流式工具在 model_request_completed 前开始调度,模型终态到达时也会反向校验已调度工具。 旧 schema 不会被静默推断或升级;必须先离线迁移到 v4,当前 runtime 才会恢复。

Store deadline 与协作取消 ​

SessionOptions.durableStoreTimeoutMs 默认对每次 Journal、subscription 和 execution lease Store 调用施加 15 秒 deadline。SDK 同时会通过 append、 read、getHeadSequence 及可选 lease 方法传递 AbortSignal。自定义 Store 应在 signal 中止后停止等待,并且不得再开始新的提交型写入。

即使 Store 忽略取消,SDK 的 host watchdog 仍是权威边界。append timeout 会被 视为 command outcome unknown:原 commandId 完成对账前,Journal 会拒绝其他 command。lease 获取、heartbeat、校验、fenced operation 或释放超时会抛出 DurableExecutionLeaseTimeoutError;heartbeat 与活动 fenced operation 超时还会 中止 lease.signal,阻止新的副作用。活动 lease 调用的实际 timeout 还会被限制 在 heartbeat 到 lease expiry 的剩余安全窗口内;即使 heartbeat 调度延迟,基于 单调时钟的本地 expiry watchdog 也会 fail-closed。Session 与 lease 同时配置 Store deadline 时,以更严格的值为准。

lease acquire 超时且结果未知后,同一进程内针对相同 Store、Session 和 owner 的 重试会自动复用之前生成的 leaseId。timeout error 也会暴露该 ID,跨进程重试可 显式传入它,对账同一次 acquire。在不确定窗口完成对账或过期前,同一 identity 必须继续使用相同的 TTL、heartbeat interval 和 Store deadline。

独立使用 Journal、subscription 和 lease API 时可设置 storeTimeoutMs。 JsonlDurableEventStore 提供 operationTimeoutMs;默认值为 Math.min(MAX_DURABLE_STORE_TIMEOUT_MS, lockTimeoutMs + 15000)。取消会移除 进程内排队的锁 waiter 并停止跨进程锁轮询;已经开始的 callback 仍持有锁,直到 清理完成。

执行租约与 fencing ​

DurableExecutionLease 在 Journal CAS 之上提供 opt-in worker 所有权:

ts
const lease = await DurableExecutionLease.acquire(store, sessionId, {
  ownerId: WorkerId('worker-a'),
  ttlMs: 30_000,
  heartbeatIntervalMs: 10_000,
  storeTimeoutMs: 15_000,
});
const journal = await DurableSessionJournal.open(store, sessionId, {
  executionLease: lease,
  executionLeaseStore: store,
  storeTimeoutMs: 15_000,
});

支持租约的 Store 实现 DurableExecutionLeaseStore。获取、续租、释放和 append(..., { executionFence }) 必须通过同一个事务边界串行化。接管会递增 FencingToken;stale writer 收到 DURABLE_EXECUTION_LEASE_LOST,活动租约 期间未携带 fence 的 append 收到 DURABLE_EXECUTION_LEASE_REQUIRED。fencing 要求是粘性的:Store 一旦为 Session 创建过 lease 状态,即使当前 lease 已过期或 释放,后续 append 和 Journal/Recovery Coordinator open 仍必须携带新的活动 lease。调用方通过 executionLeaseStore 显式提供 requiresExecutionLease() 检查端口;append 内的事务校验才是最终权威边界。短时内部持久化可通过 withExecutionLease() 在同一所有权锁内执行,避免 transcript 写入与 lease 接管交错;不要用它包裹模型或工具等长耗时外部 I/O。

进程内 lease handle 会自动 heartbeat。任何续租或校验失败都会中止 lease.signal 并保持 fail-closed。配置 SessionOptions.durableExecutionLeaseStore 与 executionLease 后,Session 会集成该 handle:模型调用与工具副作用在 I/O 前立即校验所有权,Journal commit 携带 fence,subagent 状态与 output 写入也在同一所有权边界内执行;失租时本地 执行关闭,但不会写入伪造的 durable 终态。

fence 保护 SDK 生命周期提交。工具修改其他共享资源时,还必须让下游比较 ExecutionContext.executionFence.fencingToken;通用 SDK 无法强制 fence 已经启动且下游不校验 token 的操作。SDK 会在当前 worker 内终止受管 shell 的 完整进程组并等待退出,但它不是跨进程 supervisor,不能替代共享资源上的 token 校验。

受控 worker handoff ​

session.suspendForHandoff() 在恢复前提供显式的源 worker 屏障。它会封闭新的 后台子 Agent 准入,协作取消活动执行,等待模型/工具收敛与 transcript 持久化, 关闭本地 Runtime,并返回刷新后的 journal head 和 recovery plan。

与 abort() 和 close() 不同,handoff 会保留未完成的 durable Request/Turn, 不会写入 turn_aborted、request_interrupted 或 session_closed;尚未收敛的 模型或工具边界会在方法返回前按保守语义完成记录。继任 worker 必须先按返回的 plan 调用 DurableSessionRecoveryCoordinator,只在 plan 允许后调用 resumeSession()。

仍有后台子 Agent 或归属该 Session 的后台 shell 运行时,该屏障会在取消主 Request 前拒绝;同时要求 durable journal 与 transcript storage 都已配置。 未配置 executionLease 时,handoff 只协调已知的源 worker 和继任 worker。 配置后,源 worker 会持有租约直到清理完成,并在返回前释放,继任 worker 随后可 获取下一个 token。

JSONL 持久化 ​

文件位于:

text
{storageRoot}/durable-events/{base64url(sessionId)}.jsonl
{storageRoot}/durable-events/{base64url(sessionId)}.lease.json

每行是一个完整 append batch,而不是单个事件。该布局保证进程在写入中途崩溃 时,恢复逻辑可以忽略未以换行结束的尾批次,不会接受半个事务。

每次提交:

  1. 在进程级 mutex 内排队,再获取 Session 级跨进程文件锁并重新读取、校验当前 head。
  2. 验证 compare-and-append 前置条件。
  3. 以一次 batch 写入分配连续 sequence。
  4. 调用文件 fsync 后才返回成功。

read() 和 getHeadSequence() 使用同一把锁,因此不会读取另一进程正在截断或 追加的中间状态。本进程 mutex 排队与跨进程锁获取共用默认 10 秒总预算。完整的 直接 Store 调用还受 operationTimeoutMs 限制;默认值为 Math.min(MAX_DURABLE_STORE_TIMEOUT_MS, lockTimeoutMs + 15000),显式 operation timeout 不能小于 lock timeout。跨进程锁使用操作系统 advisory lock: 进程退出或崩溃时由内核立即释放,暂停但仍存活的进程会继续持锁,不会因 wall-clock 超时被另一个进程夺取。lockTimeoutMs: 0 表示只立即尝试一次;锁已 被占用时不会排队或重试。

每个事件文件旁会保留一个 *.jsonl.lock sidecar。它的存在不表示锁当前被占用; 锁状态属于打开的文件描述符。只要仍有进程使用该 Store,就不能手动删除、替换或 移动事件文件及其 lock sidecar。

Native lock addon 仅在首次 Store 操作时加载,不影响其他 SDK API 的导入。在 addon 无法加载或平台不支持 advisory lock 时,Store 以 DURABLE_EVENT_LOCK_FAILED fail closed。当前预构建支持 macOS、glibc Linux 和 Windows 的 x64/arm64;Alpine/musl 及其他目标不受 JsonlDurableEventStore 支持。

事件文件与 lease sidecar 使用 0600 权限,Session ID 经过 base64url 编码, 不会成为文件路径。 事件会保存原始请求输入、完整模型响应、工具输入和模型侧工具结果;调用方必须将 Store 视为敏感数据存储,并自行配置加密、保留期限和访问控制。

一致性边界 ​

JsonlDurableEventStore 将事件日志和 execution lease sidecar 分成独立内部 组件,但二者共享同一 Session advisory lock。它保证同一主机上多个 Node.js 进程针对同一 Session 的 互斥读写、原子 compare-and-append,以及带单调 fencing token 的执行租约。租约 状态和 event append 共用同一把 Session 锁。该保证要求本地文件系统正确实现 advisory lock,不适用于 NFS 等共享网络文件系统。多副本服务必须实现 DurableExecutionLeaseStore,并通过数据库事务让每次 append 与 fence 校验 保持原子;不能依赖未提供 fencing 的超时 lease。

DURABLE_EVENT_WRITE_FAILED 不代表 batch 一定没有写入:底层写入成功后 fsync、解锁或关闭锁文件失败时,提交结果都可能未知。调用方重试前必须重新 读取 head,并通过 commandId 等关联字段核对结果。当前 Store 尚不提供 command 自动去重。

Store 不持久化 token delta、工具 progress 等高频 UI 事件。只有会影响恢复 决策的 domain event 应进入 durable journal。

错误 ​

错误说明
DurableCommandConflictError相同 commandId 对应不同内容或非连续事件区间
DurableCommandOutcomeUnknownError写入失败后无法确认 command 是否提交,Journal 已 fenced
DurableSessionJournalErrorcommand 输入、Store page 或 commit 返回值违反契约
DurableEventSubscriptionError订阅配置、cursor 锚点或 Store page 不满足续读契约
DurableSessionRecoveryRequiredErrorSession 存在未完成工作,必须先执行恢复或对账
DurableExecutionLeaseErrorlease 获取、heartbeat、所有权或持久化状态失败
DURABLE_EXECUTION_LEASE_CONFLICT另一个未过期的 worker lease 正持有 Session
DURABLE_EXECUTION_LEASE_REQUIRED活动 lease 存在,但 append 未携带对应 fence
DURABLE_EXECUTION_LEASE_LOSTlease 已过期、释放或被更高 token 替换
DURABLE_EXECUTION_LEASE_TIMEOUTlease Store 调用超过 deadline
SessionDurableRecorderErrorSession runtime 观察到非法 durable 生命周期状态
DurableSessionRecoveryErrorRecovery Coordinator 状态非法、目标不存在或 rollover 不安全(DURABLE_RECOVERY_INVALID_STATE、DURABLE_RECOVERY_TARGET_NOT_FOUND、DURABLE_RECOVERY_UNSAFE_ROLLOVER)
DurableEventProjectionErrorschema、事件顺序或关联关系不满足生命周期约束
DurableEventSequenceConflictErrorcompare-and-append 前置条件失败
DURABLE_EVENT_INVALID_OPTIONSJSONL Store 构造参数无效
DURABLE_EVENT_LOCK_FAILEDSession 文件锁初始化或获取失败
DURABLE_EVENT_LOCK_TIMEOUT在 lockTimeoutMs 内未能获得 Session 文件锁
DURABLE_EVENT_IO_TIMEOUTdurable Store append、read 或 head 查询超过 deadline
DurableEventStoreError参数、cursor、读写或日志完整性错误

完整、已换行但 schema 错误、事件 ID 重复或 sequence 不连续的记录被视为损坏日志;Store 不会跳过后继续执行。

Released under the MIT License.