Message Persistence

Message Persistence 是将 AI 对话消息持久化到数据库的模块,包括用户消息、AI 回复、工具调用记录和元数据的存储与查询。

#type / concept #status / evergreen #tech / backend #tech / architecture

[!info] related notes

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)
}

常见坑

  1. 半截消息当完整消息用: 流式中断后 status 仍是 streaming,下轮对话把它当历史传给 LLM
  2. 不做增量保存: 每个 token 都 UPDATE 整条消息,性能差
  3. 不记录 token 用量: 无法统计成本
  4. 消息和会话不一致: 会话删除了消息还在

参考资料

创建于 2026/6/30 更新于 2026/7/15