Message Persistence
Message Persistence 是将 AI 对话消息持久化到数据库的模块,包括用户消息、AI 回复、工具调用记录和元数据的存储与查询。
#type / concept
#status / evergreen
#tech / backend
#tech / architecture
[!info] related notes
- 所属 MOC: AI Agent Application MOC
- 相关: Session Management, Conversation Persistence
- 上游: Context Builder — 从数据库加载消息
Message Persistence
一句话定义
Message Persistence 是将 AI 对话消息持久化到数据库的模块。它不只是”保存聊天记录”,还要处理流式消息的增量保存、工具调用记录、消息状态管理和高效查询。
它解决什么问题
AI 对话的消息比传统聊天复杂:
- 流式生成: AI 回复是逐 token 到达的,需要增量保存
- 工具调用: 一条消息可能包含多个工具调用记录
- 消息状态: pending / streaming / completed / failed / cancelled
- 半截消息: 流式中断时可能只有半条回复
核心原理
消息数据结构
type Message struct {
ID string `json:"id"`
SessionID string `json:"session_id"`
TurnID string `json:"turn_id"` // 一轮对话的标识
Role string `json:"role"` // user, assistant, system
Content string `json:"content"`
Status string `json:"status"` // pending, streaming, completed, failed, cancelled
Type string `json:"type"` // text, tool_call, tool_result, ui_event
// 工具调用
ToolCalls []ToolCall `json:"tool_calls,omitempty"`
// 元数据
TokenUsage *TokenUsage `json:"token_usage,omitempty"`
Model string `json:"model,omitempty"`
Duration int `json:"duration_ms,omitempty"`
// 时间
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
CompletedAt *time.Time `json:"completed_at,omitempty"`
}
流式消息的增量保存
// 用户消息:一次性保存
func (m *MessageRepo) SaveUserMessage(ctx context.Context, sessionID, content string) (*Message, error) {
msg := &Message{
ID: generateID(),
SessionID: sessionID,
Role: "user",
Content: content,
Status: "completed",
CreatedAt: time.Now(),
}
return msg, m.db.Save(msg)
}
// AI 回复:增量保存
func (m *MessageRepo) BeginAssistantMessage(ctx context.Context, sessionID string) (*Message, error) {
msg := &Message{
ID: generateID(),
SessionID: sessionID,
Role: "assistant",
Content: "",
Status: "streaming",
CreatedAt: time.Now(),
}
return msg, m.db.Save(msg)
}
func (m *MessageRepo) AppendContent(ctx context.Context, msgID string, delta string) error {
return m.db.Exec(
"UPDATE messages SET content = content || ?, updated_at = NOW() WHERE id = ?",
delta, msgID,
)
}
func (m *MessageRepo) CompleteMessage(ctx context.Context, msgID string, tokenUsage *TokenUsage) error {
return m.db.Exec(
"UPDATE messages SET status = 'completed', completed_at = NOW(), token_usage = ? WHERE id = ?",
tokenUsage, msgID,
)
}
消息过滤(给 Context Builder 用)
func (m *MessageRepo) GetMessagesForContext(ctx context.Context, sessionID string, currentTurnID string) ([]Message, error) {
return m.db.Query(`
SELECT * FROM messages
WHERE session_id = ?
AND turn_id != ? -- 排除当前 turn
AND status = 'completed' -- 只取已完成的
AND type != 'ui_event' -- 排除 UI 事件
ORDER BY created_at ASC
`, sessionID, currentTurnID)
}
常见坑
- 半截消息当完整消息用: 流式中断后 status 仍是 streaming,下轮对话把它当历史传给 LLM
- 不做增量保存: 每个 token 都 UPDATE 整条消息,性能差
- 不记录 token 用量: 无法统计成本
- 消息和会话不一致: 会话删除了消息还在