schoober-ai-sdk:大模型信息持久化与断点续传
持久化与断点续传
github:Schoober AI SDK GitHub 仓库
各位看官求🌟一下,小的先在此谢过
Agent 任务不同于一次性的 LLM 请求——它可能涉及数十轮 ReAct 循环、多次工具调用、甚至父子任务编排。一次网络中断、服务重启或进程崩溃,如果没有持久化机制,之前的所有进度都会丢失。
schoober-ai-sdk 的持久化设计目标是:任务可以在任何时刻中断,并在之后完整恢复到中断前的状态继续执行。
1. 需要持久化什么
一个运行中的 Agent 任务,其状态分散在多个模块中。要实现断点续传,需要持久化的数据有三类:
| 数据类型 | 内容 | 恢复时的作用 |
|---|---|---|
| TaskState | 任务状态(RUNNING/PAUSED/…)、配置、上下文、子任务ID列表、错误信息 | 重建 TaskExecutor 的内存状态 |
| ApiMessage[] | LLM 的完整对话历史(user/assistant/tool_result) | ReAct 循环从中断处继续推理 |
| UserMessage[] | 前端展示的消息(文本、工具卡片、错误提示) | 恢复 UI 展示 |
| TaskInput | 任务的原始输入 | 重新启动时获取初始输入 |
其中 ApiMessage[] 是最关键的——它就是 LLM 的"记忆"。恢复后,这些消息会作为 context 传入下一次 LLM 请求,LLM 能看到之前所有的推理和工具调用结果,从断点处继续工作。
2. PersistenceManager:存储的抽象
PersistenceManager 是一个纯接口,不绑定任何具体存储实现:
interface PersistenceManager {
// 任务状态
saveTaskState(taskId: string, state: TaskState): Promise<void>;
loadTaskState(taskId: string): Promise<TaskState | null>;
updateTaskState(taskId: string, updates: Partial<TaskState>): Promise<void>;
deleteTaskState(taskId: string): Promise<void>;
listTaskStates(): Promise<TaskState[]>;
// API 消息(LLM 对话历史)
saveApiMessages(taskId: string, messages: ApiMessage[]): Promise<void>;
loadApiMessages(taskId: string): Promise<ApiMessage[]>;
appendApiMessage(taskId: string, message: ApiMessage): Promise<void>;
deleteApiMessages(taskId: string): Promise<void>;
// 用户消息(UI 展示)
saveUserMessages(taskId: string, messages: UserMessage[]): Promise<void>;
loadUserMessages(taskId: string): Promise<UserMessage[]>;
deleteUserMessages(taskId: string): Promise<void>;
// 任务输入
saveTaskInput(taskId: string, input: TaskInput): Promise<void>;
loadTaskInput(taskId: string): Promise<TaskInput | null>;
deleteTaskInput(taskId: string): Promise<void>;
// 初始化(可选)
initialize?(): Promise<void>;
}
几个设计决策:
为什么状态、消息、输入分开存储? 因为它们的读写模式完全不同。TaskState 频繁更新(每次状态变化都写)但读取少;ApiMessage 是追加写入(每轮 ReAct 循环追加一条)但恢复时全量读取;TaskInput 只写一次、读一次。分开存储让实现者可以针对不同模式做优化——比如 TaskState 用 Redis,ApiMessage 用数据库的 append-only 表。
为什么 API 消息同时有 saveApiMessages(全量写入)和 appendApiMessage(追加写入)? 全量写入用于初始化和回滚场景,追加写入用于正常运行中每轮循环结束后的增量保存。实现者可以根据存储特性选择最优策略——比如文件系统用全量写入更简单,数据库用追加更高效。
为什么没有事务支持? SDK 有意不引入事务语义。每种数据独立保存,最坏情况下(如状态保存成功但消息保存失败),恢复时能检测到不一致并处理。引入事务会大幅增加接口的实现复杂度,而 Agent 任务对一致性的要求不像支付系统那么严格——对话历史少一条消息,LLM 通常能通过上下文推断出来。
3. StateManager:防抖持久化
任务状态更新非常频繁——每次工具执行结束、每次子任务状态变化、每次重试计数增加都会触发。如果每次更新都立即写入存储,IO 开销会非常大。StateManager 通过防抖机制解决这个问题。
3.1 内存先行,异步持久化
updateState 的核心逻辑是先更新内存,再异步持久化:
updateState(updates: Partial<TaskState>): void {
const oldStatus = this.state.status;
// 1. 立即更新内存(调用者无需等待)
this.state = { ...this.state, ...updates };
// 2. 状态变化时触发回调
if (updates.status && !oldStatus.equals(this.state.status)) {
this.onStatusChange?.({ oldStatus, newStatus: this.state.status });
}
// 3. 触发防抖持久化(异步,不阻塞)
if (this.autoPersist && this.persistenceManager) {
this.scheduleSave();
}
}
调用者调用 updateState 后立即返回,内存中的状态已经是最新的。持久化在后台异步进行。
3.2 防抖 + 最大等待
scheduleSave 实现了经典的防抖(debounce)+ 最大等待(maxWait)策略:
private scheduleSave(): void {
this.hasPendingChanges = true;
// 记录首次待保存时间(用于 maxWait 计算)
if (!this.firstPendingSaveTime) {
this.firstPendingSaveTime = Date.now();
}
// 清除之前的定时器
if (this.saveTimer) {
clearTimeout(this.saveTimer);
}
const waitedTime = Date.now() - (this.firstPendingSaveTime || Date.now());
if (waitedTime >= this.maxWaitMs) {
// 达到最大等待时间(默认 5 秒),强制保存
this.performSave();
} else {
// 设置新的防抖定时器(默认 1 秒)
const delay = Math.min(this.debounceMs, this.maxWaitMs - waitedTime);
this.saveTimer = setTimeout(() => this.performSave(), delay);
}
}
用时间线来理解:
t=0ms updateState({retryCount: 1}) → 启动 1s 定时器
t=200ms updateState({context: ...}) → 重置定时器,从 200ms 开始计 1s
t=800ms updateState({error: ...}) → 重置定时器,从 800ms 开始计 1s
t=1800ms 定时器触发 → performSave() → 一次 IO,保存了三次更新的合并结果
如果更新持续不断:
t=0ms updateState(...) → 启动定时器
t=500ms updateState(...) → 重置定时器
t=1000ms updateState(...) → 重置定时器
...
t=4900ms updateState(...) → 重置定时器
t=5000ms 达到 maxWait → 强制 performSave()
防抖的好处是:在 ReAct 循环的一轮中(通常持续数秒),多次状态更新只触发一次 IO。maxWait 保证即使更新非常密集,也不会无限延迟保存。
3.3 保存的并发安全
performSave 通过状态标志确保不会并发保存:
private performSave(): void {
if (this.isSaving || !this.hasPendingChanges) {
return;
}
this.isSaving = true;
this.hasPendingChanges = false;
// 捕获当前状态快照(避免保存过程中状态被修改)
const stateSnapshot = { ...this.state };
this.saveStateInternal(stateSnapshot)
.catch(error => {
// 保存失败:标记为有待保存的变更,下次触发时重试
this.hasPendingChanges = true;
})
.finally(() => {
this.isSaving = false;
// 保存期间有新变更,继续调度
if (this.hasPendingChanges) {
this.scheduleSave();
}
});
}
关键细节:保存前先拍状态快照({ ...this.state }),保存过程中内存状态可以继续被修改,不会出现脏写。保存失败后通过 hasPendingChanges 标记,在下一次 updateState 时自动重试。
3.4 关键时刻的强制保存
防抖策略适合运行中的频繁更新,但在某些关键时刻(任务完成、中止、失败),必须立即保存,不能等防抖:
async saveStateNow(): Promise<void> {
// 取消待定的定时器
if (this.saveTimer) {
clearTimeout(this.saveTimer);
}
// 等待当前保存完成(自旋等待)
while (this.isSaving) {
await new Promise(resolve => setTimeout(resolve, 10));
}
// 立即保存
if (this.hasPendingChanges) {
await this.saveStateInternal({ ...this.state });
this.hasPendingChanges = false;
}
}
在 LifecycleManager 的 complete()、abort()、fail() 方法中,都会调用 saveStateNow() 确保终态被持久化。
4. MessageManager:消息的双队列管理
MessageManager 管理两个独立的消息队列,每个队列有自己的防抖保存逻辑:
API 消息的保存还有一个额外保障——通过 Promise 队列串行化:
private saveApiQueue: Promise<void> = Promise.resolve();
async saveApiMessagesNow(): Promise<void> {
// 串入队列,保证前一次保存完成后再执行下一次
this.saveApiQueue = this.saveApiQueue.then(async () => {
await this.persistenceManager.saveApiMessages(
this.taskId,
this.apiMessages
);
});
await this.saveApiQueue;
}
串行化保证了消息的保存顺序——在并发场景下(如多个工具并行执行后同时触发保存),不会出现后一条消息的保存覆盖了前一条。
消息保存的时机
除了防抖触发外,有几个关键时刻会强制立即保存:
| 时机 | 方法 | 原因 |
|---|---|---|
| 每轮 ReAct 循环结束 | MessageCoordinator.finalizeApiMessage() |
LLM 的完整响应必须保存,避免中断后丢失 |
| 任务完成/中止/失败 | MessageCoordinator.saveAllMessagesNow() |
终态前确保所有消息落盘 |
| 暂停时 | MessageCoordinator.saveAllMessagesNow() |
暂停后可能长时间不操作,必须保证状态完整 |
5. 任务恢复:从存储到运行
5.1 恢复入口:Agent.loadTask()
恢复的入口是 Agent.loadTask(),它做的事情和 createTask() 类似,但数据来源是持久化层而非用户输入:
async loadTask(taskId: string, context?: TaskContext, callbacks?: TaskCallbacks): Promise<Task> {
// 1. 从持久化层加载任务状态
const taskState = await this.config.persistence.loadTaskState(taskId);
// 2. 加载任务输入
const taskInput = await this.config.persistence.loadTaskInput(taskId);
// 3. 重建任务配置
const taskConfig: StartTaskConfig = {
...taskState.config,
id: taskId,
input: taskInput || undefined,
};
// 4. 合并上下文(传入的 context 覆盖持久化的 context)
const taskContext: TaskContext = deepMerge(taskState.context, context);
// 5. 创建 TaskExecutor 实例
const taskExecutor = new TaskExecutor(this, taskConfig, taskContext, callbacks, taskId);
// 6. 从持久化状态恢复
await taskExecutor.restoreFromState(taskState);
return taskExecutor;
}
第 4 步的 deepMerge 值得注意——恢复时传入的 context 会与持久化的 context 深度合并,而非简单覆盖。这允许调用者在恢复时注入新的上下文(如更新后的环境信息),同时保留持久化的上下文(如 token 使用统计)。
5.2 TaskRestoreManager:恢复的编排者
TaskExecutor.restoreFromState() 委托给 TaskRestoreManager 执行具体的恢复流程:
async restoreFromState(state: TaskState): Promise<void> {
// 1. 恢复任务状态(处理反序列化)
this.stateManager.restoreState(state);
// 2. 恢复消息历史
await this.restoreMessages();
// 3. 如果任务正在运行,恢复完成 Promise
if (this.stateManager.getState().status.equals(TaskStatus.RUNNING)) {
this.restoreCompletionPromise();
}
}
5.3 状态的反序列化
从存储中恢复的状态需要反序列化处理。JSON 不能直接表示 Date 对象和自定义值对象(如 TaskStatus),StateManager.restoreState() 处理这些转换:
restoreState(state: TaskState): void {
// TaskStatus:JSON 中是字符串 "running",需要转回值对象
let normalizedStatus: TaskStatus;
if (typeof state.status === 'string') {
normalizedStatus = TaskStatus.fromString(state.status);
} else if (state.status instanceof TaskStatus) {
normalizedStatus = state.status;
} else {
normalizedStatus = state.status as TaskStatus;
}
// Date:JSON 中是 ISO 字符串,需要转回 Date 对象
const normalizedState: TaskState = {
...state,
status: normalizedStatus,
startTime: state.startTime
? (state.startTime instanceof Date
? state.startTime
: new Date(state.startTime as any))
: undefined,
endTime: state.endTime
? (state.endTime instanceof Date
? state.endTime
: new Date(state.endTime as any))
: undefined,
};
this.state = normalizedState;
}
typeof state.status === 'string' 的判断兼容了两种来源:直接从 JSON 解析的字符串,和已经经过某层转换的 TaskStatus 对象。
5.4 消息历史的恢复
消息恢复通过 MessageCoordinator.loadAllMessages() → MessageManager.loadAllMessages() 完成:
async loadAllMessages(): Promise<void> {
if (!this.persistenceManager) return;
// 并行加载两种消息
const [apiMessages, userMessages] = await Promise.all([
this.persistenceManager.loadApiMessages(this.taskId),
this.persistenceManager.loadUserMessages(this.taskId),
]);
this.apiMessages = apiMessages;
this.userMessages = userMessages;
}
加载完成后,apiMessages 会被填充到内存中。下一次 ReAct 循环启动时,ReActEngine.executeStep() 通过 this.callbacks.getMessages() 获取到这些恢复的消息,LLM 就能"看到"之前的完整对话历史。
5.5 恢复后的 ReAct 循环
恢复后调用 task.start() 重新启动 ReAct 循环。此时的状态是:
如果状态是 RUNNING,start() 会直接进入 ReAct 循环:
如果状态是 WAITING_FOR_SUBTASK,start() 会路由到子任务:
6. 回滚:不完整轮次的清理
暂停或错误可能发生在 ReAct 循环的中间——LLM 已经返回了部分响应,但工具还没执行完。此时消息历史中可能包含不完整的 assistant 消息,如果直接恢复继续执行,LLM 会看到一条残缺的消息,导致推理混乱。
MessageCoordinator 通过轮次追踪解决这个问题:
startStreamingMessage(): void {
// 记录本轮开始前的消息数量
this.roundStartApiMessageCount = this.messageManager.getApiMessages().length;
this.isRoundInProgress = true;
// ... 插入占位消息
}
rollbackCurrentRoundApiMessages(): void {
if (!this.isRoundInProgress) return;
const currentApiMessages = this.messageManager.getApiMessages();
const rollbackCount = currentApiMessages.length - this.roundStartApiMessageCount;
// 从后向前删除本轮新增的所有 API 消息
for (let i = currentApiMessages.length - 1; i >= this.roundStartApiMessageCount; i--) {
this.messageManager.removeApiMessageById(currentApiMessages[i].id);
}
}
在 LifecycleManager.pause() 中,如果指定了 needRollback = true,会先回滚本轮不完整的消息,再保存:
注意:只回滚 API 消息,保留 User 消息。因为 API 消息是 LLM 的对话历史(影响后续推理质量),而 User 消息是 UI 展示记录(用户已经看到了,删除反而造成困惑)。
7. 数据流全景
将持久化贯穿到整个任务生命周期:
8. 小结
持久化系统的设计原则:
- 接口与实现分离:
PersistenceManager是纯接口,不绑定存储引擎,开发者可以自由选择 Redis / 数据库 / 文件系统 - 数据分类存储:状态、API 消息、User 消息、输入各自独立,适配不同的读写模式
- 防抖 + 强制保存:运行中高频更新用防抖减少 IO,关键时刻(完成/暂停/中止)用强制保存确保数据不丢
- 快照隔离:保存前拍快照,保存期间内存可以继续修改,不会相互干扰
- 回滚保护:暂停时可回滚不完整的 API 消息,避免 LLM 看到残缺的对话历史
- 懒恢复:子任务按需恢复,不全量加载,减少不必要的 IO 和内存开销
更多推荐
所有评论(0)