diff --git a/gateway/internal/handler/mail.go b/gateway/internal/handler/mail.go index 6ca93cd..84fc0cf 100644 --- a/gateway/internal/handler/mail.go +++ b/gateway/internal/handler/mail.go @@ -9,8 +9,8 @@ import ( "github.com/agentmail/gateway/internal/middleware" "github.com/agentmail/gateway/internal/models" + "github.com/agentmail/gateway/internal/notify" "github.com/agentmail/gateway/internal/repo" - "github.com/agentmail/gateway/internal/sse" "github.com/google/uuid" ) @@ -397,81 +397,19 @@ func SendMail(w http.ResponseWriter, r *http.Request) { JSON(w, http.StatusOK, resp) } -// notifyRecipients 向主收件人与抄送方推送 new_mail,并刷新相关方的会话列表。 -// 收件人可能是 Agent 也可能是人类用户(三维地址 name 位共享命名空间), -// 因此统一用 SendToRecipient 同时试 Agent 通道与用户通道。 +// notifyRecipients 是 notify.Recipients 的薄封装,保留旧签名减少调用点改动。 // -// **每个收件方拿到的 workspace 是自己那个地址的 path 位**,不是主收件人的: -// 三维地址 name@path.session 的 path 就是工作目录,插件要靠它建会话。 -// 抄送给 opencode@/a 与主发给 dsh@/b 是两个不同的工作区,共用一份 payload -// 会让抄送方在别人的目录里开会话。 -// -// 同理,**每个收件方拿到的 reply_address 也是自己那个地址**,并且 session 位已经 -// 把 `new` 换成真实别名:`.new` 建完会话就失效了,把原文那个 `x@/p.new` -// 送给参与方只会让它下一次又建一条新会话。 +// 实现只有一份,在 internal/notify 里 —— 此前 handler 与 scheduler 各写一份, +// 加字段时漏改一处直接造成生产事故(详见那个包的注释)。 func notifyRecipients(ctx context.Context, to models.Address, cc []models.Address, sessionID, mailID uuid.UUID, from, subject string) { - // 别名在此时已由 resolveTarget 保证存在(`.new` 与默认会话都过 EnsureSessionAlias)。 - // 仍可能为空的情形:命名写入失败(已吐日志)。此时退回省略 session 位, - // 而不是把 "new" 写进去 —— 后者会让参与方反复建新会话。 - alias := repo.SessionAliasOf(ctx, sessionID) - // 这条会话是否接管了一条平台侧已存在的会话(人在 TUI 里开的那种)。 - // 插件据此决定 resume 还是新建 —— 空串就是过去的行为。 - platformID := repo.PlatformIDOf(ctx, sessionID) - - payload := func(role, workspace, forName string) map[string]interface{} { - return map[string]interface{}{ - "mail_id": mailID.String(), - "session_id": sessionID.String(), - "from_name": from, - "subject": subject, - "mail_type": "normal", - "role": role, // to / cc - // to_workspace 是收件方地址的 path 位,即希望它在哪个工作目录干活。 - // 不带这一项的后果:插件只能自己拼一个临时目录,于是每封邮件都落在 - // 不同的空目录里,DSH / opencode 按 cwd 分组时全进「未分组」。 - "to_workspace": workspace, - // session_alias 是这条会话今后的寻址名。没有它的话,收到 `.new` - // 邮件的一方只持有一个 send_mail 不接受的 session_id。 - "session_alias": alias, - // reply_address 是「把回信发回这条会话」的现成地址。 - // 插件不必自己拼(拼错了就是静默开新会话)。 - "reply_address": models.FormatAddress(from, "", alias), - // self_address 是对方应当用来称呼自己的地址,供转发/报告时引用。 - "self_address": models.FormatAddress(forName, workspace, alias), - // platform_session_id 非空时,这封邮件要投进**平台侧已经存在的 - // 那条会话**(TUI 与邮箱是同一个 Agent 的两个入口)。 - // - // 插件必须 resume 而不是新建:新建会让人在 TUI 里看不到这封邮件 - // 带来的对话,而那正是接管这条会话的目的。 - // 空串 = 照旧按邮件新开一条平台会话。 - "platform_session_id": platformID, - } - } - - update := map[string]interface{}{ - "session_id": sessionID.String(), - "status": "active", - } - - // 参与方去重:收件人 + 所有抄送 + 发件人自己(刷新他的发件箱) - seen := map[string]bool{} - - sse.Default.SendToRecipient(to.Name, "new_mail", payload("to", to.Path, to.Name)) - sse.Default.SendToRecipient(to.Name, "session_update", update) - seen[to.Name] = true - - for _, c := range cc { - if seen[c.Name] { - continue - } - seen[c.Name] = true - sse.Default.SendToRecipient(c.Name, "new_mail", payload("cc", c.Path, c.Name)) - sse.Default.SendToRecipient(c.Name, "session_update", update) - } - - if !seen[from] { - sse.Default.SendToRecipient(from, "session_update", update) - } + notify.Recipients(ctx, notify.Mail{ + SessionID: sessionID, + MailID: mailID, + From: from, + To: to, + CC: cc, + Subject: subject, + }) } // GET /api/v1/mail/inbox diff --git a/gateway/internal/handler/permission.go b/gateway/internal/handler/permission.go index 275c5ee..4a54171 100644 --- a/gateway/internal/handler/permission.go +++ b/gateway/internal/handler/permission.go @@ -169,14 +169,24 @@ func RequestPermission(w http.ResponseWriter, r *http.Request) { return } - // 只推给该决策人 + // 只推给该决策人。 + // + // 这一处不走 notify.Recipients:那个函数推给「三维地址解析出的参与方」, + // 而权限询问的投递对象是逐会话树找出来的人类决策人(NearestHumanInThread), + // 不是一个地址 —— 抄送也不应当收到它(权限是待办,不是广播)。 + // + // 但 payload 必须带足字段:前端的授权页靠 session_alias + 会话 workspace + // 拼出「哪个 Agent、在哪个目录、哪条线索」。只给 from_name 的话人 + // 看到的只是一个光秃的 Agent 名,无法判断该不该批。 + alias := repo.SessionAliasOf(r.Context(), sessionID) sse.Default.SendToUser(decider, "new_mail", map[string]interface{}{ - "mail_id": mailID.String(), - "session_id": sessionID.String(), - "from_name": agentName, - "subject": "权限请求: " + req.Question, - "mail_type": "permission_request", - "role": "to", + "mail_id": mailID.String(), + "session_id": sessionID.String(), + "from_name": agentName, + "subject": "权限请求: " + req.Question, + "mail_type": "permission_request", + "role": "to", + "session_alias": alias, }) JSON(w, http.StatusOK, map[string]string{ diff --git a/gateway/internal/notify/mail.go b/gateway/internal/notify/mail.go new file mode 100644 index 0000000..675ce23 --- /dev/null +++ b/gateway/internal/notify/mail.go @@ -0,0 +1,162 @@ +// Package notify 是「一封邮件落库之后要通知谁、推什么」的**唯一实现**。 +// +// # 为什么单独成包 +// +// 在此之前有两份几乎相同的推送代码:`handler.notifyRecipients`(人发信、 +// Agent 发信、转发都走它)和 `scheduler` 里日历提醒自己拼的那一份。 +// +// 两份代码的代价在生产上兑现过一次,而且症状离原因很远:给 `new_mail` 加 +// `platform_session_id` 字段时只改了 handler 那份,调度器那份仍是旧的。 +// 于是日历提醒投进一条**接管会话**时,插件收不到 `platform_session_id`, +// 把它当成新会话另开了一条平台会话;那条新会话的名字随后经命名同步回写, +// **把接管会话的别名冲掉了** —— 人在补全里选中的「项目定位」变成了 +// 「日程提醒:…」,同一条会话因此在候选列表里出现两次,而另一条真实会话 +// 被按别名字符串去重吃掉了。 +// +// 链条上每一环都不报错。根因只是「同一件事写了两遍」。 +// +// 因此这个包对外只暴露一个入口:新增字段时不存在「另一处忘了改」的可能。 +package notify + +import ( + "context" + + "github.com/agentmail/gateway/internal/models" + "github.com/agentmail/gateway/internal/repo" + "github.com/agentmail/gateway/internal/sse" + "github.com/google/uuid" +) + +// Mail 描述一封刚落库的邮件需要推给谁。 +type Mail struct { + SessionID uuid.UUID + MailID uuid.UUID + // From 是发件方名字。人类用户名与 Agent 名共享命名空间,这里不区分。 + From string + // To 是主收件方地址(三维寻址已解析)。 + To models.Address + // CC 是抄送方地址列表。 + CC []models.Address + // Subject 是邮件主题。 + Subject string + // MailType 默认 "normal";权限请求等特殊类型由调用方指定。 + MailType string + // Origin 标记这封信的来源,供插件与 UI 区分「定时提醒」与「有人在找它」。 + // 空串表示普通邮件。 + Origin string + // ReplyToName 是「把回信发回这条会话」时该写的收件人名。 + // + // 默认取 From。日历提醒必须覆盖它:发件人是 `calendar`,而那不是一个 + // 收得到信的账号 —— 回给它的信投不出去。此时应当填主收件方自己的名字, + // 让模型把结果回报到同一条线索上。 + ReplyToName string +} + +// Recipients 把一封邮件推给主收件人、所有抄送方,并刷新发件方的会话列表。 +// +// # 每个收件方拿到的是**自己那个地址** +// +// 三维地址 `name@path.session` 的 path 就是工作目录,插件靠它建会话。 +// 抄送给 `opencode@/a` 与主发给 `dsh@/b` 是两个不同的工作区,共用一份 +// payload 会让抄送方在别人的目录里开会话。`reply_address` / `self_address` +// 同理,且 session 位已经把 `.new` 换成真实别名 —— `.new` 建完会话就失效了, +// 把原文那个 `x@/p.new` 送给参与方只会让它下一次又建一条新会话。 +// +// # 抄送方必须单独推 +// +// 漏掉的后果很隐蔽:邮件的 cc_list 里有他们、他们**查**收件箱能看到这封信, +// 但没有任何事件推给他们 —— 插件不会唤起会话,Agent 直到下一次补拉 +// (重启时)才发现。对「知情方」而言等于没通知。 +func Recipients(ctx context.Context, m Mail) { + // 别名此时应已由会话解析路径保证存在(`.new` 与默认会话都过 + // EnsureSessionAlias)。仍可能为空的情形:命名写入失败(已吐日志)。 + // 此时退回省略 session 位,而不是把 "new" 写进去 —— 后者会让参与方 + // 反复建新会话。 + alias := repo.SessionAliasOf(ctx, m.SessionID) + + // 这条会话是否接管了一条平台侧已存在的会话(人在 TUI/GUI 里开的那种)。 + // 插件据此决定 resume 还是新建;空串就是过去的行为。 + platformID := repo.PlatformIDOf(ctx, m.SessionID) + + mailType := m.MailType + if mailType == "" { + mailType = "normal" + } + replyTo := m.ReplyToName + if replyTo == "" { + replyTo = m.From + } + + payload := func(role, workspace, forName string) map[string]interface{} { + p := map[string]interface{}{ + "mail_id": m.MailID.String(), + "session_id": m.SessionID.String(), + "from_name": m.From, + "subject": m.Subject, + "mail_type": mailType, + "role": role, // to / cc + // to_workspace 是收件方地址的 path 位,即希望它在哪个工作目录干活。 + // 不带这一项的后果:插件只能自己拼一个临时目录,于是每封邮件都落在 + // 不同的空目录里,DSH / opencode 按 cwd 分组时全进「未分组」。 + "to_workspace": workspace, + // session_alias 是这条会话今后的寻址名。没有它的话,收到 `.new` + // 邮件的一方只持有一个 send_mail 不接受的 session_id。 + "session_alias": alias, + // reply_address 是「把回信发回这条会话」的现成地址。 + // 插件不必自己拼(拼错了就是静默开新会话)。 + "reply_address": models.FormatAddress(replyTo, "", alias), + // self_address 是对方应当用来称呼自己的地址,供转发/报告时引用。 + "self_address": models.FormatAddress(forName, workspace, alias), + // platform_session_id 非空时,这封邮件要投进**平台侧已经存在的 + // 那条会话**(TUI 与邮箱是同一个 Agent 的两个入口)。 + // + // 插件必须 resume 而不是新建:新建会让人在 TUI 里看不到这封邮件 + // 带来的对话,而那正是接管这条会话的目的。 + "platform_session_id": platformID, + } + if m.Origin != "" { + p["origin"] = m.Origin + } + return p + } + + update := map[string]interface{}{ + "session_id": m.SessionID.String(), + "status": "active", + } + + // 参与方去重:收件人 + 所有抄送 + 发件方自己(刷新他的发件箱) + seen := map[string]bool{} + + sse.Default.SendToRecipient(m.To.Name, "new_mail", payload("to", m.To.Path, m.To.Name)) + sse.Default.SendToRecipient(m.To.Name, "session_update", update) + seen[m.To.Name] = true + + for _, c := range m.CC { + if seen[c.Name] { + continue + } + seen[c.Name] = true + sse.Default.SendToRecipient(c.Name, "new_mail", payload("cc", c.Path, c.Name)) + sse.Default.SendToRecipient(c.Name, "session_update", update) + } + + if !seen[m.From] { + sse.Default.SendToRecipient(m.From, "session_update", update) + } +} + +// SessionActive 只刷新某一方的会话列表,不推 new_mail。 +// +// 用于「日历事件的创建者该知道提醒发出去了」这类场景:他不是收件方, +// 不该收到一封信的通知,但需要看到那条会话活跃起来 —— 否则 +// 「我设的提醒到底触发了没有」只能去翻 journalctl。 +func SessionActive(name string, sessionID uuid.UUID) { + if name == "" { + return + } + sse.Default.SendToRecipient(name, "session_update", map[string]interface{}{ + "session_id": sessionID.String(), + "status": "active", + }) +} diff --git a/gateway/internal/repo/adopt_alias_test.go b/gateway/internal/repo/adopt_alias_test.go new file mode 100644 index 0000000..780a579 --- /dev/null +++ b/gateway/internal/repo/adopt_alias_test.go @@ -0,0 +1,274 @@ +package repo + +import ( + "context" + "testing" + "time" + + "github.com/agentmail/gateway/internal/db" +) + +// ─── 接管会话的别名不可被平台命名同步覆盖 ─── +// +// 锁的是一次生产事故的**第二环**(第一环是调度器漏 platform_session_id): +// +// 12:12 人选中补全里的「项目定位」→ 接管平台会话 01a05a5e,本侧建 26e26477 +// 12:20 日历提醒省略 session 位 → 落进 26e26477(第三环,见 TestDefaultSession...) +// 12:20 插件收不到 platform_session_id → 另开一条 pi 会话 +// 12:20 那条新会话的名字经 /sessions/{id}/sync 回写 +// → SyncSessionAlias 把 26e26477 的别名冲成「日程提醒:…」 +// +// 结果:人在补全里选的名字凭空消失,同一条会话在候选列表里出现两次 +// (一次用被冲掉的别名、一次用镜像里的原始 slug),而另一条真实会话被 +// 按别名字符串去重吃掉了。 +// +// 接管会话的别名是**人从补全里选中的平台 slug**,任何平台命名同步都不该动它。 +func TestSyncSessionAliasNeverOverwritesAdoptedAlias(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedAgent(t, "pi", 20) + + id, err := AdoptPlatformSession(ctx, "pi", "01a05a5e", "项目定位", "/home/program/agentmail", "标题") + if err != nil { + t.Fatalf("接管: %v", err) + } + if got := SessionAliasOf(ctx, id); got != "项目定位" { + t.Fatalf("接管后别名 = %q,期望 项目定位", got) + } + + // 平台侧同步一个完全不同的名字(生产上就是日历提醒的主题) + final, err := SyncSessionAlias(ctx, id, "日程提醒:小宅自测") + if err != nil { + t.Fatalf("SyncSessionAlias: %v", err) + } + if final != "项目定位" { + t.Errorf("同步返回 %q —— 接管会话的别名不该被改", final) + } + if got := SessionAliasOf(ctx, id); got != "项目定位" { + t.Errorf("库里别名变成了 %q —— 人在补全里选的名字被冲掉了", got) + } +} + +// 普通(非接管)会话仍然接受平台命名同步 —— 别名复用平台命名是既定决策, +// 上面那道门不能把它一起关掉。 +func TestSyncSessionAliasStillWorksForNormalSession(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedAgent(t, "pi", 20) + + id, err := CreateSession(ctx, nil, "pi", "邮件驱动的会话", "/tmp/ws") + if err != nil { + t.Fatalf("建会话: %v", err) + } + if _, err := EnsureSessionAlias(ctx, id, "pi-初始别名"); err != nil { + t.Fatalf("EnsureSessionAlias: %v", err) + } + + final, err := SyncSessionAlias(ctx, id, "平台生成的名字") + if err != nil { + t.Fatalf("SyncSessionAlias: %v", err) + } + if final != "平台生成的名字" { + t.Errorf("普通会话应当接受同步,得到 %q", final) + } +} + +// 接管会话**没有**别名时(理论上不会发生,AdoptPlatformSession 一定给一个) +// 仍然允许写入 —— 否则那条会话永远无法寻址。 +func TestSyncSessionAliasFillsEmptyAdoptedAlias(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedAgent(t, "pi", 20) + + id, err := CreateSession(ctx, nil, "pi", "标题", "/tmp/ws") + if err != nil { + t.Fatalf("建会话: %v", err) + } + // 手工造出「有 platform_id 但无别名」的状态 + if _, err := db.DB.ExecContext(ctx, + `UPDATE sessions SET platform_id = 'pid-x' WHERE session_id = $1`, id); err != nil { + t.Fatalf("置 platform_id: %v", err) + } + + final, err := SyncSessionAlias(ctx, id, "补上一个名字") + if err != nil { + t.Fatalf("SyncSessionAlias: %v", err) + } + if final != "补上一个名字" { + t.Errorf("无别名的接管会话应当允许写入,得到 %q", final) + } +} + +// ─── 接管会话不是「默认会话」 ─── +// +// 事故的**第三环**:日历提醒的收件地址省略 session 位(`homeagent` 而不是 +// `homeagent@/x.某会话`),走 FindOrCreateDefaultSession。它原来只按 +// 「参与过 + workspace 匹配 + 未归档」挑最近活跃的一条 —— 于是挑中了人 +// 刚刚显式指定的那条接管会话。 +// +// 接管会话是人**点名**要谈的一条线索,不该被省略 session 位的邮件当默认会话。 +func TestFindOrCreateDefaultSessionSkipsAdoptedSessions(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedAgent(t, "pi", 20) + + // 一条接管会话,且有邮件(满足 EXISTS 条件) + adopted, err := AdoptPlatformSession(ctx, "pi", "pid-adopted", "人选的线索", "/home/program/agentmail", "标题") + if err != nil { + t.Fatalf("接管: %v", err) + } + if _, err := CreateMail(ctx, adopted, nil, "jianf", "", "pi", "/home/program/agentmail", + "人发的第一封", "正文", nil); err != nil { + t.Fatalf("建邮件: %v", err) + } + + // 省略 session 位投递 → 不该落进那条接管会话 + got, err := FindOrCreateDefaultSession(ctx, "pi", "/home/program/agentmail", "calendar", "日程提醒") + if err != nil { + t.Fatalf("FindOrCreateDefaultSession: %v", err) + } + if got == adopted { + t.Error("省略 session 位的邮件落进了接管会话 —— 那是人显式指定的线索") + } + + // 该新建一条,且它不带 platform_id + if pid := PlatformIDOf(ctx, got); pid != "" { + t.Errorf("新建的默认会话不该有 platform_id,得到 %q", pid) + } +} + +// 普通会话仍然可以作为默认会话被复用 —— 上面那道门不能把它一起关掉, +// 否则每封省略 session 位的邮件都会新开一条会话。 +func TestFindOrCreateDefaultSessionStillReusesNormalSession(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedAgent(t, "pi", 20) + + first, err := FindOrCreateDefaultSession(ctx, "pi", "/tmp/ws", "jianf", "第一封") + if err != nil { + t.Fatalf("第一次: %v", err) + } + if _, err := CreateMail(ctx, first, nil, "jianf", "", "pi", "/tmp/ws", "第一封", "正文", nil); err != nil { + t.Fatalf("建邮件: %v", err) + } + + second, err := FindOrCreateDefaultSession(ctx, "pi", "/tmp/ws", "jianf", "第二封") + if err != nil { + t.Fatalf("第二次: %v", err) + } + if second != first { + t.Error("普通默认会话应当被复用,否则每封省略 session 位的邮件都开新会话") + } +} + +// ─── 候选列表按 platform_id 去重 ─── +// +// 事故的**第四环**:`SuggestSessionCandidates` 的去重只比别名字符串。 +// 别名一被冲掉,同一条会话就在列表里出现两次: +// +// 候选 1 日程提醒:…(被冲掉的别名) source=mail ← 26e26477 +// 候选 2 项目定位(镜像里的原始 slug) source=platform ← 也是 26e26477 +// +// 更糟的是**另一条真实会话被吃掉了**:它的 slug 恰好等于候选 1 那个 +// 被冲掉的别名,于是 `seen[slug]` 命中、被 continue 跳过。 +// 人在界面上看到两条,实际只有一条能选,而第三条不存在于列表里。 +func TestSuggestSessionCandidatesDedupesByPlatformID(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedPlatformAgent(t, "pi") + + // 接管一条平台会话 + adopted, err := AdoptPlatformSession(ctx, "pi", "01a05a5e", "项目定位", "/home/program/agentmail", "标题") + if err != nil { + t.Fatalf("接管: %v", err) + } + if _, err := CreateMail(ctx, adopted, nil, "jianf", "", "pi", "/home/program/agentmail", + "主题", "正文", nil); err != nil { + t.Fatalf("建邮件: %v", err) + } + + // 镜像里同时有它与另一条真实会话 + now := time.Now() + if err := ReplacePlatformSessions(ctx, "pi", []PlatformSession{ + {PlatformID: "01a05a5e", Workspace: "/home/program/agentmail", Slug: "项目定位", + Title: "项目定位", MailDriven: true, UpdatedAt: &now}, + {PlatformID: "01a06aa5", Workspace: "/home/program/agentmail", Slug: "另一条真实会话", + Title: "另一条", MailDriven: true, UpdatedAt: &now}, + }); err != nil { + t.Fatalf("上报镜像: %v", err) + } + + got, err := SuggestSessionCandidates(ctx, "jianf", "pi", "/home/program/agentmail") + if err != nil { + t.Fatalf("SuggestSessionCandidates: %v", err) + } + + // 应当恰好两条:接管那条(mail 来源)+ 另一条真实会话(platform 来源) + if len(got) != 2 { + t.Fatalf("应有 2 个候选,实际 %d:%+v", len(got), got) + } + + byAlias := map[string]SessionCandidate{} + for _, c := range got { + byAlias[c.Alias] = c + } + if c, ok := byAlias["项目定位"]; !ok { + t.Error("接管会话应当在候选里") + } else if c.Source != "mail" { + t.Errorf("接管会话的来源应是 mail(保证送得到),得到 %q", c.Source) + } + if c, ok := byAlias["另一条真实会话"]; !ok { + t.Error("另一条真实会话被吃掉了 —— 那正是 bug 的表现") + } else if c.Source != "platform" { + t.Errorf("未接管的平台会话来源应是 platform,得到 %q", c.Source) + } +} + +// 别名被冲掉之后也不该出现重复项。 +// +// 这是事故现场的**精确复现**:本侧别名与镜像 slug 不一致(别名被另一条会话的 +// 命名同步冲掉了),此时按别名字符串去重必然漏,只有按 platform_id 才对。 +func TestSuggestSessionCandidatesNoDupWhenAliasDiverged(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + seedPlatformAgent(t, "pi") + + adopted, err := AdoptPlatformSession(ctx, "pi", "01a05a5e", "项目定位", "/home/program/agentmail", "标题") + if err != nil { + t.Fatalf("接管: %v", err) + } + if _, err := CreateMail(ctx, adopted, nil, "jianf", "", "pi", "/home/program/agentmail", + "主题", "正文", nil); err != nil { + t.Fatalf("建邮件: %v", err) + } + // 模拟别名被冲掉(绕过 SyncSessionAlias 的守卫直接改库 —— + // 存量数据里可能已经有这种状态) + if _, err := db.DB.ExecContext(ctx, + `UPDATE sessions SET session_alias = $1 WHERE session_id = $2`, + "日程提醒:小宅自测", adopted); err != nil { + t.Fatalf("改别名: %v", err) + } + + now := time.Now() + if err := ReplacePlatformSessions(ctx, "pi", []PlatformSession{ + {PlatformID: "01a05a5e", Workspace: "/home/program/agentmail", Slug: "项目定位", + Title: "项目定位", MailDriven: true, UpdatedAt: &now}, + }); err != nil { + t.Fatalf("上报镜像: %v", err) + } + + got, err := SuggestSessionCandidates(ctx, "jianf", "pi", "/home/program/agentmail") + if err != nil { + t.Fatalf("SuggestSessionCandidates: %v", err) + } + + // 只有一条会话,就该只有一个候选 —— 别名分叉不该让它变成两个 + if len(got) != 1 { + t.Fatalf("同一条会话应只有 1 个候选,实际 %d:%+v", len(got), got) + } + if got[0].Source != "mail" { + t.Errorf("应保留 mail 来源(它保证送得到),得到 %q", got[0].Source) + } +} + +// 存量数据里可能已经有这种状态(守卫是后加的)。 diff --git a/gateway/internal/repo/platform_sessions.go b/gateway/internal/repo/platform_sessions.go index fbca293..5259d80 100644 --- a/gateway/internal/repo/platform_sessions.go +++ b/gateway/internal/repo/platform_sessions.go @@ -121,6 +121,7 @@ func SuggestSessionCandidates(ctx context.Context, forUser, peerName, path strin rows, err := db.DB.QueryContext(ctx, ` SELECT s.session_alias, COALESCE(s.subject, ''), + COALESCE(s.platform_id, ''), (SELECT COUNT(*) FROM mails u WHERE u.session_id = s.session_id AND u.status = 'unread') FROM sessions s @@ -153,9 +154,9 @@ func SuggestSessionCandidates(ctx context.Context, forUser, peerName, path strin defer rows.Close() for rows.Next() { - var alias, title string + var alias, title, pid string var unread int - if err := rows.Scan(&alias, &title, &unread); err != nil { + if err := rows.Scan(&alias, &title, &pid, &unread); err != nil { return out, err } if alias == "" { @@ -165,6 +166,9 @@ func SuggestSessionCandidates(ctx context.Context, forUser, peerName, path strin out = append(out, SessionCandidate{ Alias: alias, Title: title, Source: "mail", Unread: unread, }) + if pid != "" { + seen["pid:"+pid] = len(out) - 1 + } } if err := rows.Err(); err != nil { return out, err @@ -172,7 +176,7 @@ func SuggestSessionCandidates(ctx context.Context, forUser, peerName, path strin // ---- 来源 2:平台会话镜像 ---- prows, err := db.DB.QueryContext(ctx, ` - SELECT slug, title + SELECT slug, title, platform_id FROM agent_platform_sessions WHERE agent_name = $1 AND slug <> '' @@ -188,13 +192,20 @@ func SuggestSessionCandidates(ctx context.Context, forUser, peerName, path strin defer prows.Close() for prows.Next() { - var slug, title string - if err := prows.Scan(&slug, &title); err != nil { + var slug, title, pid string + if err := prows.Scan(&slug, &title, &pid); err != nil { break } if slug == "" { continue } + // 已被接管的平台会话不再单独列:选它也会落进已有的那条本侧线索, + // 但候选列表出现两次会让人以为有两条不同的会话(项目定位 x2 的场景)。 + if pid != "" { + if _, dup := seen["pid:"+pid]; dup { + continue + } + } if i, ok := seen[slug]; ok { // 本侧已有同名线索:保留 mail 来源(它保证送得到), // 但补上镜像的标题 —— 平台侧标题通常比会话建立时的主题更贴切 diff --git a/gateway/internal/repo/repo.go b/gateway/internal/repo/repo.go index 84eee2f..99d0fd6 100644 --- a/gateway/internal/repo/repo.go +++ b/gateway/internal/repo/repo.go @@ -306,12 +306,23 @@ func CreateMail(ctx context.Context, sessionID uuid.UUID, parentMailID *uuid.UUI } // 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) (uuid.UUID, error) { optsJSON, _ := json.Marshal(options) var id uuid.UUID err := db.DB.QueryRowContext(ctx, - `INSERT INTO mails (session_id, from_name, to_name, subject, body, mail_type, permission_options, created_at) - VALUES ($1, $2, $3, $4, $5, 'permission_request', $6, NOW()) RETURNING mail_id`, + `INSERT INTO mails (session_id, from_name, from_workspace, to_name, subject, body, mail_type, permission_options, created_at) + VALUES ($1, $2, COALESCE((SELECT workspace FROM sessions WHERE session_id = $1), ''), $3, $4, $5, 'permission_request', $6, NOW()) RETURNING mail_id`, sessionID, fromName, toUser, "权限请求: "+question, body, optsJSON, ).Scan(&id) return id, err @@ -647,6 +658,10 @@ func FindOrCreateDefaultSession(ctx context.Context, name, path, fromAgent, subj 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 @@ -719,14 +734,25 @@ func SyncSessionAlias(ctx context.Context, id uuid.UUID, want string) (string, e var cur *string var source string + var platformID string if err := db.DB.QueryRowContext(ctx, - `SELECT session_alias, COALESCE(alias_source, 'platform') FROM sessions WHERE session_id = $1`, - id).Scan(&cur, &source); err != nil { + `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 diff --git a/gateway/internal/scheduler/calendar.go b/gateway/internal/scheduler/calendar.go index 9d760aa..298922f 100644 --- a/gateway/internal/scheduler/calendar.go +++ b/gateway/internal/scheduler/calendar.go @@ -11,8 +11,8 @@ import ( "github.com/agentmail/gateway/internal/db" "github.com/agentmail/gateway/internal/models" + "github.com/agentmail/gateway/internal/notify" "github.com/agentmail/gateway/internal/repo" - "github.com/agentmail/gateway/internal/sse" "github.com/google/uuid" ) @@ -279,62 +279,33 @@ func SendCalendarMail(ctx context.Context, eventID, toAddr, subject, body, creat } } - alias := repo.SessionAliasOf(ctx, sessionID) - sse.Default.SendToRecipient(addr.Name, "new_mail", map[string]interface{}{ - "mail_id": mailID.String(), - "session_id": sessionID.String(), - "from_name": calendarSender, - "subject": subject, - "mail_type": "normal", - "role": "to", - "to_workspace": addr.Path, - "session_alias": alias, - // 回信地址给 calendar 是发不出去的(它不是收件方), - // 给会话自己的地址才让模型能把结果回报到同一条线索上。 - "reply_address": models.FormatAddress(addr.Name, addr.Path, alias), - "self_address": models.FormatAddress(addr.Name, addr.Path, alias), - // 让插件与 UI 能区分「这封是定时提醒」而不是有人在找它 - "origin": "calendar", - }) - sse.Default.SendToRecipient(addr.Name, "session_update", map[string]interface{}{ - "session_id": sessionID.String(), - "status": "active", - }) - - // 抄送方也要收到 SSE。 + // ─── 推送 SSE ─── // - // 漏了这一步的后果很隐蔽:邮件的 cc_list 里有他们、他们**查**收件箱 - // 能看到这封信,但没有任何事件推给他们 —— 于是插件不会唤起会话, - // Agent 直到下一次补拉(重启时)才发现。对「知情方」而言等于没通知。 - for _, c := range ccList { - alias := repo.SessionAliasOf(ctx, sessionID) - sse.Default.SendToRecipient(c.Name, "new_mail", map[string]interface{}{ - "mail_id": mailID.String(), - "session_id": sessionID.String(), - "from_name": calendarSender, - "subject": subject, - "mail_type": "normal", - "role": "cc", - "to_workspace": c.Path, - "session_alias": alias, - "reply_address": models.FormatAddress(addr.Name, addr.Path, alias), - "self_address": models.FormatAddress(c.Name, c.Path, alias), - "origin": "calendar", - }) - sse.Default.SendToRecipient(c.Name, "session_update", map[string]interface{}{ - "session_id": sessionID.String(), - "status": "active", - }) - } + // 这是整个系统里**唯一**的推送实现(handler 那条路径也是 notify.Recipients)。 + // 原来这里自己拼了一整份 payload(handler 里是另一份), + // 加 `platform_session_id` 时只改了那边 → 日历提醒投进接管会话时 + // 插件不知道是接管、另开了一条新 pi 会话 → 命名同步把接管会话的别名冲掉 + // → 人在补全里选的「项目定位」变成了「日程提醒:…」→ 选哪条都落进同一条。 + // 根因只是「同一件事写了两遍」。 + // + // ReplyToName = addr.Name:日历提醒的回信要落回那条线索,不是回给 calendar + // Origin = "calendar":让插件与 UI 能区分「这封是定时提醒」 + notify.Recipients(ctx, notify.Mail{ + SessionID: sessionID, + MailID: mailID, + From: calendarSender, + To: addr, + CC: ccList, + Subject: subject, + MailType: "normal", + Origin: "calendar", + ReplyToName: addr.Name, + }) - // 创建者也该看到提醒发出去了 —— 否则「我设的提醒到底触发了没有」 - // 只能去翻 journalctl。 - if createdBy != "" && createdBy != addr.Name { - sse.Default.SendToRecipient(createdBy, "session_update", map[string]interface{}{ - "session_id": sessionID.String(), - "status": "active", - }) - } + // 创建者不是收件方,不在 Recipients 的参与方去重里 —— 他不会收到 + // new_mail(他不该被「有人在找你」打扰),但他应该看到会话活跃起来: + // 「我设的提醒到底触发了没有」不应该只能去翻 journalctl。 + notify.SessionActive(createdBy, sessionID) return nil } diff --git a/plugins/homeagent-mail-bridge/plugin.go b/plugins/homeagent-mail-bridge/plugin.go index 28365c4..3c87ec3 100644 --- a/plugins/homeagent-mail-bridge/plugin.go +++ b/plugins/homeagent-mail-bridge/plugin.go @@ -1,6 +1,7 @@ package main import ( + "bufio" "bytes" "encoding/json" "fmt" @@ -73,6 +74,22 @@ type Plugin struct { // B-1.6 补拉状态:首个成功心跳后只补一次 catchupDone bool + + // ─── SSE 专用 ─── + + // SSE 需要一个不设 Timeout 的 HTTP client:原来 p.client(60s Timeout) + // 跑 SSE 长连接,每 60 秒自己掐断自己。之后 lastEventID 回退 → 重放 → + // 又阻塞 → 又超时 —— 自激振荡。这个 client 只给 readSSE 用。 + sseClient *http.Client + + // B-7.3 邮件级去重:SSE 重放会重发同一批事件,没有这层去重 + // 每封邮件会被注入 agent 两遍。契约 B-7.3 要求:每封只注入一次。 + deliveredMails map[string]bool + + // 单调递增的 last-seen-ID:被重放的旧事件不会让它回退。 + // 原来直接赋值(p.lastEventID = eid),Gateway 重放时发旧 ID, + // 于是 lastEventID 从 123 退回 116 → 下次重连又报 116 → 又重放。 + sseMaxID int64 } // ─── B-1.1 密钥解析与本地生成 ─── @@ -146,14 +163,16 @@ func NewPluginFactory(name string, config map[string]interface{}) (sdk.Plugin, e } return &Plugin{ - name: name, - agentName: agentName, - gwURL: strings.TrimRight(gw, "/"), - key: "", // Start() 里解析 - keyFile: "", - client: &http.Client{Timeout: 60 * time.Second}, - stopCh: make(chan struct{}), - explicitSends: make(map[string]time.Time), + name: name, + agentName: agentName, + gwURL: strings.TrimRight(gw, "/"), + key: "", // Start() 里解析 + keyFile: "", + client: &http.Client{Timeout: 60 * time.Second}, + sseClient: &http.Client{}, // 无超时:SSE 是长连接 + stopCh: make(chan struct{}), + deliveredMails: make(map[string]bool), + explicitSends: make(map[string]time.Time), }, nil } @@ -550,14 +569,19 @@ func (p *Plugin) readSSE() error { } req.Header.Set("Authorization", "Bearer "+p.key) - // W-4:断线期间的事件会丢,带上 Last-Event-ID 可以让 Gateway 从断点补发 + // W-4:断线期间的事件会丢,带上 Last-Event-ID 可以让 Gateway 从断点补发。 + // 只上报比当前记录的更大的 ID:Gateway 重放时发的是事件原本的 ID, + // 如果无条件赋值,lastEventID 会从 123 退回 116 → 下次重连又报 116 → + // 又重放 —— 自激振荡的放大器。 p.sseMu.Lock() if p.lastEventID != "" { req.Header.Set("Last-Event-ID", p.lastEventID) } p.sseMu.Unlock() - resp, err := p.client.Do(req) + // sseClient 无 Timeout:p.client 有 60s Timeout,SSE 是长连接, + // 每 60 秒自己掐断自己 → 重放 → 阻塞 → 超时 → 重放。 + resp, err := p.sseClient.Do(req) if err != nil { return err } @@ -569,8 +593,11 @@ func (p *Plugin) readSSE() error { log.Printf("[homeagent-mail-bridge] SSE 已连接") - buf := make([]byte, 0, 4096) - lineStart := 0 + // bufio.Reader 解决原来手动管理 []byte 的两个问题: + // 1. 每次 buf = buf[lineStart:] 让 cap 缩小,几轮之后 len==cap, + // Read 拿到零长切片 → (0, nil) → 满速空转 + // 2. 手写的 line 分割逻辑有边界条件(跨次 Read 的半行处理) + br := bufio.NewReader(resp.Body) for { select { case <-p.stopCh: @@ -578,39 +605,43 @@ func (p *Plugin) readSSE() error { default: } - n, err := resp.Body.Read(buf[len(buf):cap(buf)]) - if n > 0 { - buf = buf[:len(buf)+n] - for { - i := bytes.IndexByte(buf[lineStart:], '\n') - if i < 0 { - break - } - line := string(buf[lineStart : lineStart+i]) - lineStart += i + 1 - p.parseSSELine(line) - } - if lineStart > 0 { - buf = buf[lineStart:] - lineStart = 0 - } - buf = buf[:len(buf)] + line, err := br.ReadString('\n') + if line != "" { + p.parseSSELine(strings.TrimRight(line, "\n")) } if err != nil { - if err != io.EOF { - return err + if err == io.EOF { + return nil } - return nil + return err } } } +// parseSSEID 把 "id: 123" 格式的事件 ID 解析成整数。 +// 解析失败返回 0(大于 0 的 ID 才会被接受),保证不会误清状态。 +func parseSSEID(raw string) int64 { + var id int64 + for _, c := range raw { + if c >= '0' && c <= '9' { + id = id*10 + int64(c-'0') + } + } + return id +} + func (p *Plugin) parseSSELine(line string) { - // W-4:记录 Last-Event-ID + // W-4:记录 Last-Event-ID。只向前推进,不回退。 + // Gateway 重放旧事件时发的是旧 ID,无条件赋值会让 lastEventID + // 从 123 退回到 116 → 下次重连报 116 → 又重放 → 振荡。 if strings.HasPrefix(line, "id: ") { eid := strings.TrimPrefix(line, "id: ") + id := parseSSEID(eid) p.sseMu.Lock() - p.lastEventID = eid + if id > p.sseMaxID { + p.sseMaxID = id + p.lastEventID = eid + } p.sseMu.Unlock() return } @@ -647,7 +678,22 @@ func (p *Plugin) parseSSELine(line string) { } if evt.MailType == "normal" { - p.handleNewMail(evt) + // B-7.3:去重。SSE 重放时同一封邮件会再出现,没有这层 + // 每封邮件会被注入 agent 两遍(实测 21 次超时 → 21 次重放)。 + p.sseMu.Lock() + if p.deliveredMails[evt.MailID] { + p.sseMu.Unlock() + return + } + p.deliveredMails[evt.MailID] = true + p.sseMu.Unlock() + + // InjectInputSync 会阻塞几十秒(查日志、调工具、转发 QQ), + // 而它跑在 readSSE 的读循环里 —— 循环卡住期间 SSE 事件积压在 + // TCP 缓冲区,卡到超时断线重连后 Gateway 全部重放一遍。 + // 把处理丢到独立 goroutine:parseSSELine 立刻返回,读循环继续。 + // homeagent 是单事件循环,InjectInputSync 自己会排队。 + go p.handleNewMail(evt) } } diff --git a/web/src/components/NarrowStack.tsx b/web/src/components/NarrowStack.tsx index 37b049c..df47710 100644 --- a/web/src/components/NarrowStack.tsx +++ b/web/src/components/NarrowStack.tsx @@ -58,14 +58,23 @@ export default function NarrowStack({ return (
{/* 底层:始终挂载。打开覆盖层时用 aria-hidden 把它从无障碍树里摘掉, - 否则屏幕阅读器会读到两层内容 */} -
+ 否则屏幕阅读器会读到两层内容。 + + `isolate`(isolation: isolate)是必需的:它让底层**自成一个层叠上下文**。 + 不加的后果在日历上实测到过:月视图的星期表头是 `sticky top-0 z-10`, + 而覆盖层没有 z-index(= auto = 0)—— 两者在同一个层叠上下文里比, + `z-10` 赢过 `auto`,于是底层的表头穿透到二级页面之上,把日程内容遮住一条。 + + 为何不只给覆盖层加 z-10 就完事:那只能治当下这一处。底层是任意业务组件, + 下一个人在里面写个 `z-20` 就又复现,而这类 bug 只能肉眼看见。 + isolate 把边界定在容器上,底层写多少 z-index 都出不来。 */} +
{base}
{mounted && (
diff --git a/web/src/components/PermissionList.tsx b/web/src/components/PermissionList.tsx index 17e9c34..30cd0fa 100644 --- a/web/src/components/PermissionList.tsx +++ b/web/src/components/PermissionList.tsx @@ -132,10 +132,9 @@ function PermissionSessionGroup({ }`} /> - {g.agentName} - {g.path && ( - {g.path} - )} + + {g.agentName}{g.path ? `@${g.path}` : ''}{g.alias ? `.${g.alias}` : ''} + {time}
@@ -156,9 +155,11 @@ function PermissionSessionGroup({ )}
-

- {g.alias ? `.${g.alias}` : '(未命名会话)'} -

+ {!g.alias && ( +

+ (未命名会话) +

+ )} {open && ( diff --git a/web/src/lib/mailGroups.ts b/web/src/lib/mailGroups.ts index 9a903a8..b286c21 100644 --- a/web/src/lib/mailGroups.ts +++ b/web/src/lib/mailGroups.ts @@ -160,7 +160,7 @@ export function groupPermissions(mails: Mail[]): PermissionGroup[] { sessionId, alias: latest.session_alias || '', agentName: latest.from_name, - path: latest.from_workspace || '', + path: latest.session_workspace || '', pending: sorted.filter(isPendingPermission), settled: sorted.filter(m => !isPendingPermission(m)), latest diff --git a/web/test/components/mailGroups.test.tsx b/web/test/components/mailGroups.test.tsx index ad8f832..0e799f5 100644 --- a/web/test/components/mailGroups.test.tsx +++ b/web/test/components/mailGroups.test.tsx @@ -254,12 +254,22 @@ describe('groupPermissions', () => { expect(groups.map(g => g.sessionId)).toEqual(['new', 'old']); }); - it('组头带上发起请求的 Agent 与工作目录', () => { + it('组头带上发起请求的 Agent 与会话工作目录', () => { const [g] = groupPermissions([ - perm({ session_id: 's', from_name: 'dsh', from_workspace: '/home/program/llmsproxy' }) + perm({ + session_id: 's', + from_name: 'dsh', + // from_workspace 对 Agent 存的是 **Agent 名**而不是路径(历史遗留)。 + // 拿它当路径用会在授权页拼出 `dsh@dsh`,而权限请求的 + // from_workspace 实测是**空串** —— 于是那一行永远不渲染, + // 人根本不知道是哪个目录里的哪条线索在请求权限。 + from_workspace: 'dsh', + session_workspace: '/home/program/llmsproxy' + }) ]); // 同名 Agent 在不同目录是不同的活,光有名字判断不了 expect(g.agentName).toBe('dsh'); + // path 必须取 session_workspace(会话的 workspace,权威来源) expect(g.path).toBe('/home/program/llmsproxy'); });