四个各自独立的生产缺陷,共同的根源都是「本该属于会话的属性没有存在会话上」。 ## 1. dsh 指定工作目录完全失效(所有会话落进「未分组」) 插件建会话时用的 cwd 是自己拼的 `~/.dsh/mail-sessions/mail-<uuid>` —— 每封邮件一个全新的空目录。DSH 与 opencode 都按 cwd 给会话分组,于是所有 邮件会话既不属于任何项目、彼此也不同组。 而 Gateway 从来没把地址里的 path 位发给插件:`notifyRecipients` 的 payload 只有 mail_id/session_id/from_name/subject,`to_workspace` 虽然入库了却不在 SSE 事件里,插件即使想用也拿不到。 - SSE `new_mail` 事件加 `to_workspace`。**每个收件方拿到自己那个地址的 path**, 不是主收件人的 —— 抄送给 opencode@/a 与主发给 dsh@/b 是两个工作区 - 两个插件的 cwd 都改为取寻址的 path 位;不存在的目录**不创建**而是回退到 兜底目录(一个笔误不该在磁盘上落下真目录,Agent 会在里面一无所获地干活) - 拒绝相对路径:cwd 的相对基准是 harness 进程的启动目录,systemd 下通常是 `/` ## 2. 会话别名列不出工作区下的历史会话(无法选择) workspace 只存在于 `mails.to_workspace` 上,「这个工作区下有哪些会话」必须 JOIN mails 再从收发双方的 workspace 里猜。而 Agent 回信时 from_workspace 填的是 **Agent 名**而不是路径,旧条件 `to_workspace = $p OR from_workspace = $p` 在只剩 Agent 回信可匹配时两边都对不上。 - `sessions.workspace` 新列,`CreateSession` 从地址的 path 位带入 - `SuggestSessionCandidates` 取代 `SuggestSessionsFor`:以会话自己的 workspace 为权威,历史会话(该列为空)回退到 mails 反推 —— 升级后老会话不该消失 - `FindOrCreateDefaultSession` 同步改用会话的 workspace ## 3. 平台侧会话在补全里根本不存在 人直接在 opencode/DSH 界面上开的会话,Gateway 一无所知。 新增 `agent_platform_sessions` 镜像表,插件在心跳里上报快照。 **上报而非 Gateway 反向拉取**:当前架构是单向的(Agent 持密钥主动连 Gateway, Gateway 从不外呼),反向拉取需要它保存各平台的地址与凭证,那是另一套信任模型。 - 与 sessions 表分开存:镜像里是别人家的会话,id 属于平台的 id 空间,没有 本侧的 owner/预算/邮件。混进 sessions 会让每一处「按会话鉴权」都要先判断 这条到底是不是真的本侧会话 - **整表替换而非增量合并**:平台侧删掉的会话必须从候选里消失 —— session 位是 三态语义,指向不存在的会话直接 404 - **nil 与空数组语义不同**:插件拉不到列表时省略该字段(保留镜像), 而不是传空数组把镜像抹掉 - **subagent 子会话不上报**:实测 DSH 的 list 里混着 49 条子会话,标题就是 派活的提示词前缀(九条都叫 "You are auditing ONE file"),slug 全撞名; 它们是父 agent 内部的工作单元,人往里发邮件毫无意义 - **slug 撞名只留最近那条**:服务端只能取其中一条,上报同名项只会让补全里 出现几个点哪个都不确定的候选 - DSH 插件此前**完全没有心跳** —— Gateway 靠 last_seen 判在线,一直靠注册撑着 补全候选带标题与来源:`suggestions` 保留纯字符串数组(不打破已部署的前端与 第三方客户端),新增同序的 `candidates`。过滤时标题也参与匹配 —— 人记得的是 「缓存选型」而不是 brisk-harbor 这种随机短名。 ## 4. 对话树看不见抄送与转发产生的分支 旧实现从锚点分「祖先链 + 子树」两路展开,而**兄弟节点既不是锚点的祖先也不是 它的子孙**:一封抄送给两个 Agent 的邮件收到两个回复,从其中一个看树永远看不到 另一个;挂在原件上的转发分支同理。 改为先 `ThreadRootOf` 上溯到线索根,再从根整树 BFS。只剩一个加载方向, 因此不再需要滚动位置补偿。前端补上抄送人列表与转发标记 —— 树上两个兄弟节点 为什么并列,唯一的解释就是父邮件抄送给了两个人。 ## 5. DSH 插件(Phase 7.7) 卡了一下午的 `Cannot read properties of undefined (reading 'kind')` 根因是 `followup()` 的参数形状:DSH 要完整的 UserMessage(content + source), 而我照抄了 opencode 的 parts 数组。错误抛在 agent-loop 内部,不指向调用点。 - `agent/status` → idle 时自动转发最后一条 assistant 消息(对应 opencode 的 session.idle),复用 relay-dedup 让位于模型的主动回信,走免配额通道 - `approval/request` 权限询问转邮件问人。与 opencode 的差异:那边的 permission.ask 是同步钩子只能立即返回 ask,DSH 这边是异步 waterfall, 可以真的等人 —— 拆插件时未决询问一律 fail closed,否则 await 永不返回 - 会话别名由模型标题派生(保留中文,去掉 `.` `@` `/` 等寻址分隔符 —— 留在别名里会让它自己被解析器切开) - 逻辑放 lib/ 下的纯函数并加测试:三类约定都是「错了不当场报错、只在深处 炸一个无关错误」 ## 其他 - `deploy/reset-demo.sh`:清空演示邮件数据,保留账号与密钥。备份用 `.backup` 而非 cp(WAL 下 cp 拿到的是缺尾巴的库);手工按依赖顺序删(SQLite 的 foreign_keys 默认关,声明了 REFERENCES 也不级联);只在目标是默认库时才碰 systemd(演练时误停过一次生产服务) - 插件 dist/ 不进版本库,install.sh 负责构建 - `permission_decision` 事件补 session_id:插件重启丢了待决映射时要靠它定位会话
244 lines
8.9 KiB
Go
244 lines
8.9 KiB
Go
package repo
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
|
||
"github.com/agentmail/gateway/internal/db"
|
||
"github.com/agentmail/gateway/internal/models"
|
||
"github.com/google/uuid"
|
||
)
|
||
|
||
// 对话树。
|
||
//
|
||
// **不另建 tree_nodes 表**:`mails.parent_mail_id` 已经完整编码了树结构 ——
|
||
// 回复指向来信,转发指向被转发的原件。再维护一张 tree_nodes 就是第二份真相,
|
||
// 两处不一致时无法判断谁对。这里直接用递归 CTE 在 mails 上查。
|
||
//
|
||
// 树可以跨会话:转发把线索引到新会话,但 parent 仍指向原件。这正是「对话树」比
|
||
// 「会话内平铺」更有价值的地方 —— 能看出一条线索分叉去了哪里。
|
||
// 也正因如此,读取时必须按会话逐个鉴权(见 handler):
|
||
// A 转发给 B 之后,B 与 C 在新会话里的往来不能回流给 A。
|
||
//
|
||
// **从根展开,而不是从锚点展开**:曾经的实现是「锚点的祖先链 + 锚点的子树」,
|
||
// 于是兄弟节点整条分支都在盲区里 —— 一封抄送给两个 Agent 的邮件,两个回复
|
||
// 互为兄弟,从其中一个看树看不到另一个;挂在原件上的转发同理。
|
||
// 兄弟既不是锚点的祖先也不是它的子孙,只有先上溯到根、再整棵 BFS 才能覆盖。
|
||
//
|
||
// **分块加载而非截断**:线索可以有几百封,一次全取要把几 MB 预览塞给前端。
|
||
// 从根 BFS 后只剩一个方向,游标就是「已取到的节点数」。
|
||
|
||
// TreeMail 是树里的一个节点。正文只带预览:整棵线索带全文可能几百 KB,
|
||
// 前端点开某封时再单取全文与附件清单。
|
||
type TreeMail struct {
|
||
models.Mail
|
||
// Depth 是**距线索根**的层级:0 = 根,1 = 它的直接回复。
|
||
// 从根展开后根一定在结果里,绝对深度因此总是可知的(早先按相对锚点算,
|
||
// 是因为那时根可能还没取到)。
|
||
Depth int `json:"depth"`
|
||
AttachmentCount int `json:"attachment_count"`
|
||
}
|
||
|
||
// descendantDepthCap 只是数据损坏时的兜底。
|
||
//
|
||
// parent_mail_id 正常不成环(新邮件只能指向已存在的旧邮件),但一旦被外部工具改坏,
|
||
// 无上限的递归 CTE 会把进程拖死。取得足够大,正常数据碰不到。
|
||
const descendantDepthCap = 10000
|
||
|
||
const threadCols = `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,
|
||
m.status, m.created_at, s.session_alias,
|
||
(SELECT COUNT(*) FROM attachments a WHERE a.mail_id = m.mail_id) AS attach_count`
|
||
|
||
// ThreadRootOf 沿 parent_mail_id 上溯到线索的根,返回根的 mail_id 与锚点到根的层数。
|
||
//
|
||
// 「根」= 链条最上面那封:parent_mail_id 为 NULL,或者指向一封已被删掉的邮件
|
||
// (JOIN 断掉,递归自然停在这一层)。锚点自己没有父时返回它自己、depth 0。
|
||
//
|
||
// **不做可见性过滤**:不可见的中间段必须能穿过 —— 转发把线索引进别人的会话,
|
||
// 再往上却可能仍是自己参与的往来。只返回 id 与层数,不泄露任何内容。
|
||
func ThreadRootOf(ctx context.Context, anchorID uuid.UUID) (uuid.UUID, int, error) {
|
||
var rootID uuid.UUID
|
||
var lvl int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
WITH RECURSIVE up(mail_id, parent_mail_id, lvl) AS (
|
||
SELECT mail_id, parent_mail_id, 0 FROM mails WHERE mail_id = $1
|
||
UNION ALL
|
||
SELECT m.mail_id, m.parent_mail_id, up.lvl + 1
|
||
FROM mails m JOIN up ON m.mail_id = up.parent_mail_id
|
||
WHERE up.lvl < $2
|
||
)
|
||
SELECT mail_id, lvl FROM up ORDER BY lvl DESC LIMIT 1
|
||
`, anchorID, descendantDepthCap).Scan(&rootID, &lvl)
|
||
if err != nil {
|
||
return uuid.Nil, 0, err
|
||
}
|
||
return rootID, lvl, nil
|
||
}
|
||
|
||
// AncestorsRaw 沿 parent_mail_id 上溯,取第 offset+1 .. offset+limit 层的祖先。
|
||
// 层号 1 = 父,2 = 祖父;返回的 Depth 为负数(相对锚点)。
|
||
//
|
||
// 从根 BFS 之后这个函数只在一处还有用:巨型线索里锚点没落在 BFS 首页时,
|
||
// 用它把「根到锚点」这条路径单独补齐,保证点开的那封一定看得见。
|
||
// 调用方需要自己把负 depth 换算成绝对深度(锚点绝对深度由 ThreadRootOf 给出)。
|
||
//
|
||
// **不做可见性过滤**,理由同 ThreadRootOf。过滤放在 handler 层(那里知道调用者是谁)。
|
||
//
|
||
// 第二个返回值表示 offset+limit 层之上还有节点。
|
||
func AncestorsRaw(ctx context.Context, anchorID uuid.UUID, offset, limit int) ([]TreeMail, bool, error) {
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
WITH RECURSIVE up(mail_id, parent_mail_id, lvl) AS (
|
||
SELECT mail_id, parent_mail_id, 0 FROM mails WHERE mail_id = $1
|
||
UNION ALL
|
||
SELECT m.mail_id, m.parent_mail_id, up.lvl + 1
|
||
FROM mails m JOIN up ON m.mail_id = up.parent_mail_id
|
||
WHERE up.lvl < $2
|
||
)
|
||
SELECT `+threadCols+`, u.lvl
|
||
FROM up u
|
||
JOIN mails m ON m.mail_id = u.mail_id
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE u.lvl > $3
|
||
ORDER BY u.lvl ASC
|
||
`, anchorID, offset+limit+1, offset)
|
||
if err != nil {
|
||
return nil, false, err
|
||
}
|
||
// 多取一层用来判断「上面还有没有」,不返回给调用方
|
||
out, err := scanTreeRows(rows, true)
|
||
if err != nil {
|
||
return nil, false, err
|
||
}
|
||
hasMore := len(out) > limit
|
||
if hasMore {
|
||
out = out[:limit]
|
||
}
|
||
return out, hasMore, nil
|
||
}
|
||
|
||
// DescendantsRaw 取给定节点及其全部子孙,BFS 顺序(同层按时间),按节点数分页。
|
||
//
|
||
// 传线索的根(见 ThreadRootOf)就能覆盖整棵树:兄弟、抄送产生的平行回复、
|
||
// 挂在原件上的转发分支,全都是根的子孙。offset = 0 时结果第一个是起点自己(Depth 0)。
|
||
//
|
||
// 同样不做可见性过滤:不可见的子节点下面可能挂着可见的孙节点
|
||
// (别人把线索转走又转回来给我)。
|
||
//
|
||
// 注意 CTE 每次都会走完整棵子树,LIMIT 只截断输出。一条邮件线索通常几十封,
|
||
// 这个代价可以接受;真出现巨型线索时再加物化。
|
||
func DescendantsRaw(ctx context.Context, anchorID uuid.UUID, offset, limit int) ([]TreeMail, bool, error) {
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
WITH RECURSIVE down(mail_id, lvl) AS (
|
||
SELECT mail_id, 0 FROM mails WHERE mail_id = $1
|
||
UNION ALL
|
||
SELECT m.mail_id, down.lvl + 1
|
||
FROM mails m JOIN down ON m.parent_mail_id = down.mail_id
|
||
WHERE down.lvl < $2
|
||
)
|
||
SELECT `+threadCols+`, d.lvl
|
||
FROM down d
|
||
JOIN mails m ON m.mail_id = d.mail_id
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
ORDER BY d.lvl ASC, m.created_at ASC, m.mail_id ASC
|
||
LIMIT $3 OFFSET $4
|
||
`, anchorID, descendantDepthCap, limit+1, offset)
|
||
if err != nil {
|
||
return nil, false, err
|
||
}
|
||
out, err := scanTreeRows(rows, false)
|
||
if err != nil {
|
||
return nil, false, err
|
||
}
|
||
hasMore := len(out) > limit
|
||
if hasMore {
|
||
out = out[:limit]
|
||
}
|
||
return out, hasMore, nil
|
||
}
|
||
|
||
// TreeMailByID 取单封邮件的树节点形式,深度由调用方给定。
|
||
//
|
||
// 补齐「根 → 锚点」路径时用得上:AncestorsRaw 从父开始,不含锚点自己。
|
||
// 同样不做可见性过滤,由 handler 负责。
|
||
func TreeMailByID(ctx context.Context, id uuid.UUID, depth int) (*TreeMail, error) {
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT `+threadCols+`, $2
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.mail_id = $1
|
||
`, id, depth)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
out, err := scanTreeRows(rows, false)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(out) == 0 {
|
||
return nil, sql.ErrNoRows
|
||
}
|
||
return &out[0], nil
|
||
}
|
||
|
||
// scanTreeRows 读出节点。negate 为真时把层号取负(祖先方向)。
|
||
func scanTreeRows(rows interface {
|
||
Next() bool
|
||
Scan(...interface{}) error
|
||
Err() error
|
||
Close() error
|
||
}, negate bool) ([]TreeMail, error) {
|
||
defer rows.Close()
|
||
|
||
out := []TreeMail{}
|
||
for rows.Next() {
|
||
var t TreeMail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
var lvl int
|
||
if err := rows.Scan(&t.ID, &t.SessionID, &t.ParentMailID,
|
||
&t.FromName, &t.FromWorkspace, &t.ToName, &t.ToWorkspace,
|
||
&ccJSON, &t.Subject, &t.Body, &t.MailType, &t.PermResult,
|
||
&t.Status, &t.CreatedAt, &alias, &t.AttachmentCount, &lvl); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &t.CCList)
|
||
}
|
||
if t.CCList == nil {
|
||
t.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
t.SessionAlias = *alias
|
||
}
|
||
if negate {
|
||
t.Depth = -lvl
|
||
} else {
|
||
t.Depth = lvl
|
||
}
|
||
t.BodyPreview = preview(t.Body, 240)
|
||
t.Body = "" // 树视图只要预览,全文按需单取
|
||
out = append(out, t)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
// preview 按 UTF-8 边界截断正文。
|
||
// 直接切字节会把多字节字符切成半个,前端渲染出 U+FFFD 替换符。
|
||
func preview(s string, max int) string {
|
||
if len(s) <= max {
|
||
return s
|
||
}
|
||
cut := max
|
||
for cut > 0 && !utf8Start(s[cut]) {
|
||
cut--
|
||
}
|
||
return s[:cut] + "..."
|
||
}
|
||
|
||
// utf8Start 判断某字节是否为一个 UTF-8 序列的首字节
|
||
func utf8Start(b byte) bool { return b&0xC0 != 0x80 }
|