Files
MailUI4Agents/gateway/internal/repo/attachments.go
JianFeeeee 79c4171c9d feat: L0 线协议冻结 + 附件链路修复 + 人/Agent 区分
L0 核心:
- 严格解码 Decode(DisallowUnknownFields) 全覆盖 29 个 DecodeBody 调用点
- DecodeLenient 心跳专用:容忍新字段但回报 unknown_fields
- 400 消息列出本端点接受的全部字段(jsonFieldNames 反射 tag)
- 日历 status 校验(create 补字段 + update 拦非法值)
- 新增 strictdecode_test.go 10 例 + blob/list_test.go 6 例

A-4 附件挂载回滚:checkAttachable 在 CreateMail 前校验,失败按
解挂→释放 relay→删邮件→退预算回滚,幽灵邮件这条路堵住了

A-5 反向 GC:blob.Store.List() 枚举磁盘(跳 .upload-*),
SweepUnreferencedBlobs 按 attachments + calendar_attachments 反查,
48h 年龄下限兜上传窗口。已接进每小时 sweep 循环

C 人/Agent 区分:四个读路径 + threadCols 补 from_human / to_human
(EXISTS users 判定),models.Mail 加 ToHuman。前端判据从
workspace 启发式改成显式布尔,mailCounterpart/sessionCounterpart
从 session_workspace 取 path(修 dsh@dsh 拼接 bug)

契约文档:SSE new_mail 补 4 字段(in_reply_to/from_human/
permission_mode/permission_enforcement),B-5 加 B-5.6
(Agent→Agent 不转发),B-3.4 MUST 改条件式,心跳补 mode_enforcement
+ unknown_fields,demo 死链修复 + from_human 检查
验收清单加 Agent→Agent 负向对照项
2026-09-06 15:18:06 +08:00

382 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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`,
} {
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()
}