Files
MailUI4Agents/plugins/homeagent-mail-bridge/plugin.go
JianFeeeee e8583ecd41 fix: 事故全链路修复 — 调度器合流 + 接管保护 + 补全去重 + 权限显示 + homeagent SSE 振荡
## 事故现场

用户选中补全里的「项目定位」→ 邮件投进另一条会话,界面显示的名字也不是
自己选的那个。授权页只显示 Agent 名,看不出哪个目录哪条线索。

## 四处因果链

**① 调度器自己的 new_mail payload(起点)。** `notifyRecipients`(handler)
与 `SendCalendarMail`(scheduler)是两份代码。加 `platform_session_id` 时只改了
handler 那份 → 日历提醒投进接管会话时插件不知道是接管 → 另开一条新会话 →
命名同步冲掉接管会话的别名。

修法:抽出 `internal/notify` 包,唯一入口 `notify.Recipients`。
handler / scheduler / permission.go 都走它。新增字段时不存在「另一处忘了改」。

**② SyncSessionAlias 覆盖接管别名。** 别名是人从补全里选中的平台 slug,
任何平台命名同步都不该动它。加守卫 `platform_id <> ''` → 有绑定就返回当前值。

**③ SuggestSessionCandidates 按别名字符串去重。** 别名一被冲掉,同一条会话
出现两次(一次被冲的名字、一次镜像 slug),而另一条真实会话被吃掉。
改按 `platform_id` 去重。 mail 侧查 `s.platform_id`,镜像侧查 `platform_id`。

**④ FindOrCreateDefaultSession 不排除接管会话。** 日历提醒省略 session 位 →
FindOrCreateDefaultSession 挑中人显式指定的接管会话。加 `platform_id = ''` 条件。

## 权限页

**CreatePermissionMail 不写 from_workspace。** `from_workspace` 存空串 →
前端 `g.path && ...` 不渲染 → 人只看到光秃的 Agent 名,不知道哪个目录
哪条线索在请求权限。修法:INSERT 时从 sessions.workspace 取。

**SSE payload 缺 session_alias。** permission.go 的 SSE 不走 notify 包(决策人
不是地址解析出的参与方),但 payload 也要带 `session_alias` → 前端拼出
`pi@/home/program/agentmail.别名`,而不是光秃的 `pi`。

**mailGroups.ts:path ← session_workspace。** `from_workspace` 对 Agent 存的是
Agent 名(历史遗留),不能当路径用。PermissionList 显示完整三段地址
`agent@path.alias`。

## NarrowStack z-index

窄屏日历的星期表头(`sticky top-0 z-10`)穿透到二级页面之上。覆盖层
auto z-index 输给 z-10 → 底层组件的层叠穿透到覆盖层。

修法:底层容器加 `isolate`(isolation: isolate),自成层叠上下文;
覆盖层加 `z-10`。只给覆盖层加 z-index 只能治当前一处,底层再写更大的
z-index 又会复现。

## homeagent SSE 自激振荡

根因:五处缺陷叠加,SSE 每 60 秒断一次 → Gateway 全量重放 → 再断 → 再重放。

1. `p.client`(60s Timeout)跑 SSE 长连接 → 新增 `sseClient`(无超时)
2. `InjectInputSync` 在读循环里同步调用 → 改为 `go p.handleNewMail(evt)`
3. `lastEventID` 无条件赋值,Gateway 重放时发旧 ID → 单调递增 `sseMaxID`
4. 无邮件级去重 → 补 `deliveredMails map[string]bool`
5. 手动 `[]byte` 管理:每次 `buf[lineStart:]` 缩小 cap → 最终 len==cap
   → Read 零长切片 → 满速空转。改 `bufio.Reader`。

## 清库

保留 jianf + 4 个 Agent 密钥 + 模型范围配置。清掉 mails/sessions/
calendar_events/attachments/agent_platform_sessions/relayed_mails/
permission_requests/rate_limits。测试数据已全部清零。

## 测试

- gateway 7 包全过;repo + 7 例(adopt_alias_test.go)
- web 182 例(mailGroups 新增 session_workspace 断言)
- 前端构建通过
2026-09-04 15:35:48 +08:00

1245 lines
42 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

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

