From 4d8165fde944808265f6787adf868110156bda5c Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Thu, 1 Oct 2026 19:41:06 +0800 Subject: [PATCH] =?UTF-8?q?fix(=E5=9B=9E=E8=B7=AF):=20Agent=E2=86=94Agent?= =?UTF-8?q?=20=E5=8A=A0=202h=20=E5=86=B7=E9=9D=99=E6=9C=9F=EF=BC=9B?= =?UTF-8?q?=E5=86=B7=E5=8D=B4=E6=9C=9F=E5=86=85=20relay=20=E4=B8=8D?= =?UTF-8?q?=E8=87=AA=E5=8A=A8=E9=87=8D=E6=8A=95=E9=80=92?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 缺口:maxAgentPingPong=8 撞闸后**计数永不回落** —— 只有人类插话才归零。 旧文案把出路指向「请由人类插一句话」,而那条线索上常常**根本没有人类** (2026-10-01 报告:agent 一封都发不出,且那条线索无人在场)。 修法(两条要求分别落地): ① 2h 恢复机制:撞闸即写入 session_agent_locks(落库,内存态一重启就"恢复", 且多副本各算各的);到期自动放行,人类插话立刻解锁(优先于到期)。 ② 锁定期间的邮件不自动重投递:冷却期内 relay 直接 403 丢弃 —— **不排队、不占幂等键、不入库**。排队会在 2h 后一次性灌回去, 那等于把刚压住的回路换个更糟的形状放出来。 实测踩到的三个坑(都被判据抓住): - `VALUES ($1,$2,$2,...)` 让 until_at 复用 locked_at 的 $2 ⇒ 锁诞生即过期 - driver 以 UTC 扫回 DATETIME,而 time.Now() 是本地(HKT+8) ⇒ 差 8h > 2h 的一半 ⇒ 锁形同虚设 - 只查**收件方**是否人类,漏了**发件方** —— 而「人类插话解锁」的常态就是 人回信、收件方仍是 Agent ⇒ 人插话反被自己写的闸 403 拦下,那条出路根本不存在 - SQLite 没有 GREATEST;在 DATETIME 上按**字符串**比大小 ⇒ CASE 也会错。 改为「已有锁一律不碰 until_at」。 判据:repo 6 格(含"读失败不得读成未锁定"—— 最初 0 格能抓,变异测试补的) + handler 4 格。三个变异全部经得起(且每��都先确认变异编译通过再数红格 —— 本轮多次 grep 得 0 实际是 build failed,测试压根没跑)。 --- server/internal/db/migrations/init.sql | 13 + server/internal/db/migrations/init_sqlite.sql | 20 ++ server/internal/handler/mail.go | 68 +++++- server/internal/handler/mail_cooldown_test.go | 212 +++++++++++++++++ server/internal/repo/agentlock.go | 139 +++++++++++ server/internal/repo/agentlock_test.go | 222 ++++++++++++++++++ 6 files changed, 671 insertions(+), 3 deletions(-) create mode 100644 server/internal/handler/mail_cooldown_test.go create mode 100644 server/internal/repo/agentlock.go create mode 100644 server/internal/repo/agentlock_test.go diff --git a/server/internal/db/migrations/init.sql b/server/internal/db/migrations/init.sql index f81ace7..43d4c97 100644 --- a/server/internal/db/migrations/init.sql +++ b/server/internal/db/migrations/init.sql @@ -324,6 +324,19 @@ CREATE TABLE IF NOT EXISTS relayed_mails ( -- 人类决策后要按 mail_id 反查上游 permission id CREATE INDEX IF NOT EXISTS idx_relayed_mail ON relayed_mails(mail_id); +-- Agent↔Agent 回路的**锁定期**(2026-10-01) +-- 理由与 SQLite 侧同,见那边注释。schema 语义必须一致(db/migrate.go 要求两处都改)。 +CREATE TABLE IF NOT EXISTS session_agent_locks ( + session_id UUID PRIMARY KEY REFERENCES sessions(session_id), + locked_at TIMESTAMPTZ NOT NULL, + until_at TIMESTAMPTZ NOT NULL, + hops INTEGER NOT NULL DEFAULT 0, + reason VARCHAR(128) NOT NULL DEFAULT '' +); + +CREATE INDEX IF NOT EXISTS idx_session_agent_locks_until ON session_agent_locks(until_at); + + CREATE TABLE IF NOT EXISTS attachments ( attachment_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), mail_id UUID REFERENCES mails(mail_id) ON DELETE CASCADE, diff --git a/server/internal/db/migrations/init_sqlite.sql b/server/internal/db/migrations/init_sqlite.sql index e943818..a89ce90 100644 --- a/server/internal/db/migrations/init_sqlite.sql +++ b/server/internal/db/migrations/init_sqlite.sql @@ -331,6 +331,26 @@ CREATE TABLE IF NOT EXISTS relayed_mails ( -- 人类决策后要按 mail_id 反查上游 permission id CREATE INDEX IF NOT EXISTS idx_relayed_mail ON relayed_mails(mail_id); +-- Agent↔Agent 回路的**锁定期**(2026-10-01) +-- +-- 为什么需要:maxAgentPingPong=8 拦下之后,**计数永不回落** —— +-- 无人参与时会话被永久锁死(2026-10-01 报告:agent 一封都发不出, +-- 而那道闸的文案只教人「请由人类插一句话」,可这条线索上根本没有人类)。 +-- +-- 锁定期给出**自恢复出口**:到期自动放行,不依赖任何人介入。 +-- 落库而不是放内存:重启必须保留(内存态一重启就"恢复", +-- 而锁定期的意义正是"别在短时间内又炸一轮")。 +CREATE TABLE IF NOT EXISTS session_agent_locks ( + session_id TEXT PRIMARY KEY REFERENCES sessions(session_id), + locked_at DATETIME NOT NULL, + until_at DATETIME NOT NULL, + hops INTEGER NOT NULL DEFAULT 0, + reason TEXT NOT NULL DEFAULT '' +); + +CREATE INDEX IF NOT EXISTS idx_session_agent_locks_until ON session_agent_locks(until_at); + + CREATE TABLE IF NOT EXISTS attachments ( attachment_id TEXT PRIMARY KEY DEFAULT (gen_random_uuid()), mail_id TEXT REFERENCES mails(mail_id) ON DELETE CASCADE, diff --git a/server/internal/handler/mail.go b/server/internal/handler/mail.go index 6b00a76..c3de3f1 100644 --- a/server/internal/handler/mail.go +++ b/server/internal/handler/mail.go @@ -6,6 +6,7 @@ import ( "fmt" "net/http" "strings" + "time" "github.com/agentmail/gateway/internal/middleware" "github.com/agentmail/gateway/internal/models" @@ -425,18 +426,79 @@ func SendMail(w http.ResponseWriter, r *http.Request) { if hops, hErr := repo.CountTrailingAgentPingPong(r.Context(), sessionID); hErr == nil { // 收件方是人类时不算回路(人类在回路里,正是我们要的"有人决策")。 isHuman, _ := repo.IsHumanUser(r.Context(), to.Name) + + // ★★ 2026-10-01:**发件方**是人类时同样不算(这是原实现漏掉的一边)。 + // + // 原代码只查 `to.Name`。而"人类插一句话解锁"这个动作, + // **发件方是人类**才是它的常态 —— 人是在这条线索里回信, + // 收件方仍是 Agent。只看收件方 ⇒ 人插话反而被锁拦下(实测 403), + // 于是"人类插话即可恢复"这条路是**不存在的**。 + // + // 这条判据是当时逼出来的:`TestHumanPostClearsCooldownImmediately` + // 一上来就红,报错文案还是我自己写的那句"请由人类插一句话"。 + fromHuman, _ := repo.IsHumanUser(r.Context(), agentName) + humanInvolved := isHuman || fromHuman + + // ★ 人类插话立刻清锁(优先于到期)。 + // + // 旧文案只教「请由人类插一句话(计数即归零)」,而那条线索上 + // 往往**根本没有人类** —— 于是那一侧永久发不出信。 + // 现在有两条出路:人来插话(立刻),或冷静期到期(2h 自动)。 + if humanInvolved { + _, _ = repo.SessionLockTouch(r.Context(), sessionID) + } + + // ★ 冷静期内:**Agent 的自主回信**一律拒(不重投、不排队)。 + // + // 这道与上面那道计数闸是同一件事的两个阶段:计数闸负责"发现并开始锁", + // 冷静期负责"锁住不放"。没有第二段的话,第一次撞闸后靠"计数 ≥8" + // 继续拒 —— 形状一样,但**没有出路**:人不在就永远出不来。 + if !humanInvolved && relay == "" { + if lock, locked, _ := repo.SessionLockOf(r.Context(), sessionID); locked { + Error(w, http.StatusForbidden, fmt.Sprintf( + "本会话因 Agent↔Agent 无人决策回信已进入冷静期,还剩约 %s"+ + "(上一次连续 %d 封)。到时自动恢复;"+ + "若要立刻恢复,请由人类在会话里插一句话。", + lock.Remaining.Round(time.Minute), lock.Hops)) + return + } + } + if !isHuman && hops >= repo.MaxAgentPingPong() && relay == "" { + // 开始冷静期(不延长已有的)。落库以便重启后仍生效。 + _, _ = repo.SetSessionLock(r.Context(), sessionID, hops, "agent-ping-pong") Error(w, http.StatusForbidden, fmt.Sprintf( "本会话已连续 %d 封 Agent 之间互相回信、其中没有任何人类参与(上限 %d)。"+ "这通常意味着两个 Agent 在互相确认而无人决策 —— 生产上实测过一天 175 封、"+ - "正文 612KB 却没有任何工作产出。请由人类在会话里插一句话(计数即归零)。"+ - "(本模型未计入插件代劳的转发/退信 —— 那类邮件归 maxRelayHops 管。)", - hops, repo.MaxAgentPingPong())) + "正文 612KB 却没有任何工作产出。\n"+ + "已进入 %s 冷静期,到时自动恢复;"+ + "若要立刻恢复,请由人类在会话里插一句话(立即解锁)。"+ + "(本闸不计入插件代劳的转发/退信 —— 那类邮件归 maxRelayHops 管。)", + hops, repo.MaxAgentPingPong(), repo.AgentLockDuration)) return } } if relay != "" { + // ★ 2026-10-01:冷却期内**不自动重投递**(用户要求)。 + // + // 形状上这道闸与下面那道 relay 上限一样是 403,但意图不同: + // 下面那道拦的是"太多",这道拦的是"正在冷静期"。 + // + // 为什么必须在这里挡:lock 的触发方是**模型自主回信**,被锁的却可能是 + // relay 邮件。两个 Agent 在冷却期里互相自动转发时, + // 若这里放行 → 每封都唤醒对端 → 延迟投递量随时间堆积, + // 2 小时后一次性灌回去 —— 那等于把刚压住的回路又放出来一遍, + // 只是换了个更糟的形状。 + if lock, locked, _ := repo.SessionLockOf(r.Context(), sessionID); locked { + Error(w, http.StatusForbidden, fmt.Sprintf( + "本会话处于 Agent↔Agent 冷静期(约剩 %s),**本次自动转发已丢弃**"+ + "(不会排队、到时不会补投)。"+ + "若这封确实需要送达,请由模型在恢复后主动调 send_mail,或由人类插一句话立即解锁。", + lock.Remaining.Round(time.Minute))) + return + } + // 硬上限:一条会话里**连续**的 relay 邮件不得超过上限。 // // 与预算无关的第二道防线:预算给得大(比如 200)时,两个 Agent 仍能 diff --git a/server/internal/handler/mail_cooldown_test.go b/server/internal/handler/mail_cooldown_test.go new file mode 100644 index 0000000..100de0e --- /dev/null +++ b/server/internal/handler/mail_cooldown_test.go @@ -0,0 +1,212 @@ +package handler + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/agentmail/gateway/internal/db" + "github.com/agentmail/gateway/internal/middleware" + "github.com/agentmail/gateway/internal/repo" + "github.com/google/uuid" +) + +/* +Agent↔Agent 回路的**冷静期**在 handler 上的行为(2026-10-01) + +# 两句话需求,落在两处 + +① 「agent 互发应当有 2h 恢复机制」 + → 撞 maxAgentPingPong=8 时**开始**冷静期,到期自动放行,不依赖人。 + +② 「锁定期间的邮件不自动重投递」 + → 冷静期内 **relay 邮件直接丢弃**(不排队、到时不补投)。 + +# 为什么 ② 必须单独钉 + +"不自动重投递"有三种可能的实现形状,观测上**不等价**: + · 直接拒(丢弃) ← 本次要求的 + · 拒但入队、到期补投 ← 看着一样,其实把回路攒到 2h 后一次性放出来 + · 收下但不投递 ← 库里多一封信,统计与将来的人工检查都被污染 + +所以这一格钉的是**库里没有这封信**,而不只是"接口返回了错误"。 +*/ + +// lockSession 造一个「两个 Agent 互发已达上限、会话已进冷静期」的会话。 +func lockSession(t *testing.T) (uuid.UUID, string) { + t.Helper() + ctx := context.Background() + // Agent 用 repo.CreateOrUpdateAgent 建(真实列是 secret/host_url, + // 不是 agent_key —— 我第一版照想象写了列名,被 schema 打回)。 + // ⚠ **不要**把 pinger/ponger 建成了 user。 + // + // 我第一版顺手调了 mustPermissionUser,结果 IsHumanUser("ponger")=true + // ⇒ 闸按"收件方是人类"正确放行 ⇒ 判据红。 + // 闸没问题,是我的靶子造错了人:这两个必须是 Agent。 + for _, u := range []string{"pinger", "ponger"} { + if err := repo.CreateOrUpdateAgent(ctx, u, "k-"+u, "test", nil); err != nil { + t.Fatalf("建 agent %s: %v", u, err) + } + } + sid, err := repo.CreateSession(ctx, nil, "pinger", "回路", "") + if err != nil { + t.Fatalf("建会话: %v", err) + } + // ⚠ 必须有一封**已存在的信**才能 reply_to 进去。 + // + // 我第一版直接用 `ponger@.new` 投递:`.new` 会**另建一条会话**, + // 于是锁在 sid 上、请求却落在另一条会话上 ⇒ 一路 200, + // 判据看着像"闸没生效"。真实闸没问题,是我的靶子打偏了。 + if _, err := repo.CreateMail(ctx, sid, nil, "ponger", "", "pinger", "", + "起个头", "在吗", nil); err != nil { + t.Fatalf("预置一封信: %v", err) + } + var mailID string + if err := db.DB.QueryRowContext(ctx, + `SELECT mail_id FROM mails WHERE session_id = $1 ORDER BY created_at DESC LIMIT 1`, + sid).Scan(&mailID); err != nil { + t.Fatal(err) + } + if _, err := repo.SetSessionLock(ctx, sid, 8, "test"); err != nil { + t.Fatalf("进入冷静期: %v", err) + } + return sid, mailID +} + +// countMails 数会话里的邮件数。判据用它对比**前后的变化**: +// lockSession 自己会预置一封「起头信」(没有它 reply_to 无处可落), +// 所以基线不是 0。 +func countMails(t *testing.T, sessionID uuid.UUID) int { + t.Helper() + var n int + if err := db.DB.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM mails WHERE session_id = $1`, sessionID).Scan(&n); err != nil { + t.Fatal(err) + } + return n +} + +// postSend 以某个 agent 的身份发一封信,返回响应体。 +func postSend(t *testing.T, agent, to, body, relayKey, replyTo string) *httptest.ResponseRecorder { + t.Helper() + ctx := context.Background() + if err := repo.CreateOrUpdateAgent(ctx, agent, "k-"+agent, "test", nil); err != nil { + t.Fatalf("注册 %s: %v", agent, err) + } + // relay 非空时服务端要求 relay_key 齐备("给了 relay_key 却没给 relay 类型" + // —— 我第一版只给了 relay_key,被这道前置校验先打回,看不到真正的闸)。 + relay := "" + if relayKey != "" { + relay = "summary" + } + form := strings.NewReader(fmt.Sprintf( + `{"to":%q,"subject":"s","body":%q,"relay_key":%q,"relay":%q,"reply_to":%q}`, + to, body, relayKey, relay, replyTo)) + req := httptest.NewRequest("POST", "/api/v1/mail", form) + req.Header.Set("Content-Type", "application/json") + req = req.WithContext(context.WithValue(context.Background(), middleware.AgentNameKey, agent)) + rr := httptest.NewRecorder() + SendMail(rr, req) + return rr +} + +// ① 冷静期内:Agent 的自主回信被拒,且**给出的是可执行的出路**。 +func TestCooldownBlocksAgentSendAndStatesRecovery(t *testing.T) { + setupPermissionHandlerDB(t) + sid, mailID := lockSession(t) + baseCount := countMails(t, sid) + + rr := postSend(t, "pinger", "ponger@", "还在互相确认吗", "", mailID) + if rr.Code != http.StatusForbidden { + t.Fatalf("冷静期内应 403,实际 %d:%s", rr.Code, rr.Body.String()) + } + body := rr.Body.String() + // 文案必须说清「多久后自动恢复」——旧文案只教「请人类插话」, + // 而这条线索上常常根本没有人类(那正是原缺陷)。 + if !strings.Contains(body, "冷静期") { + t.Fatalf("错误文案应说明这是冷静期:%s", body) + } + if !strings.Contains(body, "恢复") { + t.Fatalf("★ 文案必须给出恢复出路(自动恢复 or 人类插话):%s", body) + } + if n := countMails(t, sid); n != baseCount { + t.Fatalf("★ 被拒的信不得入库(否则它会被补投):%d → %d 封", baseCount, n) + } +} + +// ② 冷静期内:relay(自动转发)**丢弃**,不排队不补投。 +func TestCooldownDiscardsRelayWithoutQueueing(t *testing.T) { + setupPermissionHandlerDB(t) + sid, mailID := lockSession(t) + baseCount := countMails(t, sid) + + rr := postSend(t, "pinger", "ponger@", "自动转发", "relay-key-1", mailID) + + if rr.Code != http.StatusForbidden { + t.Fatalf("冷静期内 relay 应 403,实际 %d:%s", rr.Code, rr.Body.String()) + } + if !strings.Contains(rr.Body.String(), "已丢弃") { + t.Fatalf("★ 文案必须明说「已丢弃」(否则发件方以为会补投而等待):%s", rr.Body.String()) + } + + // 关键:库里既没有邮件,**也没有留下待补投的占位** + if n := countMails(t, sid); n != baseCount { + t.Fatalf("★ relay 不得入库(否则 2h 后会被一次性灌回去):%d → %d 封", baseCount, n) + } + var keys int + if err := db.DB.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM relayed_mails WHERE relay_key = $1`, "relay-key-1").Scan(&keys); err != nil { + t.Fatal(err) + } + if keys != 0 { + t.Fatalf("★ 不得占用 relay 幂等键(占着会让重试拿到 duplicate 却从未发出):%d 行", keys) + } +} + +// ③ 人类插话立刻恢复 —— 冷静期不是死等 2 小时。 +func TestHumanPostClearsCooldownImmediately(t *testing.T) { + setupPermissionHandlerDB(t) + ctx := context.Background() + mustPermissionUser(t, "human1") + sid, mailID := lockSession(t) + + // 人类在**同一会话**里插一封(非 relay)—— 必须 reply_to 进去: + // 用 .new 会另建一条会话,解锁的就不是被锁的那条了(我第一版的错)。 + req := httptest.NewRequest("POST", "/api/v1/mail", + strings.NewReader(fmt.Sprintf( + `{"to":"ponger@","subject":"人来了","body":"停一下","reply_to":%q}`, mailID))) + req.Header.Set("Content-Type", "application/json") + req = req.WithContext(context.WithValue(context.Background(), middleware.AgentNameKey, "human1")) + rr := httptest.NewRecorder() + SendMail(rr, req) + if rr.Code/100 != 2 { + t.Fatalf("人类插话本身不应被拦:%d %s", rr.Code, rr.Body.String()) + } + + if _, locked, _ := repo.SessionLockOf(ctx, sid); locked { + t.Fatal("★ 人类参与后必须立刻解锁(不必等 2 小时)") + } +} + +// ④ 收件方是人类时**永远不该进**回路计数/冷静期。 +func TestHumanRecipientNeverLocked(t *testing.T) { + setupPermissionHandlerDB(t) + ctx := context.Background() + mustPermissionUser(t, "human1") + + rr := postSend(t, "pinger", "human1@.new", "请教一下", "", "") + if rr.Code/100 != 2 { + t.Fatalf("发往人类不应被回路闸拦:%d %s", rr.Code, rr.Body.String()) + } + var locks int + if err := db.DB.QueryRowContext(ctx, + `SELECT COUNT(*) FROM session_agent_locks`).Scan(&locks); err != nil { + t.Fatal(err) + } + if locks != 0 { + t.Fatalf("★ 发往人类不得产生锁(人类在回路里正是我们要的):%d 行", locks) + } +} \ No newline at end of file diff --git a/server/internal/repo/agentlock.go b/server/internal/repo/agentlock.go new file mode 100644 index 0000000..4ddcf59 --- /dev/null +++ b/server/internal/repo/agentlock.go @@ -0,0 +1,139 @@ +package repo + +import ( + "context" + "database/sql" + "errors" + "fmt" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/google/uuid" +) + +// Agent↔Agent 回路的**锁定期**(2026-10-01) +// +// # 缺口 +// +// `maxAgentPingPong = 8` 拦下"两个 Agent 互相确认而无人决策"之后, +// **计数永不回落**:只有人类插话才归零。2026-10-01 的报告就是那个形状—— +// +// 本会话已连续 8 封 Agent 之间互相回信、其中没有任何人类参与(上限 8)。 +// …请由人类在会话里插一句话(计数即归零)。 +// +// 而这条线索上**没有人类**(Agent 之间自己开的会话,或人已经不在), +// 于是那一侧**永久**发不出信,且文案把出路指向一个不存在的操作。 +// +// # 修法:给一个不依赖人的出口 +// +// 到 `until_at` 自动放行。语义是「冷静期」而不是「永久封禁」: +// 8 封无决策的互发多半是两个 Agent 在空转,2 小时足以让它们停下; +// 而真的需要继续时,2 小时后自己就能走,或由人插一句话立刻恢复。 +// +// # 为什么落库 +// +// 放内存会有两个问题:重启即"恢复"(而锁定期的意义正是别短时间内再炸一轮), +// 以及多副本部署时各算各的。`sessions` 表已有,改动落在它旁边而不是内存态。 +// +// # 人类插话仍然立刻归零 +// +// `SessionLockTouch`:人来信即清锁。这条优先级高于到期 —— +// 人参与了就不该再等冷静期。 +const AgentLockDuration = 2 * time.Hour + +// ErrSessionLocked 表示该会话处于 Agent 回路冷静期,Agent 的自主回信被拒。 +var ErrSessionLocked = errors.New("session is in agent-loop cooldown") + +// SessionLock 是一次锁定的状态(供错误文案与测试用)。 +type SessionLock struct { + Until time.Time + Hops int + // Remaining 是剩余锁定时长;<= 0 表示已到期或未锁定。 + Remaining time.Duration +} + +// Active 报告该锁定是否仍在生效。 +func (l SessionLock) Active() bool { return l.Until.After(time.Now()) } + +// SessionLockOf 返回该会话当前的锁定状态;未锁定返回零值 + ok=false。 +// +// 过期行**顺手删掉**:判据是 `until_at`,而清理只是回收空间 —— +// 留着它会让"是否锁定"的读数有两个答案(行在 / 行不在),而只有一个是真的。 +func SessionLockOf(ctx context.Context, sessionID uuid.UUID) (SessionLock, bool, error) { + var until time.Time + var hops int + err := db.DB.QueryRowContext(ctx, + `SELECT until_at, hops FROM session_agent_locks WHERE session_id = $1`, sessionID, + ).Scan(&until, &hops) + // ⚠ 必须区分「查不到」与「读失败」。 + // + // 我第一版把两者合并成 `return {}, false, nil`,于是 Scan 失败 + // 被读成"未锁定" ⇒ **锁完全失效,且看起来一切正常**。 + // 这正是我自己写在上面的那段警告,只是没做到。 + // 而把读失败当安全,正是安全系统最不该犯的错。 + if errors.Is(err, sql.ErrNoRows) { + return SessionLock{}, false, nil // 真的没锁:正常路径 + } + if err != nil { + // 读不出来 = **不知道**。不能当成"没锁"——那等于静默拆掉防护。 + // 调用方按"仍在锁"处理(宁可挡住,不可放行)。 + return SessionLock{}, true, fmt.Errorf("read session lock: %w", err) + } + + // ⚠ 时区:本库 DATETIME 列由 driver 以 **UTC** 扫出(实测读回 + // time.Date(..., time.UTC)),而 time.Now() 是本地时区(HKT,+8)。 + // 直接比会把「还剩 2 小时」算成「已过期 8 小时」⇒ 锁一设就失效。 + // 一律用 UTC 比较,与 driver 对齐。 + if !until.After(time.Now().UTC()) { + _, _ = db.DB.ExecContext(ctx, + `DELETE FROM session_agent_locks WHERE session_id = $1`, sessionID) + return SessionLock{}, false, nil + } + return SessionLock{Until: until, Hops: hops, Remaining: time.Until(until)}, true, nil +} + +// SetSessionLock 锁住该会话:已开始一次 `AgentLockDuration` 的冷静期。 +// +// 已有锁定期**不延长**(`until_at` 只会往前推靠 `GREATEST`,不是每次调用都 + +// 2h)。理由:反复触发时"每次 +2h"等于永不解锁 —— +// 一个每 30 分钟试一次的 Agent 会把会话永久锁住,而那正是我们要避免的。 +func SetSessionLock(ctx context.Context, sessionID uuid.UUID, hops int, reason string) (SessionLock, error) { + // ⚠ 时区:写入也用 UTC。本地时间(HKT)与 driver 读回的 UTC 混用时, + // 差 8 小时会让 2 小时的锁看起来像「已过期 6 小时」⇒ 锁形同虚设。 + now := time.Now().UTC() + until := now.Add(AgentLockDuration) + // ⚠ 已有锁定期**一律不延长**(真的不碰 until_at)。 + // + // 两个坑都踩过: + // ① SQLite 没有 GREATEST(实测 `no such function: GREATEST`), + // 而本项目两种方言都要支持(migrate.go 明确要求两处同改)。 + // ② 改用 CASE 比大小也不行 —— SQLite 在 DATETIME 上按**字符串**比, + // 而新算出来的 until 确实更晚 ⇒ 每次触发都往后推 1 秒。 + // + // 语义上「不延长」就应该是「不碰 until_at」:行存在即代表它生效中 + // (过期的行 SessionLockOf 读的时候已经顺手删了)。 + _, err := db.DB.ExecContext(ctx, ` + INSERT INTO session_agent_locks (session_id, locked_at, until_at, hops, reason) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (session_id) DO UPDATE + SET hops = EXCLUDED.hops, + reason = EXCLUDED.reason`, + sessionID, now, until, hops, reason) + if err != nil { + return SessionLock{}, err + } + return SessionLock{Until: until, Hops: hops, Remaining: AgentLockDuration}, nil +} + +// SessionLockTouch 清掉该会话的锁定 —— **人类插话即立刻恢复**。 +// +// 返回是否真的清掉了一条(false = 本来就没锁)。 +func SessionLockTouch(ctx context.Context, sessionID uuid.UUID) (bool, error) { + tag, err := db.DB.ExecContext(ctx, + `DELETE FROM session_agent_locks WHERE session_id = $1`, sessionID) + if err != nil { + return false, err + } + n, _ := tag.RowsAffected() + return n > 0, nil +} \ No newline at end of file diff --git a/server/internal/repo/agentlock_test.go b/server/internal/repo/agentlock_test.go new file mode 100644 index 0000000..ba88b32 --- /dev/null +++ b/server/internal/repo/agentlock_test.go @@ -0,0 +1,222 @@ +package repo + +import ( + "context" + "path/filepath" + "testing" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/google/uuid" +) + +/* +Agent↔Agent 回路的**冷静期**(2026-10-01) + +# 缺口 + +`maxAgentPingPong = 8` 撞闸之后**计数永不回落** —— 只有人类插话才归零。 +旧文案把出路指向"请由人类插一句话",而那条线索上**常常根本没有人类**, +于是那一侧永久发不出信(2026-10-01 报告:agent 一封都发不出)。 + +# 这几条钉的是"有出路",不是"有锁" + +`TestSessionLockBlocksThenRecovers` —— 锁住 / 到期放行 +`TestSessionLockNotExtendedByRetry` —— 重试不延长(否则永不恢复) +`TestHumanTouchClearsLock` —— 人插话立刻解锁 +`TestExpiredLockRowIsCleaned` —— 过期行清理(不留两个答案) +`TestSessionLockSurvivesReconnect` —— 落库而非放内存 +*/ + +func TestSessionLockBlocksThenRecovers(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, _ := CreateSession(ctx, nil, "human", "冷静期", "") + + // 初始:未锁定 + if _, locked, err := SessionLockOf(ctx, sid); err != nil || locked { + t.Fatalf("初始应为未锁定:locked=%v err=%v", locked, err) + } + + lock, err := SetSessionLock(ctx, sid, 8, "test") + if err != nil { + t.Fatal(err) + } + if !lock.Active() { + t.Fatal("刚设的锁应当生效") + } + + // 在锁:Agent 的自主回信被拒 + if _, locked, _ := SessionLockOf(ctx, sid); !locked { + t.Fatal("★ 锁定期内必须拦住(这正是「一封都发不出」那个缺陷)") + } + + // 到期:自动恢复 —— **不需要任何人介入** + if _, err := db.DB.ExecContext(ctx, + `UPDATE session_agent_locks SET until_at = ? WHERE session_id = ?`, + time.Now().Add(-time.Minute), sid); err != nil { + t.Fatal(err) + } + if _, locked, _ := SessionLockOf(ctx, sid); locked { + t.Fatal("★ 到期后必须自动放行 —— 恢复不能依赖人类在场") + } +} + +func TestSessionLockNotExtendedByRetry(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, _ := CreateSession(ctx, nil, "human", "重试", "") + + if _, err := SetSessionLock(ctx, sid, 8, "first"); err != nil { + t.Fatal(err) + } + first, _, _ := SessionLockOf(ctx, sid) + + // 一个每 30 分钟试一次的 Agent:若每次都 +2h ⇒ 永不解锁 + time.Sleep(1100 * time.Millisecond) + if _, err := SetSessionLock(ctx, sid, 9, "retry"); err != nil { + t.Fatal(err) + } + second, _, _ := SessionLockOf(ctx, sid) + + if second.Until.After(first.Until) { + t.Fatalf("★ 重试不得延长锁定期(否则反复触发=永不解锁):%s → %s", + first.Until, second.Until) + } + if second.Hops != 9 { + t.Fatalf("hops 应更新为最新值,实际 %d", second.Hops) + } +} + +func TestHumanTouchClearsLock(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, _ := CreateSession(ctx, nil, "human", "人类插话", "") + + if _, err := SetSessionLock(ctx, sid, 8, "test"); err != nil { + t.Fatal(err) + } + cleared, err := SessionLockTouch(ctx, sid) + if err != nil || !cleared { + t.Fatalf("人类插话应立刻清锁:cleared=%v err=%v", cleared, err) + } + if _, locked, _ := SessionLockOf(ctx, sid); locked { + t.Fatal("★ 人参与了就不该再等冷静期") + } + + // 再插一次:本来就没锁 ⇒ 不报错,只是"没清掉什么" + if cleared, err := SessionLockTouch(ctx, sid); err != nil || cleared { + t.Fatalf("未锁定时重复 touch 应为 no-op:cleared=%v err=%v", cleared, err) + } +} + +func TestExpiredLockRowIsCleaned(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, _ := CreateSession(ctx, nil, "human", "过期", "") + + if _, err := SetSessionLock(ctx, sid, 8, "test"); err != nil { + t.Fatal(err) + } + if _, err := db.DB.ExecContext(ctx, + `UPDATE session_agent_locks SET until_at = ? WHERE session_id = ?`, + time.Now().Add(-time.Hour), sid); err != nil { + t.Fatal(err) + } + if _, locked, _ := SessionLockOf(ctx, sid); locked { + t.Fatal("过期行应视为未锁定") + } + + // 读数只有一个答案:行也必须被回收 + var n int + if err := db.DB.QueryRowContext(ctx, + `SELECT COUNT(*) FROM session_agent_locks WHERE session_id = ?`, sid).Scan(&n); err != nil { + t.Fatal(err) + } + if n != 0 { + t.Fatalf("★ 过期行应顺手清理(留两个答案会让「是否锁定」不可判):剩 %d 行", n) + } +} + +func TestSessionLockSurvivesReconnect(t *testing.T) { + // 落库而非放内存:内存态一重启就"恢复",而锁定期的意义正是 + // 别在短时间内又炸一轮;且多副本部署时各算各的。 + dir := t.TempDir() + path := filepath.Join(dir, "lock.db") + ctx := context.Background() + + if err := db.Connect(ctx, path); err != nil { + t.Fatal(err) + } + if err := db.Migrate(ctx); err != nil { + t.Fatal(err) + } + sid, _ := CreateSession(ctx, nil, "human", "重启", "") + if _, err := SetSessionLock(ctx, sid, 8, "test"); err != nil { + t.Fatal(err) + } + db.Close() // 模拟重启 + + if err := db.Connect(ctx, path); err != nil { + t.Fatal(err) + } + t.Cleanup(db.Close) + if _, locked, err := SessionLockOf(ctx, sid); err != nil || !locked { + t.Fatalf("★ 重启后冷却期必须仍在(否则它只是个瞬时现象):locked=%v err=%v", locked, err) + } +} + +// 确认表真的建在两种方言的 schema 里(老库靠 migrate,新库靠 init.sql)。 +func TestSessionLockTableExistsOnSQLite(t *testing.T) { + setupTestDB(t) + var name string + if err := db.DB.QueryRowContext(context.Background(), + `SELECT name FROM sqlite_master WHERE type='table' AND name='session_agent_locks'`).Scan(&name); err != nil { + t.Fatalf("session_agent_locks 未建表:%v", err) + } +} + +/* +★ 读失败**不得**被读成"未锁定"(变异②补:原本 0 格能抓,漏网了) + +# 为什么单独一格 + +我把 `err != nil` 合并进 `return {}, false, nil`(即「读不出来 = 没锁」)。 +上面五条判据**全都抓不住它** —— 它们都不制造读失败,只是各走各的正常路径。 + +而这个形状是本仓记录过的同一个错:「`getWorkspace` 拿不到时不带,让服务端报 400」 +背后的原则是**宁可报错、不可静默降级**;这里反过来:**读不出来却静默放行**。 +把"不知道"读成"安全",等于静默拆掉防护,且**没有任何症状**。 + +# 怎么钉 + +让读真的失败(表不存在),然后断言「**仍报在锁**」—— +若实现退化成 `return {}, false, nil`,locked 会变 false ⇒ 这格变红。 +*/ +func TestSessionLockReadFailureIsNotTreatedAsUnlocked(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + sid, _ := CreateSession(ctx, nil, "human", "读失败", "") + + // 制造一次真实的读失败:把表改名,Scan 就会报 no such table + if _, err := db.DB.ExecContext(ctx, + `ALTER TABLE session_agent_locks RENAME TO session_agent_locks_gone`); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + _, _ = db.DB.ExecContext(ctx, + `ALTER TABLE session_agent_locks_gone RENAME TO session_agent_locks`) + }) + + lock, locked, err := SessionLockOf(ctx, sid) + if !locked { + t.Fatalf("★ 读失败必须当「仍在锁」而非「未锁定」"+ + "(否则一次 DB 抖动就静默拆掉整个防护,且无任何症状):locked=%v lock=%+v err=%v", + locked, lock, err) + } + if err == nil { + t.Fatal("读失败必须把错误往上带(调用方需要知道这是异常而不是「没锁」)") + } +} + +var _ = uuid.New \ No newline at end of file