上一提交(1619399)把已读改成按读者记录后,回填只把历史 `status='read'` 记到**主收件人**
名下 —— 对抄送方等于"突然多出一批未读旧邮件"。这不是理论风险,**当天就在野外发生了**:
opencode 桥(部署后 46 分钟):
16:06:48 [mail-bridge] 已接入 http://127.0.0.1:8180,身份 opencode(密钥认证)
16:06:49 [mail-bridge] 补投 2 封离线期间的邮件(共 2 封未读)
→ 它对 05:42 那封「打个招呼」**又回了两次信**(08:07:21Z / 08:08:37Z)
即桥的 `pending_mails = CountUnread` 因迁移变大 ⇒ 桥一重启就把旧信当漏投重放并再次回信。
两个人工探针当时都只覆盖主收件人,恰好绕过这个面("同一封被多人共享"的坑,
判据必须站到每个收件人各自的位置上)。
修法(`backfillMailReadsCC`):迁移前的邮件(`created_at <` 切换时刻)凡 `status='read'`,
给它的**所有收件人**(主 + 抄送)各补一行 —— 与旧模型下"所有人看到的都是已读"完全一致;
迁移后的邮件一律不碰(那条界线是判据核心:越界就会把"某个人读过"错写成"所有收件人都读过")。
切换时刻:迁移时写进 `app_meta(read_model_switchover_at)`;老库没有这个键时退化成
`MIN(mail_reads.read_at)`(那张表的第一笔写入就是回填批次)。
判据 `internal/db/migrate_reads_test.go`:迁移前的老邮件必须补到抄送方、**迁移后的不能碰**、
重复执行不重复插。扰动验证:去掉时间界线 → 判据红(补记 2 行,期望 1)。
实测收口:
- 迁移日志「再给 4 个抄送方补记历史已读」;"抄送方仍算未读(已读邮件)" 计数 **0**。
- **重放反证**:重启 opencode / pi 的桥 → 无"补投"行、3 分钟内 0 封新邮件 ✓
(对比修复前 opencode 重启即补投并回信)。
- 清掉那 2 封由这次迁移产生的误回信(happy-pixel 回到 6 封)。
- 全量 server 10 包 + client/electron vitest 239 + 五 Agent 演练 20/20 全绿。
教训:**语义迁移必须让"可观测状态"保持不变**,新语义只对迁移后新增的对象生效 ——
否则用户会看到一批凭空冒出来的未读,而下游(这里是桥的补投)会把它当真实信号动作。
366 lines
16 KiB
Go
366 lines
16 KiB
Go
package db
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
_ "embed"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
|
||
"github.com/google/uuid"
|
||
"strings"
|
||
"time"
|
||
)
|
||
|
||
//go:embed migrations/init.sql
|
||
var initSQLPostgres string
|
||
|
||
//go:embed migrations/init_sqlite.sql
|
||
var initSQLSQLite string
|
||
|
||
// Migrate 建表建索引。两种方言各有一份 schema,语义保持一致。
|
||
func Migrate(ctx context.Context) error {
|
||
switch D {
|
||
case Postgres:
|
||
// PG 侧含 DO $$ … $$ 迁移块,必须整体提交
|
||
if _, err := DB.ExecContext(ctx, initSQLPostgres); err != nil {
|
||
return fmt.Errorf("migrate postgres: %w", err)
|
||
}
|
||
case SQLite:
|
||
// modernc.org/sqlite 的 Exec 不接受多语句,逐条执行
|
||
for i, stmt := range splitStatements(initSQLSQLite) {
|
||
if _, err := DB.ExecContext(ctx, stmt); err != nil {
|
||
return fmt.Errorf("migrate sqlite (语句 #%d: %.60s): %w", i+1, stmt, err)
|
||
}
|
||
}
|
||
// CREATE TABLE IF NOT EXISTS 不会给**已存在**的表补列,而 SQLite 又没有
|
||
// ADD COLUMN IF NOT EXISTS。已部署的库靠这一步补齐新列。
|
||
if err := addMissingColumns(ctx); err != nil {
|
||
return err
|
||
}
|
||
default:
|
||
return fmt.Errorf("migrate: 未初始化的方言")
|
||
}
|
||
|
||
if err := backfillMailReads(ctx); err != nil {
|
||
return err
|
||
}
|
||
if err := backfillMailReadsCC(ctx); err != nil {
|
||
return err
|
||
}
|
||
|
||
fmt.Printf("数据库迁移完成(%s)\n", D)
|
||
return nil
|
||
}
|
||
|
||
/*
|
||
backfillMailReads 把已读模型从"邮件级"迁到"读者级"时补一次历史数据。
|
||
|
||
背景:`mails.status='read'` 原先表示"有人读过",但**没记是谁读的**。新的
|
||
mail_reads 表按读者记录,因此老数据只能推断 —— 取主收件人(to_name)作为默认读者:
|
||
绝大多数已读发生在主收件人身上,而抄送方"被代读"的情况本来就是要修掉的错。
|
||
|
||
**只能跑一次**:每次启动都跑的话,它会把"某个抄送方读过"的邮件按主收件人写成已读 ——
|
||
正是这次要修的语义错误。所以用 app_meta 里的标记守住(见 init_sqlite.sql 的注释)。
|
||
*/
|
||
func backfillMailReads(ctx context.Context) error {
|
||
const marker = "read_model_per_recipient_v1"
|
||
var v string
|
||
err := DB.QueryRowContext(ctx, `SELECT value FROM app_meta WHERE key = $1`, marker).Scan(&v)
|
||
if err == nil {
|
||
return nil // 已迁过
|
||
}
|
||
if !errors.Is(err, sql.ErrNoRows) {
|
||
return fmt.Errorf("backfill mail_reads: 读标记失败: %w", err)
|
||
}
|
||
|
||
// INSERT ... SELECT ... WHERE NOT EXISTS 两种方言都认;WHERE 只是防御
|
||
// (标记保证只跑一次,但万一上次中断在半路,这条能安全续上)。
|
||
res, err := DB.ExecContext(ctx, `
|
||
INSERT INTO mail_reads (mail_id, reader_name)
|
||
SELECT m.mail_id, m.to_name
|
||
FROM mails m
|
||
WHERE m.status = 'read'
|
||
AND NOT EXISTS (
|
||
SELECT 1 FROM mail_reads r WHERE r.mail_id = m.mail_id AND r.reader_name = m.to_name
|
||
)`)
|
||
if err != nil {
|
||
return fmt.Errorf("backfill mail_reads: %w", err)
|
||
}
|
||
n, _ := res.RowsAffected()
|
||
|
||
if _, err := DB.ExecContext(ctx, `INSERT INTO app_meta (key, value) VALUES ($1, $2)`, marker, fmt.Sprintf("done rows=%d", n)); err != nil {
|
||
return fmt.Errorf("backfill mail_reads: 写标记失败: %w", err)
|
||
}
|
||
// 记下切换时刻:抄送方的回填(backfillMailReadsCC)需要它划"迁移前的老邮件"这条线。
|
||
// 记不进去不算致命 —— 那条路会退化成用 MIN(read_at) 推断(见 switchoverAt)。
|
||
if _, err := DB.ExecContext(ctx,
|
||
`INSERT INTO app_meta (key, value) VALUES ($1, $2)`,
|
||
"read_model_switchover_at", time.Now().UTC().Format("2006-01-02 15:04:05")); err != nil {
|
||
return fmt.Errorf("backfill mail_reads: 写切换时刻失败: %w", err)
|
||
}
|
||
if n > 0 {
|
||
fmt.Printf("已读模型迁移:把 %d 封历史已读邮件记到主收件人名下(一次性)\n", n)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// switchoverAt 返回已读语义的切换时刻(UTC):
|
||
// 迁移时写过就用它;没写过(例如切换发生在这个键出现之前)退化成 mail_reads 里最早的一行
|
||
// —— 那张表的第一笔写入就是回填批次,所以 MIN(read_at) 恰好等于切换时刻。
|
||
func switchoverAt(ctx context.Context) (time.Time, error) {
|
||
var v string
|
||
err := DB.QueryRowContext(ctx, `SELECT value FROM app_meta WHERE key = 'read_model_switchover_at'`).Scan(&v)
|
||
if err == nil {
|
||
if t, perr := time.Parse("2006-01-02 15:04:05", v); perr == nil {
|
||
return t.UTC(), nil
|
||
}
|
||
} else if !errors.Is(err, sql.ErrNoRows) {
|
||
return time.Time{}, err
|
||
}
|
||
var raw sql.NullString
|
||
if err := DB.QueryRowContext(ctx, `SELECT MIN(read_at) FROM mail_reads`).Scan(&raw); err != nil {
|
||
return time.Time{}, err
|
||
}
|
||
if !raw.Valid {
|
||
return time.Now().UTC(), nil // 表还是空的:本机没有"老邮件"要照顾
|
||
}
|
||
for _, layout := range []string{"2006-01-02 15:04:05.000", "2006-01-02 15:04:05", time.RFC3339} {
|
||
if t, perr := time.Parse(layout, raw.String); perr == nil {
|
||
return t.UTC(), nil
|
||
}
|
||
}
|
||
return time.Now().UTC(), nil
|
||
}
|
||
|
||
/*
|
||
backfillMailReadsCC 把**抄送方**在旧模型下的状态也补上,让这次迁移对使用者是"行为不变"的。
|
||
|
||
为什么必须补(2026-09-13 实测的回归):旧模型里 `mails.status='read'` 对**所有人**都算已读。
|
||
第一版回填只把它记到主收件人名下 ⇒ 抄送方在新模型下突然看到一批"未读"的旧邮件 ⇒
|
||
桥的补投判据 `pending_mails = CountUnread` 跟着变大 ⇒ **桥一重启就把旧信当漏投重放并再次回信**。
|
||
现场证据(opencode 桥,部署后 46 分钟):
|
||
|
||
16:06:48 [mail-bridge] 已接入 http://127.0.0.1:8180,身份 opencode
|
||
16:06:49 [mail-bridge] 补投 2 封离线期间的邮件(共 2 封未读)
|
||
|
||
随后它对 05:42 那封"打个招呼"又回了两封信 —— 纯粹是我这次迁移造成的新邮件噪声。
|
||
|
||
判据:**迁移前的邮件**(created_at < 切换时刻)凡是 `status='read'`,就给它的**所有收件人**
|
||
(主收件人 + 抄送)各记一行 —— 与旧模型下"所有人看到的都是已读"完全一致。
|
||
迁移后的邮件不走这条路(那时已是按读者记录,只给真正读过的人记)—— 这条边界是本函数的判据核心。
|
||
*/
|
||
func backfillMailReadsCC(ctx context.Context) error {
|
||
const marker = "read_model_cc_backfill_v1"
|
||
var v string
|
||
err := DB.QueryRowContext(ctx, `SELECT value FROM app_meta WHERE key = $1`, marker).Scan(&v)
|
||
if err == nil {
|
||
return nil
|
||
}
|
||
if !errors.Is(err, sql.ErrNoRows) {
|
||
return fmt.Errorf("backfill mail_reads(cc): 读标记失败: %w", err)
|
||
}
|
||
switchover, err := switchoverAt(ctx)
|
||
if err != nil {
|
||
return fmt.Errorf("backfill mail_reads(cc): 取切换时刻失败: %w", err)
|
||
}
|
||
n, err := backfillMailReadsCCSince(ctx, switchover)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if _, err := DB.ExecContext(ctx, `INSERT INTO app_meta (key, value) VALUES ($1, $2)`, marker, fmt.Sprintf("done rows=%d", n)); err != nil {
|
||
return fmt.Errorf("backfill mail_reads(cc): 写标记失败: %w", err)
|
||
}
|
||
if n > 0 {
|
||
fmt.Printf("已读模型迁移:再给 %d 个抄送方补记历史已读(迁移行为保持不变)\n", n)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// backfillMailReadsCCSince 是上面那条的一次性逻辑本体(switchover 显式传入,便于判据直接调)。
|
||
func backfillMailReadsCCSince(ctx context.Context, switchover time.Time) (int, error) {
|
||
rows, err := DB.QueryContext(ctx, `
|
||
SELECT m.mail_id, m.cc_list
|
||
FROM mails m
|
||
WHERE m.status = 'read'
|
||
AND m.created_at < $1`, switchover)
|
||
if err != nil {
|
||
return 0, fmt.Errorf("backfill mail_reads(cc): 选区失败: %w", err)
|
||
}
|
||
type ccEntry struct {
|
||
Name string `json:"name"`
|
||
}
|
||
type pending struct {
|
||
mailID uuid.UUID
|
||
name string
|
||
}
|
||
var todo []pending
|
||
for rows.Next() {
|
||
var id uuid.UUID
|
||
var ccRaw []byte
|
||
if err := rows.Scan(&id, &ccRaw); err != nil {
|
||
rows.Close()
|
||
return 0, err
|
||
}
|
||
var ccs []ccEntry
|
||
if len(ccRaw) > 0 {
|
||
if err := json.Unmarshal(ccRaw, &ccs); err != nil {
|
||
continue // cc_list 形状不认识就跳过这封,别让整个迁移挂掉
|
||
}
|
||
}
|
||
for _, c := range ccs {
|
||
if c.Name != "" {
|
||
todo = append(todo, pending{id, c.Name})
|
||
}
|
||
}
|
||
}
|
||
rows.Close()
|
||
if err := rows.Err(); err != nil {
|
||
return 0, err
|
||
}
|
||
|
||
inserted := 0
|
||
for _, p := range todo {
|
||
res, err := DB.ExecContext(ctx,
|
||
`INSERT INTO mail_reads (mail_id, reader_name)
|
||
SELECT $1, $2
|
||
WHERE NOT EXISTS (SELECT 1 FROM mail_reads WHERE mail_id = $1 AND reader_name = $2)`,
|
||
p.mailID, p.name)
|
||
if err != nil {
|
||
return inserted, fmt.Errorf("backfill mail_reads(cc): 插行失败: %w", err)
|
||
}
|
||
if k, _ := res.RowsAffected(); k > 0 {
|
||
inserted += int(k)
|
||
}
|
||
}
|
||
return inserted, nil
|
||
}
|
||
|
||
// splitStatements 按分号切分 SQL 脚本并剔除注释行。
|
||
// 本项目的 SQLite schema 只有 CREATE 语句,不含字符串字面量里的分号,
|
||
// 因此按分号朴素切分是安全的;若将来加入含分号的字面量需改用真正的词法切分。
|
||
func splitStatements(script string) []string {
|
||
var out []string
|
||
for _, raw := range strings.Split(script, ";") {
|
||
var lines []string
|
||
for _, line := range strings.Split(raw, "\n") {
|
||
if t := strings.TrimSpace(line); t == "" || strings.HasPrefix(t, "--") {
|
||
continue
|
||
}
|
||
lines = append(lines, line)
|
||
}
|
||
if stmt := strings.TrimSpace(strings.Join(lines, "\n")); stmt != "" {
|
||
out = append(out, stmt)
|
||
}
|
||
}
|
||
return out
|
||
}
|
||
|
||
// sqliteAddColumns 声明 SQLite 侧需要在已存在的表上补齐的列。
|
||
//
|
||
// 新库由 init_sqlite.sql 的 CREATE TABLE 一次建全,这里只服务**已部署的库**。
|
||
// PG 侧用 ALTER TABLE ... ADD COLUMN IF NOT EXISTS 就够,SQLite 没有这个语法,
|
||
// 只能先查 pragma 再决定加不加。
|
||
//
|
||
// 新增列时同时改两处:init_sqlite.sql 的 CREATE TABLE(给新库)与这张表(给老库)。
|
||
var sqliteAddColumns = []struct{ table, column, ddl string }{
|
||
{"mails", "rename_alias", "ALTER TABLE mails ADD COLUMN rename_alias TEXT"},
|
||
{"mails", "rename_reason", "ALTER TABLE mails ADD COLUMN rename_reason TEXT"},
|
||
{"sessions", "rename_dismissed", "ALTER TABLE sessions ADD COLUMN rename_dismissed TEXT"},
|
||
{"sessions", "alias_source", "ALTER TABLE sessions ADD COLUMN alias_source TEXT NOT NULL DEFAULT 'platform'"},
|
||
// 会话级往返预算(0 = 不限)。旧库默认 0:引入预算不应该把已在进行的会话卡死。
|
||
{"sessions", "max_rounds", "ALTER TABLE sessions ADD COLUMN max_rounds INTEGER NOT NULL DEFAULT 0"},
|
||
{"sessions", "used_rounds", "ALTER TABLE sessions ADD COLUMN used_rounds INTEGER NOT NULL DEFAULT 0"},
|
||
// 会话所属的工作目录。旧库默认空串:历史会话的 workspace 无法可靠反推
|
||
// (Agent 回信的 from_workspace 存的是 Agent 名而不是路径),强行回填只会
|
||
// 造出一批看起来有值实际是错的数据。
|
||
{"sessions", "workspace", "ALTER TABLE sessions ADD COLUMN workspace TEXT NOT NULL DEFAULT ''"},
|
||
// 日历多收件人。旧库默认 '[]':读的时候由 EffectiveRecipients() 退回
|
||
// to_address / agent_name,历史事件因此继续工作,不需要数据迁移。
|
||
{"calendar_events", "recipients", "ALTER TABLE calendar_events ADD COLUMN recipients TEXT NOT NULL DEFAULT '[]'"},
|
||
{"calendar_events", "delivery_mode", "ALTER TABLE calendar_events ADD COLUMN delivery_mode TEXT NOT NULL DEFAULT 'separate'"},
|
||
// 日历事件已触发的 occurrence。旧库为 NULL:等价于「从未触发」,
|
||
// 于是已过期的一次性事件会补发一次提醒 —— 这是可接受的,
|
||
// 而反过来(默认成 event_time)会让正在等的提醒永远发不出去。
|
||
{"calendar_events", "fired_for", "ALTER TABLE calendar_events ADD COLUMN fired_for DATETIME"},
|
||
// 日历事件的权限档位(plan / workspace / full)。事件触发时若新建会话,
|
||
// 用这一列定死档位;复用已有会话则取「会话现档 与 事件档」中更严那个。
|
||
//
|
||
// 旧库默认 'workspace':历史事件补发提醒不该静默升到 full(提权路径
|
||
// 会被 P1 calendar 投递接线堵住,但这里默认值也得守住)。
|
||
// 与 sessions.permission_mode 的默认取向一致。
|
||
{"calendar_events", "permission_mode", "ALTER TABLE calendar_events ADD COLUMN permission_mode TEXT NOT NULL DEFAULT 'workspace'"},
|
||
// 本侧会话接管的平台会话 id。旧库默认空串 = 「不是接管来的」,
|
||
// 与新建会话的语义一致,不需要数据迁移。
|
||
{"sessions", "platform_id", "ALTER TABLE sessions ADD COLUMN platform_id TEXT NOT NULL DEFAULT ''"},
|
||
// 派给该 Agent 的新任务默认多少个来回。
|
||
// 旧库也给 20:之前的 max_rounds 默认是 10 但那是终身额度,语义不同,
|
||
// 不能直接搬过来当单任务预算。
|
||
{"agents", "default_rounds", "ALTER TABLE agents ADD COLUMN default_rounds INTEGER NOT NULL DEFAULT 20"},
|
||
// 会话级权限档位(plan / workspace / full)。
|
||
//
|
||
// 旧库默认 'workspace' 而不是 'full':已在进行的会话大多是「在这个目录里干活」,
|
||
// 给 workspace 与它们的实际形态一致。默认 full 则等于给所有历史会话追授全权,
|
||
// 而「我忘了收紧」与「我确实需要全权」在数据上从此无法区分。
|
||
{"sessions", "permission_mode", "ALTER TABLE sessions ADD COLUMN permission_mode TEXT NOT NULL DEFAULT 'workspace'"},
|
||
// 接收平台实际做到的强制力(native / advisory),由插件心跳自报后落到会话上。
|
||
//
|
||
// 旧库默认 'advisory':没自报过的插件,我们不能替它宣称「档位在这里是被强制的」。
|
||
// 保守方向是承认做不到,而不是假装做到了。
|
||
{"sessions", "permission_enforcement", "ALTER TABLE sessions ADD COLUMN permission_enforcement TEXT NOT NULL DEFAULT 'advisory'"},
|
||
// Agent 自报的档位强制能力(native / advisory),随心跳更新。
|
||
// 与 sessions.permission_enforcement 的区别:这里是平台的能力,那里是
|
||
// 某条会话建立时的事实快照 —— 插件升级后能力会变,已结束的会话不该被改写。
|
||
{"agents", "mode_enforcement", "ALTER TABLE agents ADD COLUMN mode_enforcement TEXT NOT NULL DEFAULT 'advisory'"},
|
||
// 待决请求类型(permission/question)与多选语义,供 ask_user_question 桥接使用。
|
||
// 旧库默认 'permission'/0:历史请求是危险工具审批,语义不变。
|
||
{"permission_requests", "kind", "ALTER TABLE permission_requests ADD COLUMN kind TEXT NOT NULL DEFAULT 'permission'"},
|
||
{"permission_requests", "multi_select", "ALTER TABLE permission_requests ADD COLUMN multi_select INTEGER NOT NULL DEFAULT 0"},
|
||
// 邮件上的请求类型与多选标记(与 permission_requests 表一致)。
|
||
// 旧库默认 ''/0:历史权限邮件按单选审批渲染。
|
||
{"mails", "permission_kind", "ALTER TABLE mails ADD COLUMN permission_kind TEXT NOT NULL DEFAULT ''"},
|
||
{"mails", "permission_multi_select", "ALTER TABLE mails ADD COLUMN permission_multi_select INTEGER NOT NULL DEFAULT 0"},
|
||
}
|
||
|
||
// sqliteAddIndexes 是建表后才能建的索引(依赖上面补的列)。
|
||
// CREATE INDEX IF NOT EXISTS 天然幂等,直接执行即可。
|
||
var sqliteAddIndexes = []string{
|
||
// 接管平台会话时按 platform_id 反查(依赖上面补的列)
|
||
"CREATE INDEX IF NOT EXISTS idx_sessions_platform ON sessions(platform_id) WHERE platform_id <> ''",
|
||
// 人类决策后要按 mail_id 反查上游 permission id
|
||
"CREATE INDEX IF NOT EXISTS idx_relayed_mail ON relayed_mails(mail_id)",
|
||
}
|
||
|
||
func addMissingColumns(ctx context.Context) error {
|
||
for _, c := range sqliteAddColumns {
|
||
has, err := columnExists(ctx, c.table, c.column)
|
||
if err != nil {
|
||
return fmt.Errorf("migrate sqlite: 检查 %s.%s: %w", c.table, c.column, err)
|
||
}
|
||
if has {
|
||
continue
|
||
}
|
||
if _, err := DB.ExecContext(ctx, c.ddl); err != nil {
|
||
return fmt.Errorf("migrate sqlite: 补列 %s.%s: %w", c.table, c.column, err)
|
||
}
|
||
fmt.Printf("补列 %s.%s\n", c.table, c.column)
|
||
}
|
||
for _, ddl := range sqliteAddIndexes {
|
||
if _, err := DB.ExecContext(ctx, ddl); err != nil {
|
||
return fmt.Errorf("migrate sqlite: 建索引 %.60s: %w", ddl, err)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func columnExists(ctx context.Context, table, column string) (bool, error) {
|
||
// pragma_table_info 是表函数形式的 PRAGMA,可以直接当表查(比解析 PRAGMA 输出干净)。
|
||
// table 与 column 都来自上面的硬编码常量表,不存在注入面。
|
||
var n int
|
||
err := DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM pragma_table_info(?) WHERE name = ?`,
|
||
table, column).Scan(&n)
|
||
return n > 0, err
|
||
}
|