package main
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"log"
"mime/multipart"
"net/http"
"os"
"path/filepath"
"strings"
"sync"
"time"
"gitcode.com/JianFeeeee/homeagent-sdk/sdk"
)
const (
defaultGateway = "http://127.0.0.1:8180"
heartbeatIntval = 30 * time.Second // 契约 B-2 要求 30 秒
sseRetry = 3 * time.Second
// B-5.3 explicitSends 记忆窗口:模型在一轮里发过的信,
// 在该窗口内不再自动 relay 同一封。超过窗口说明那轮已经结束。
dedupWindow = 10 * time.Minute
// B-1.6 补投上限(与 DSH/pi 保持一致)
catchupLimit = 5
// 连续 relay 跳数上限(与 Gateway 常量一致)
maxRelayHops = 5
)
// Plugin 实现 sdk.Plugin。
//
// TrueAgent 是单一常驻 agent 事件循环,没有「每个对话一个 session」的概念。
// 因此本插件不做会话映射 —— 所有邮件注入同一个事件循环,像 QQ 插件一样。
// 邮件的 context 完全靠中断消息的文本传递,不靠 platform_sessions 上报。
type Plugin struct {
// name 是**插件名**homed 注册用,如 homeagent-mail-bridge
name string
// agentName 是**AgentMail 身份**(如 homeagent
//
// 两者必须分开:密钥绑定的是 AgentMail 身份,拿插件名去注册会被拒
// 403 该密钥已绑定到 Agent "homeagent",不能用于注册 "homeagent-mail-bridge"
// 这不是 Gateway 太严格 —— 名字与人类用户名共用命名空间,
// 让一把密钥能注册任意名字等于让它能冒充任何人。
agentName string
sdk *sdk.PluginSDK
gwURL string
key string
keyFile string // B-1.1 密钥文件路径
client *http.Client
stopCh chan struct{}
stopOnce sync.Once
// W-4 SSE 重连 —— 断线期间的事件会丢,带上 Last-Event-ID 可以补回
lastEventID string
sseMu sync.Mutex // 保护 lastEventID
// B-5.3 explicitSends模型在当前 turn 里通过 send_mail/output_send 发过的
// 邮件 IDrelay_key 格式)。自动 relay 前查这个表,已有则让位。
//
// 为什么不是按 session 隔离TrueAgent 是单事件循环,所有邮件共享一个 turn。
// 模型如果调了 send_mail 回给发件人,那就是它自己的回复,不该再 relay。
explicitSends map[string]time.Time
explicitSendsMu sync.Mutex
// B-2.2 模型目录缓存
modelCatalog []string
// B-1.6 补拉状态:首个成功心跳后只补一次
catchupDone bool
// ─── SSE 专用 ───
// SSE 需要一个不设 Timeout 的 HTTP client原来 p.client60s 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 = eidGateway 重放时发旧 ID
// 于是 lastEventID 从 123 退回 116 → 下次重连又报 116 → 又重放。
sseMaxID int64
}
// ─── B-1.1 密钥解析与本地生成 ───
//
// 契约要求:环境变量 → ~/.agentmail/agent.key → 本地生成一把。
// 本地生成时打印到 stderr进 journalctl落盘到 key 文件0600
// 这样管理员拿到日志里的密钥全文去后台登记,下次重启就不再需要环境变量。
func resolveKey() (key, keyFile string) {
keyFile = os.Getenv("AGENTMAIL_CONFIG_DIR")
if keyFile == "" {
home, _ := os.UserHomeDir()
keyFile = filepath.Join(home, ".agentmail")
}
keyFile = filepath.Join(keyFile, "agent.key")
// 1. 环境变量
if k := strings.TrimSpace(os.Getenv("AGENTMAIL_AGENT_KEY")); k != "" {
return k, keyFile
}
// 2. 本地文件
if data, err := os.ReadFile(keyFile); err == nil {
k := strings.TrimSpace(string(data))
if k != "" {
return k, keyFile
}
}
// 3. 本地生成ak_ 前缀 + 24 字节 hex与 dsh 保持一致)
b := make([]byte, 24)
for i := range b {
b[i] = "0123456789abcdef"[time.Now().UnixNano()%16]
time.Sleep(1)
}
k := "ak_" + fmt.Sprintf("%x", b)
dir := filepath.Dir(keyFile)
os.MkdirAll(dir, 0700)
os.WriteFile(keyFile, []byte(k), 0600)
// 契约 9.8:密钥打印到 stderr 进 journalctl不走平台 logger
fmt.Fprintf(os.Stderr, "[homeagent-mail-bridge] 本地生成密钥,请让管理员在 AgentMail 后台「Agent 密钥」中登记:\n%s\n文件%s\n", k, keyFile)
return k, keyFile
}
// ─── 工厂 ───
func NewPluginFactory(name string, config map[string]interface{}) (sdk.Plugin, error) {
gw := ""
if v, ok := config["gateway_url"].(string); ok {
gw = v
}
if gw == "" {
gw = os.Getenv("AGENTMAIL_GATEWAY_URL")
}
if gw == "" {
gw = defaultGateway
}
// AgentMail 身份config > 环境变量 > 从插件名去掉 -mail-bridge 后缀
agentName := ""
if v, ok := config["agent_name"].(string); ok {
agentName = strings.TrimSpace(v)
}
if agentName == "" {
agentName = strings.TrimSpace(os.Getenv("AGENTMAIL_AGENT_NAME"))
}
if agentName == "" {
agentName = strings.TrimSuffix(name, "-mail-bridge")
}
return &Plugin{
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
}
func (p *Plugin) Name() string { return p.name }
func (p *Plugin) Start(s *sdk.PluginSDK) error {
p.sdk = s
s.SetAutoRestart(true)
// B-1.1 解析密钥(环境变量 → 本地文件 → 生成)
p.key, p.keyFile = resolveKey()
// ─── 注册工具 ───
// registerTool 包一层只为计数:日志里的工具数必须与实际注册数一致。
registeredToolCount := 0
registerTool := func(name string, def sdk.ToolDef, h func(map[string]interface{}) (interface{}, error)) {
registeredToolCount++
s.RegisterTool(name, def, h)
}
registerTool("read_inbox", sdk.ToolDef{
Name: "read_inbox",
Description: "查阅收件箱中的邮件。收到新邮件通知后应立即调用此工具。每封含 mail_id、发件人、主题、正文与附件清单。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"status": map[string]interface{}{"type": "string", "description": "过滤条件 unread|all默认 unread"},
"limit": map[string]interface{}{"type": "number", "description": "返回数量,默认 5"},
},
},
}, p.handleReadInbox)
registerTool("read_mail", sdk.ToolDef{
Name: "read_mail",
Description: "读一封邮件的完整内容,含收件人、抄送清单、附件与每个参与方的可投递地址。",
Parameters: oneStringParam("mail_id", "邮件 ID", true),
}, p.handleReadMail)
registerTool("send_mail", sdk.ToolDef{
Name: "send_mail",
Description: "发送邮件。三维地址 name@path.session。回复来信请传 reply_to。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"to": map[string]interface{}{"type": "string", "description": "收件人三维地址"},
"subject": map[string]interface{}{"type": "string", "description": "邮件主题"},
"body": map[string]interface{}{"type": "string", "description": "邮件正文Markdown"},
"cc": map[string]interface{}{"type": "string", "description": "抄送"},
"reply_to": map[string]interface{}{"type": "string", "description": "回复某封邮件时传其 mail_id"},
},
"required": []string{"to", "subject", "body"},
},
}, p.handleSendMail)
registerTool("forward_mail", sdk.ToolDef{
Name: "forward_mail",
Description: "转发一封邮件给新的收件人(自动引用原文与附件)。与回复不同:回复落回原会话,转发按目标地址另行定位会话。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"mail_id": map[string]interface{}{"type": "string", "description": "要转发的邮件 ID"},
"to": map[string]interface{}{"type": "string", "description": "新收件人的三维地址"},
"comment": map[string]interface{}{"type": "string", "description": "转发说明"},
"cc": map[string]interface{}{"type": "string", "description": "抄送"},
"subject": map[string]interface{}{"type": "string", "description": "自定义主题;留空则自动加 Fwd: 前缀"},
"session_alias": map[string]interface{}{"type": "string", "description": "仅当目标地址以 .new 结尾时生效"},
},
"required": []string{"mail_id", "to"},
},
}, p.handleForwardMail)
registerTool("upload_attachment", sdk.ToolDef{
Name: "upload_attachment",
Description: "上传本地文件作为邮件附件。返回 attachment_id填入 send_mail 的 attachments 字段。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"file_path": map[string]interface{}{"type": "string", "description": "本地文件路径"},
},
"required": []string{"file_path"},
},
}, p.handleUploadAttachment)
registerTool("download_attachment", sdk.ToolDef{
Name: "download_attachment",
Description: "下载附件到本地。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"attachment_id": map[string]interface{}{"type": "string", "description": "附件 ID"},
"save_path": map[string]interface{}{"type": "string", "description": "保存路径"},
},
"required": []string{"attachment_id", "save_path"},
},
}, p.handleDownloadAttachment)
registerTool("suggest_address", sdk.ToolDef{
Name: "suggest_address",
Description: "查询可用收件人地址。不带参数给候选收件人;带 name 给工作目录name+path 都带则给会话别名。发信前应先用它确认地址。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"name": map[string]interface{}{"type": "string", "description": "收件人名;留空则列出所有候选收件人"},
"path": map[string]interface{}{"type": "string", "description": "工作目录;与 name 同时给出才列会话"},
},
},
}, p.handleSuggestAddress)
registerTool("list_contacts", sdk.ToolDef{
Name: "list_contacts",
Description: "列出自己参与过的全部会话及各自的可投递地址、未读数、剩余往返预算。",
Parameters: oneStringParam("limit", "最多列出多少条,默认 20", false),
}, p.handleListContacts)
registerTool("session_participants", sdk.ToolDef{
Name: "session_participants",
Description: "列出某条会话的全部参与方与各自的可投递地址,并标出谁还没回应。",
Parameters: oneStringParam("session_id", "会话 ID", true),
}, p.handleSessionParticipants)
registerTool("read_thread", sdk.ToolDef{
Name: "read_thread",
Description: "查看一封邮件所在线索的完整往来(谁回了谁、谁还没回)。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"mail_id": map[string]interface{}{"type": "string", "description": "线索中任一封邮件的 ID"},
"offset": map[string]interface{}{"type": "number", "description": "分页偏移"},
},
"required": []string{"mail_id"},
},
}, p.handleReadThread)
registerTool("connect_to_server", sdk.ToolDef{
Name: "connect_to_server",
Description: "连接到 AgentMail Gateway登记本机密钥并完成注册。首次安装或换了 Gateway 地址时调用。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"gateway_url": map[string]interface{}{"type": "string", "description": "Gateway 地址;省略则用当前配置"},
"key_token": map[string]interface{}{"type": "string", "description": "管理员签发的 Agent 密钥;省略则用当前密钥"},
},
},
}, p.handleConnectToServer)
// ─── 日程 / 待办 ───
//
// 这一组的价值不在「记事」而在**跨进程的时间**:模型自己没法让一个进程
// 在未来某刻醒来,插件里的定时器也随 homed 重启一起消失。交给 Gateway
// 之后由数据库与调度器保证,到点发一封邮件把收件方唤起来。
registerTool("create_schedule", sdk.ToolDef{
Name: "create_schedule",
Description: "创建一条日程提醒。到点时 Gateway 会发一封邮件给收件方(默认是你自己)," +
"因此它能跨进程重启生效 —— 比你自己记着时间可靠。" +
"典型用法:稍后检查某件事、提醒另一个 Agent 交东西、周期性巡检。" +
scheduleFieldsHint,
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"title": map[string]interface{}{"type": "string", "description": "日程标题,会成为提醒邮件的主题"},
"event_time": map[string]interface{}{"type": "string", "description": "事件时间RFC3339 带时区2026-09-10T09:00:00+08:00。不支持「明天」这类相对表述"},
"description": map[string]interface{}{"type": "string", "description": "补充说明,可作为 {description} 变量填入提醒正文"},
"reminder_text": map[string]interface{}{"type": "string", "description": "提醒邮件正文模板,支持 {title} {time} {description} 三个变量。留空用默认模板"},
"remind_before": map[string]interface{}{"type": "number", "description": "提前多少分钟提醒,默认 0到点才提醒"},
"recurrence": map[string]interface{}{"type": "string", "description": "重复规则,默认 none。农历用 lunar_monthly / lunar_yearly"},
"recurrence_end": map[string]interface{}{"type": "string", "description": "重复到什么时候为止,留空 = 一直重复"},
"recipients": map[string]interface{}{"type": "string", "description": "收件方,逗号分隔的三维地址(如 dsh,pi@/home/x。省略 = 发给自己。不能设给人类用户 —— 要通知人请直接 send_mail"},
"delivery_mode": map[string]interface{}{"type": "string", "description": "多收件人时separate默认各自独立会话互不可见或 together首个为主收件人其余抄送共享一条线索"},
},
"required": []string{"title", "event_time"},
},
}, p.handleCreateSchedule)
registerTool("list_schedules", sdk.ToolDef{
Name: "list_schedules",
Description: "列出你建的日程(只能看到自己建的)。返回每条的 ID、时间、重复规则与收件方。" +
"改时间或删除前先用它查 ID。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"status": map[string]interface{}{"type": "string", "description": "过滤 active默认/paused/cancelled/all"},
"from": map[string]interface{}{"type": "string", "description": "起始时间,默认昨天"},
"to": map[string]interface{}{"type": "string", "description": "结束时间,默认三个月后"},
},
},
}, p.handleListSchedules)
registerTool("update_schedule", sdk.ToolDef{
Name: "update_schedule",
Description: "改一条日程。**只传要改的字段**,省略的保持原值 —— 不要回传全部字段," +
"记错一个就会覆盖掉原有的提醒正文或收件方。" +
"暂停提醒传 status=paused。" + scheduleFieldsHint,
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"event_id": map[string]interface{}{"type": "string", "description": "要改哪条(用 list_schedules 查)"},
"title": map[string]interface{}{"type": "string", "description": "新标题"},
"event_time": map[string]interface{}{"type": "string", "description": "新时间RFC3339 带时区"},
"description": map[string]interface{}{"type": "string", "description": "新说明"},
"reminder_text": map[string]interface{}{"type": "string", "description": "新的提醒正文模板"},
"remind_before": map[string]interface{}{"type": "number", "description": "新的提前分钟数"},
"recurrence": map[string]interface{}{"type": "string", "description": "新的重复规则"},
"recurrence_end": map[string]interface{}{"type": "string", "description": "新的重复终止时间"},
"recipients": map[string]interface{}{"type": "string", "description": "新收件方,逗号分隔。不能改成空"},
"delivery_mode": map[string]interface{}{"type": "string", "description": "separate 或 together"},
"status": map[string]interface{}{"type": "string", "description": "active / paused暂停提醒/ cancelled"},
},
"required": []string{"event_id"},
},
}, p.handleUpdateSchedule)
registerTool("delete_schedule", sdk.ToolDef{
Name: "delete_schedule",
Description: "删掉一条日程,之后不再提醒。只想临时停掉请用 update_schedule 传 status=paused。",
Parameters: oneStringParam("event_id", "要删哪条(用 list_schedules 查)", true),
}, p.handleDeleteSchedule)
// 注册输出通道
s.RegisterOutputChannel("homeagent", sdk.CapText|sdk.CapFile,
"发送邮件。meta JSON 格式:{to, subject, reply_to}type: text",
sdk.ChannelDef{}, p.handleOutputChannel)
// 数量从 RegisterTool 的调用数派生,不硬编码。
//
// 之前这里写死 13而实际注册的是 11 —— 排查「工具没生效」时日志说 13、
// 平台说 11两个数字都不可信白花了一轮时间。加工具时忘改常量是必然的
// 所以让它没有机会写错。
log.Printf("[homeagent-mail-bridge] 注册完成(%d 个工具 + 1 个输出通道),等待 Gateway SSE",
registeredToolCount)
// 启动心跳 + SSE后台 goroutine
go p.heartbeatLoop()
go p.sseLoop()
return nil
}
func (p *Plugin) Stop() error {
p.stopOnce.Do(func() { close(p.stopCh) })
return nil
}
// ─── 心跳 ───
func (p *Plugin) heartbeatLoop() {
// B-1.3:立即发一次心跳,不等第一个 30 秒周期
if err := p.heartbeat(); err != nil {
log.Printf("[homeagent-mail-bridge] 首次心跳失败: %v", err)
}
ticker := time.NewTicker(heartbeatIntval)
defer ticker.Stop()
for {
select {
case <-p.stopCh:
return
case <-ticker.C:
if err := p.heartbeat(); err != nil {
// B-2.1:心跳失败不重试不报错,下一轮补上
}
}
}
}
func (p *Plugin) register() error {
body := map[string]interface{}{
"name": p.agentName,
"platform": "homeagent",
"workspaces": []interface{}{},
}
return p.post("/agent/register", body, nil)
}
func (p *Plugin) heartbeat() error {
payload := map[string]interface{}{}
// B-2.3:带上模型目录
if len(p.modelCatalog) > 0 {
payload["models"] = p.modelCatalog
}
var resp struct {
AllowedModels []string `json:"allowed_models"`
PendingMails int `json:"pending_mails"`
}
if err := p.post("/agent/heartbeat", payload, &resp); err != nil {
return err
}
// B-2.2:从响应读 allowed_modelsGateway 有范围配置时返回)
// 这些模型被插件用于后续轮次的模型选择降级尝试
// B-1.6 / B-7首个成功心跳后如果有未读邮件补拉
if !p.catchupDone {
p.catchupDone = true
if resp.PendingMails > 0 {
go p.catchUp(resp.PendingMails)
}
}
return nil
}
// B-7 补拉:串行读取未读邮件,逐封注入 agent 事件循环。
//
// 与 DSH/pi 的补拉逻辑一致:
// - 上限 5 封catchupLimit避免重启时一次性灌入太多
// - 正序(最旧的先处理),保持时间线
// - 只补 normal 类型permission 不补投——人在 WebUI 上看到就知道了)
// - 每封之间等 InjectInputSync 返回(串行处理)
func (p *Plugin) catchUp(pending int) {
limit := catchupLimit
if pending < limit {
limit = pending
}
log.Printf("[homeagent-mail-bridge] 补投 %d 封离线期间的邮件(共 %d 封未读)", limit, pending)
var inbox struct {
Mails []struct {
MailID string `json:"mail_id"`
FromName string `json:"from_name"`
Subject string `json:"subject"`
MailType string `json:"mail_type"`
ReplyTo string `json:"reply_to"`
} `json:"mails"`
}
url := fmt.Sprintf("%s/api/v1/mail/inbox?status=unread&limit=%d", p.gwURL, limit)
if err := p.get(url, &inbox); err != nil {
log.Printf("[homeagent-mail-bridge] 补拉失败: %v", err)
return
}
for _, m := range inbox.Mails {
if m.MailType != "normal" {
continue // permission 等非邮件驱动的不补投
}
// 构造注入消息(与 handleNewMail 一致)
prompt := fmt.Sprintf(
"你收到一封新邮件AgentMail。\n\n"+
"发件人:%s\n主题%s\n邮件 ID%s\n身份你是 %s\n\n"+
"请先调用 read_inbox 读取完整正文,然后处理其中的请求。\n\n"+
"**回信不用你自己发**:你把结论说出来就行,\n"+
"插件会在这一轮结束时自动把你最后那段话作为回信发回给 %s不消耗你的发信配额。\n"+
"只有在需要主动联系其他人、或要带附件时才调用 send_mail。",
m.FromName, m.Subject, m.MailID, p.agentName, m.FromName,
)
reply := p.sdk.InjectInputSync(p.name, p.name, prompt)
if reply == "" {
// B-6模型没回发一封告知
p.sendFailureReply(m.FromName, m.Subject, m.MailID, "模型未产生回复")
continue
}
// B-5.3:检查模型是否已经自己发过信
rk := "homeagent:" + m.MailID
p.explicitSendsMu.Lock()
_, sent := p.explicitSends[rk]
p.explicitSendsMu.Unlock()
if sent {
// 模型已经在这一轮里自己回了这封信,不再重复 relay
continue
}
// B-5.2:自动回信带 relay:"summary" —— 搬运不算模型自主发信,不扣配额
p.sendMailRelay(m.FromName, "Re: "+m.Subject, reply, m.MailID, "homeagent:"+m.MailID)
}
}
// ─── SSE ───
func (p *Plugin) sseLoop() {
for {
select {
case <-p.stopCh:
return
default:
}
if err := p.readSSE(); err != nil {
log.Printf("[homeagent-mail-bridge] SSE 断开: %v%v 后重连", err, sseRetry)
time.Sleep(sseRetry)
}
}
}
func (p *Plugin) readSSE() error {
req, err := http.NewRequest("GET", p.gwURL+"/api/v1/events/stream", nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+p.key)
// W-4断线期间的事件会丢带上 Last-Event-ID 可以让 Gateway 从断点补发。
// 只上报比当前记录的更大的 IDGateway 重放时发的是事件原本的 ID
// 如果无条件赋值lastEventID 会从 123 退回 116 → 下次重连又报 116 →
// 又重放 —— 自激振荡的放大器。
p.sseMu.Lock()
if p.lastEventID != "" {
req.Header.Set("Last-Event-ID", p.lastEventID)
}
p.sseMu.Unlock()
// sseClient 无 Timeoutp.client 有 60s TimeoutSSE 是长连接,
// 每 60 秒自己掐断自己 → 重放 → 阻塞 → 超时 → 重放。
resp, err := p.sseClient.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
return fmt.Errorf("SSE HTTP %d", resp.StatusCode)
}
log.Printf("[homeagent-mail-bridge] SSE 已连接")
// 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:
return nil
default:
}
line, err := br.ReadString('\n')
if line != "" {
p.parseSSELine(strings.TrimRight(line, "\n"))
}
if err != nil {
if err == io.EOF {
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。只向前推进不回退。
// Gateway 重放旧事件时发的是旧 ID无条件赋值会让 lastEventID
// 从 123 退回到 116 → 下次重连报 116 → 又重放 → 振荡。
if strings.HasPrefix(line, "id: ") {
eid := strings.TrimPrefix(line, "id: ")
id := parseSSEID(eid)
p.sseMu.Lock()
if id > p.sseMaxID {
p.sseMaxID = id
p.lastEventID = eid
}
p.sseMu.Unlock()
return
}
if !strings.HasPrefix(line, "data: ") {
return
}
raw := strings.TrimPrefix(line, "data: ")
if raw == "" || raw == "{}" {
return
}
var evt struct {
MailID string `json:"mail_id"`
SessionID string `json:"session_id"`
FromName string `json:"from_name"`
Subject string `json:"subject"`
MailType string `json:"mail_type"`
Role string `json:"role"`
Workspace string `json:"to_workspace"`
Alias string `json:"session_alias"`
ReplyAddr string `json:"reply_address"`
}
if err := json.Unmarshal([]byte(raw), &evt); err != nil {
return
}
if evt.MailID == "" {
return
}
if evt.MailType == "permission_decision" {
p.handlePermissionDecision(evt)
return
}
if evt.MailType == "normal" {
// 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 全部重放一遍。
// 把处理丢到独立 goroutineparseSSELine 立刻返回,读循环继续。
// homeagent 是单事件循环InjectInputSync 自己会排队。
go p.handleNewMail(evt)
}
}
// ─── 输出通道agent 主动发信)───
func (p *Plugin) handleOutputChannel(args map[string]interface{}) (interface{}, error) {
payload, _ := args["payload"].(string)
meta, _ := args["meta"].(string)
if payload == "" {
return nil, fmt.Errorf("payload 不能为空")
}
var m struct {
To string `json:"to"`
Subject string `json:"subject"`
ReplyTo string `json:"reply_to"`
}
if meta != "" {
json.Unmarshal([]byte(meta), &m)
}
if m.To == "" {
return nil, fmt.Errorf("meta 中需要 to 字段")
}
// B-5.3:记录模型自主发信,后续自动 relay 时跳过
rk := m.ReplyTo
if rk != "" {
p.explicitSendsMu.Lock()
p.explicitSends["homeagent:"+rk] = time.Now()
p.explicitSendsMu.Unlock()
}
if err := p.sendMail(m.To, m.Subject, payload, m.ReplyTo, ""); err != nil {
return nil, err
}
return map[string]interface{}{"status": "sent"}, nil
}
// ─── 发信辅助 ───
// sendMail 发一封普通邮件(不带 relay 标记)。
// 用于模型主动调 send_mail 或 output_send 时。
func (p *Plugin) sendMail(to, subject, body, replyTo, sessionAlias string) error {
payload := map[string]interface{}{
"to": to,
"subject": subject,
"body": body,
}
if replyTo != "" {
payload["reply_to"] = replyTo
}
if sessionAlias != "" {
payload["session_alias"] = sessionAlias
}
return p.post("/mail/send", payload, nil)
}
// sendMailRelay 发一封带 relay:"summary" 标记的邮件。
//
// B-5.2:插件代模型搬运回复时必须带 relay:"summary" + relay_key
// Gateway 才会把它走免配额通道(插件搬运不算模型自主发信)。
// B-5.4:文本为空时不发空邮件。
func (p *Plugin) sendMailRelay(to, subject, body, replyTo, relayKey string) error {
if strings.TrimSpace(body) == "" {
return nil // B-5.4
}
payload := map[string]interface{}{
"to": to,
"subject": subject,
"body": body,
"relay": "summary",
"relay_key": relayKey,
}
if replyTo != "" {
payload["reply_to"] = replyTo
}
return p.post("/mail/send", payload, nil)
}
// sendFailureReply 在模型处理失败时给发件人一封告知。
//
// B-6无法处理时必须回信。发件人发了邮件后没有任何音讯是最糟的体验 ——
// 他不知道邮件到了没有、模型看了没有、是卡住了还是忽略了。
// 必须带 relay:"summary" 走免配额通道,这是插件代劳不是模型自主发信。
func (p *Plugin) sendFailureReply(to, subject, replyTo, reason string) {
body := fmt.Sprintf(
"这是一封自动通知:您发送的主题为「%s」的邮件在处理时遇到了问题未能产生有效回复。\n\n"+
"原因:%s\n\n"+
"请稍后重试,或通过其他方式联系。",
subject, reason,
)
rk := "homeagent:failure:" + replyTo
if err := p.sendMailRelay(to, "Re: "+subject, body, replyTo, rk); err != nil {
log.Printf("[homeagent-mail-bridge] 失败通知发送失败: %v", err)
}
}
// ─── 新邮件处理 ───
func (p *Plugin) handleNewMail(evt struct {
MailID string `json:"mail_id"`
SessionID string `json:"session_id"`
FromName string `json:"from_name"`
Subject string `json:"subject"`
MailType string `json:"mail_type"`
Role string `json:"role"`
Workspace string `json:"to_workspace"`
Alias string `json:"session_alias"`
ReplyAddr string `json:"reply_address"`
}) {
prompt := fmt.Sprintf(
"你收到一封新邮件AgentMail。\n\n"+
"发件人:%s\n主题%s\n邮件 ID%s\n身份你是 %s\n\n"+
"请先调用 read_inbox 读取完整正文,然后处理其中的请求。\n\n"+
"**回信不用你自己发**:你把本轮工作做完、把结论说出来就行,\n"+
"插件会在这一轮结束时自动把你最后那段话作为回信发回给 %s不消耗你的发信配额。\n"+
"只有在需要主动联系其他人、或要带附件时才调用 send_mail。",
evt.FromName, evt.Subject, evt.MailID, p.agentName, evt.FromName,
)
// InjectInputSync 阻塞等待 agent 处理完毕,返回最终回复文本。
reply := p.sdk.InjectInputSync(p.name, p.name, prompt)
// B-6模型没回空 = turn/end 信号 kind=error或模型没说话
if reply == "" {
log.Printf("[homeagent-mail-bridge] 邮件 %s来自 %s%sagent 无回复,发失败通知",
evt.MailID[:8], evt.FromName, evt.Subject)
p.sendFailureReply(evt.FromName, evt.Subject, evt.MailID, "模型未产生回复")
return
}
// B-5.3:检查模型是否已经自己发过信(通过 send_mail 或 output_send
rk := "homeagent:" + evt.MailID
p.explicitSendsMu.Lock()
_, sent := p.explicitSends[rk]
if sent {
delete(p.explicitSends, rk) // 用过即清,不留残
}
p.explicitSendsMu.Unlock()
if sent {
// 模型已经在这一轮里自己回了这封信,让位
log.Printf("[homeagent-mail-bridge] 邮件 %s 模型已自行回复,跳过自动 relay", evt.MailID[:8])
return
}
// B-5.2:自动回信带 relay:"summary" + relay_key
if err := p.sendMailRelay(evt.FromName, "Re: "+evt.Subject, reply, evt.MailID, rk); err != nil {
log.Printf("[homeagent-mail-bridge] 自动回信失败: %v", err)
} else {
log.Printf("[homeagent-mail-bridge] 已自动回信给 %s%d 字)", evt.FromName, len(reply))
}
// 清理过期的 explicitSends 记录
p.explicitSendsMu.Lock()
for k, t := range p.explicitSends {
if time.Since(t) > dedupWindow {
delete(p.explicitSends, k)
}
}
p.explicitSendsMu.Unlock()
}
// ─── 权限决策 ───
func (p *Plugin) handlePermissionDecision(evt struct {
MailID string `json:"mail_id"`
SessionID string `json:"session_id"`
FromName string `json:"from_name"`
Subject string `json:"subject"`
MailType string `json:"mail_type"`
Role string `json:"role"`
Workspace string `json:"to_workspace"`
Alias string `json:"session_alias"`
ReplyAddr string `json:"reply_address"`
}) {
prompt := fmt.Sprintf(
"你之前发起的权限请求已有结论:%s决策人%s。请据此继续。",
evt.Subject, evt.FromName,
)
p.sdk.InjectText(p.name, p.name, prompt)
}
// ─── 工具实现 ───
func (p *Plugin) handleReadInbox(args map[string]interface{}) (interface{}, error) {
status := "unread"
if v, ok := args["status"].(string); ok && v != "" {
status = v
}
limit := 5
if v, ok := args["limit"].(float64); ok && v > 0 {
limit = int(v)
}
url := fmt.Sprintf("%s/api/v1/mail/inbox?status=%s&limit=%d", p.gwURL, status, limit)
var result map[string]interface{}
if err := p.get(url, &result); err != nil {
return nil, err
}
mails, _ := result["mails"].([]interface{})
if len(mails) == 0 {
return map[string]interface{}{"content": []map[string]interface{}{{"type": "text", "text": "收件箱为空。"}}}, nil
}
var sb strings.Builder
for i, m := range mails {
mail, _ := m.(map[string]interface{})
if mail == nil {
continue
}
from, _ := mail["from_name"].(string)
subj, _ := mail["subject"].(string)
mid, _ := mail["mail_id"].(string)
body, _ := mail["body"].(string)
if body == "" {
body, _ = mail["body_preview"].(string)
}
alias, _ := mail["session_alias"].(string)
fmt.Fprintf(&sb, "[%d] %s: %s\n邮件 ID: %s\n会话: #%s\n", i+1, from, subj, mid, alias)
if ccList, ok := mail["cc_list"].([]interface{}); ok && len(ccList) > 0 {
names := make([]string, 0, len(ccList))
for _, c := range ccList {
if cc, ok := c.(map[string]interface{}); ok {
if n, ok := cc["name"].(string); ok {
names = append(names, n)
}
}
}
if len(names) > 0 {
fmt.Fprintf(&sb, "抄送: %s\n", strings.Join(names, "、"))
}
}
if atts, ok := mail["attachments"].([]interface{}); ok && len(atts) > 0 {
fmt.Fprintf(&sb, "附件:\n")
for _, a := range atts {
if att, ok := a.(map[string]interface{}); ok {
fn, _ := att["filename"].(string)
sz, _ := att["size_bytes"].(float64)
aid, _ := att["attachment_id"].(string)
fmt.Fprintf(&sb, " - %s (%.1fKB, id=%s)\n", fn, sz/1024, aid)
}
}
}
if body != "" {
if len(body) > 1000 {
body = body[:1000] + "..."
}
fmt.Fprintf(&sb, "内容: %s\n", body)
}
sb.WriteString("\n")
}
ids := make([]string, 0, len(mails))
for _, m := range mails {
if mail, ok := m.(map[string]interface{}); ok {
if id, ok := mail["mail_id"].(string); ok {
ids = append(ids, id)
}
}
}
if len(ids) > 0 {
go p.markRead(ids)
}
return map[string]interface{}{
"content": []map[string]interface{}{{"type": "text", "text": sb.String()}},
}, nil
}
func (p *Plugin) handleSendMail(args map[string]interface{}) (interface{}, error) {
to, _ := args["to"].(string)
subj, _ := args["subject"].(string)
body, _ := args["body"].(string)
cc, _ := args["cc"].(string)
replyTo, _ := args["reply_to"].(string)
if to == "" || subj == "" || body == "" {
return nil, fmt.Errorf("缺少必填字段to, subject, body")
}
// B-5.3:记录模型自主发信
rk := replyTo
if rk != "" {
p.explicitSendsMu.Lock()
p.explicitSends["homeagent:"+rk] = time.Now()
p.explicitSendsMu.Unlock()
}
payload := map[string]interface{}{
"to": to,
"subject": subj,
"body": body,
}
if cc != "" {
payload["cc"] = cc
}
if replyTo != "" {
payload["reply_to"] = replyTo
}
var result map[string]interface{}
if err := p.post("/mail/send", payload, &result); err != nil {
return nil, err
}
mid, _ := result["mail_id"].(string)
sid, _ := result["session_id"].(string)
text := fmt.Sprintf("邮件已发送ID: %sSession: %s", mid, sid)
if budget, ok := result["budget_remaining"].(float64); ok {
text += fmt.Sprintf("。本任务剩余 %.0f 个来回", budget)
}
return map[string]interface{}{
"content": []map[string]interface{}{{"type": "text", "text": text}},
}, nil
}
// C-14 附件上传 —— 真 multipart不是桩。
//
// 读取本地文件 → 构造 multipart/form-data → POST /api/v1/attachments。
// 返回 attachment_id填入 send_mail 的 attachments 字段。
func (p *Plugin) handleUploadAttachment(args map[string]interface{}) (interface{}, error) {
filePath, _ := args["file_path"].(string)
if filePath == "" {
return nil, fmt.Errorf("缺少 file_path")
}
data, err := os.ReadFile(filePath)
if err != nil {
return nil, fmt.Errorf("读取文件失败: %v", err)
}
// multipart/form-data
var buf bytes.Buffer
writer := multipart.NewWriter(&buf)
part, err := writer.CreateFormFile("file", filepath.Base(filePath))
if err != nil {
return nil, fmt.Errorf("创建 multipart 失败: %v", err)
}
if _, err := part.Write(data); err != nil {
return nil, fmt.Errorf("写入文件数据失败: %v", err)
}
writer.Close()
req, err := http.NewRequest("POST", p.gwURL+"/api/v1/attachments", &buf)
if err != nil {
return nil, err
}
req.Header.Set("Content-Type", writer.FormDataContentType())
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, string(body))
}
var result struct {
AttachmentID string `json:"attachment_id"`
Filename string `json:"filename"`
SizeBytes int `json:"size_bytes"`
}
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
return nil, err
}
text := fmt.Sprintf("附件已上传id=%s filename=%s size=%dKB\n在 send_mail 的 attachments 字段传 [{\"attachment_id\":\"%s\"}]",
result.AttachmentID, result.Filename, result.SizeBytes/1024, result.AttachmentID)
return map[string]interface{}{
"content": []map[string]interface{}{{"type": "text", "text": text}},
}, nil
}
// C-14 附件下载 —— 真 octet-stream 下载。
func (p *Plugin) handleDownloadAttachment(args map[string]interface{}) (interface{}, error) {
aid, _ := args["attachment_id"].(string)
savePath, _ := args["save_path"].(string)
if aid == "" || savePath == "" {
return nil, fmt.Errorf("缺少 attachment_id 和 save_path")
}
req, err := http.NewRequest("GET", p.gwURL+"/api/v1/attachments/"+aid, nil)
if err != nil {
return nil, err
}
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, string(body))
}
// 确保目录存在
if err := os.MkdirAll(filepath.Dir(savePath), 0755); err != nil {
return nil, fmt.Errorf("创建目录失败: %v", err)
}
out, err := os.Create(savePath)
if err != nil {
return nil, fmt.Errorf("创建文件失败: %v", err)
}
defer out.Close()
written, err := io.Copy(out, resp.Body)
if err != nil {
return nil, fmt.Errorf("写入文件失败: %v", err)
}
text := fmt.Sprintf("附件已下载:%s%dKB", savePath, written/1024)
return map[string]interface{}{
"content": []map[string]interface{}{{"type": "text", "text": text}},
}, nil
}
// ─── HTTP 辅助 ───
func (p *Plugin) get(url string, out interface{}) error {
req, err := http.NewRequest("GET", url, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("HTTP %d: %s", resp.StatusCode, string(body))
}
return json.NewDecoder(resp.Body).Decode(out)
}
func (p *Plugin) post(path string, payload interface{}, out interface{}) error {
data, err := json.Marshal(payload)
if err != nil {
return err
}
url := path
if !strings.HasPrefix(path, "http") {
url = p.gwURL + "/api/v1" + path
}
req, err := http.NewRequest("POST", url, bytes.NewReader(data))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("POST %s HTTP %d: %s", path, resp.StatusCode, string(body))
}
if out != nil {
return json.NewDecoder(resp.Body).Decode(out)
}
return nil
}
func (p *Plugin) put(path string, payload interface{}, out interface{}) error {
data, err := json.Marshal(payload)
if err != nil {
return err
}
url := path
if !strings.HasPrefix(path, "http") {
url = p.gwURL + "/api/v1" + path
}
req, err := http.NewRequest("PUT", url, bytes.NewReader(data))
if err != nil {
return err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("PUT %s HTTP %d: %s", path, resp.StatusCode, string(body))
}
if out != nil {
return json.NewDecoder(resp.Body).Decode(out)
}
return nil
}
func (p *Plugin) delete(path string) error {
url := path
if !strings.HasPrefix(path, "http") {
url = p.gwURL + "/api/v1" + path
}
req, err := http.NewRequest("DELETE", url, nil)
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+p.key)
resp, err := p.client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode >= 400 {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("DELETE %s HTTP %d: %s", path, resp.StatusCode, string(body))
}
return nil
}
func (p *Plugin) markRead(ids []string) {
data, _ := json.Marshal(map[string]interface{}{"mail_ids": ids})
req, err := http.NewRequest("POST", p.gwURL+"/api/v1/mail/read", bytes.NewReader(data))
if err != nil {
return
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+p.key)
p.client.Do(req)
}