fix(回路): Agent↔Agent 加 2h 冷静期;冷却期内 relay 不自动重投递

缺口: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,测试压根没跑)。
This commit is contained in:
2026-10-01 19:41:06 +08:00
parent 660bbd983a
commit 4d8165fde9
6 changed files with 671 additions and 3 deletions

View File

@ -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,

View File

@ -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,

View File

@ -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 仍能

View File

@ -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)
}
}

View File

@ -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
}

View File

@ -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