pi 报告:离线补投的邮件把 full 档会话当 workspace 档申请审批 → 服务端 409 → 桥按「永久失败」当场 block → 这一轮 bash/write/edit 全被拦(SSE 实时送达不受影响)。 jianf 让 pi 把这件转给我,我这边定位后**先跑变异再改**。 ## 1 根因:`mailToEvent` 少搬字段(不是服务端不给) `lib/catchup.js` 的 `mailToEvent()` 只搬了 8 个字段,没有 `permission_mode` / `permission_enforcement`,于是 worker 的 `msg.data?.permission_mode || 'workspace'` 落到默认档。**pi 以为收件箱行不含档位、于是建议"要么动服务端载荷要么另取一次"—— 实测不成立**:服务端一直就给了(`repo.ListInboxScoped` 的 SQL 里有 `JOIN sessions s` + `COALESCE(NULLIF(s.permission_mode,''),'workspace')`, `models.Mail.PermissionMode` 的注释写明"补拉路径必须有它们")。所以修因只在插件侧: 补上这两个字段,键名与 SSE 逐字一致;缺字段时给空串(**不猜档**,猜宽了就是提权)。 `lib/catchup.js` 在四个桥里**逐字节相同**,一次改动四边同步(改后 md5 仍为一份)。 ## 2 兜底:409 带档位时按档位处置 服务端在"档位不该问人"时也回 409,并在回包里带 `permission_mode`。两种 409 的正确反应 **相反**:无人可问 → 拦;**full 档 → 放行**(本档无需审批,拦了就是把能干的活干死)。 `src/worker.mjs` 的 409 分支先认 `permission_mode === MODE_FULL` 放行, plan 档与"链上没有人类"照旧 fail closed —— 只有服务端明说 full 才放行。 ## 3 顺出的同族字段(用"配对"扫出来的,不是猜的) 把四个桥**读投递事件的字段**与 `mailToEvent` 的产出对了一遍,邮件类字段还缺三个: - `from_human`:dsh 的提示词靠它决定说不说"回信不用你自己发"。缺了它, **人发来的信在补投路径上被当成 Agent 来信、失去自动回信**(服务端注释早写明)。 - `in_reply_to`:SSE 那边等于 `ParentMailID`。缺了它,"这封是对我的回复"被当成新派的活, 两边互相客套到撞 hop 上限(生产实测 6 轮)。行里叫 `parent_mail_id`,**只改名不推算**。 - `session_alias`:缺了它插件只有 session_id,而 `send_mail` 不接受 session_id。 `reply_address` 是**唯一**行里真的没有的字段(SSE 在 notify 里按收件人现算)。 服务端注释明确说"插件不必自己拼(拼错了就是静默开新会话)",所以由服务端补: `models.Mail.ReplyAddress` + `ListInboxScoped` 填 `FormatAddress(from_name,"",alias)`, 插件只搬运。 ## 4 判据(这次事故**单独看任何一个桥的测试都发现不了** —— 缺口在接口上) - `test/catchup.test.mjs`:补投必须带档位(缺字段给空串而非猜档); ★ **四桥配对**:把 dsh/pi/zcode 读的邮件字段与补投产出配对,缺了就红 (非邮件事件字段走显式 ALLOW 并各写理由,白名单不许膨胀)。这条正是本次缺口的形状。 - `test/permission-mode-409.test.mjs`:409 + full 必须放行且放行分支在 block 之前, 非 full 仍拦;带**判据自检**(拿掉放行分支后必须判红)。 - `server/internal/repo/session_scope_test.go`:收件箱行带 `permission_mode`(含"没设过 回落 workspace"的反向对照)与 `reply_address`(与 `FormatAddress` 同形、path 位为空)。 变异验证:mailToEvent 去掉档位 → 2 条红;worker 新读一个补投没产的字段 → 配对判据红**并点名该字段**; 409 分支拿掉 full 放行 → 自检红;SQL 把档位写死成 workspace → Go 判据红。 ## 验证 `go test ./...` 全绿(新增 2 条);四个桥套件全绿(pi 439 / dsh 381 / opencode 331 / zcode 385)。 **未部署**:`/opt/agentmail` 与 `sudo ./deploy/install.sh` 都在我的工作区之外(本会话文件策略 workspace-write,放宽需审批而这条链上没有人类),所以修复已进仓但**线上仍是有缺陷的版本** —— 需要有人跑一次 `sudo ./deploy/install.sh`(脚本自己会跑齐各套件)。
1706 lines
66 KiB
Go
1706 lines
66 KiB
Go
package repo
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/agentmail/gateway/internal/db"
|
||
"github.com/agentmail/gateway/internal/models"
|
||
"github.com/google/uuid"
|
||
)
|
||
|
||
// ---------- Agent ----------
|
||
|
||
// ErrAgentDisabled 表示该 Agent 已被管理员停用。
|
||
//
|
||
// 停用是可逆的「归档」:邮件、会话、权限记录全部保留,只是不再接受新任务。
|
||
// 与删除分开是因为往来邮件里有一半是人自己写的 —— 停用 Agent 不该删掉
|
||
// 用户的东西;而 Agent 名与人类用户名共用命名空间,历史邮件里的 from_name
|
||
// 指向一个已删除的名字时,下一个同名注册者会看起来像是当初的发信人。
|
||
var ErrAgentDisabled = errors.New("agent disabled")
|
||
|
||
func CreateOrUpdateAgent(ctx context.Context, name, secret, platform string, workspaces []models.Workspace) error {
|
||
// 已停用的 Agent 不得靠重新注册复活。
|
||
//
|
||
// 少了这一步的后果:管理员停用后,那个平台的插件下次启动就会重新注册
|
||
// (注册是插件启动流程的一部分),status 被写回 online,停用等于没做。
|
||
// 必须让插件收到一个明确的错误,而不是静默成功。
|
||
disabled, dErr := AgentDisabled(ctx, name)
|
||
if dErr != nil {
|
||
return dErr
|
||
}
|
||
if disabled {
|
||
return ErrAgentDisabled
|
||
}
|
||
|
||
wsJSON, _ := json.Marshal(workspaces)
|
||
// 注意 DO UPDATE 里【不】碰 default_rounds:
|
||
// 那是管理员配的值,Agent 重启重新注册不应该把它冲回默认。
|
||
_, err := db.DB.ExecContext(ctx, `
|
||
INSERT INTO agents (agent_name, secret, workspaces, platform, status, last_seen)
|
||
VALUES ($1, $2, $3, $4, 'online', NOW())
|
||
ON CONFLICT (agent_name) DO UPDATE SET
|
||
secret = EXCLUDED.secret,
|
||
workspaces = EXCLUDED.workspaces,
|
||
platform = EXCLUDED.platform,
|
||
status = 'online',
|
||
last_seen = NOW()
|
||
`, name, secret, wsJSON, platform)
|
||
return err
|
||
}
|
||
|
||
func HeartbeatAgent(ctx context.Context, agentName string) (int, error) {
|
||
// 只把【非停用】的 Agent 标成在线。
|
||
//
|
||
// 不加这个条件的话,停用后插件的心跳会把 status 从 disabled 改回 online
|
||
// —— 而心跳是每 30 秒一次的,停用最多维持半分钟。
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE agents SET last_seen = NOW(), status = 'online'
|
||
WHERE agent_name = $1 AND status <> 'disabled'`,
|
||
agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
return CountUnread(ctx, agentName)
|
||
}
|
||
|
||
// AgentDisabled 该 Agent 是否已被停用。Agent 不存在时返回 false ——
|
||
// 「还没注册」与「被停用」是两件事,前者应当能正常注册。
|
||
func AgentDisabled(ctx context.Context, agentName string) (bool, error) {
|
||
var status string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT status FROM agents WHERE agent_name = $1`, agentName).Scan(&status)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return false, nil
|
||
}
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
return status == "disabled", nil
|
||
}
|
||
|
||
// ErrRecipientUnknown 收件人既不是人类用户也不是在册 Agent。
|
||
var ErrRecipientUnknown = errors.New("recipient unknown")
|
||
|
||
// ErrRecipientDisabled 收件 Agent 已被管理员停用。
|
||
var ErrRecipientDisabled = errors.New("recipient disabled")
|
||
|
||
// RecipientDeliverable 校验一个收件人名当前能不能收信。
|
||
//
|
||
// 为什么必须有这道检查:发信路径原来只校验地址语法、调用权限与会话别名,
|
||
// 从不问「这个名字存在吗」。于是发给已删除或已停用的 Agent 一律返回 200 ——
|
||
// 邮件入库、分配预算、建好会话,而那一端永远不会有人读。发件人看到 200
|
||
// 和一个 session_id,以为送出去了。这是静默丢件,比报错严重:
|
||
// 报错能立刻改,静默丢件要等对方追问才发现。
|
||
//
|
||
// 三种放行/拒绝:
|
||
// - 人类用户 → 放行(人的收件箱一直在)
|
||
// - Agent 在册且未停用 → 放行
|
||
// - Agent 不存在 → ErrRecipientUnknown(对应 404)
|
||
// - Agent 已停用 → ErrRecipientDisabled(对应 409)
|
||
//
|
||
// 停用选择「当场拒收」而不是「入库等恢复后补投」:停用的语义就是这个 Agent
|
||
// 现在不干活,让发件人以为信已送达更坏 —— 它会照常等回信。
|
||
func RecipientDeliverable(ctx context.Context, name string) error {
|
||
if name == "" {
|
||
return nil // 空名由上层的地址解析负责
|
||
}
|
||
|
||
isHuman, err := IsHumanUser(ctx, name)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if isHuman {
|
||
return nil
|
||
}
|
||
|
||
var status string
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT status FROM agents WHERE agent_name = $1`, name).Scan(&status)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return ErrRecipientUnknown
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if status == "disabled" {
|
||
return ErrRecipientDisabled
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// SetAgentDisabled 停用或恢复一个 Agent。
|
||
//
|
||
// 停用时连带撤销它的全部密钥:留着密钥的话,那个平台的插件仍然能用它调
|
||
// /mail/send —— 停用的语义是「这个 Agent 不再参与工作」,不只是「不出现在
|
||
// 补全列表里」。恢复时不会把密钥变回来,管理员需要重新签发。
|
||
//
|
||
// 返回撤销的密钥数,供界面提示。
|
||
func SetAgentDisabled(ctx context.Context, agentName string, disabled bool) (int, error) {
|
||
if !disabled {
|
||
// 恢复:回到 offline 而不是 online —— 它是否真的在线由下一次心跳决定,
|
||
// 直接写 online 会让界面显示一个其实没在跑的 Agent 为在线。
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE agents SET status = 'offline' WHERE agent_name = $1 AND status = 'disabled'`,
|
||
agentName)
|
||
return 0, err
|
||
}
|
||
|
||
tx, err := db.DB.BeginTx(ctx, nil)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
res, err := tx.ExecContext(ctx,
|
||
`UPDATE agents SET status = 'disabled' WHERE agent_name = $1`, agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
if n, _ := res.RowsAffected(); n == 0 {
|
||
return 0, sql.ErrNoRows
|
||
}
|
||
|
||
keyRes, err := tx.ExecContext(ctx,
|
||
`DELETE FROM agent_keys WHERE agent_name = $1`, agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
revoked, _ := keyRes.RowsAffected()
|
||
|
||
if err := tx.Commit(); err != nil {
|
||
return 0, err
|
||
}
|
||
return int(revoked), nil
|
||
}
|
||
|
||
// ListAgents 列出 Agent。
|
||
//
|
||
// statusFilter 为空时默认**排除已停用的** —— 这个函数的三个调用点
|
||
// (地址补全、GET /agents、可授权范围)都是在回答「现在能派活给谁」,
|
||
// 而停用的 Agent 不该出现在那里。要连停用一起看,传 statusFilter="all"。
|
||
func ListAgents(ctx context.Context, statusFilter string) ([]models.Agent, error) {
|
||
// 带上 default_rounds:前端补全收件人时要显示「派给它的任务默认几个来回」,
|
||
// 否则人得先去管理员页查一遍才敢派活。
|
||
//
|
||
// 也要带 mode_enforcement:它是插件自报的**档位强制能力**
|
||
// (native / partial / advisory)。"这个 Agent 到底能不能真的拦住危险操作"
|
||
// 是使用者在派活前必须知道的事。
|
||
//
|
||
// 该字段曾在列表接口**静默丢失**(SELECT 里没有,Scan 也就没扫),
|
||
// 而库里四个 Agent 的值一直是正确的 —— 表现为 /api/v1/agents 全部返回空串。
|
||
// 与列数不匹配不同,这种"少取一列"不会报错,只会安静地少一个事实。
|
||
// 回归测试:repo/agent_mode_list_test.go。
|
||
q := `SELECT agent_id, agent_name, workspaces, platform, status,
|
||
COALESCE(default_rounds, 0),
|
||
COALESCE(NULLIF(mode_enforcement, ''), 'advisory') FROM agents`
|
||
args := []any{}
|
||
switch statusFilter {
|
||
case "":
|
||
q += ` WHERE status <> 'disabled'`
|
||
case "all":
|
||
// 不加条件
|
||
default:
|
||
q += ` WHERE status = $1`
|
||
args = append(args, statusFilter)
|
||
}
|
||
q += ` ORDER BY agent_name`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
agents := []models.Agent{}
|
||
for rows.Next() {
|
||
var a models.Agent
|
||
var wsJSON []byte
|
||
if err := rows.Scan(&a.ID, &a.Name, &wsJSON, &a.Platform, &a.Status,
|
||
&a.DefaultRounds, &a.ModeEnforcement); err != nil {
|
||
return nil, err
|
||
}
|
||
if wsJSON != nil {
|
||
json.Unmarshal(wsJSON, &a.Workspaces)
|
||
}
|
||
agents = append(agents, a)
|
||
}
|
||
return agents, nil
|
||
}
|
||
|
||
func VerifyAgent(ctx context.Context, name, secret string) (*models.Agent, error) {
|
||
var a models.Agent
|
||
var wsJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT agent_id, agent_name, workspaces, platform, status
|
||
FROM agents WHERE agent_name = $1 AND secret = $2`,
|
||
name, secret,
|
||
).Scan(&a.ID, &a.Name, &wsJSON, &a.Platform, &a.Status)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if wsJSON != nil {
|
||
json.Unmarshal(wsJSON, &a.Workspaces)
|
||
}
|
||
return &a, nil
|
||
}
|
||
|
||
// ---------- Session ----------
|
||
|
||
// CreateSession 建会话。
|
||
//
|
||
// 显式传了 alias(发信时的 session_alias 参数)= 调用方亲自命名,标为 manual,
|
||
// 平台后续自动同步不得覆盖;未传则等待平台命名,标为 platform。
|
||
func CreateSession(ctx context.Context, alias *string, fromAgent, subject, workspace string) (uuid.UUID, error) {
|
||
source := "platform"
|
||
if alias != nil && *alias != "" {
|
||
source = "manual"
|
||
}
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO sessions (session_alias, from_agent, subject, alias_source, workspace)
|
||
VALUES ($1, $2, $3, $4, $5) RETURNING session_id`,
|
||
alias, fromAgent, subject, source, strings.TrimSpace(workspace),
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
func GetSessionByID(ctx context.Context, id uuid.UUID) (*models.Session, error) {
|
||
var s models.Session
|
||
var dismissed *string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id, session_alias, from_agent, subject, status, owner_user_id,
|
||
created_at, updated_at, rename_dismissed, COALESCE(alias_source, 'platform'),
|
||
COALESCE(max_rounds, 0), COALESCE(used_rounds, 0),
|
||
COALESCE(NULLIF(permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(permission_enforcement, ''), 'advisory')
|
||
FROM sessions WHERE session_id = $1`, id,
|
||
).Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status, &s.OwnerUserID,
|
||
&s.CreatedAt, &s.UpdatedAt, &dismissed, &s.AliasSource,
|
||
&s.MaxRounds, &s.UsedRounds,
|
||
&s.PermissionMode, &s.PermissionEnforcement)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if dismissed != nil {
|
||
s.RenameDismissed = *dismissed
|
||
}
|
||
return &s, nil
|
||
}
|
||
|
||
func TouchSession(ctx context.Context, id uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET updated_at = NOW(), status = 'active' WHERE session_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
// UpdateSessionAlias 手工改名(人显式指定)。
|
||
//
|
||
// 同时把 alias_source 标为 'manual':人的选择优先于平台自动命名。
|
||
// 否则平台下一次 session.updated 会把人刚定的名字冲掉,
|
||
// 人上一秒记住的寻址地址下一秒失效。
|
||
func UpdateSessionAlias(ctx context.Context, id uuid.UUID, alias string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET session_alias = $1, alias_source = 'manual', updated_at = NOW()
|
||
WHERE session_id = $2`,
|
||
alias, id)
|
||
return err
|
||
}
|
||
|
||
// ---------- Mail ----------
|
||
|
||
func CreateMail(ctx context.Context, sessionID uuid.UUID, parentMailID *uuid.UUID,
|
||
fromName, fromWorkspace, toName, toWorkspace, subject, body string, ccList []models.Address) (uuid.UUID, error) {
|
||
if ccList == nil {
|
||
ccList = []models.Address{}
|
||
}
|
||
ccJSON, _ := json.Marshal(ccList)
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
// created_at 显式给 NOW():SQLite 的 DEFAULT CURRENT_TIMESTAMP 只有秒精度,
|
||
// 同秒插入的多封邮件排序不确定(「会话里最早/最后那封」都会取错行)。
|
||
// 改 schema 的默认值只对新库生效 —— CREATE TABLE IF NOT EXISTS 不改已存在的表,
|
||
// 而 SQLite 没有 ALTER COLUMN,因此这里显式传。
|
||
`INSERT INTO mails (session_id, parent_mail_id, from_name, from_workspace,
|
||
to_name, to_workspace, subject, body, cc_list, created_at)
|
||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, NOW()) RETURNING mail_id`,
|
||
sessionID, parentMailID, fromName, fromWorkspace, toName, toWorkspace, subject, body, ccJSON,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// CreatePermissionMail 创建权限请求邮件,toUser 为目标人类用户名
|
||
// CreatePermissionMail 创建一封权限询问邮件(Agent → 人类决策人)。
|
||
//
|
||
// **from_workspace 必须写入发起方的工作目录。**
|
||
//
|
||
// 不写的后果在授权页上很具体:那一列存空串,而前端拿 `from_workspace`
|
||
// 当「发起方在哪个目录干活」渲染 —— 于是那一行永远不显示,
|
||
// 人只看到一个光秃的 Agent 名,不知道是哪个目录里的哪条线索在请求权限。
|
||
// 同名 Agent 在不同目录是不同的活,那正是决策时最需要的信息。
|
||
//
|
||
// 取会话的 workspace 而不是传参:会话的工作目录在它建立时就定下了,
|
||
// 而询问发起于那条会话里。
|
||
func CreatePermissionMail(ctx context.Context, sessionID uuid.UUID, fromName, toUser, question, body string, options []string, kind string, multiSelect bool) (uuid.UUID, error) {
|
||
optsJSON, _ := json.Marshal(options)
|
||
var multiSelectInt int
|
||
if multiSelect {
|
||
multiSelectInt = 1
|
||
}
|
||
var id uuid.UUID
|
||
// 主题按待办类型区分:「权限请求」意味着有人要放行一个危险操作,
|
||
// 「需要回答」只是模型缺信息。授权页与收件箱都靠这一行措辞判断该做什么。
|
||
subject := "权限请求: " + question
|
||
if kind == "question" {
|
||
subject = "需要回答: " + question
|
||
}
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO mails (session_id, from_name, from_workspace, to_name, subject, body, mail_type, permission_options, permission_kind, permission_multi_select, created_at)
|
||
VALUES ($1, $2, COALESCE((SELECT workspace FROM sessions WHERE session_id = $1), ''), $3, $4, $5, 'permission_request', $6, $7, $8, NOW()) RETURNING mail_id`,
|
||
sessionID, fromName, toUser, subject, body, optsJSON, kind, multiSelectInt,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// CreateDecisionMail 创建人类决策邮件(fromUser → toAgent)
|
||
func CreateDecisionMail(ctx context.Context, sessionID uuid.UUID, parentMailID uuid.UUID, fromUser, toAgent, decision, note string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
body := decision
|
||
if note != "" {
|
||
body = fmt.Sprintf("%s\n\n备注: %s", decision, note)
|
||
}
|
||
// `mail_type` 必须与普通邮件区分开:这封不是"新任务",而是**控制面回执**。
|
||
// 它的内容(决策 + 备注)已经随 SSE 的 permission_decision 直接交给了发起询问的
|
||
// worker,桥若再按"新邮件"起一轮,同一件事就被处理两次 —— 2026-09-13 实测:
|
||
// 人类的更正邮件被挤在队列后面 8 分钟才被看到,而 agent 期间一直在重问同一条命令。
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO mails (session_id, parent_mail_id, from_name, to_name, subject, body, mail_type, created_at)
|
||
VALUES ($1, $2, $3, $4, $5, $6, 'permission_decision', NOW()) RETURNING mail_id`,
|
||
sessionID, parentMailID, fromUser, toAgent, "Re: 权限请求 - "+decision, body,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// DeleteMailByID 删一封邮件。
|
||
//
|
||
// **只用于回滚一次刚失败的发信**,不是给人用的「删邮件」功能 ——
|
||
// 邮件是不可篡改的历史记录,没有任何人面入口能删它。
|
||
//
|
||
// 场景:附件挂载在建邮件之后才发现冲突(竞态窗口),此时这封邮件不应存在:
|
||
// 发件方收到的是 4xx,它会重试,而一封无附件的残余邮件会让收件方收到两封。
|
||
//
|
||
// attachments 表的外键是 ON DELETE CASCADE,所以已经挂上去的那几条会跟着消失;
|
||
// relayed_mails 的 mail_id 无 CASCADE,由调用方用 ReleaseRelay 归还幂等键。
|
||
func DeleteMailByID(ctx context.Context, id uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx, `DELETE FROM mails WHERE mail_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
func GetMailByID(ctx context.Context, id uuid.UUID) (*models.Mail, error) {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
var renameAlias, renameReason *string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace, m.rename_alias, m.rename_reason,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.mail_id = $1`, id,
|
||
).Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &renameAlias, &renameReason,
|
||
&m.FromHuman, &m.ToHuman)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
// 改名提议随单封返回,让「谁在哪一封里提了什么」可追溯;
|
||
// 【待处理】的提议另有专用端点(GET /sessions/:id/rename-proposal)。
|
||
if renameAlias != nil {
|
||
m.RenameAlias = *renameAlias
|
||
}
|
||
if renameReason != nil {
|
||
m.RenameReason = *renameReason
|
||
}
|
||
AttachPermissionDeadline(&m)
|
||
return &m, nil
|
||
}
|
||
|
||
/*
|
||
─── 已读语义:按**读者**记录,而不是邮件行上的一个全局列 ───────────────
|
||
|
||
2026-09-13 实测缺陷:`mails.status` 是邮件级的,任何收件人读掉,对所有收件人
|
||
(含抄送)都变成已读。后果:Agent 的 read_inbox(默认 unread)拿不到信(dsh 明确
|
||
回报"收件箱列表未展示它,直接按 mail_id 读取成功");人类的未读被抄送的 Agent
|
||
读掉;桥的补投判据 CountUnread 归零 ⇒ 那封信不再补投(静默丢信)。
|
||
|
||
现在的判据:**对某个读者未读 = mail_reads 里没有 (mail_id, reader_name) 这一行。**
|
||
`mails.status` 保留为"有人读过 / 已归档"的冗余,不再作为未读判据。
|
||
|
||
下列两个片段把这件事收在一处,避免每条 SQL 各写一遍判据(写岔了就是又一次语义漂移)。
|
||
*/
|
||
|
||
// unreadFor 返回"$n 这个读者看这封邮件是未读"的谓词;`m` 必须是 mails 的别名。
|
||
func unreadFor(arg string) string {
|
||
return `(m.status <> 'archived' AND NOT EXISTS (
|
||
SELECT 1 FROM mail_reads r WHERE r.mail_id = m.mail_id AND r.reader_name = ` + arg + `))`
|
||
}
|
||
|
||
// readStateFor 返回给前端的 status 值(archived 是全局的,read/unread 按读者算)。
|
||
func readStateFor(arg string) string {
|
||
return `CASE WHEN m.status = 'archived' THEN 'archived'
|
||
WHEN EXISTS (SELECT 1 FROM mail_reads r WHERE r.mail_id = m.mail_id AND r.reader_name = ` + arg + `)
|
||
THEN 'read' ELSE 'unread' END`
|
||
}
|
||
|
||
// markReadFor 批量记下"某个读者读过哪些邮件"。
|
||
//
|
||
// `where` 是**只用于 INSERT 的**筛选片段,里面可以(也应当)用 `m.` 前缀引用 mails ——
|
||
// 调用方各自负责随后刷新 mails.status 那列冗余(那个语句没有 `m` 别名)。
|
||
// 之前我把同一个 where 复用到 UPDATE 上,直接 SQL 报 "no such column: m.mail_id"
|
||
// (测试当场抓到)。
|
||
func markReadFor(ctx context.Context, reader string, where string, args ...any) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`INSERT INTO mail_reads (mail_id, reader_name)
|
||
SELECT m.mail_id, $1 FROM mails m
|
||
WHERE `+where+`
|
||
AND NOT EXISTS (SELECT 1 FROM mail_reads r WHERE r.mail_id = m.mail_id AND r.reader_name = $1)`,
|
||
append([]any{reader}, args...)...)
|
||
return err
|
||
}
|
||
|
||
// MarkMailRead 记下**这个读者**读过这封邮件(人类端点)。
|
||
func MarkMailRead(ctx context.Context, id uuid.UUID, reader string) error {
|
||
if _, err := db.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)`,
|
||
id, reader); err != nil {
|
||
return err
|
||
}
|
||
_, err := db.DB.ExecContext(ctx, `UPDATE mails SET status = 'read' WHERE mail_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
/*
|
||
─── 会话维度(2026-09-14)───────────────────────────────────────────────
|
||
|
||
用户报的缺陷:「不同 session 的 agent 都可以看到全部邮件」。
|
||
|
||
原先 `read_inbox` 是**按 Agent** 的:列的是该 Agent 的全部未读(含别的会话的来信),
|
||
并且按契约把列出来的都标成已读 ⇒ A 会话的 worker 会把 B 会话的未读标掉。
|
||
平时看不出来(SSE 事件在途时队列兜着),但桥重启/漏事件后的补投判据是
|
||
`?status=unread` —— 被别人标掉的那封**再也不会补投** ⇒ 静默丢信。
|
||
|
||
修法:列表与"全部标已读"都支持按 `session_id` 收窄,桥把自己的会话传进来。
|
||
原函数保持原语义(不带会话 = 整个 Agent 的收件箱),新增带会话的变体 ——
|
||
老调用点一个都不用改。
|
||
*/
|
||
func ListInbox(ctx context.Context, agentName, status string, limit int) ([]models.Mail, error) {
|
||
return ListInboxScoped(ctx, agentName, status, limit, uuid.Nil)
|
||
}
|
||
|
||
// ListInboxScoped 与 ListInbox 相同,但 `sessionID` 非零时只列该会话的邮件。
|
||
func ListInboxScoped(ctx context.Context, agentName, status string, limit int, sessionID uuid.UUID) ([]models.Mail, error) {
|
||
q := `SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
` + readStateFor("$1") + ` AS status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human,
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace') AS permission_mode,
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory') AS permission_enforcement
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE (m.to_name = $1 OR ` + db.CCHas("m.cc_list", 1) + `)
|
||
AND s.status <> 'archived'`
|
||
args := []any{agentName}
|
||
if sessionID != uuid.Nil {
|
||
// 会话收窄:只列这条线索里的邮件(见上面「会话维度」的说明)
|
||
args = append(args, sessionID)
|
||
q += fmt.Sprintf(` AND m.session_id = $%d`, len(args))
|
||
}
|
||
|
||
if status != "" && status != "all" {
|
||
// 未读/已读都按**这个读者**算(原先直接比 m.status,于是被抄送方读掉
|
||
// 别人的未读也跟着变 —— 这就是要修的那条)
|
||
if status == "unread" {
|
||
q += ` AND ` + unreadFor("$1")
|
||
} else if status == "read" {
|
||
q += ` AND NOT ` + unreadFor("$1") + ` AND m.status <> 'archived'`
|
||
} else {
|
||
q += ` AND m.status = $2`
|
||
args = append(args, status)
|
||
}
|
||
}
|
||
q += ` ORDER BY m.created_at DESC, m.mail_id DESC`
|
||
if limit > 0 {
|
||
q += fmt.Sprintf(` LIMIT %d`, limit)
|
||
}
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman,
|
||
&m.PermissionMode, &m.PermissionEnforcement); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
// 回信地址与 SSE 载荷同形(见 models.Mail.ReplyAddress 的说明):
|
||
// 收件方回信 = 寄回本条会话里的**发件人**,path 位留空(与 notify 一致)。
|
||
m.ReplyAddress = models.FormatAddress(m.FromName, "", m.SessionAlias)
|
||
// Body preview
|
||
if len(m.Body) > 200 {
|
||
m.BodyPreview = m.Body[:200] + "..."
|
||
} else {
|
||
m.BodyPreview = m.Body
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
AttachPermissionDeadlines(mails)
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
func CountUnread(ctx context.Context, agentName string) (int, error) {
|
||
var count int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*)
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE (m.to_name = $1 OR `+db.CCHas("m.cc_list", 1)+`)
|
||
AND `+unreadFor("$1")+`
|
||
AND s.status <> 'archived'
|
||
`, agentName).Scan(&count)
|
||
return count, err
|
||
}
|
||
|
||
func GetSessionMails(ctx context.Context, sessionID uuid.UUID) ([]models.Mail, error) {
|
||
rows, err := db.DB.QueryContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.session_id = $1
|
||
ORDER BY m.created_at ASC, m.mail_id ASC`, sessionID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
AttachPermissionDeadlines(mails)
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
// ---------- Permission ----------
|
||
|
||
func CreatePermissionRequest(ctx context.Context, mailID, sessionID uuid.UUID, agentName, question string, options []string, contextStr string, kind string, multiSelect bool) error {
|
||
optsJSON, _ := json.Marshal(options)
|
||
var multiSelectInt int
|
||
if multiSelect {
|
||
multiSelectInt = 1
|
||
}
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`INSERT INTO permission_requests (mail_id, session_id, agent_name, question, options, context, kind, multi_select)
|
||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`,
|
||
mailID, sessionID, agentName, question, optsJSON, contextStr, kind, multiSelectInt)
|
||
return err
|
||
}
|
||
|
||
// AttachPermissionDeadline 给仍未决策的权限请求邮件补上失效时刻。
|
||
//
|
||
// 为什么是「推导」而不是查询时算完落库:
|
||
// - 它是 CreatedAt 的纯函数,存下来就会随时间失真(存的是派生值,不是事实);
|
||
// - 只有**仍未决策**的 permission_request 才有意义 —— 已决策的不再是待办,
|
||
// 给它一个「失效时刻」只会让界面把历史记录也标成过期。
|
||
//
|
||
// 为什么只在邮件上给「时刻」而不给「是否失效」:
|
||
//
|
||
// 布尔值是发送那一刻的快照,经 SSE 缓存到客户端后会永久停在旧值。
|
||
// 时刻是持久事实,任何客户端在任何时候都能自己比出现在过没过期。
|
||
func AttachPermissionDeadline(m *models.Mail) {
|
||
if m == nil || m.MailType != "permission_request" || m.PermResult != "" {
|
||
return
|
||
}
|
||
d := models.PermissionDeadline(m.CreatedAt)
|
||
m.PermissionExpiresAt = &d
|
||
}
|
||
|
||
// AttachPermissionDeadlines 是切片版本,供列表读路径一次处理。
|
||
func AttachPermissionDeadlines(mails []models.Mail) {
|
||
for i := range mails {
|
||
AttachPermissionDeadline(&mails[i])
|
||
}
|
||
}
|
||
|
||
// decider 是**做出决策的人**(用户名):权限邮件对他也应当变成已读 —— 但只对他一个人,
|
||
// 而不是像原先那样把邮件行的全局 status 一改(抄送的其他 Agent 会因此看不到它)。
|
||
func DecidePermission(ctx context.Context, mailID uuid.UUID, decider, decision string) (*models.PermissionRequest, error) {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`UPDATE permission_requests SET result = $1, decided_at = NOW()
|
||
WHERE mail_id = $2
|
||
RETURNING request_id, mail_id, session_id, agent_name, question, options, context, result, decided_at, created_at`,
|
||
decision, mailID,
|
||
).Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Result, &pr.DecidedAt, &pr.CreatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
|
||
// Also update the mail
|
||
_, _ = db.DB.ExecContext(context.Background(),
|
||
`UPDATE mails SET permission_result = $1, status = 'read' WHERE mail_id = $2`,
|
||
decision, mailID)
|
||
if decider != "" {
|
||
// 记到决策人名下(按读者记,见 repo.markReadFor)
|
||
_, _ = db.DB.ExecContext(context.Background(),
|
||
`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)`,
|
||
mailID, decider)
|
||
}
|
||
|
||
return &pr, nil
|
||
}
|
||
|
||
func ListPendingPermissions(ctx context.Context) ([]models.PermissionRequest, error) {
|
||
rows, err := db.DB.QueryContext(ctx,
|
||
`SELECT request_id, mail_id, session_id, agent_name, question, options, context, kind, multi_select, result, decided_at, created_at
|
||
FROM permission_requests WHERE result IS NULL
|
||
ORDER BY created_at DESC`)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
reqs := []models.PermissionRequest{}
|
||
for rows.Next() {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
if err := rows.Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
pr.ExpiresAt = models.PermissionDeadline(pr.CreatedAt)
|
||
reqs = append(reqs, pr)
|
||
}
|
||
return reqs, nil
|
||
}
|
||
|
||
// ---------- Check permission ownership ----------
|
||
|
||
func GetPermissionByMailID(ctx context.Context, mailID uuid.UUID) (*models.PermissionRequest, error) {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT request_id, mail_id, session_id, agent_name, question, options, context, kind, multi_select, result, decided_at, created_at
|
||
FROM permission_requests WHERE mail_id = $1`, mailID,
|
||
).Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
return &pr, nil
|
||
}
|
||
|
||
func GetSessionMailByID(ctx context.Context, sessionID, mailID uuid.UUID) (*models.Mail, error) {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.session_id = $1 AND m.mail_id = $2`, sessionID, mailID,
|
||
).Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
AttachPermissionDeadline(&m)
|
||
return &m, nil
|
||
}
|
||
|
||
func FindSessionByAlias(ctx context.Context, alias string) (*models.Session, error) {
|
||
var s models.Session
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id, session_alias, from_agent, subject, status, owner_user_id, created_at, updated_at
|
||
FROM sessions WHERE session_alias = $1 AND status <> 'archived'`, alias,
|
||
).Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status, &s.OwnerUserID, &s.CreatedAt, &s.UpdatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &s, nil
|
||
}
|
||
|
||
// ErrSessionNotFound 表示三维地址里指定的 session 别名不存在(或不属于该收件人)。
|
||
// 调用方应据此回 404「无法送达」,而不是悄悄新建一个会话。
|
||
var ErrSessionNotFound = errors.New("session not found")
|
||
|
||
// FindNamedSessionFor 查找收件人 name@path 名下别名为 alias 的会话。
|
||
// 严格匹配:会话必须存在、未归档,且该收件人确实参与过该会话,否则返回 ErrSessionNotFound。
|
||
// FindNamedSessionFor 实现 session 位给具体别名时的语义:必须已存在。
|
||
//
|
||
// 不限定 workspace:别名全局唯一且本身就承担寻址职责,
|
||
// 再叠一层工作区校验只会让「名字对上了却送不到」变成一种难查的失败。
|
||
func FindNamedSessionFor(ctx context.Context, name, path, alias string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT s.session_id
|
||
FROM sessions s
|
||
WHERE s.session_alias = $1
|
||
AND s.status <> 'archived'
|
||
AND EXISTS (
|
||
SELECT 1 FROM mails m
|
||
WHERE m.session_id = s.session_id
|
||
AND (m.to_name = $2 OR m.from_name = $2 OR `+db.CCHas("m.cc_list", 2)+`)
|
||
)
|
||
ORDER BY s.updated_at DESC
|
||
LIMIT 1
|
||
`, alias, name).Scan(&id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return uuid.Nil, ErrSessionNotFound
|
||
}
|
||
return id, err
|
||
}
|
||
|
||
// FindOrCreateDefaultSession 实现 session 位省略时的「默认会话」语义:
|
||
// 复用 name@path 最近一次活跃的非归档会话;从未通过信则建立一个新的作为默认会话。
|
||
//
|
||
// 匹配工作区优先看 sessions.workspace(权威来源),旧会话那列为空时回退到
|
||
// mails.to_workspace 反推。只看 to_workspace:Agent 回信时 from_workspace 存的是
|
||
// Agent 名而不是路径,拿它比路径永远匹配不上(旧实现就挂在这里)。
|
||
func FindOrCreateDefaultSession(ctx context.Context, name, path, fromAgent, subject string) (uuid.UUID, error) {
|
||
id, _, err := FindOrCreateDefaultSessionCreated(ctx, name, path, fromAgent, subject)
|
||
return id, err
|
||
}
|
||
|
||
// FindOrCreateDefaultSessionCreated 与 FindOrCreateDefaultSession 相同,但额外返回
|
||
// **这次调用是否真的新建了会话**。
|
||
//
|
||
// 为什么需要这个返回值:调用方此前用 `parentMailID == nil` 判断「是不是新建会话」,
|
||
// 而复用已有默认会话时 parentMailID 也是 nil —— 于是「仅在新建时生效」的字段
|
||
// (往返预算、权限档位)在每一封省略 session 位的信上都被重写了。
|
||
// 实测:第一封 max_rounds=7 → 第二封省略该字段 → 预算被静默改成默认的 20。
|
||
func FindOrCreateDefaultSessionCreated(ctx context.Context, name, path, fromAgent, subject string) (uuid.UUID, bool, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT s.session_id
|
||
FROM sessions s
|
||
WHERE s.status <> 'archived'
|
||
-- 接管会话是人显式指定的线索,不该被省略 session 位的邮件当「默认会话」吃掉。
|
||
-- 日历提醒省略 session 位后落进了人选的那条会话 → 那条会话的别名被另开的
|
||
-- 会话的命名同步冲掉(项目定位 → 日程提醒:…)的链条,起点就在这里。
|
||
AND (s.platform_id IS NULL OR s.platform_id = '')
|
||
AND EXISTS (
|
||
SELECT 1 FROM mails m
|
||
WHERE m.session_id = s.session_id
|
||
AND (m.to_name = $1 OR m.from_name = $1 OR `+db.CCHas("m.cc_list", 1)+`)
|
||
)
|
||
AND (s.workspace = $2
|
||
OR (s.workspace = '' AND EXISTS (
|
||
SELECT 1 FROM mails w
|
||
WHERE w.session_id = s.session_id
|
||
AND COALESCE(w.to_workspace,'') = $2
|
||
)))
|
||
ORDER BY s.updated_at DESC
|
||
LIMIT 1
|
||
`, name, path).Scan(&id)
|
||
if err == nil {
|
||
TouchSession(ctx, id)
|
||
return id, false, nil
|
||
}
|
||
if !errors.Is(err, sql.ErrNoRows) {
|
||
return uuid.Nil, false, err
|
||
}
|
||
newID, cErr := CreateSession(ctx, nil, fromAgent, subject, path)
|
||
return newID, cErr == nil, cErr
|
||
}
|
||
|
||
// SessionAliasOf 返回会话别名,未命名或查询失败时返回空串。
|
||
// 仅用于响应体回显,不影响投递路径,所以吞错是可接受的。
|
||
func SessionAliasOf(ctx context.Context, id uuid.UUID) string {
|
||
var alias *string
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_alias FROM sessions WHERE session_id = $1`, id).Scan(&alias); err != nil {
|
||
return ""
|
||
}
|
||
if alias == nil {
|
||
return ""
|
||
}
|
||
return *alias
|
||
}
|
||
|
||
// AgentCanAccessSession 判断 Agent 是否参与过该会话(发件/收件/被抄送)。
|
||
// Agent 只能改自己参与的会话的别名,避免跨会话改名。
|
||
func AgentCanAccessSession(ctx context.Context, agentName string, sessionID uuid.UUID) (bool, error) {
|
||
var n int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*) FROM mails m
|
||
WHERE m.session_id = $1
|
||
AND (m.from_name = $2 OR m.to_name = $2 OR `+db.CCHas("m.cc_list", 2)+`)
|
||
`, sessionID, agentName).Scan(&n)
|
||
return n > 0, err
|
||
}
|
||
|
||
// SyncSessionTitle 更新会话主题(Agent 平台生成的摘要标题)。
|
||
func SyncSessionTitle(ctx context.Context, id uuid.UUID, title string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET subject = $1, updated_at = NOW() WHERE session_id = $2`,
|
||
title, id)
|
||
return err
|
||
}
|
||
|
||
// SyncSessionAlias 把 Agent 平台侧的会话标识写为本侧别名。
|
||
// 平台侧标识(如 opencode 的 slug)在平台内不保证全局唯一,而本侧别名负责寻址必须唯一,
|
||
// 因此撞名时自动追加 -2、-3… 后缀而不是报错——同步是后台行为,不该因撞名失败。
|
||
// 返回最终落库的别名。该会话已持有目标别名时直接返回,不做无谓写入。
|
||
//
|
||
// **人显式定过的别名不覆盖**(alias_source = 'manual'):
|
||
// 用户刚接受了 Agent 的改名提议,或手工敲了一个名字,平台下一次 session.updated
|
||
// 不该把它冲掉 —— 那会让人上一秒记住的寻址地址下一秒失效。
|
||
// 此时返回当前别名,调用方据此知道同步未生效。
|
||
func SyncSessionAlias(ctx context.Context, id uuid.UUID, want string) (string, error) {
|
||
const maxAttempts = 50
|
||
|
||
var cur *string
|
||
var source string
|
||
var platformID string
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_alias, COALESCE(alias_source, 'platform'), COALESCE(platform_id, '')
|
||
FROM sessions WHERE session_id = $1`,
|
||
id).Scan(&cur, &source, &platformID); err != nil {
|
||
return "", err
|
||
}
|
||
if source == "manual" && cur != nil && *cur != "" {
|
||
return *cur, nil
|
||
}
|
||
// 接管会话的别名是人从补全里选中的平台 slug,任何平台命名同步
|
||
// 都不该动它。不守这道门的话,事件日历经由桥另开会话 → 命名同步
|
||
// 回写 → 把接管会话的别名冲掉 → 人在补全里选的名字凭空消失。
|
||
// 生产上已经兑现过一次(项目定位 → 日程提醒:…)。
|
||
if platformID != "" {
|
||
if cur != nil && *cur != "" {
|
||
return *cur, nil
|
||
}
|
||
}
|
||
|
||
for i := 0; i < maxAttempts; i++ {
|
||
candidate := want
|
||
if i > 0 {
|
||
candidate = fmt.Sprintf("%s-%d", want, i+1)
|
||
}
|
||
|
||
owner, err := aliasOwner(ctx, candidate)
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
if owner != nil {
|
||
if *owner == id {
|
||
return candidate, nil // 已经是这个别名,无需写入
|
||
}
|
||
continue // 被别人占用,试下一个后缀
|
||
}
|
||
|
||
// 只在仍是 platform 来源时写入:并发下用户可能刚好接受了改名提议,
|
||
// 条件放进 WHERE 才能保证「检查」与「写入」不被插进来的手工改名割开
|
||
res, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET session_alias = $1, updated_at = NOW()
|
||
WHERE session_id = $2 AND COALESCE(alias_source, 'platform') <> 'manual'`,
|
||
candidate, id)
|
||
if err == nil {
|
||
if n, _ := res.RowsAffected(); n == 0 {
|
||
// 期间变成 manual 了,尊重人的选择
|
||
return SessionAliasOf(ctx, id), nil
|
||
}
|
||
return candidate, nil
|
||
}
|
||
// 并发下另一个请求刚占走该别名(唯一索引拦下),继续试下一个后缀
|
||
if db.IsUniqueViolation(err) {
|
||
continue
|
||
}
|
||
return "", err
|
||
}
|
||
return "", fmt.Errorf("alias %q: 连同 -2..-%d 后缀均被占用", want, maxAttempts)
|
||
}
|
||
|
||
// aliasOwner 返回持有该别名的会话 ID;无人持有时返回 nil。
|
||
func aliasOwner(ctx context.Context, alias string) (*uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id FROM sessions WHERE session_alias = $1`, alias).Scan(&id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return nil, nil
|
||
}
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &id, nil
|
||
}
|
||
|
||
func SessionMailCount(ctx context.Context, sessionID uuid.UUID) (int, error) {
|
||
var count int
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM mails WHERE session_id = $1`, sessionID).Scan(&count)
|
||
return count, err
|
||
}
|
||
|
||
// ---------- Contacts / Archive ----------
|
||
|
||
// Contact 是「联系人」= 一条 name@path.session 三维地址
|
||
type Contact struct {
|
||
SessionID uuid.UUID `json:"session_id"`
|
||
AgentName string `json:"agent_name"`
|
||
Path string `json:"path"`
|
||
SessionAlias string `json:"session_alias"`
|
||
Address string `json:"address"` // name@path.session
|
||
Status string `json:"status"`
|
||
MailCount int `json:"mail_count"`
|
||
UnreadCount int `json:"unread_count"`
|
||
LastActivity time.Time `json:"last_activity"`
|
||
|
||
// Subject 是会话主题(多由 Agent 平台的模型生成的摘要)。
|
||
// 卡片视图要靠它回答「这条线索在干什么」—— 光有 name@path.alias
|
||
// 只能看出跟谁在聊,看不出在聊什么。
|
||
Subject string `json:"subject"`
|
||
|
||
// MaxRounds/UsedRounds 是本任务的往返预算(0 = 不限)。
|
||
// 列表上直接可见,才不用点进每条会话去查哪件事快跑满了。
|
||
MaxRounds int `json:"max_rounds"`
|
||
UsedRounds int `json:"used_rounds"`
|
||
|
||
// PermissionMode 与 PermissionEnforcement 必须成对出现在列表上:
|
||
// 前者是「这条任务要求什么」,后者是「对方平台实际做到了什么」。
|
||
// 只显示前者会让人以为 plan 档把 homeagent 管住了(它没有拦截点)。
|
||
PermissionMode string `json:"permission_mode"`
|
||
PermissionEnforcement string `json:"permission_enforcement"`
|
||
|
||
// LastFrom/LastPreview 是最后一封邮件的发件人与正文摘要,
|
||
// 卡片视图用它显示「最新进展」——列表视图只显示地址时,
|
||
// 人必须逐条点开才知道哪条有新动静。
|
||
LastFrom string `json:"last_from"`
|
||
LastPreview string `json:"last_preview"`
|
||
}
|
||
|
||
// ListContactsFor 按 (agent, path, session) 聚合出联系人清单。
|
||
// forUser 非空时只列该用户参与的会话(owner / 收发 / 抄送);空表示不限(管理员全局视图)。
|
||
// archived=false 只列活跃会话,true 只列归档会话。
|
||
func ListContactsFor(ctx context.Context, forUser string, archived bool) ([]Contact, error) {
|
||
op := "<>"
|
||
if archived {
|
||
op = "="
|
||
}
|
||
scope := ""
|
||
args := []any{}
|
||
if forUser != "" {
|
||
scope = ` AND (s.owner_user_id = (SELECT user_id FROM users WHERE username = $1)
|
||
OR EXISTS (
|
||
SELECT 1 FROM mails mm
|
||
WHERE mm.session_id = s.session_id
|
||
AND (mm.from_name = $1 OR mm.to_name = $1
|
||
OR ` + db.CCHas("mm.cc_list", 1) + `)
|
||
))`
|
||
args = append(args, forUser)
|
||
}
|
||
// 取会话里最早那封邮件作为联系人身份。
|
||
// PG 用 LATERAL 子查询;SQLite 无 LATERAL,改用关联子查询逐列取值
|
||
// (同一个 min(created_at) 子句,四列取自同一行)。
|
||
var firstMail string
|
||
if db.D == db.Postgres {
|
||
firstMail = `
|
||
JOIN LATERAL (
|
||
SELECT to_name, to_workspace, from_name, from_workspace
|
||
FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at ASC
|
||
LIMIT 1
|
||
) m ON TRUE`
|
||
} else {
|
||
firstMail = `
|
||
JOIN mails m ON m.mail_id = (
|
||
SELECT mail_id FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at ASC, mail_id ASC
|
||
LIMIT 1
|
||
)`
|
||
}
|
||
|
||
// 最后一封邮件用关联子查询取,不再 JOIN 一次:
|
||
// 两个 JOIN(最早一封 + 最新一封)在 SQLite 下要写两段方言分支,
|
||
// 而这里每个会话只多两次索引查找(idx_mails_session 已有)。
|
||
lastMail := `(
|
||
SELECT mail_id FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at DESC, mail_id DESC
|
||
LIMIT 1
|
||
)`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT s.session_id,
|
||
COALESCE(NULLIF(m.to_name, 'human'), m.from_name) AS agent_name,
|
||
COALESCE(NULLIF(m.to_workspace, ''), m.from_workspace) AS path,
|
||
COALESCE(s.session_alias, '') AS alias,
|
||
s.status,
|
||
(SELECT COUNT(*) FROM mails x WHERE x.session_id = s.session_id),
|
||
(SELECT COUNT(*) FROM mails x WHERE x.session_id = s.session_id
|
||
AND x.status <> 'archived'
|
||
AND NOT EXISTS (SELECT 1 FROM mail_reads r WHERE r.mail_id = x.mail_id AND r.reader_name = $1)),
|
||
s.updated_at,
|
||
s.subject,
|
||
COALESCE(s.max_rounds, 0),
|
||
COALESCE(s.used_rounds, 0),
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory'),
|
||
COALESCE((SELECT from_name FROM mails WHERE mail_id = `+lastMail+`), ''),
|
||
COALESCE((SELECT body FROM mails WHERE mail_id = `+lastMail+`), '')
|
||
FROM sessions s`+firstMail+`
|
||
WHERE s.status `+op+` 'archived'`+scope+`
|
||
ORDER BY s.updated_at DESC
|
||
`, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
contacts := []Contact{}
|
||
for rows.Next() {
|
||
var c Contact
|
||
if err := rows.Scan(&c.SessionID, &c.AgentName, &c.Path, &c.SessionAlias,
|
||
&c.Status, &c.MailCount, &c.UnreadCount, &c.LastActivity,
|
||
&c.Subject, &c.MaxRounds, &c.UsedRounds,
|
||
&c.PermissionMode, &c.PermissionEnforcement,
|
||
&c.LastFrom, &c.LastPreview); err != nil {
|
||
return nil, err
|
||
}
|
||
c.Address = c.AgentName + "@" + c.Path
|
||
if c.SessionAlias != "" {
|
||
c.Address += "." + c.SessionAlias
|
||
}
|
||
// 正文只留一段摘要:卡片上放不下全文,而整个列表带全文可能几百 KB。
|
||
// 按 rune 截断而非字节 —— 中文 3 字节/字,裸切会产生 U+FFFD。
|
||
c.LastPreview = previewRunes(c.LastPreview, 90)
|
||
contacts = append(contacts, c)
|
||
}
|
||
return contacts, rows.Err()
|
||
}
|
||
|
||
// previewRunes 按字符数截断,附省略号。
|
||
// 不按字节切:中文一字三字节,裸切会在末尾留半个字符(渲染成 U+FFFD)。
|
||
func previewRunes(s string, n int) string {
|
||
s = strings.TrimSpace(s)
|
||
rs := []rune(s)
|
||
if len(rs) <= n {
|
||
return s
|
||
}
|
||
return string(rs[:n]) + "…"
|
||
}
|
||
|
||
// ArchiveSession 归档一个会话(邮箱界面不再展示,数据保留)
|
||
func ArchiveSession(ctx context.Context, sessionID uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET status = 'archived', updated_at = NOW() WHERE session_id = $1`,
|
||
sessionID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 同时把该会话下的邮件标记为已归档,收件箱不再列出
|
||
_, err = db.DB.ExecContext(ctx,
|
||
`UPDATE mails SET status = 'archived' WHERE session_id = $1 AND status <> 'archived'`,
|
||
sessionID)
|
||
return err
|
||
}
|
||
|
||
// FindSessionByAddress 按 name@path.session 定位会话
|
||
func FindSessionByAddress(ctx context.Context, name, path, alias string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT s.session_id
|
||
FROM sessions s
|
||
JOIN mails m ON m.session_id = s.session_id
|
||
WHERE COALESCE(s.session_alias, '') = $1
|
||
AND (m.to_name = $2 OR m.from_name = $2)
|
||
AND (COALESCE(m.to_workspace,'') = $3 OR COALESCE(m.from_workspace,'') = $3)
|
||
LIMIT 1
|
||
`, alias, name, path).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// SuggestPaths 给出某个收件方可用的工作目录候选(三段式补全的 path 位)。
|
||
//
|
||
// 三个来源并集,**历史优先**:
|
||
//
|
||
// 1. mails.to_workspace 里真实投递过的目录(按最近使用倒序)
|
||
// 2. agent_platform_sessions.workspace 里心跳上报的目录(镜像)
|
||
// 3. agents.workspaces 里注册时自报的目录
|
||
//
|
||
// 早先只看第 1、3 项,而第 1 项在清库后为空,第 3 项对两个官方插件
|
||
// **永远是空的** —— 契约(W-2)明确要求 `workspaces: []`,
|
||
// 因为工作目录由每封邮件的 to_workspace 决定而不是注册时固定。
|
||
// 于是补全的第二段在清库后对所有 Agent 都给空列表。
|
||
//
|
||
// 加入第 2 项:镜像是心跳上报的平台会话快照,其中 workspace 字段
|
||
// 就是各平台的真实工作目录。它在清库后仍然存在(心跳重新上报)。
|
||
//
|
||
// 按最近使用倒序而非字典序:同一个 Agent 常在几个项目间切,刚用过的那个
|
||
// 几乎总是下一封想用的那个。
|
||
func SuggestPaths(ctx context.Context, agentName string) ([]string, error) {
|
||
out := []string{}
|
||
seen := map[string]bool{}
|
||
|
||
// 来源 1:真实投递历史。即使 Agent 未注册(人类收件方)也能给出候选。
|
||
//
|
||
// 只 SELECT 路径一列,排序用的 MAX(created_at) 不进结果集:
|
||
// SQLite 把时间戳存为 TEXT,把它 Scan 进 time.Time 会失败,
|
||
// 而那个值除了排序之外无用 —— 取回来只是多一个能静默失败的环节。
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT to_workspace
|
||
FROM mails
|
||
WHERE to_name = $1 AND COALESCE(to_workspace, '') <> ''
|
||
GROUP BY to_workspace
|
||
ORDER BY MAX(created_at) DESC
|
||
`, agentName)
|
||
if err == nil {
|
||
defer rows.Close()
|
||
for rows.Next() {
|
||
var p string
|
||
if rows.Scan(&p) != nil {
|
||
continue
|
||
}
|
||
if !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
// 来源 2:平台会话镜像。心跳上报的平台会话快照,workspace 就是
|
||
// 各平台的真实工作目录。清库后仍然存在(心跳重新上报)。
|
||
// 没有这一步,清库后路径补全对所有 Agent 都空。
|
||
mrows, err := db.DB.QueryContext(ctx, `
|
||
SELECT workspace, COUNT(*) AS n
|
||
FROM agent_platform_sessions
|
||
WHERE agent_name = $1 AND COALESCE(workspace, '') <> ''
|
||
GROUP BY workspace
|
||
ORDER BY n DESC
|
||
`, agentName)
|
||
if err == nil {
|
||
defer mrows.Close()
|
||
for mrows.Next() {
|
||
var p string
|
||
var n int
|
||
if mrows.Scan(&p, &n) != nil {
|
||
continue
|
||
}
|
||
if !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
// 来源 3:注册时自报。报了但还没收过信的目录靠这一步进入候选,
|
||
// 否则新接入的 Agent 在第一封邮件之前仍然无路径可选。
|
||
var wsJSON []byte
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`SELECT workspaces FROM agents WHERE agent_name = $1`, agentName).Scan(&wsJSON); err == nil {
|
||
var ws []models.Workspace
|
||
if len(wsJSON) > 0 {
|
||
json.Unmarshal(wsJSON, &ws)
|
||
}
|
||
for _, w := range ws {
|
||
// 三维地址的 path 位是**路径**,不是工作区的展示名。
|
||
// 这里取 Path 而不是 Name:remotebot 报的是
|
||
// {name:"demo", path:"/tmp/remotebot-ws"},取 Name 会给出 "demo"
|
||
// —— 一个拉起会话时不存在的目录。
|
||
p := w.Path
|
||
if p == "" {
|
||
p = w.Name
|
||
}
|
||
if p != "" && !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
return out, nil
|
||
}
|
||
|
||
// ListSentBy 列出某发件人发出的邮件(发件箱),排除已归档会话
|
||
func ListSentBy(ctx context.Context, fromName string, limit int) ([]models.Mail, error) {
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.from_name = $1 AND s.status <> 'archived'
|
||
ORDER BY m.created_at DESC, m.mail_id DESC
|
||
LIMIT $2
|
||
`, fromName, limit)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
if len(m.Body) > 200 {
|
||
m.BodyPreview = m.Body[:200] + "..."
|
||
} else {
|
||
m.BodyPreview = m.Body
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
AttachPermissionDeadlines(mails)
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
// ListPendingPermissionsFor 列出待决权限请求;forUser 非空时只列发给该用户的
|
||
func ListPendingPermissionsFor(ctx context.Context, forUser string) ([]models.PermissionRequest, error) {
|
||
q := `SELECT pr.request_id, pr.mail_id, pr.session_id, pr.agent_name, pr.question,
|
||
pr.options, pr.context, COALESCE(pr.kind, '') AS kind, pr.multi_select, pr.result, pr.decided_at, pr.created_at
|
||
FROM permission_requests pr
|
||
JOIN mails m ON m.mail_id = pr.mail_id
|
||
JOIN sessions s ON s.session_id = pr.session_id
|
||
WHERE pr.result IS NULL AND s.status <> 'archived'`
|
||
args := []any{}
|
||
if forUser != "" {
|
||
q += ` AND m.to_name = $1`
|
||
args = append(args, forUser)
|
||
}
|
||
q += ` ORDER BY pr.created_at DESC`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
reqs := []models.PermissionRequest{}
|
||
for rows.Next() {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
if err := rows.Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
pr.ExpiresAt = models.PermissionDeadline(pr.CreatedAt)
|
||
reqs = append(reqs, pr)
|
||
}
|
||
return reqs, nil
|
||
}
|
||
|
||
// ListSessionsFor 列出某人类用户参与的会话(owner / 收发 / 抄送);forUser 空表示不限
|
||
func ListSessionsFor(ctx context.Context, forUser string, limit int) ([]models.Session, error) {
|
||
q := `SELECT s.session_id, s.session_alias, s.from_agent, s.subject, s.status,
|
||
s.owner_user_id, s.created_at, s.updated_at,
|
||
COALESCE(s.max_rounds, 0), COALESCE(s.used_rounds, 0),
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory'),
|
||
(SELECT COUNT(*) FROM mails m WHERE m.session_id = s.session_id)
|
||
FROM sessions s
|
||
WHERE s.status <> 'archived'`
|
||
args := []any{}
|
||
if forUser != "" {
|
||
q += ` AND (s.owner_user_id = (SELECT user_id FROM users WHERE username = $1)
|
||
OR EXISTS (
|
||
SELECT 1 FROM mails mm
|
||
WHERE mm.session_id = s.session_id
|
||
AND (mm.from_name = $1 OR mm.to_name = $1
|
||
OR ` + db.CCHas("mm.cc_list", 1) + `)
|
||
))`
|
||
args = append(args, forUser)
|
||
}
|
||
q += ` ORDER BY s.updated_at DESC`
|
||
if limit > 0 {
|
||
q += fmt.Sprintf(` LIMIT %d`, limit)
|
||
}
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
sessions := []models.Session{}
|
||
for rows.Next() {
|
||
var s models.Session
|
||
// 列数必须与上面的 SELECT 一一对应 —— 预算两列曾经只加进了查询而没加进这里,
|
||
// 结果 /me/sessions 整个 500,联系人栅拉不到任何数据。
|
||
if err := rows.Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status,
|
||
&s.OwnerUserID, &s.CreatedAt, &s.UpdatedAt,
|
||
&s.MaxRounds, &s.UsedRounds,
|
||
&s.PermissionMode, &s.PermissionEnforcement, &s.MailCount); err != nil {
|
||
return nil, err
|
||
}
|
||
sessions = append(sessions, s)
|
||
}
|
||
return sessions, rows.Err()
|
||
}
|
||
|
||
// CountUnreadInSession 统计某人在某会话内的未读数(含被抄送)
|
||
func CountUnreadInSession(ctx context.Context, name string, sessionID uuid.UUID) (int, error) {
|
||
var n int
|
||
// 抄送判定必须走 db.CCHas:这条路原先写的是 PG 专有的 `cc_list @> $3::jsonb`,
|
||
// 而**线上是 SQLite** —— 那条 SQL 直接语法错误(unrecognized token: "@"),
|
||
// 调用点又是 `unread, _ :=`(吞错),于是会话列表的未读数一直显示 0。
|
||
// 用的是与 CountUnread 同一个助手,两种方言都正确。
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*) FROM mails m
|
||
WHERE m.session_id = $1
|
||
AND `+unreadFor("$2")+`
|
||
AND (m.to_name = $2 OR `+db.CCHas("m.cc_list", 2)+`)
|
||
`, sessionID, name).Scan(&n)
|
||
return n, err
|
||
}
|
||
|
||
// SetMailRenameProposal 记录某封邮件里 Agent 提议的新会话别名。
|
||
//
|
||
// 单独一条 UPDATE 而不是塞进 CreateMail 的参数表:提议是可选的旁支信息,
|
||
// 让三个调用点都多传两个几乎总是空串的参数不值当。
|
||
func SetMailRenameProposal(ctx context.Context, mailID uuid.UUID, alias, reason string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE mails SET rename_alias = $1, rename_reason = $2 WHERE mail_id = $3`,
|
||
alias, reason, mailID)
|
||
return err
|
||
}
|
||
|
||
// PendingRenameProposal 返回某会话里**最新一条尚未处理**的改名提议。
|
||
//
|
||
// 「尚未处理」= 提议的别名既不是当前别名(已接受),也不在驳回记录里。
|
||
// 无提议时返回 ("", "", nil)。
|
||
//
|
||
// 只看最新一条:Agent 干活过程中可能多次提议,最后那条才是它现在的结论。
|
||
func PendingRenameProposal(ctx context.Context, sessionID uuid.UUID) (alias, reason string, err error) {
|
||
var cur, dismissed *string
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT session_alias, rename_dismissed FROM sessions WHERE session_id = $1`,
|
||
sessionID).Scan(&cur, &dismissed)
|
||
if err != nil {
|
||
return "", "", err
|
||
}
|
||
|
||
var a, rs *string
|
||
err = db.DB.QueryRowContext(ctx, `
|
||
SELECT rename_alias, rename_reason FROM mails
|
||
WHERE session_id = $1 AND rename_alias IS NOT NULL AND rename_alias <> ''
|
||
ORDER BY created_at DESC, mail_id DESC
|
||
LIMIT 1
|
||
`, sessionID).Scan(&a, &rs)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return "", "", nil
|
||
}
|
||
if err != nil {
|
||
return "", "", err
|
||
}
|
||
if a == nil || *a == "" {
|
||
return "", "", nil
|
||
}
|
||
// 已经改成这个名字了 = 提议已被接受,不必再提示
|
||
if cur != nil && *cur == *a {
|
||
return "", "", nil
|
||
}
|
||
// 用户驳回过这个建议
|
||
if dismissed != nil && *dismissed == *a {
|
||
return "", "", nil
|
||
}
|
||
if rs != nil {
|
||
reason = *rs
|
||
}
|
||
return *a, reason, nil
|
||
}
|
||
|
||
// DismissRenameProposal 记下用户驳回了哪个建议,好让提示条不再反复弹。
|
||
//
|
||
// 只存最后驳回的那一个而不是一张列表:Agent 每次提的名字都不同,
|
||
// 攒一张历史表除了占地方没有别的用处 —— 需要判断的只是「当前这条提议是否被否过」。
|
||
func DismissRenameProposal(ctx context.Context, sessionID uuid.UUID, alias string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET rename_dismissed = $1 WHERE session_id = $2`,
|
||
alias, sessionID)
|
||
return err
|
||
}
|
||
|
||
// MarkMailsReadFor 把一批邮件标记为某收件人已读,返回实际影响的行数。
|
||
//
|
||
// **鉴权写进 WHERE 而不是先查后改**:`to_name = $1 OR cc 含 $1` 直接放在
|
||
// UPDATE 条件里,于是「不是发给我的邮件」根本改不动 —— 既省掉一次查询,
|
||
// 也没有「查完到改之间邮件被转走」的时间窗。
|
||
//
|
||
// 已经是 read 的不计入影响行数(`status = 'unread'` 条件),
|
||
// 调用方据此知道这次真正标掉了几封。
|
||
func MarkMailsReadFor(ctx context.Context, recipient string, ids []uuid.UUID) (int, error) {
|
||
if len(ids) == 0 {
|
||
return 0, nil
|
||
}
|
||
|
||
// IN 子句的占位符按方言编号(SQLite 与 PG 都认 $N)。
|
||
// 不用一条条 UPDATE:一次网络往返 + 一次事务,SQLite 单写者下差别明显。
|
||
ph := make([]string, len(ids))
|
||
args := make([]any, 0, len(ids)+1)
|
||
args = append(args, recipient)
|
||
for i, id := range ids {
|
||
ph[i] = fmt.Sprintf("$%d", i+2)
|
||
args = append(args, id)
|
||
}
|
||
|
||
if err := markReadFor(ctx, recipient,
|
||
`m.mail_id IN (`+strings.Join(ph, ",")+`) AND (m.to_name = $1 OR `+db.CCHas("m.cc_list", 1)+`)`,
|
||
args...); err != nil {
|
||
return 0, err
|
||
}
|
||
// 冗余列:mails.status 只表示"有人读过 / 已归档",**不再作为未读判据**
|
||
// (判据是 mail_reads,见 unreadFor)。保留它是为了兼容仍在读这一列的老路径,
|
||
// 以及让 SQL 层面"邮件是否被任何人读过"仍可一眼看出。
|
||
res, err := db.DB.ExecContext(ctx, `
|
||
UPDATE mails SET status = 'read'
|
||
WHERE mail_id IN (`+strings.Join(ph, ",")+`)
|
||
AND status = 'unread'
|
||
AND (to_name = $1 OR `+db.CCHas("cc_list", 1)+`)
|
||
`, args...)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
n, _ := res.RowsAffected()
|
||
return int(n), nil
|
||
}
|
||
|
||
// MarkAllInboxReadFor 把某收件人收件箱里全部未读标为已读,返回影响行数。
|
||
//
|
||
// 排除已归档会话:那些邮件在收件箱里根本看不到,
|
||
// 标掉它们只会让「标记了 N 封」这个数字与用户看到的对不上。
|
||
func MarkAllInboxReadFor(ctx context.Context, recipient string) (int, error) {
|
||
return MarkAllInboxReadForSession(ctx, recipient, uuid.Nil)
|
||
}
|
||
|
||
// MarkAllInboxReadForSession 只标掉某条会话里发给 recipient 的未读。
|
||
//
|
||
// 为什么需要:Agent 的「不给 mail_ids,全部标掉」在会话驱动的 worker 里会跨会话
|
||
// 误伤(见 ListInboxScoped 上面那段说明)。不带 sessionID(uuid.Nil)时是旧语义。
|
||
func MarkAllInboxReadForSession(ctx context.Context, recipient string, sessionID uuid.UUID) (int, error) {
|
||
scope := `(m.to_name = $1 OR ` + db.CCHas("m.cc_list", 1) + `)
|
||
AND m.session_id IN (SELECT session_id FROM sessions WHERE status <> 'archived')`
|
||
args := []any{recipient}
|
||
if sessionID != uuid.Nil {
|
||
scope += ` AND m.session_id = $2`
|
||
args = append(args, sessionID)
|
||
}
|
||
if err := markReadFor(ctx, recipient, scope, args...); err != nil {
|
||
return 0, err
|
||
}
|
||
// 同上:刷新冗余列,未读判据在 mail_reads。
|
||
// ★ 这条 UPDATE 也必须跟着同一个 scope —— 只给上面的 INSERT 收窄、漏掉它,
|
||
// 返回的计数与"实际标掉多少"都会跨会话(测试当场抓到:应当 1 封、实际 2 封)。
|
||
upd := `UPDATE mails SET status = 'read'
|
||
WHERE status = 'unread'
|
||
AND (to_name = $1 OR ` + db.CCHas("cc_list", 1) + `)
|
||
AND session_id IN (SELECT session_id FROM sessions WHERE status <> 'archived')`
|
||
updArgs := []any{recipient}
|
||
if sessionID != uuid.Nil {
|
||
upd += ` AND session_id = $2`
|
||
updArgs = append(updArgs, sessionID)
|
||
}
|
||
res, err := db.DB.ExecContext(ctx, upd, updArgs...)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
n, _ := res.RowsAffected()
|
||
return int(n), nil
|
||
}
|
||
|
||
// ─── Agent 删除 ───
|
||
|
||
// DeleteAgent 删除 Agent 的全部运行态,保留邮件历史,名字进保留名单。
|
||
//
|
||
// 删除范围:
|
||
// - agent_keys(全部撤销)
|
||
// - agent_platform_sessions(清掉镜像)
|
||
// - models_scope(模型范围配置)
|
||
// - rate_limits(速率限制计数器)
|
||
// - calendar_events:cancelled(不触发、不静默空转)
|
||
// - agents 行本身
|
||
//
|
||
// 保留范围:
|
||
// - mails(历史是审计凭据,不能删)
|
||
// - sessions(与 mails 一起组成线索)
|
||
//
|
||
// 取消邮件:Agent 名在 mails 里的字段(from_name / to_name)不改 ——
|
||
// 那是历史记录的固有属性。未来要"查这封信是谁发的"仍然能查到。
|
||
func DeleteAgent(ctx context.Context, name string) (int, error) {
|
||
tx, err := db.DB.BeginTx(ctx, nil)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
// 撤销全部密钥
|
||
r1, _ := tx.ExecContext(ctx, `DELETE FROM agent_keys WHERE agent_name = $1`, name)
|
||
keysRevoked, _ := r1.RowsAffected()
|
||
|
||
// 清镜像
|
||
tx.ExecContext(ctx, `DELETE FROM agent_platform_sessions WHERE agent_name = $1`, name)
|
||
|
||
// 清模型范围
|
||
tx.ExecContext(ctx, `DELETE FROM agent_allowed_models WHERE agent_name = $1`, name)
|
||
tx.ExecContext(ctx, `DELETE FROM agent_model_catalog WHERE agent_name = $1`, name)
|
||
|
||
// 清速率限制
|
||
tx.ExecContext(ctx, `DELETE FROM rate_limits WHERE agent_name = $1`, name)
|
||
|
||
// 日历事件置 cancelled:Agent 创建的提醒不该继续触发并投给一个不存在的收件人
|
||
tx.ExecContext(ctx,
|
||
`UPDATE calendar_events SET status = 'cancelled', updated_at = $2
|
||
WHERE created_by = $1 AND status = 'active'`, name, time.Now())
|
||
|
||
// 删 agents 行
|
||
if _, err := tx.ExecContext(ctx, `DELETE FROM agents WHERE agent_name = $1`, name); err != nil {
|
||
return 0, err
|
||
}
|
||
|
||
return int(keysRevoked), tx.Commit()
|
||
}
|
||
|
||
// IsRetiredAgentName 检查一个名字是否已被删除(用于注册时拒绝同名重建)。
|
||
func IsRetiredAgentName(ctx context.Context, name string) (bool, error) {
|
||
// 如果 agents 表里没有这个名字,且没有任何邮件引用它,就当「已退役」。
|
||
// 更严格的做法是建一张 retired_agents 表,但现有数据量下这个查询够了。
|
||
var cnt int
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM agents WHERE agent_name = $1`, name).Scan(&cnt)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
if cnt > 0 {
|
||
return false, nil // 还活着
|
||
}
|
||
// 检查有没有历史邮件用这个名字
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM mails WHERE from_name = $1 OR to_name = $1`, name).Scan(&cnt)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
return cnt > 0, nil
|
||
}
|