package repo import ( "context" "database/sql" "errors" "fmt" "strings" "time" "github.com/agentmail/gateway/internal/db" "github.com/agentmail/gateway/internal/models" "github.com/google/uuid" ) // ---------- 附件 ---------- // // 元数据在库、内容在磁盘(internal/blob)。两者的一致性由调用顺序保证: // 先落盘再入库 —— 反过来会出现「库里有记录但文件不存在」的下载 500。 // 落盘成功但入库失败时最多留下一个无引用的文件,由 GC 回收,不影响正确性。 var ( // ErrAttachmentNotFound 附件不存在 ErrAttachmentNotFound = errors.New("attachment not found") // ErrAttachmentNotOwned 附件不属于该上传者 ErrAttachmentNotOwned = errors.New("attachment not owned by uploader") // ErrAttachmentAlreadyAttached 附件已挂到别的邮件上 ErrAttachmentAlreadyAttached = errors.New("attachment already attached") ) const attachmentCols = `attachment_id, mail_id, uploader, filename, content_type, size_bytes, sha256, created_at` func scanAttachment(sc interface{ Scan(...any) error }) (*models.Attachment, error) { var a models.Attachment if err := sc.Scan(&a.ID, &a.MailID, &a.Uploader, &a.Filename, &a.ContentType, &a.SizeBytes, &a.SHA256, &a.CreatedAt); err != nil { return nil, err } return &a, nil } // CreateAttachment 登记一条待挂载的附件(mail_id 为空)。 func CreateAttachment(ctx context.Context, uploader, filename, contentType string, size int64, sum string) (*models.Attachment, error) { a := &models.Attachment{ Uploader: uploader, Filename: filename, ContentType: contentType, SizeBytes: size, SHA256: sum, } err := db.DB.QueryRowContext(ctx, ` INSERT INTO attachments (uploader, filename, content_type, size_bytes, sha256) VALUES ($1, $2, $3, $4, $5) RETURNING attachment_id, created_at `, uploader, filename, contentType, size, sum).Scan(&a.ID, &a.CreatedAt) if err != nil { return nil, err } return a, nil } // GetAttachment 读取一条附件元数据。 func GetAttachment(ctx context.Context, id uuid.UUID) (*models.Attachment, error) { a, err := scanAttachment(db.DB.QueryRowContext(ctx, `SELECT `+attachmentCols+` FROM attachments WHERE attachment_id = $1`, id)) if errors.Is(err, sql.ErrNoRows) { return nil, ErrAttachmentNotFound } return a, err } // ListAttachmentsFor 列出某封邮件的附件。 func ListAttachmentsFor(ctx context.Context, mailID uuid.UUID) ([]models.Attachment, error) { rows, err := db.DB.QueryContext(ctx, `SELECT `+attachmentCols+` FROM attachments WHERE mail_id = $1 ORDER BY created_at`, mailID) if err != nil { return nil, err } defer rows.Close() out := []models.Attachment{} for rows.Next() { a, err := scanAttachment(rows) if err != nil { return nil, err } out = append(out, *a) } return out, rows.Err() } // EnsureAttachable 只做**读取校验**:这批附件是否存在、属于该上传者、且尚未挂载。 // // # 为什么要有一个「只查不改」的版本 // // 原先只有 AttachToMail,而它在 CreateMail **之后**调用。于是附件不合法时 // (不属于我 / 已随别的邮件发出)请求返回 403/409,但那封邮件**已经入库、已经 // 通知了收件人、已经扣掉了会话预算**。实测两封探针邮件(403 与 409)都躺在库里, // used_rounds 也涨了。发件方看到 4xx 会重试,收件方于是收到两封。 // // 纯输入校验必须在产生任何副作用之前做完 —— 与「400 之后会话已建好」是同一个教训。 // // 它不能取代 AttachToMail 里的原子判断:两次调用之间仍有竞态窗口 // (另一个请求把同一个附件挂走了)。那条路径靠调用方回滚,见 handler.attachAll。 func EnsureAttachable(ctx context.Context, ids []uuid.UUID, uploader string) error { for _, id := range ids { a, err := GetAttachment(ctx, id) if err != nil { return err // ErrAttachmentNotFound 或库错误 } if a.Uploader != uploader { return ErrAttachmentNotOwned } if a.MailID != nil { return ErrAttachmentAlreadyAttached } } return nil } // AttachToMail 把一批待挂载附件绑到某封邮件上。 // // 每条都要求:存在、属于该上传者、且尚未挂载。 // 用 WHERE mail_id IS NULL AND uploader = ? 一条 UPDATE 完成判断与写入, // 避免「先查后改」在并发下把同一个附件挂到两封邮件上。 func AttachToMail(ctx context.Context, mailID uuid.UUID, ids []uuid.UUID, uploader string) error { for _, id := range ids { tag, err := db.DB.ExecContext(ctx, ` UPDATE attachments SET mail_id = $1 WHERE attachment_id = $2 AND uploader = $3 AND mail_id IS NULL `, mailID, id, uploader) if err != nil { return err } if n, _ := tag.RowsAffected(); n > 0 { continue } // 没改到:查明原因,给调用方一个能照着修的错误 a, gErr := GetAttachment(ctx, id) if gErr != nil { return gErr } if a.Uploader != uploader { return ErrAttachmentNotOwned } return ErrAttachmentAlreadyAttached } return nil } // CopyAttachmentsTo 把源邮件的附件复制到目标邮件(转发时用)。 // // 内容寻址下「复制」只是新增一条指向同一 sha256 的元数据,不拷磁盘文件。 // uploader 记为转发人:附件随新邮件重新分发,其可见范围由新邮件的参与方决定, // 而不是沿用原上传者。返回复制的数量。 func CopyAttachmentsTo(ctx context.Context, srcMailID, dstMailID uuid.UUID, forwarder string) (int, error) { src, err := ListAttachmentsFor(ctx, srcMailID) if err != nil { return 0, err } for _, a := range src { _, err := db.DB.ExecContext(ctx, ` INSERT INTO attachments (mail_id, uploader, filename, content_type, size_bytes, sha256) VALUES ($1, $2, $3, $4, $5, $6) `, dstMailID, forwarder, a.Filename, a.ContentType, a.SizeBytes, a.SHA256) if err != nil { return 0, err } } return len(src), nil } // DeleteAttachment 删除一条附件元数据,返回它的 sha256 以及该内容是否已无人引用。 // 内容寻址下多条记录可能共享同一个文件,只有最后一条引用消失才能删磁盘文件。 func DeleteAttachment(ctx context.Context, id uuid.UUID) (sum string, orphaned bool, err error) { a, err := GetAttachment(ctx, id) if err != nil { return "", false, err } if _, err = db.DB.ExecContext(ctx, `DELETE FROM attachments WHERE attachment_id = $1`, id); err != nil { return "", false, err } var refs int if err = db.DB.QueryRowContext(ctx, `SELECT COUNT(*) FROM attachments WHERE sha256 = $1`, a.SHA256).Scan(&refs); err != nil { return "", false, err } return a.SHA256, refs == 0, nil } // SweepOrphanAttachments 清理超过 age 仍未挂载到邮件的附件记录, // 返回可以从磁盘删除的 sha256 列表(已确认无任何记录引用)。 // // 上传后没走完发信流程(用户取消、Agent 崩溃)会留下这类记录, // 不清理的话磁盘只会单调增长。 func SweepOrphanAttachments(ctx context.Context, age time.Duration) ([]string, error) { cutoff := time.Now().Add(-age) rows, err := db.DB.QueryContext(ctx, `SELECT attachment_id, sha256 FROM attachments WHERE mail_id IS NULL AND created_at < $1`, cutoff) if err != nil { return nil, err } type orphan struct { id uuid.UUID sum string } var found []orphan for rows.Next() { var o orphan if err := rows.Scan(&o.id, &o.sum); err != nil { rows.Close() return nil, err } found = append(found, o) } rows.Close() if err := rows.Err(); err != nil { return nil, err } var removable []string for _, o := range found { if _, err := db.DB.ExecContext(ctx, `DELETE FROM attachments WHERE attachment_id = $1`, o.id); err != nil { return nil, err } var refs int if err := db.DB.QueryRowContext(ctx, `SELECT COUNT(*) FROM attachments WHERE sha256 = $1`, o.sum).Scan(&refs); err != nil { return nil, err } if refs == 0 { removable = append(removable, o.sum) } } return removable, nil } // SweepUnreferencedBlobs 删掉磁盘上没有任何库记录指向的内容文件。 // // # 为什么 SweepOrphanAttachments 不够 // // 那个函数走的是 `SELECT … FROM attachments WHERE mail_id IS NULL` —— 它只能看见 // **库里还有记录**的孤儿。一旦记录本身消失(清库、手工 DELETE、迁移), // 对应的文件就永远脱离了 GC 的视野:本机实测磁盘 8 个 blob 里 7 个没有任何库记录, // 全部来自 09-03 那次清库,之后一直躺在那里。 // // 这个反向清理从**磁盘**出发:枚举全部内容文件,凡是 attachments 与 // calendar_attachments 都不引用的就删。返回删掉的数量。 // // # 为什么要 minAge // // 上传是「先落盘、再入库」(顺序不能反,否则会出现「库里有记录、磁盘没文件」的 // 下载 500)。那两步之间有一个窗口,此刻文件确实没有任何库记录 —— 不设年龄下限 // 会把正在上传的文件删掉。取一个远大于单次上传耗时的值。 func SweepUnreferencedBlobs(ctx context.Context, blobs BlobLister, minAge time.Duration) (int, error) { if blobs == nil { return 0, nil } sums, err := blobs.List() if err != nil { return 0, err } if len(sums) == 0 { return 0, nil } // 一次查回全部被引用的 sha256。逐个文件查一次库是 N 次往返, // 而这两张表加起来通常只有几百行。 referenced := map[string]struct{}{} for _, q := range []string{ `SELECT sha256 FROM attachments`, `SELECT sha256 FROM calendar_attachments`, // 用户壁纸也是 blob 存储里的内容文件。漏掉这一张表,用户的壁纸会在 // 下一次 GC 时被当成"没人引用的孤儿"删掉 —— 而库里那行还在, // 表现为"图片 404、设置却显示已设置"(2026-09-14 新增功能时特意先查了 GC)。 `SELECT image_sha256 FROM user_appearance WHERE image_sha256 <> ''`, } { rows, qErr := db.DB.QueryContext(ctx, q) if qErr != nil { return 0, qErr } for rows.Next() { var s string if sErr := rows.Scan(&s); sErr != nil { rows.Close() return 0, sErr } referenced[s] = struct{}{} } rows.Close() if rErr := rows.Err(); rErr != nil { return 0, rErr } } cutoff := time.Now().Add(-minAge) removed := 0 for sum, mod := range sums { if _, ok := referenced[sum]; ok { continue } if mod.After(cutoff) { continue // 可能正在上传(落盘与入库之间的窗口) } if rErr := blobs.Remove(sum); rErr != nil { continue // 删不掉就下一轮再试,不该让整次清理中断 } removed++ } return removed, nil } // BlobLister 是 SweepUnreferencedBlobs 需要的存储能力。 // // 用 map[string]time.Time 而不是自定义结构体:那样 blob 包就不必 import repo // (底层存储依赖上层仓储会很怪),而 Go 的接口是结构化匹配的,签名一致即可。 type BlobLister interface { // List 返回 sha256 → 该内容文件的修改时间。 List() (map[string]time.Time, error) Remove(sum string) error } // AttachmentAccessible 判断某人是否有权读取某附件: // 已挂载的看邮件所属会话的参与关系,未挂载的只有上传者本人能看。 func AttachmentAccessible(ctx context.Context, a *models.Attachment, name string) (bool, error) { if a.MailID == nil { return a.Uploader == name, nil } var n int err := db.DB.QueryRowContext(ctx, ` SELECT COUNT(*) FROM mails m WHERE m.mail_id = $1 AND (m.from_name = $2 OR m.to_name = $2 OR `+db.CCHas("m.cc_list", 2)+`) `, *a.MailID, name).Scan(&n) return n > 0, err } // ListAttachmentsForMails 批量取多封邮件的附件,返回 mail_id → 附件列表。 // // 为什么要批量:会话线程与收发件箱都是「一批邮件」,逐封调 ListAttachmentsFor // 就是 N+1 —— 一个 200 封的会话打开一次要打 200 次库。 // 用 IN (...) 一次取回后在内存里分组。 // // 占位符手工拼而非用数组参数:SQLite 驱动不支持 PG 的 = ANY($1), // 而这里的元素是已解析的 uuid.UUID,不存在注入面。 func ListAttachmentsForMails(ctx context.Context, mailIDs []uuid.UUID) (map[uuid.UUID][]models.Attachment, error) { out := map[uuid.UUID][]models.Attachment{} if len(mailIDs) == 0 { return out, nil } ph := make([]string, len(mailIDs)) args := make([]any, len(mailIDs)) for i, id := range mailIDs { ph[i] = fmt.Sprintf("$%d", i+1) args[i] = id } rows, err := db.DB.QueryContext(ctx, `SELECT `+attachmentCols+` FROM attachments WHERE mail_id IN (`+strings.Join(ph, ",")+`) ORDER BY created_at`, args...) if err != nil { return nil, err } defer rows.Close() for rows.Next() { a, err := scanAttachment(rows) if err != nil { return nil, err } if a.MailID == nil { continue // WHERE 已排除,只是防御 } out[*a.MailID] = append(out[*a.MailID], *a) } return out, rows.Err() }