用户 12 天前就提过(`552fbc7` 只修了 session_id 那一维),这轮才真修。 用户原话:「难道让一个不在项目工作区的 agentsession 去修工程吗?」 # 缺陷(生产实测,2026-09-26) 在 `mc` 工作区干活的 pi 读收件箱拿到 **200 封,其中 191 封属于 `/home/program/agentmail`** —— 它照着那些信里的断言去改 agentmail 的代码, 把手上的 mc 活丢在一边。用户当场问它「你怎么干着干着修 agentmail 去了?」 (这条对话就在 mc 会话的 jsonl 里) 根因:`ListInboxScoped` 的 WHERE 只有 `m.to_name = $1`(+ 可选 session_id), **没有任何 workspace 条件**。三维地址 `name@path.session` 的 path 位 在收件箱侧从未生效 —— 那不是"另一种语义",是没兑现契约。 # 三条守卫全部只覆盖自动转发,防不住这个 | 守卫 | 只覆盖 | 为何无效 | | --- | --- | --- | | 会话预算 | `relay != ""` 才扣 | 这批信 relay=0(模型主动发)⇒ 不扣 | | maxRelayHops=5 | 同上,只数 relay | 同上 ⇒ 不进那个分支 | | 插件自动转发守卫 | 插件代劳时 | 日志明说"本轮不自动转发" ⇒ 模型自己发的不受管 | # 服务端 · `ListInboxScoped` / `CountUnreadScoped` / `MarkAllInboxReadForSession` 三处统一加 `s.workspace = $N`(用会话的 workspace,不用 mails.to_workspace: 后者是信封字段、可能是抄送或历史遗留;"线索属于哪个工作区"是会话属性)。 ★ 三处必须是**同一个谓词** —— 列表看不到的信却被"全部标掉"标掉就是静默丢信 (session_scope_test.go 记过这个形状)。 · **workspace 在 Agent 侧必需,缺了 400**(用户裁定:「不带 workspace 是错误 发件格式,直接退回!」)。旧语义(不带=全部)正是缺陷本身,不留兼容回退。 · 人类侧**不过滤**(一个人跨工作区,WebUI 按 session_workspace 分组显示)—— 所以"必需"这条约束放在 Handler 而不是 repo 层:它是接口契约,不是数据层不变量。 · 新增 `UnreadWorkspaces`:心跳是**进程级**(一个桥服务所有工作区),没有 "我的工作区"可言;但只有总数桥不知道去哪个工作区补投 ⇒ 心跳回 `pending_workspaces` 清单,桥逐个消费。 · 决策载荷补 `workspace`(服务端知道 session→workspace,插件重启后推不出来)。 · `TouchAgentLastSeen` 从 HeartbeatAgent 拆出:middleware 在每个认证请求上都调它, 而那时工作区还没解析(请求体没读),原来在白算一次 CountUnread。 # 三个插件(pi / opencode / dsh) · 读类工具带 `workspace`;补投从"读一次全局收件箱"改为**逐工作区**读。 · pi:worker 信封的 `to_workspace` 经闭包递进工具(不是会话文件 header 的 cwd —— 后者是"会话上次落在哪",前者是"这封信寄到哪个工作区")。 · opencode/dsh:插件常驻、信封在 deliverMail 那刻就消费掉了 ⇒ 新增 `sessionWorkspace` 映射(键与既有 reverseMap 同一把)。 · 修一处真 bug:`UnreadWorkspaces` 原先会返回相对路径工作区(历史库里有 `workspace='root'`),桥侧实测撞 400(`补投工作区 root 失败`)⇒ 只报可寻址的。 # 实测凭据 · 改前:`pi` 的收件箱 200 封混 3 个工作区(agentmail 191 / TrueAgent 7 / huawei 2) · 改后:agentmail=100(total 228)、mc=16、TrueAgent=7 —— 各工作区独立 · 不带 workspace ⇒ **HTTP 400**,话术给出可执行步骤 · 桥日志:`rw=/home/newqqagent/plugindev/editdoc-upgrade` —— 终于是别的工作区了 (改前 78 次 worker 启动**全部**是 `/home/program/agentmail`) # 判据 · `server/internal/repo/workspace_scope_test.go`(3 条): 两向收窄 + **反向对照**(不带时两条都看得到 ⇒ 证明是收窄不是清空)+ 未读数同口径 + 相对路径必须报错 · `plugins/pi-mail-bridge/test/inbox-workspace-scope.test.mjs`(4 条):接线 + 取信封而非 cwd + 补投逐工作区 + 判据自检 · dsh 那条 `取不到会话时退回整体收件箱` **改了**:它钉的"退回整体"正是缺陷, 现在钉"两维各自缺席时各自不带、服务端 400 让错误可见" · 变异验证:服务端 2 处 + 插件 3 处,全部判红后恢复回绿 全量:server `go test ./...` 绿;三插件 513+340+403 全绿。
2104 lines
85 KiB
Go
2104 lines
85 KiB
Go
package repo
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"path/filepath"
|
||
"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
|
||
}
|
||
|
||
/*
|
||
UnreadWorkspaces 列出"这个 Agent 在哪些工作区有未读"。
|
||
|
||
# 为什么要它,而不是让心跳收一个 workspace 参数
|
||
|
||
心跳是**进程级**的(一个桥进程同时服务所有工作区),而收件箱是**worker 级**的
|
||
(每个 worker 手上只有一封信,信封上有明确的 path 位)。在进程级强制要求
|
||
workspace 是概念错配 —— 它没有一个"我的工作区"可言。
|
||
|
||
但 `pending_mails` 是桥的补投判据:若它是一个跨工作区的总数,桥就不知道该去
|
||
**哪个工作区**补投。⇒ 心跳返回这个清单,桥逐个工作区去读(见 catchUp)。
|
||
|
||
这同时修掉一个隐蔽问题:补投原先调 `/mail/inbox?status=unread`(不带收窄),
|
||
按当时的语义会列出**所有工作区**的未读并逐封重放 —— 在 mc 干活时会去补投
|
||
agentmail 的信。
|
||
*/
|
||
func UnreadWorkspaces(ctx context.Context, agentName string) ([]string, error) {
|
||
if err := requireReader(agentName); err != nil {
|
||
return nil, err
|
||
}
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT s.workspace, COUNT(*) AS n
|
||
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'
|
||
-- ★ 只报**可寻址**的工作区(绝对路径)。
|
||
--
|
||
-- 历史库里存在 workspace 为相对路径的行(早期以 "pi@root" 寻址留下的
|
||
-- 测试会话)。收件箱接口要求绝对路径,把它们放进清单只会在桥侧撞 400
|
||
-- —— 我实测就撞到了:"补投工作区 root 失败: HTTP 400"。
|
||
--
|
||
-- 过滤放在这里而不是让桥去试错:这个清单的语义是"**能去补投**的工作区",
|
||
-- 列出不可寻址的等于给调用方递一个注定失败的任务。
|
||
AND s.workspace LIKE '/%'
|
||
GROUP BY s.workspace
|
||
ORDER BY n DESC`, agentName)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
out := []string{}
|
||
for rows.Next() {
|
||
var ws string
|
||
var n int
|
||
if err := rows.Scan(&ws, &n); err != nil {
|
||
return nil, err
|
||
}
|
||
out = append(out, ws)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
/*
|
||
TouchAgentLastSeen 只刷新"我还活着",**不**算未读数。
|
||
|
||
# 为什么要把它拆出来
|
||
|
||
原先 middleware 在**每一个**已认证请求上都调 HeartbeatAgent(它会顺手算
|
||
`CountUnread` 并丢掉返回值)—— 那是白算一次全表扫描,而且现在 `CountUnread`
|
||
还需要工作区,而 middleware 那一层拿不到(请求体还没解析)。
|
||
|
||
⇒ 拆成两件事:心跳副作用(只更新 last_seen)留在 middleware;
|
||
"这个工作区还有多少未读"由心跳 handler 按请求体里的 workspace 算。
|
||
*/
|
||
func TouchAgentLastSeen(ctx context.Context, agentName string) 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)
|
||
return err
|
||
}
|
||
|
||
// HeartbeatAgent 刷新在线状态并返回**该 Agent 的未读总数**。
|
||
//
|
||
// ★ 这里是**全局**口径(跨工作区),与 ListInboxScoped 不同 —— 原因见
|
||
//
|
||
// UnreadWorkspaces 上面那段:心跳是进程级,收件箱是 worker 级。
|
||
// 桥拿到这个总数后,用 `pending_workspaces` 清单逐工作区去补投,
|
||
// 两边合起来才是"有没有信、在哪个工作区"。
|
||
func HeartbeatAgent(ctx context.Context, agentName string) (int, error) {
|
||
if err := TouchAgentLastSeen(ctx, agentName); 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 permOptsJSON []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,
|
||
COALESCE(m.permission_options,'[]') AS permission_options,
|
||
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, &permOptsJSON,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &renameAlias, &renameReason,
|
||
&m.FromHuman, &m.ToHuman)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
/*
|
||
* ★★ 2026-09-23 补:`permission_options` 从 INSERT 起就写进 `mails` 表,
|
||
* 但**从来没有任何读路径选过它** ⇒ 详情端点永远返回空数组。
|
||
*
|
||
* WebUI 的 `MailView.tsx` 决策面板读 `mail.permission_options ?? []`,
|
||
* 于是提问型的预设选项在两端**全部落空**(审批型靠 `['同意','拒绝']`
|
||
* 兜底蒙混过去,提问型是直接没有)。
|
||
*
|
||
* 修法与 `cc_list` 同款:JSON 列 → 字节 → `json.Unmarshal` 进 `[]string`,
|
||
* 空则保底 `[]`(与前端 `?? []` 同义,但**在服务端**完成,
|
||
* 免得每个读方各写一遍兜底)。
|
||
*/
|
||
if len(permOptsJSON) > 0 && string(permOptsJSON) != "[]" {
|
||
json.Unmarshal(permOptsJSON, &m.PermOptions)
|
||
}
|
||
if m.PermOptions == nil {
|
||
m.PermOptions = []string{}
|
||
}
|
||
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 各写一遍判据(写岔了就是又一次语义漂移)。
|
||
*/
|
||
|
||
/*
|
||
* requireReader:**读侧**的同一条规则(pi 2026-09-14 裁定 §2)。
|
||
*
|
||
* 写侧加了"空 reader 必须报错"之后,读侧仍然是洞 —— 而且更隐蔽:`reader` 在查询里是
|
||
* **过滤条件**,空串不会写坏数据、也不会报错,只会**算出一个错误的数**:
|
||
* `unreadFor('')` 的 `NOT EXISTS(... reader_name = '')` 恒真 ⇒ 于是
|
||
* `CountUnread(ctx, "")` 把**所有**邮件都算成未读(用户看到的是"全都没读"),
|
||
* 而 `ListInbox(ctx, "", "read")` 恒空。没有异常、没有坏数据,只有一个错数字。
|
||
*
|
||
* 二选一(pi 要求显式定,不许落在"没人知道"):① 报错;② 明确定义"空 reader = 汇总裁剪语义"并钉住。
|
||
* 这里选 **①报错** —— 因为当前没有任何调用方需要"汇总"语义(HTTP 层传的都是登录用户名),
|
||
* 而②会立刻需要一条判据去定义"汇总"到底是什么意思(那是一个还没有需求的功能)。
|
||
* 将来真需要汇总,就新增一个**名字里带汇总**的函数,而不是让空串偷偷兼职。
|
||
*/
|
||
func requireReader(reader string) error {
|
||
if strings.TrimSpace(reader) == "" {
|
||
return fmt.Errorf("reader 不能为空:未读/已读是**按读者**算的,空读者会静默算出一个错误的数")
|
||
}
|
||
return nil
|
||
}
|
||
|
||
/*
|
||
checkWorkspace 校验 workspace 的形状(**允许为空**)。
|
||
|
||
# 空与非空各是什么语义(两条不同的入口,别混)
|
||
|
||
· **非空** = 只列该工作区的信。**Agent 侧必须非空**(Handler 层强制)。
|
||
· **空** = 不过滤,列该名字的全部。**人类侧就是这个**:一个人跨工作区,
|
||
WebUI 把结果按 `session_workspace` **分组显示**(`MailList.tsx:167`),
|
||
而不是只给一个工作区 —— 强行让人也带工作区,等于让人在多个工作区之间反复切。
|
||
|
||
⇒「必需」这条约束放在 **Handler**(`GetInbox` / `MarkInboxRead` / `HeartbeatAgent`),
|
||
|
||
不是这里:它是**接口契约**,不是数据层不变量。放 repo 会让人类那条路也没法用,
|
||
而人类侧并没有"我处在哪个工作区"这个概念。
|
||
|
||
# 为什么 Agent 必须带工作区
|
||
|
||
Agent 的收件箱原先只按名字过滤(`WHERE m.to_name = $1`),于是 `pi` 这个名字下
|
||
**所有工作区**的信混成一个池子。生产实测:在 `mc` 工作区干活的 pi 读收件箱拿到
|
||
200 封,其中 191 封属于 `/home/program/agentmail` —— 它照着那些信里的断言去改
|
||
agentmail 的代码,把手上 mc 的活丢在一边(用户当场问「你怎么干着干着修
|
||
agentmail 去了?」)。
|
||
|
||
三维地址是 `name@path.session` —— **path 位本来就该参与寻址**。收件箱侧此前
|
||
完全没用它,那不是"另一种语义",是没兑现契约。
|
||
*/
|
||
func checkWorkspace(workspace string) error {
|
||
w := strings.TrimSpace(workspace)
|
||
if w == "" {
|
||
return nil // 人类侧:不过滤,由 WebUI 按工作区分组显示
|
||
}
|
||
// 绝对路径:相对路径在服务端无法解释,且不同调用方 cwd 不同 ⇒ 拼出来必然对不上。
|
||
if !strings.HasPrefix(w, "/") {
|
||
return fmt.Errorf("workspace 必须是绝对路径,收到 %q", w)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// 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"
|
||
// (测试当场抓到)。
|
||
//
|
||
// ★★ reader 的占位符**必须排在调用方实参之后**(2026-09-20 修,一个静默少行的真 bug)。
|
||
//
|
||
// 原实现把 reader **前置**(`append([]any{reader}, args...)`)当 `$1`,而两个调用方
|
||
// 的 `where` 早就把 `$1` 用成了 recipient:
|
||
//
|
||
// MarkMailsReadFor : $1=recipient,$2..$(N+1)=ids,args=[recipient, ids...]
|
||
// MarkAllInboxReadForSession : $1=recipient,$2=sessionID,args=[recipient, …]
|
||
//
|
||
// 前置之后整张绑定表**右移一格** ⇒ `$2` 拿到的是 recipient 而不是 `id1`:
|
||
// - `IN ($2 …, $(N+1))` 实际是 `(recipient, id1 … idN-1)` ⇒ **最后一封永远不插**;
|
||
// 只传 1 封时 `$2=recipient` ⇒ **一封都不插**。
|
||
// - `MarkAllInboxReadForSession` 带 session 时 `$2=recipient` 被当成 `session_id` ⇒ 同样一行不插。
|
||
//
|
||
// **为什么长期零痕迹**:调用方返回的"标了几封"来自**随后那条 `UPDATE mails`**,
|
||
// 它用的是没被前置的 args ⇒ 计数正确、那列冗余值也正确,
|
||
// 只有权威列 `mail_reads` 静默少行。而未读判据是 `mail_reads`(见 `unreadFor`),
|
||
// 于是邮件**看起来已读、实际仍是未读** ⇒ 重启补投时被当新信重投。
|
||
//
|
||
// 现在把 reader 放在**最后一个**占位符,两个调用方的 `$1` 语义各自保持不变。
|
||
func markReadFor(ctx context.Context, reader string, where string, args ...any) error {
|
||
// reader 的编号紧跟在调用方实参之后,避免与它们已占用的 $N 相撞。
|
||
n := len(args) + 1
|
||
ph := fmt.Sprintf("$%d", n)
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`INSERT INTO mail_reads (mail_id, reader_name)
|
||
SELECT m.mail_id, `+ph+` 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 = `+ph+`)`,
|
||
append(append([]any{}, args...), reader)...)
|
||
return err
|
||
}
|
||
|
||
// MarkMailRead 记下**这个读者**读过这封邮件(人类端点)。
|
||
func MarkMailRead(ctx context.Context, id uuid.UUID, reader string) error {
|
||
/*
|
||
* 空读者**直接报错**,不兜底成某个读者(pi 2026-09-14 提的"把约定变成做错会红")。
|
||
*
|
||
* 这条约定("已读是**按读者**记的")原先只靠注释和调用方的自觉:HTTP 那一层取的是
|
||
* `user.Username`(不是请求体里的参数),所以线上不会传空 —— 但**函数本身**接受空串,
|
||
* 于是将来任何一个新调用方传 `""`,就会写入一行 `reader_name=''` 的垃圾:
|
||
* 它不属于任何人,却会让"某人的未读"统计出偏差,而且**没有任何东西会红**。
|
||
* 按读者记账的东西,"读者是谁"是必填语义,不是可选项。
|
||
*/
|
||
if strings.TrimSpace(reader) == "" {
|
||
return fmt.Errorf("MarkMailRead: reader 不能为空(已读按读者记录,reader 是必填语义)")
|
||
}
|
||
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, workspace string, limit int) ([]models.Mail, error) {
|
||
return ListInboxScoped(ctx, agentName, status, workspace, limit, uuid.Nil)
|
||
}
|
||
|
||
// ListInboxScoped 与 ListInbox 相同,但 `sessionID` 非零时只列该会话的邮件。
|
||
//
|
||
// `workspace` 非空时收窄到该工作区(Agent 侧 Handler 强制必填;人类侧为空):只列属于该工作区的会话里的信。
|
||
// 判据用 `s.workspace`(会话的权威工作区)而不是 `m.to_workspace`:
|
||
// 后者是**这封信**的信封字段,可能是抄送、可能是历史遗留;而"这条线索属于哪个
|
||
// 工作区"是会话的属性,一处定死才不会两种答案。
|
||
func ListInboxScoped(ctx context.Context, agentName, status, workspace string, limit int, sessionID uuid.UUID) ([]models.Mail, error) {
|
||
if err := requireReader(agentName); err != nil {
|
||
return nil, err
|
||
}
|
||
if err := checkWorkspace(workspace); err != nil {
|
||
return nil, err
|
||
}
|
||
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'`
|
||
// ★ 工作区收窄(必需):`pi` 在两个工作区各有一条收件箱,互不可见。
|
||
//
|
||
// 用 EXISTS 而不是再 JOIN 一次 sessions:s 已经在上面 JOIN 过了,
|
||
// 这里直接把条件写进 WHERE 即可(同一条 s)。
|
||
args := []any{agentName}
|
||
if strings.TrimSpace(workspace) != "" {
|
||
args = append(args, workspace)
|
||
q += fmt.Sprintf(` AND s.workspace = $%d`, len(args))
|
||
}
|
||
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 {
|
||
// ★ 占位符必须是**动态序号**:$2 现在被 workspace 占了。
|
||
q += fmt.Sprintf(` AND m.status = $%d`, len(args)+1)
|
||
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, workspace string) (int, error) {
|
||
return CountUnreadScoped(ctx, agentName, workspace, uuid.Nil)
|
||
}
|
||
|
||
// CountUnreadScoped 与 CountUnread 相同,但 `sessionID` 非零时只数那条会话。
|
||
//
|
||
// ★ 必须与 ListInboxScoped **同一个收窄口径**:桥的补投判据是
|
||
//
|
||
// `pending_mails = CountUnread` —— 两者口径不一致时,列表看不到的信会一直
|
||
// 被算成"还有未读",桥每次心跳都重放一遍(这是设计文档里记过的那个坑)。
|
||
func CountUnreadScoped(ctx context.Context, agentName, workspace string, sessionID uuid.UUID) (int, error) {
|
||
if err := requireReader(agentName); err != nil {
|
||
return 0, err
|
||
}
|
||
if err := checkWorkspace(workspace); err != nil {
|
||
return 0, err
|
||
}
|
||
var count int
|
||
q := `
|
||
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'`
|
||
args := []any{agentName}
|
||
if strings.TrimSpace(workspace) != "" {
|
||
args = append(args, workspace)
|
||
q += fmt.Sprintf(` AND s.workspace = $%d`, len(args))
|
||
}
|
||
if sessionID != uuid.Nil {
|
||
args = append(args, sessionID)
|
||
q += fmt.Sprintf(` AND m.session_id = $%d`, len(args))
|
||
}
|
||
err := db.DB.QueryRowContext(ctx, q, args...).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:别名全局唯一且本身就承担寻址职责,
|
||
// 再叠一层工作区校验只会让「名字对上了却送不到」变成一种难查的失败。
|
||
// ErrSessionAmbiguous:地址里没给 path,而该别名在多个工作目录下都存在。
|
||
//
|
||
// 不给"随便挑一条"的兜底:那等于把信随机投进某个工作区 —— 2026-09-15 那次错投
|
||
// (工作区在 TrueAgent 的 pi 发的信进了 workspace=/home 的会话/客户端)就是这么来的。
|
||
var ErrSessionAmbiguous = errors.New("会话别名在多个工作目录下都存在,请在地址里写明 path")
|
||
|
||
// FindNamedSessionFor 按**会话身份**找会话:地址 `name@path.<别名>`。
|
||
//
|
||
// # path 是身份的一半,不是提示(用户 2026-09-15 订正)
|
||
//
|
||
// 原话:「agent 平台的 session 是和 path 绑定的,path+session 才能指定到准确的 agent,
|
||
// 而授权也是对 session 授权,而不是整个 agent」。所以:
|
||
//
|
||
// - 地址里**给了** path → 必须 path 与别名**同时命中**。对不上就是"没有这条会话",
|
||
// 而不是"退回按别名找"—— 后者正是错投的根因(path 被忽略 ⇒ 两条不相干的线索
|
||
// 落进同一条会话,甚至同一个 pi 会话文件)。
|
||
// - 地址里**没给** path(`name@.别名`,人类回信常用)→ 只有该别名**唯一**时才认;
|
||
// 同名多条时返回 ErrSessionAmbiguous,让调用方要求补 path。
|
||
//
|
||
// name 仍要在这条线索里出现过(from/to/cc):否则任何会话都能被叫任意名字。
|
||
func FindNamedSessionFor(ctx context.Context, name, path, alias string) (uuid.UUID, error) {
|
||
path = strings.TrimSpace(path)
|
||
if path != "" {
|
||
path = filepath.Clean(path)
|
||
}
|
||
|
||
// 两个变体共用同一个"我参与过"判据。
|
||
var participation = `
|
||
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) + `)
|
||
)`
|
||
|
||
if path != "" {
|
||
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 COALESCE(s.workspace, '') = $3`+participation+`
|
||
ORDER BY s.updated_at DESC LIMIT 1`, alias, name, path).Scan(&id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return uuid.Nil, ErrSessionNotFound
|
||
}
|
||
return id, err
|
||
}
|
||
|
||
// 没给 path:取最多两条,用于区分"唯一"与"歧义"。
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT s.session_id FROM sessions s
|
||
WHERE s.session_alias = $1 AND s.status <> 'archived'`+participation+`
|
||
ORDER BY s.updated_at DESC LIMIT 2`, alias, name)
|
||
if err != nil {
|
||
return uuid.Nil, err
|
||
}
|
||
defer rows.Close()
|
||
var ids []uuid.UUID
|
||
for rows.Next() {
|
||
var id uuid.UUID
|
||
if err := rows.Scan(&id); err != nil {
|
||
return uuid.Nil, err
|
||
}
|
||
ids = append(ids, id)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return uuid.Nil, err
|
||
}
|
||
switch len(ids) {
|
||
case 0:
|
||
return uuid.Nil, ErrSessionNotFound
|
||
case 1:
|
||
return ids[0], nil
|
||
default:
|
||
return uuid.Nil, ErrSessionAmbiguous
|
||
}
|
||
}
|
||
|
||
// 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) {
|
||
return listContacts(ctx, forUser, nil, archived)
|
||
}
|
||
|
||
func listContacts(ctx context.Context, forUser string, onlyWorkspace *string, archived bool) ([]Contact, error) {
|
||
op := "<>"
|
||
if archived {
|
||
op = "="
|
||
}
|
||
// $1 恒为 forUser,即使它是空串。
|
||
//
|
||
// 未读计数那个子查询里 `r.reader_name = $1` **一直在引用 $1**,而原先的写法是
|
||
// 「forUser 为空就不传参」—— 那样 $1 就悬空了:Postgres 直接报 no parameter $1,
|
||
// SQLite 则静默把 `= $1` 当 `= NULL` 比(次次不成立,未读计数退化成“全部未归档”)。
|
||
// 管理员 `?all=true` 正是走的那条悬空路径(scope="")。
|
||
scope := ""
|
||
args := []any{forUser}
|
||
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) + `)
|
||
))`
|
||
}
|
||
if onlyWorkspace != nil {
|
||
args = append(args, *onlyWorkspace)
|
||
scope += fmt.Sprintf(` AND COALESCE(s.workspace, '') = $%d`, len(args))
|
||
}
|
||
// 取会话里最早那封邮件作为联系人身份。
|
||
// 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,
|
||
-- ★ 2026-09-15:path 取自**会话**,不再取第一封邮件的 from_workspace ——
|
||
-- 那一列历史上存的是 **agent 名**(写入侧 bug,已修),于是联系人地址被拼成
|
||
-- zcode@zcode.<别名>,而它正是可被复制出去的"错误地址"(用户 09-15:
|
||
-- 「错误的老数据直接清除,否则有概率发生感染」—— 这就是那条感染通道)。
|
||
-- 顺序:会话的工作区(权威)→ 该会话首封邮件的 to_workspace(寻址时写的那个)
|
||
-- → 空。**只用 to_***:from_workspace 与"对方在哪"无关。
|
||
COALESCE(NULLIF(s.workspace, ''), NULLIF(m.to_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 filterWorkspaces(out), nil
|
||
}
|
||
|
||
/*
|
||
filterWorkspaces 只留「看起来像工作目录」的候选。
|
||
|
||
为什么需要这一步(2026-09-15 用户报的「莫名其妙的 pi@/home 会话」):
|
||
建议列表是从**历史数据**学的(mails.to_workspace + 平台心跳 + 注册自报),
|
||
而历史里混进了「进程恰好所在的目录」:`/root`、`/home`、`/home/program`、
|
||
`/root/.pi/mail-sessions/<uuid>`。它们一旦进了候选,用户点第一条就发信给
|
||
`zcode@/home` → 那条会话的 workspace 就成了 `/home` → 之后这条线索里
|
||
**所有参与方**都显示 `xxx@/home`,pi 桥还会真把 worker 起在那里
|
||
(实测日志 `新建 pi 会话 …(cwd=/home)`)—— 而沙箱的 rw 只有
|
||
/home/program/agentmail,那个 worker 连文件都写不了。
|
||
|
||
★ 这是一个**自增强环**:污染的会话 → 污染的候选 → 新的污染会话。
|
||
只清理现有的 `/home` 不够,必须同时不把这类值当候选。
|
||
|
||
判据不用黑名单(那是针对已知值),而是问「它是不是一个**具体**的工作目录」:
|
||
|
||
1. 绝对路径(相对路径不可能是工作目录的地址);
|
||
2. 任何一段以 `.` 开头 → 不是(那是缓存/会话存储,如 `/root/.pi/...`);
|
||
3. 进程用户的家目录(`/root`)→ 不是(agent 的工作目录不是它的进程家目录);
|
||
4. 是**另一个候选的祖先** → 不是(`/home`、`/home/program` 都被
|
||
`/home/program/agentmail` 这条干掉:它只是容器,不是干活的地方)。
|
||
|
||
刻意**不检查目录是否存在**:注册时自报的目录可能还没建(`/tmp/remotebot-ws`
|
||
就是那种),存在性检查会把合法候选误杀。
|
||
*/
|
||
func filterWorkspaces(in []string) []string {
|
||
if len(in) == 0 {
|
||
return in
|
||
}
|
||
out := make([]string, 0, len(in))
|
||
seen := map[string]bool{}
|
||
for _, p := range in {
|
||
// 只归一化去重。**不做"这目录像不像工作目录"的判断**:
|
||
// 用户 2026-09-15 否掉了那个方向(「拒收那个目录干啥?明明是你的网关投递
|
||
// 逻辑造成的错误」)—— /home 是合法目录,错在投递,不在值。
|
||
c := filepath.Clean(p)
|
||
if seen[c] {
|
||
continue
|
||
}
|
||
seen[c] = true
|
||
out = append(out, c)
|
||
}
|
||
return out
|
||
}
|
||
|
||
func hasHiddenSegment(p string) bool {
|
||
for _, seg := range strings.Split(strings.Trim(p, "/"), "/") {
|
||
if len(seg) > 1 && strings.HasPrefix(seg, ".") {
|
||
return true
|
||
}
|
||
}
|
||
return false
|
||
}
|
||
|
||
// isAncestorOfOther 报告 p 是不是另一个候选的**严格**祖先。
|
||
// 只比给定的候选集,不碰文件系统(不用「它下面有没有目录」这种猜测)。
|
||
func isAncestorOfOther(p string, all []string) bool {
|
||
prefix := strings.TrimSuffix(p, "/") + "/"
|
||
for _, q := range all {
|
||
if q == p {
|
||
continue
|
||
}
|
||
if strings.HasPrefix(filepath.Clean(q), prefix) {
|
||
return true
|
||
}
|
||
}
|
||
return false
|
||
}
|
||
|
||
// 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) {
|
||
if err := requireReader(name); err != nil {
|
||
return 0, err
|
||
}
|
||
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, workspace string) (int, error) {
|
||
return MarkAllInboxReadForSession(ctx, recipient, workspace, uuid.Nil)
|
||
}
|
||
|
||
// MarkAllInboxReadForSession 只标掉某条会话里发给 recipient 的未读。
|
||
//
|
||
// 为什么需要:Agent 的「不给 mail_ids,全部标掉」在会话驱动的 worker 里会跨会话
|
||
// 误伤(见 ListInboxScoped 上面那段说明)。不带 sessionID(uuid.Nil)时是旧语义。
|
||
func MarkAllInboxReadForSession(ctx context.Context, recipient, workspace string, sessionID uuid.UUID) (int, error) {
|
||
if err := checkWorkspace(workspace); err != nil {
|
||
return 0, err
|
||
}
|
||
// ★ 工作区收窄:与 ListInboxScoped **同一个谓词** —— 列表看不到的信却被
|
||
// "全部标掉"标掉,就是静默丢信(session_scope_test.go 记过这个形状)。
|
||
wsFilter := ""
|
||
args := []any{recipient}
|
||
if strings.TrimSpace(workspace) != "" {
|
||
args = append(args, workspace)
|
||
wsFilter = fmt.Sprintf(" AND workspace = $%d", len(args))
|
||
}
|
||
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'` + wsFilter + `)`
|
||
if sessionID != uuid.Nil {
|
||
args = append(args, sessionID)
|
||
scope += fmt.Sprintf(` AND m.session_id = $%d`, len(args))
|
||
}
|
||
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'` + wsFilter + `)`
|
||
updArgs := []any{recipient}
|
||
if strings.TrimSpace(workspace) != "" {
|
||
updArgs = append(updArgs, workspace)
|
||
}
|
||
if sessionID != uuid.Nil {
|
||
updArgs = append(updArgs, sessionID)
|
||
upd += fmt.Sprintf(` AND session_id = $%d`, len(updArgs))
|
||
}
|
||
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
|
||
}
|