问题(P0):DSH 有两个独立的人机交互 seam —— approval/request(危险工具审批) 与 ask_user_question → ctx.userQuestions(模型主动提问)。原来只桥接了前者。 邮件驱动的会话没有本地 UI,而 ask() 的 provider 是 DSH host 注册的本地 UI 实现, 于是在那里等人点选永久等不到,那一轮工具调用**静默挂死**。 修法(不抢注全局 provider —— registerProvider 只允许一个活动实例,抢注会让 平台自己的界面失效):在 tools/execute around-dispatch 里只对**邮件驱动**的 会话接管 ask_user_question,其余原样 next()。失败一律当场报错而不是 next(): 下一个 answerer 是本地 UI,邮件会话没有兜底 UI,放过去就是挂死。 - lib/user-question.js(三桥逐字节同源,14 例测试):DSH questions[] ↔ AgentMail 单问题询问邮件的双向映射。多问题时把选项并集摊平、按 label 归属分配回各问题 (label 认不出来就不猜测放行);无选项题走自由文本 custom。 - Gateway:kind=question 且无选项时**不再**回落「同意/拒绝」(那会让自由文本 问题变成两个毫无意义的按钮);主题按类型区分「权限请求 / 需要回答」; 推送 payload 带上 permission_kind / multi_select / options。 - mails.permission_kind / permission_multi_select 此前只存在于结构体与写入路径, 五个读路径的 SELECT/Scan 都没带 —— 前端永远拿到空串,把提问渲染成批准/拒绝。 container 修正五处并加 repo 测试(含反向验证:删掉任一处字段,测试即失败)。 - 前端 PermissionPanel:question 走「勾选 + 自由文本」,多选/单选、空回答禁止提交; approval 路径不变(回归测试覆盖)。 测试:opencode 316 / dsh 349 / pi 405 / 前端 185 / Go 全量 全绿。
1524 lines
56 KiB
Go
1524 lines
56 KiB
Go
package repo
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/agentmail/gateway/internal/db"
|
||
"github.com/agentmail/gateway/internal/models"
|
||
"github.com/google/uuid"
|
||
)
|
||
|
||
// ---------- Agent ----------
|
||
|
||
// ErrAgentDisabled 表示该 Agent 已被管理员停用。
|
||
//
|
||
// 停用是可逆的「归档」:邮件、会话、权限记录全部保留,只是不再接受新任务。
|
||
// 与删除分开是因为往来邮件里有一半是人自己写的 —— 停用 Agent 不该删掉
|
||
// 用户的东西;而 Agent 名与人类用户名共用命名空间,历史邮件里的 from_name
|
||
// 指向一个已删除的名字时,下一个同名注册者会看起来像是当初的发信人。
|
||
var ErrAgentDisabled = errors.New("agent disabled")
|
||
|
||
func CreateOrUpdateAgent(ctx context.Context, name, secret, platform string, workspaces []models.Workspace) error {
|
||
// 已停用的 Agent 不得靠重新注册复活。
|
||
//
|
||
// 少了这一步的后果:管理员停用后,那个平台的插件下次启动就会重新注册
|
||
// (注册是插件启动流程的一部分),status 被写回 online,停用等于没做。
|
||
// 必须让插件收到一个明确的错误,而不是静默成功。
|
||
disabled, dErr := AgentDisabled(ctx, name)
|
||
if dErr != nil {
|
||
return dErr
|
||
}
|
||
if disabled {
|
||
return ErrAgentDisabled
|
||
}
|
||
|
||
wsJSON, _ := json.Marshal(workspaces)
|
||
// 注意 DO UPDATE 里【不】碰 default_rounds:
|
||
// 那是管理员配的值,Agent 重启重新注册不应该把它冲回默认。
|
||
_, err := db.DB.ExecContext(ctx, `
|
||
INSERT INTO agents (agent_name, secret, workspaces, platform, status, last_seen)
|
||
VALUES ($1, $2, $3, $4, 'online', NOW())
|
||
ON CONFLICT (agent_name) DO UPDATE SET
|
||
secret = EXCLUDED.secret,
|
||
workspaces = EXCLUDED.workspaces,
|
||
platform = EXCLUDED.platform,
|
||
status = 'online',
|
||
last_seen = NOW()
|
||
`, name, secret, wsJSON, platform)
|
||
return err
|
||
}
|
||
|
||
func HeartbeatAgent(ctx context.Context, agentName string) (int, error) {
|
||
// 只把【非停用】的 Agent 标成在线。
|
||
//
|
||
// 不加这个条件的话,停用后插件的心跳会把 status 从 disabled 改回 online
|
||
// —— 而心跳是每 30 秒一次的,停用最多维持半分钟。
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE agents SET last_seen = NOW(), status = 'online'
|
||
WHERE agent_name = $1 AND status <> 'disabled'`,
|
||
agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
return CountUnread(ctx, agentName)
|
||
}
|
||
|
||
// AgentDisabled 该 Agent 是否已被停用。Agent 不存在时返回 false ——
|
||
// 「还没注册」与「被停用」是两件事,前者应当能正常注册。
|
||
func AgentDisabled(ctx context.Context, agentName string) (bool, error) {
|
||
var status string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT status FROM agents WHERE agent_name = $1`, agentName).Scan(&status)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return false, nil
|
||
}
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
return status == "disabled", nil
|
||
}
|
||
|
||
// ErrRecipientUnknown 收件人既不是人类用户也不是在册 Agent。
|
||
var ErrRecipientUnknown = errors.New("recipient unknown")
|
||
|
||
// ErrRecipientDisabled 收件 Agent 已被管理员停用。
|
||
var ErrRecipientDisabled = errors.New("recipient disabled")
|
||
|
||
// RecipientDeliverable 校验一个收件人名当前能不能收信。
|
||
//
|
||
// 为什么必须有这道检查:发信路径原来只校验地址语法、调用权限与会话别名,
|
||
// 从不问「这个名字存在吗」。于是发给已删除或已停用的 Agent 一律返回 200 ——
|
||
// 邮件入库、分配预算、建好会话,而那一端永远不会有人读。发件人看到 200
|
||
// 和一个 session_id,以为送出去了。这是静默丢件,比报错严重:
|
||
// 报错能立刻改,静默丢件要等对方追问才发现。
|
||
//
|
||
// 三种放行/拒绝:
|
||
// - 人类用户 → 放行(人的收件箱一直在)
|
||
// - Agent 在册且未停用 → 放行
|
||
// - Agent 不存在 → ErrRecipientUnknown(对应 404)
|
||
// - Agent 已停用 → ErrRecipientDisabled(对应 409)
|
||
//
|
||
// 停用选择「当场拒收」而不是「入库等恢复后补投」:停用的语义就是这个 Agent
|
||
// 现在不干活,让发件人以为信已送达更坏 —— 它会照常等回信。
|
||
func RecipientDeliverable(ctx context.Context, name string) error {
|
||
if name == "" {
|
||
return nil // 空名由上层的地址解析负责
|
||
}
|
||
|
||
isHuman, err := IsHumanUser(ctx, name)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if isHuman {
|
||
return nil
|
||
}
|
||
|
||
var status string
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT status FROM agents WHERE agent_name = $1`, name).Scan(&status)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return ErrRecipientUnknown
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if status == "disabled" {
|
||
return ErrRecipientDisabled
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// SetAgentDisabled 停用或恢复一个 Agent。
|
||
//
|
||
// 停用时连带撤销它的全部密钥:留着密钥的话,那个平台的插件仍然能用它调
|
||
// /mail/send —— 停用的语义是「这个 Agent 不再参与工作」,不只是「不出现在
|
||
// 补全列表里」。恢复时不会把密钥变回来,管理员需要重新签发。
|
||
//
|
||
// 返回撤销的密钥数,供界面提示。
|
||
func SetAgentDisabled(ctx context.Context, agentName string, disabled bool) (int, error) {
|
||
if !disabled {
|
||
// 恢复:回到 offline 而不是 online —— 它是否真的在线由下一次心跳决定,
|
||
// 直接写 online 会让界面显示一个其实没在跑的 Agent 为在线。
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE agents SET status = 'offline' WHERE agent_name = $1 AND status = 'disabled'`,
|
||
agentName)
|
||
return 0, err
|
||
}
|
||
|
||
tx, err := db.DB.BeginTx(ctx, nil)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
res, err := tx.ExecContext(ctx,
|
||
`UPDATE agents SET status = 'disabled' WHERE agent_name = $1`, agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
if n, _ := res.RowsAffected(); n == 0 {
|
||
return 0, sql.ErrNoRows
|
||
}
|
||
|
||
keyRes, err := tx.ExecContext(ctx,
|
||
`DELETE FROM agent_keys WHERE agent_name = $1`, agentName)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
revoked, _ := keyRes.RowsAffected()
|
||
|
||
if err := tx.Commit(); err != nil {
|
||
return 0, err
|
||
}
|
||
return int(revoked), nil
|
||
}
|
||
|
||
// ListAgents 列出 Agent。
|
||
//
|
||
// statusFilter 为空时默认**排除已停用的** —— 这个函数的三个调用点
|
||
// (地址补全、GET /agents、可授权范围)都是在回答「现在能派活给谁」,
|
||
// 而停用的 Agent 不该出现在那里。要连停用一起看,传 statusFilter="all"。
|
||
func ListAgents(ctx context.Context, statusFilter string) ([]models.Agent, error) {
|
||
// 带上 default_rounds:前端补全收件人时要显示「派给它的任务默认几个来回」,
|
||
// 否则人得先去管理员页查一遍才敢派活。
|
||
q := `SELECT agent_id, agent_name, workspaces, platform, status,
|
||
COALESCE(default_rounds, 0) FROM agents`
|
||
args := []any{}
|
||
switch statusFilter {
|
||
case "":
|
||
q += ` WHERE status <> 'disabled'`
|
||
case "all":
|
||
// 不加条件
|
||
default:
|
||
q += ` WHERE status = $1`
|
||
args = append(args, statusFilter)
|
||
}
|
||
q += ` ORDER BY agent_name`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
agents := []models.Agent{}
|
||
for rows.Next() {
|
||
var a models.Agent
|
||
var wsJSON []byte
|
||
if err := rows.Scan(&a.ID, &a.Name, &wsJSON, &a.Platform, &a.Status,
|
||
&a.DefaultRounds); err != nil {
|
||
return nil, err
|
||
}
|
||
if wsJSON != nil {
|
||
json.Unmarshal(wsJSON, &a.Workspaces)
|
||
}
|
||
agents = append(agents, a)
|
||
}
|
||
return agents, nil
|
||
}
|
||
|
||
func VerifyAgent(ctx context.Context, name, secret string) (*models.Agent, error) {
|
||
var a models.Agent
|
||
var wsJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT agent_id, agent_name, workspaces, platform, status
|
||
FROM agents WHERE agent_name = $1 AND secret = $2`,
|
||
name, secret,
|
||
).Scan(&a.ID, &a.Name, &wsJSON, &a.Platform, &a.Status)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if wsJSON != nil {
|
||
json.Unmarshal(wsJSON, &a.Workspaces)
|
||
}
|
||
return &a, nil
|
||
}
|
||
|
||
// ---------- Session ----------
|
||
|
||
// CreateSession 建会话。
|
||
//
|
||
// 显式传了 alias(发信时的 session_alias 参数)= 调用方亲自命名,标为 manual,
|
||
// 平台后续自动同步不得覆盖;未传则等待平台命名,标为 platform。
|
||
func CreateSession(ctx context.Context, alias *string, fromAgent, subject, workspace string) (uuid.UUID, error) {
|
||
source := "platform"
|
||
if alias != nil && *alias != "" {
|
||
source = "manual"
|
||
}
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO sessions (session_alias, from_agent, subject, alias_source, workspace)
|
||
VALUES ($1, $2, $3, $4, $5) RETURNING session_id`,
|
||
alias, fromAgent, subject, source, strings.TrimSpace(workspace),
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
func GetSessionByID(ctx context.Context, id uuid.UUID) (*models.Session, error) {
|
||
var s models.Session
|
||
var dismissed *string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id, session_alias, from_agent, subject, status, owner_user_id,
|
||
created_at, updated_at, rename_dismissed, COALESCE(alias_source, 'platform'),
|
||
COALESCE(max_rounds, 0), COALESCE(used_rounds, 0),
|
||
COALESCE(NULLIF(permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(permission_enforcement, ''), 'advisory')
|
||
FROM sessions WHERE session_id = $1`, id,
|
||
).Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status, &s.OwnerUserID,
|
||
&s.CreatedAt, &s.UpdatedAt, &dismissed, &s.AliasSource,
|
||
&s.MaxRounds, &s.UsedRounds,
|
||
&s.PermissionMode, &s.PermissionEnforcement)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if dismissed != nil {
|
||
s.RenameDismissed = *dismissed
|
||
}
|
||
return &s, nil
|
||
}
|
||
|
||
func TouchSession(ctx context.Context, id uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET updated_at = NOW(), status = 'active' WHERE session_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
// UpdateSessionAlias 手工改名(人显式指定)。
|
||
//
|
||
// 同时把 alias_source 标为 'manual':人的选择优先于平台自动命名。
|
||
// 否则平台下一次 session.updated 会把人刚定的名字冲掉,
|
||
// 人上一秒记住的寻址地址下一秒失效。
|
||
func UpdateSessionAlias(ctx context.Context, id uuid.UUID, alias string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET session_alias = $1, alias_source = 'manual', updated_at = NOW()
|
||
WHERE session_id = $2`,
|
||
alias, id)
|
||
return err
|
||
}
|
||
|
||
// ---------- Mail ----------
|
||
|
||
func CreateMail(ctx context.Context, sessionID uuid.UUID, parentMailID *uuid.UUID,
|
||
fromName, fromWorkspace, toName, toWorkspace, subject, body string, ccList []models.Address) (uuid.UUID, error) {
|
||
if ccList == nil {
|
||
ccList = []models.Address{}
|
||
}
|
||
ccJSON, _ := json.Marshal(ccList)
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
// created_at 显式给 NOW():SQLite 的 DEFAULT CURRENT_TIMESTAMP 只有秒精度,
|
||
// 同秒插入的多封邮件排序不确定(「会话里最早/最后那封」都会取错行)。
|
||
// 改 schema 的默认值只对新库生效 —— CREATE TABLE IF NOT EXISTS 不改已存在的表,
|
||
// 而 SQLite 没有 ALTER COLUMN,因此这里显式传。
|
||
`INSERT INTO mails (session_id, parent_mail_id, from_name, from_workspace,
|
||
to_name, to_workspace, subject, body, cc_list, created_at)
|
||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, NOW()) RETURNING mail_id`,
|
||
sessionID, parentMailID, fromName, fromWorkspace, toName, toWorkspace, subject, body, ccJSON,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// 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, kind string, multiSelect bool) (uuid.UUID, error) {
|
||
optsJSON, _ := json.Marshal(options)
|
||
var multiSelectInt int
|
||
if multiSelect {
|
||
multiSelectInt = 1
|
||
}
|
||
var id uuid.UUID
|
||
// 主题按待办类型区分:「权限请求」意味着有人要放行一个危险操作,
|
||
// 「需要回答」只是模型缺信息。授权页与收件箱都靠这一行措辞判断该做什么。
|
||
subject := "权限请求: " + question
|
||
if kind == "question" {
|
||
subject = "需要回答: " + question
|
||
}
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO mails (session_id, from_name, from_workspace, to_name, subject, body, mail_type, permission_options, permission_kind, permission_multi_select, created_at)
|
||
VALUES ($1, $2, COALESCE((SELECT workspace FROM sessions WHERE session_id = $1), ''), $3, $4, $5, 'permission_request', $6, $7, $8, NOW()) RETURNING mail_id`,
|
||
sessionID, fromName, toUser, subject, body, optsJSON, kind, multiSelectInt,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// CreateDecisionMail 创建人类决策邮件(fromUser → toAgent)
|
||
func CreateDecisionMail(ctx context.Context, sessionID uuid.UUID, parentMailID uuid.UUID, fromUser, toAgent, decision, note string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
body := decision
|
||
if note != "" {
|
||
body = fmt.Sprintf("%s\n\n备注: %s", decision, note)
|
||
}
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`INSERT INTO mails (session_id, parent_mail_id, from_name, to_name, subject, body, created_at)
|
||
VALUES ($1, $2, $3, $4, $5, $6, NOW()) RETURNING mail_id`,
|
||
sessionID, parentMailID, fromUser, toAgent, "Re: 权限请求 - "+decision, body,
|
||
).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// DeleteMailByID 删一封邮件。
|
||
//
|
||
// **只用于回滚一次刚失败的发信**,不是给人用的「删邮件」功能 ——
|
||
// 邮件是不可篡改的历史记录,没有任何人面入口能删它。
|
||
//
|
||
// 场景:附件挂载在建邮件之后才发现冲突(竞态窗口),此时这封邮件不应存在:
|
||
// 发件方收到的是 4xx,它会重试,而一封无附件的残余邮件会让收件方收到两封。
|
||
//
|
||
// attachments 表的外键是 ON DELETE CASCADE,所以已经挂上去的那几条会跟着消失;
|
||
// relayed_mails 的 mail_id 无 CASCADE,由调用方用 ReleaseRelay 归还幂等键。
|
||
func DeleteMailByID(ctx context.Context, id uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx, `DELETE FROM mails WHERE mail_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
func GetMailByID(ctx context.Context, id uuid.UUID) (*models.Mail, error) {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
var renameAlias, renameReason *string
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace, m.rename_alias, m.rename_reason,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.mail_id = $1`, id,
|
||
).Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &renameAlias, &renameReason,
|
||
&m.FromHuman, &m.ToHuman)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
// 改名提议随单封返回,让「谁在哪一封里提了什么」可追溯;
|
||
// 【待处理】的提议另有专用端点(GET /sessions/:id/rename-proposal)。
|
||
if renameAlias != nil {
|
||
m.RenameAlias = *renameAlias
|
||
}
|
||
if renameReason != nil {
|
||
m.RenameReason = *renameReason
|
||
}
|
||
return &m, nil
|
||
}
|
||
|
||
func MarkMailRead(ctx context.Context, id uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx, `UPDATE mails SET status = 'read' WHERE mail_id = $1`, id)
|
||
return err
|
||
}
|
||
|
||
func ListInbox(ctx context.Context, agentName, status string, limit int) ([]models.Mail, error) {
|
||
q := `SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human,
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace') AS permission_mode,
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory') AS permission_enforcement
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE (m.to_name = $1 OR ` + db.CCHas("m.cc_list", 1) + `)
|
||
AND s.status <> 'archived'`
|
||
args := []any{agentName}
|
||
if status != "" && status != "all" {
|
||
q += ` AND m.status = $2`
|
||
args = append(args, status)
|
||
}
|
||
q += ` ORDER BY m.created_at DESC, m.mail_id DESC`
|
||
if limit > 0 {
|
||
q += fmt.Sprintf(` LIMIT %d`, limit)
|
||
}
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman,
|
||
&m.PermissionMode, &m.PermissionEnforcement); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
// Body preview
|
||
if len(m.Body) > 200 {
|
||
m.BodyPreview = m.Body[:200] + "..."
|
||
} else {
|
||
m.BodyPreview = m.Body
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
func CountUnread(ctx context.Context, agentName string) (int, error) {
|
||
var count int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*)
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE (m.to_name = $1 OR `+db.CCHas("m.cc_list", 1)+`)
|
||
AND m.status = 'unread'
|
||
AND s.status <> 'archived'
|
||
`, agentName).Scan(&count)
|
||
return count, err
|
||
}
|
||
|
||
func GetSessionMails(ctx context.Context, sessionID uuid.UUID) ([]models.Mail, error) {
|
||
rows, err := db.DB.QueryContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.session_id = $1
|
||
ORDER BY m.created_at ASC, m.mail_id ASC`, sessionID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
// ---------- Permission ----------
|
||
|
||
func CreatePermissionRequest(ctx context.Context, mailID, sessionID uuid.UUID, agentName, question string, options []string, contextStr string, kind string, multiSelect bool) error {
|
||
optsJSON, _ := json.Marshal(options)
|
||
var multiSelectInt int
|
||
if multiSelect {
|
||
multiSelectInt = 1
|
||
}
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`INSERT INTO permission_requests (mail_id, session_id, agent_name, question, options, context, kind, multi_select)
|
||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`,
|
||
mailID, sessionID, agentName, question, optsJSON, contextStr, kind, multiSelectInt)
|
||
return err
|
||
}
|
||
|
||
func DecidePermission(ctx context.Context, mailID uuid.UUID, decision string) (*models.PermissionRequest, error) {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`UPDATE permission_requests SET result = $1, decided_at = NOW()
|
||
WHERE mail_id = $2
|
||
RETURNING request_id, mail_id, session_id, agent_name, question, options, context, result, decided_at, created_at`,
|
||
decision, mailID,
|
||
).Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Result, &pr.DecidedAt, &pr.CreatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
|
||
// Also update the mail
|
||
_, _ = db.DB.ExecContext(context.Background(),
|
||
`UPDATE mails SET permission_result = $1, status = 'read' WHERE mail_id = $2`,
|
||
decision, mailID)
|
||
|
||
return &pr, nil
|
||
}
|
||
|
||
func ListPendingPermissions(ctx context.Context) ([]models.PermissionRequest, error) {
|
||
rows, err := db.DB.QueryContext(ctx,
|
||
`SELECT request_id, mail_id, session_id, agent_name, question, options, context, kind, multi_select, result, decided_at, created_at
|
||
FROM permission_requests WHERE result IS NULL
|
||
ORDER BY created_at DESC`)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
reqs := []models.PermissionRequest{}
|
||
for rows.Next() {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
if err := rows.Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
reqs = append(reqs, pr)
|
||
}
|
||
return reqs, nil
|
||
}
|
||
|
||
// ---------- Check permission ownership ----------
|
||
|
||
func GetPermissionByMailID(ctx context.Context, mailID uuid.UUID) (*models.PermissionRequest, error) {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT request_id, mail_id, session_id, agent_name, question, options, context, kind, multi_select, result, decided_at, created_at
|
||
FROM permission_requests WHERE mail_id = $1`, mailID,
|
||
).Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
return &pr, nil
|
||
}
|
||
|
||
func GetSessionMailByID(ctx context.Context, sessionID, mailID uuid.UUID) (*models.Mail, error) {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.session_id = $1 AND m.mail_id = $2`, sessionID, mailID,
|
||
).Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
return &m, nil
|
||
}
|
||
|
||
func FindSessionByAlias(ctx context.Context, alias string) (*models.Session, error) {
|
||
var s models.Session
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id, session_alias, from_agent, subject, status, owner_user_id, created_at, updated_at
|
||
FROM sessions WHERE session_alias = $1 AND status <> 'archived'`, alias,
|
||
).Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status, &s.OwnerUserID, &s.CreatedAt, &s.UpdatedAt)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &s, nil
|
||
}
|
||
|
||
// ErrSessionNotFound 表示三维地址里指定的 session 别名不存在(或不属于该收件人)。
|
||
// 调用方应据此回 404「无法送达」,而不是悄悄新建一个会话。
|
||
var ErrSessionNotFound = errors.New("session not found")
|
||
|
||
// FindNamedSessionFor 查找收件人 name@path 名下别名为 alias 的会话。
|
||
// 严格匹配:会话必须存在、未归档,且该收件人确实参与过该会话,否则返回 ErrSessionNotFound。
|
||
// FindNamedSessionFor 实现 session 位给具体别名时的语义:必须已存在。
|
||
//
|
||
// 不限定 workspace:别名全局唯一且本身就承担寻址职责,
|
||
// 再叠一层工作区校验只会让「名字对上了却送不到」变成一种难查的失败。
|
||
func FindNamedSessionFor(ctx context.Context, name, path, alias string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT s.session_id
|
||
FROM sessions s
|
||
WHERE s.session_alias = $1
|
||
AND s.status <> 'archived'
|
||
AND EXISTS (
|
||
SELECT 1 FROM mails m
|
||
WHERE m.session_id = s.session_id
|
||
AND (m.to_name = $2 OR m.from_name = $2 OR `+db.CCHas("m.cc_list", 2)+`)
|
||
)
|
||
ORDER BY s.updated_at DESC
|
||
LIMIT 1
|
||
`, alias, name).Scan(&id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return uuid.Nil, ErrSessionNotFound
|
||
}
|
||
return id, err
|
||
}
|
||
|
||
// FindOrCreateDefaultSession 实现 session 位省略时的「默认会话」语义:
|
||
// 复用 name@path 最近一次活跃的非归档会话;从未通过信则建立一个新的作为默认会话。
|
||
//
|
||
// 匹配工作区优先看 sessions.workspace(权威来源),旧会话那列为空时回退到
|
||
// mails.to_workspace 反推。只看 to_workspace:Agent 回信时 from_workspace 存的是
|
||
// Agent 名而不是路径,拿它比路径永远匹配不上(旧实现就挂在这里)。
|
||
func FindOrCreateDefaultSession(ctx context.Context, name, path, fromAgent, subject string) (uuid.UUID, error) {
|
||
id, _, err := FindOrCreateDefaultSessionCreated(ctx, name, path, fromAgent, subject)
|
||
return id, err
|
||
}
|
||
|
||
// FindOrCreateDefaultSessionCreated 与 FindOrCreateDefaultSession 相同,但额外返回
|
||
// **这次调用是否真的新建了会话**。
|
||
//
|
||
// 为什么需要这个返回值:调用方此前用 `parentMailID == nil` 判断「是不是新建会话」,
|
||
// 而复用已有默认会话时 parentMailID 也是 nil —— 于是「仅在新建时生效」的字段
|
||
// (往返预算、权限档位)在每一封省略 session 位的信上都被重写了。
|
||
// 实测:第一封 max_rounds=7 → 第二封省略该字段 → 预算被静默改成默认的 20。
|
||
func FindOrCreateDefaultSessionCreated(ctx context.Context, name, path, fromAgent, subject string) (uuid.UUID, bool, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
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
|
||
AND (m.to_name = $1 OR m.from_name = $1 OR `+db.CCHas("m.cc_list", 1)+`)
|
||
)
|
||
AND (s.workspace = $2
|
||
OR (s.workspace = '' AND EXISTS (
|
||
SELECT 1 FROM mails w
|
||
WHERE w.session_id = s.session_id
|
||
AND COALESCE(w.to_workspace,'') = $2
|
||
)))
|
||
ORDER BY s.updated_at DESC
|
||
LIMIT 1
|
||
`, name, path).Scan(&id)
|
||
if err == nil {
|
||
TouchSession(ctx, id)
|
||
return id, false, nil
|
||
}
|
||
if !errors.Is(err, sql.ErrNoRows) {
|
||
return uuid.Nil, false, err
|
||
}
|
||
newID, cErr := CreateSession(ctx, nil, fromAgent, subject, path)
|
||
return newID, cErr == nil, cErr
|
||
}
|
||
|
||
// SessionAliasOf 返回会话别名,未命名或查询失败时返回空串。
|
||
// 仅用于响应体回显,不影响投递路径,所以吞错是可接受的。
|
||
func SessionAliasOf(ctx context.Context, id uuid.UUID) string {
|
||
var alias *string
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_alias FROM sessions WHERE session_id = $1`, id).Scan(&alias); err != nil {
|
||
return ""
|
||
}
|
||
if alias == nil {
|
||
return ""
|
||
}
|
||
return *alias
|
||
}
|
||
|
||
// AgentCanAccessSession 判断 Agent 是否参与过该会话(发件/收件/被抄送)。
|
||
// Agent 只能改自己参与的会话的别名,避免跨会话改名。
|
||
func AgentCanAccessSession(ctx context.Context, agentName string, sessionID uuid.UUID) (bool, error) {
|
||
var n int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*) FROM mails m
|
||
WHERE m.session_id = $1
|
||
AND (m.from_name = $2 OR m.to_name = $2 OR `+db.CCHas("m.cc_list", 2)+`)
|
||
`, sessionID, agentName).Scan(&n)
|
||
return n > 0, err
|
||
}
|
||
|
||
// SyncSessionTitle 更新会话主题(Agent 平台生成的摘要标题)。
|
||
func SyncSessionTitle(ctx context.Context, id uuid.UUID, title string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET subject = $1, updated_at = NOW() WHERE session_id = $2`,
|
||
title, id)
|
||
return err
|
||
}
|
||
|
||
// SyncSessionAlias 把 Agent 平台侧的会话标识写为本侧别名。
|
||
// 平台侧标识(如 opencode 的 slug)在平台内不保证全局唯一,而本侧别名负责寻址必须唯一,
|
||
// 因此撞名时自动追加 -2、-3… 后缀而不是报错——同步是后台行为,不该因撞名失败。
|
||
// 返回最终落库的别名。该会话已持有目标别名时直接返回,不做无谓写入。
|
||
//
|
||
// **人显式定过的别名不覆盖**(alias_source = 'manual'):
|
||
// 用户刚接受了 Agent 的改名提议,或手工敲了一个名字,平台下一次 session.updated
|
||
// 不该把它冲掉 —— 那会让人上一秒记住的寻址地址下一秒失效。
|
||
// 此时返回当前别名,调用方据此知道同步未生效。
|
||
func SyncSessionAlias(ctx context.Context, id uuid.UUID, want string) (string, error) {
|
||
const maxAttempts = 50
|
||
|
||
var cur *string
|
||
var source string
|
||
var platformID string
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`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
|
||
if i > 0 {
|
||
candidate = fmt.Sprintf("%s-%d", want, i+1)
|
||
}
|
||
|
||
owner, err := aliasOwner(ctx, candidate)
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
if owner != nil {
|
||
if *owner == id {
|
||
return candidate, nil // 已经是这个别名,无需写入
|
||
}
|
||
continue // 被别人占用,试下一个后缀
|
||
}
|
||
|
||
// 只在仍是 platform 来源时写入:并发下用户可能刚好接受了改名提议,
|
||
// 条件放进 WHERE 才能保证「检查」与「写入」不被插进来的手工改名割开
|
||
res, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET session_alias = $1, updated_at = NOW()
|
||
WHERE session_id = $2 AND COALESCE(alias_source, 'platform') <> 'manual'`,
|
||
candidate, id)
|
||
if err == nil {
|
||
if n, _ := res.RowsAffected(); n == 0 {
|
||
// 期间变成 manual 了,尊重人的选择
|
||
return SessionAliasOf(ctx, id), nil
|
||
}
|
||
return candidate, nil
|
||
}
|
||
// 并发下另一个请求刚占走该别名(唯一索引拦下),继续试下一个后缀
|
||
if db.IsUniqueViolation(err) {
|
||
continue
|
||
}
|
||
return "", err
|
||
}
|
||
return "", fmt.Errorf("alias %q: 连同 -2..-%d 后缀均被占用", want, maxAttempts)
|
||
}
|
||
|
||
// aliasOwner 返回持有该别名的会话 ID;无人持有时返回 nil。
|
||
func aliasOwner(ctx context.Context, alias string) (*uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT session_id FROM sessions WHERE session_alias = $1`, alias).Scan(&id)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return nil, nil
|
||
}
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &id, nil
|
||
}
|
||
|
||
func SessionMailCount(ctx context.Context, sessionID uuid.UUID) (int, error) {
|
||
var count int
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM mails WHERE session_id = $1`, sessionID).Scan(&count)
|
||
return count, err
|
||
}
|
||
|
||
// ---------- Contacts / Archive ----------
|
||
|
||
// Contact 是「联系人」= 一条 name@path.session 三维地址
|
||
type Contact struct {
|
||
SessionID uuid.UUID `json:"session_id"`
|
||
AgentName string `json:"agent_name"`
|
||
Path string `json:"path"`
|
||
SessionAlias string `json:"session_alias"`
|
||
Address string `json:"address"` // name@path.session
|
||
Status string `json:"status"`
|
||
MailCount int `json:"mail_count"`
|
||
UnreadCount int `json:"unread_count"`
|
||
LastActivity time.Time `json:"last_activity"`
|
||
|
||
// Subject 是会话主题(多由 Agent 平台的模型生成的摘要)。
|
||
// 卡片视图要靠它回答「这条线索在干什么」—— 光有 name@path.alias
|
||
// 只能看出跟谁在聊,看不出在聊什么。
|
||
Subject string `json:"subject"`
|
||
|
||
// MaxRounds/UsedRounds 是本任务的往返预算(0 = 不限)。
|
||
// 列表上直接可见,才不用点进每条会话去查哪件事快跑满了。
|
||
MaxRounds int `json:"max_rounds"`
|
||
UsedRounds int `json:"used_rounds"`
|
||
|
||
// PermissionMode 与 PermissionEnforcement 必须成对出现在列表上:
|
||
// 前者是「这条任务要求什么」,后者是「对方平台实际做到了什么」。
|
||
// 只显示前者会让人以为 plan 档把 homeagent 管住了(它没有拦截点)。
|
||
PermissionMode string `json:"permission_mode"`
|
||
PermissionEnforcement string `json:"permission_enforcement"`
|
||
|
||
// LastFrom/LastPreview 是最后一封邮件的发件人与正文摘要,
|
||
// 卡片视图用它显示「最新进展」——列表视图只显示地址时,
|
||
// 人必须逐条点开才知道哪条有新动静。
|
||
LastFrom string `json:"last_from"`
|
||
LastPreview string `json:"last_preview"`
|
||
}
|
||
|
||
// ListContactsFor 按 (agent, path, session) 聚合出联系人清单。
|
||
// forUser 非空时只列该用户参与的会话(owner / 收发 / 抄送);空表示不限(管理员全局视图)。
|
||
// archived=false 只列活跃会话,true 只列归档会话。
|
||
func ListContactsFor(ctx context.Context, forUser string, archived bool) ([]Contact, error) {
|
||
op := "<>"
|
||
if archived {
|
||
op = "="
|
||
}
|
||
scope := ""
|
||
args := []any{}
|
||
if forUser != "" {
|
||
scope = ` AND (s.owner_user_id = (SELECT user_id FROM users WHERE username = $1)
|
||
OR EXISTS (
|
||
SELECT 1 FROM mails mm
|
||
WHERE mm.session_id = s.session_id
|
||
AND (mm.from_name = $1 OR mm.to_name = $1
|
||
OR ` + db.CCHas("mm.cc_list", 1) + `)
|
||
))`
|
||
args = append(args, forUser)
|
||
}
|
||
// 取会话里最早那封邮件作为联系人身份。
|
||
// PG 用 LATERAL 子查询;SQLite 无 LATERAL,改用关联子查询逐列取值
|
||
// (同一个 min(created_at) 子句,四列取自同一行)。
|
||
var firstMail string
|
||
if db.D == db.Postgres {
|
||
firstMail = `
|
||
JOIN LATERAL (
|
||
SELECT to_name, to_workspace, from_name, from_workspace
|
||
FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at ASC
|
||
LIMIT 1
|
||
) m ON TRUE`
|
||
} else {
|
||
firstMail = `
|
||
JOIN mails m ON m.mail_id = (
|
||
SELECT mail_id FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at ASC, mail_id ASC
|
||
LIMIT 1
|
||
)`
|
||
}
|
||
|
||
// 最后一封邮件用关联子查询取,不再 JOIN 一次:
|
||
// 两个 JOIN(最早一封 + 最新一封)在 SQLite 下要写两段方言分支,
|
||
// 而这里每个会话只多两次索引查找(idx_mails_session 已有)。
|
||
lastMail := `(
|
||
SELECT mail_id FROM mails
|
||
WHERE session_id = s.session_id
|
||
ORDER BY created_at DESC, mail_id DESC
|
||
LIMIT 1
|
||
)`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT s.session_id,
|
||
COALESCE(NULLIF(m.to_name, 'human'), m.from_name) AS agent_name,
|
||
COALESCE(NULLIF(m.to_workspace, ''), m.from_workspace) AS path,
|
||
COALESCE(s.session_alias, '') AS alias,
|
||
s.status,
|
||
(SELECT COUNT(*) FROM mails x WHERE x.session_id = s.session_id),
|
||
(SELECT COUNT(*) FROM mails x WHERE x.session_id = s.session_id AND x.status = 'unread'),
|
||
s.updated_at,
|
||
s.subject,
|
||
COALESCE(s.max_rounds, 0),
|
||
COALESCE(s.used_rounds, 0),
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory'),
|
||
COALESCE((SELECT from_name FROM mails WHERE mail_id = `+lastMail+`), ''),
|
||
COALESCE((SELECT body FROM mails WHERE mail_id = `+lastMail+`), '')
|
||
FROM sessions s`+firstMail+`
|
||
WHERE s.status `+op+` 'archived'`+scope+`
|
||
ORDER BY s.updated_at DESC
|
||
`, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
contacts := []Contact{}
|
||
for rows.Next() {
|
||
var c Contact
|
||
if err := rows.Scan(&c.SessionID, &c.AgentName, &c.Path, &c.SessionAlias,
|
||
&c.Status, &c.MailCount, &c.UnreadCount, &c.LastActivity,
|
||
&c.Subject, &c.MaxRounds, &c.UsedRounds,
|
||
&c.PermissionMode, &c.PermissionEnforcement,
|
||
&c.LastFrom, &c.LastPreview); err != nil {
|
||
return nil, err
|
||
}
|
||
c.Address = c.AgentName + "@" + c.Path
|
||
if c.SessionAlias != "" {
|
||
c.Address += "." + c.SessionAlias
|
||
}
|
||
// 正文只留一段摘要:卡片上放不下全文,而整个列表带全文可能几百 KB。
|
||
// 按 rune 截断而非字节 —— 中文 3 字节/字,裸切会产生 U+FFFD。
|
||
c.LastPreview = previewRunes(c.LastPreview, 90)
|
||
contacts = append(contacts, c)
|
||
}
|
||
return contacts, rows.Err()
|
||
}
|
||
|
||
// previewRunes 按字符数截断,附省略号。
|
||
// 不按字节切:中文一字三字节,裸切会在末尾留半个字符(渲染成 U+FFFD)。
|
||
func previewRunes(s string, n int) string {
|
||
s = strings.TrimSpace(s)
|
||
rs := []rune(s)
|
||
if len(rs) <= n {
|
||
return s
|
||
}
|
||
return string(rs[:n]) + "…"
|
||
}
|
||
|
||
// ArchiveSession 归档一个会话(邮箱界面不再展示,数据保留)
|
||
func ArchiveSession(ctx context.Context, sessionID uuid.UUID) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET status = 'archived', updated_at = NOW() WHERE session_id = $1`,
|
||
sessionID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
// 同时把该会话下的邮件标记为已归档,收件箱不再列出
|
||
_, err = db.DB.ExecContext(ctx,
|
||
`UPDATE mails SET status = 'archived' WHERE session_id = $1 AND status <> 'archived'`,
|
||
sessionID)
|
||
return err
|
||
}
|
||
|
||
// FindSessionByAddress 按 name@path.session 定位会话
|
||
func FindSessionByAddress(ctx context.Context, name, path, alias string) (uuid.UUID, error) {
|
||
var id uuid.UUID
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT s.session_id
|
||
FROM sessions s
|
||
JOIN mails m ON m.session_id = s.session_id
|
||
WHERE COALESCE(s.session_alias, '') = $1
|
||
AND (m.to_name = $2 OR m.from_name = $2)
|
||
AND (COALESCE(m.to_workspace,'') = $3 OR COALESCE(m.from_workspace,'') = $3)
|
||
LIMIT 1
|
||
`, alias, name, path).Scan(&id)
|
||
return id, err
|
||
}
|
||
|
||
// SuggestPaths 给出某个收件方可用的工作目录候选(三段式补全的 path 位)。
|
||
//
|
||
// 三个来源并集,**历史优先**:
|
||
//
|
||
// 1. mails.to_workspace 里真实投递过的目录(按最近使用倒序)
|
||
// 2. agent_platform_sessions.workspace 里心跳上报的目录(镜像)
|
||
// 3. agents.workspaces 里注册时自报的目录
|
||
//
|
||
// 早先只看第 1、3 项,而第 1 项在清库后为空,第 3 项对两个官方插件
|
||
// **永远是空的** —— 契约(W-2)明确要求 `workspaces: []`,
|
||
// 因为工作目录由每封邮件的 to_workspace 决定而不是注册时固定。
|
||
// 于是补全的第二段在清库后对所有 Agent 都给空列表。
|
||
//
|
||
// 加入第 2 项:镜像是心跳上报的平台会话快照,其中 workspace 字段
|
||
// 就是各平台的真实工作目录。它在清库后仍然存在(心跳重新上报)。
|
||
//
|
||
// 按最近使用倒序而非字典序:同一个 Agent 常在几个项目间切,刚用过的那个
|
||
// 几乎总是下一封想用的那个。
|
||
func SuggestPaths(ctx context.Context, agentName string) ([]string, error) {
|
||
out := []string{}
|
||
seen := map[string]bool{}
|
||
|
||
// 来源 1:真实投递历史。即使 Agent 未注册(人类收件方)也能给出候选。
|
||
//
|
||
// 只 SELECT 路径一列,排序用的 MAX(created_at) 不进结果集:
|
||
// SQLite 把时间戳存为 TEXT,把它 Scan 进 time.Time 会失败,
|
||
// 而那个值除了排序之外无用 —— 取回来只是多一个能静默失败的环节。
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT to_workspace
|
||
FROM mails
|
||
WHERE to_name = $1 AND COALESCE(to_workspace, '') <> ''
|
||
GROUP BY to_workspace
|
||
ORDER BY MAX(created_at) DESC
|
||
`, agentName)
|
||
if err == nil {
|
||
defer rows.Close()
|
||
for rows.Next() {
|
||
var p string
|
||
if rows.Scan(&p) != nil {
|
||
continue
|
||
}
|
||
if !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
// 来源 2:平台会话镜像。心跳上报的平台会话快照,workspace 就是
|
||
// 各平台的真实工作目录。清库后仍然存在(心跳重新上报)。
|
||
// 没有这一步,清库后路径补全对所有 Agent 都空。
|
||
mrows, err := db.DB.QueryContext(ctx, `
|
||
SELECT workspace, COUNT(*) AS n
|
||
FROM agent_platform_sessions
|
||
WHERE agent_name = $1 AND COALESCE(workspace, '') <> ''
|
||
GROUP BY workspace
|
||
ORDER BY n DESC
|
||
`, agentName)
|
||
if err == nil {
|
||
defer mrows.Close()
|
||
for mrows.Next() {
|
||
var p string
|
||
var n int
|
||
if mrows.Scan(&p, &n) != nil {
|
||
continue
|
||
}
|
||
if !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
// 来源 3:注册时自报。报了但还没收过信的目录靠这一步进入候选,
|
||
// 否则新接入的 Agent 在第一封邮件之前仍然无路径可选。
|
||
var wsJSON []byte
|
||
if err := db.DB.QueryRowContext(ctx,
|
||
`SELECT workspaces FROM agents WHERE agent_name = $1`, agentName).Scan(&wsJSON); err == nil {
|
||
var ws []models.Workspace
|
||
if len(wsJSON) > 0 {
|
||
json.Unmarshal(wsJSON, &ws)
|
||
}
|
||
for _, w := range ws {
|
||
// 三维地址的 path 位是**路径**,不是工作区的展示名。
|
||
// 这里取 Path 而不是 Name:remotebot 报的是
|
||
// {name:"demo", path:"/tmp/remotebot-ws"},取 Name 会给出 "demo"
|
||
// —— 一个拉起会话时不存在的目录。
|
||
p := w.Path
|
||
if p == "" {
|
||
p = w.Name
|
||
}
|
||
if p != "" && !seen[p] {
|
||
seen[p] = true
|
||
out = append(out, p)
|
||
}
|
||
}
|
||
}
|
||
|
||
return out, nil
|
||
}
|
||
|
||
// ListSentBy 列出某发件人发出的邮件(发件箱),排除已归档会话
|
||
func ListSentBy(ctx context.Context, fromName string, limit int) ([]models.Mail, error) {
|
||
rows, err := db.DB.QueryContext(ctx, `
|
||
SELECT m.mail_id, m.session_id, m.parent_mail_id,
|
||
m.from_name, m.from_workspace, m.to_name, m.to_workspace,
|
||
m.cc_list, m.subject, m.body, m.mail_type, COALESCE(m.permission_result,'') AS permission_result,
|
||
COALESCE(m.permission_kind,'') AS permission_kind,
|
||
COALESCE(m.permission_multi_select,0) AS permission_multi_select,
|
||
m.status, m.created_at, s.session_alias, s.workspace,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.from_name) AS from_human,
|
||
EXISTS (SELECT 1 FROM users u WHERE u.username = m.to_name) AS to_human
|
||
FROM mails m
|
||
JOIN sessions s ON m.session_id = s.session_id
|
||
WHERE m.from_name = $1 AND s.status <> 'archived'
|
||
ORDER BY m.created_at DESC, m.mail_id DESC
|
||
LIMIT $2
|
||
`, fromName, limit)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mails := []models.Mail{}
|
||
for rows.Next() {
|
||
var m models.Mail
|
||
var alias *string
|
||
var ccJSON []byte
|
||
if err := rows.Scan(&m.ID, &m.SessionID, &m.ParentMailID,
|
||
&m.FromName, &m.FromWorkspace, &m.ToName, &m.ToWorkspace,
|
||
&ccJSON, &m.Subject, &m.Body, &m.MailType, &m.PermResult,
|
||
&m.PermissionKind, &m.PermissionMulti,
|
||
&m.Status, &m.CreatedAt, &alias, &m.SessionWorkspace, &m.FromHuman, &m.ToHuman); err != nil {
|
||
return nil, err
|
||
}
|
||
if len(ccJSON) > 0 {
|
||
json.Unmarshal(ccJSON, &m.CCList)
|
||
}
|
||
if m.CCList == nil {
|
||
m.CCList = []models.Address{}
|
||
}
|
||
if alias != nil {
|
||
m.SessionAlias = *alias
|
||
}
|
||
if len(m.Body) > 200 {
|
||
m.BodyPreview = m.Body[:200] + "..."
|
||
} else {
|
||
m.BodyPreview = m.Body
|
||
}
|
||
mails = append(mails, m)
|
||
}
|
||
return mails, rows.Err()
|
||
}
|
||
|
||
// ListPendingPermissionsFor 列出待决权限请求;forUser 非空时只列发给该用户的
|
||
func ListPendingPermissionsFor(ctx context.Context, forUser string) ([]models.PermissionRequest, error) {
|
||
q := `SELECT pr.request_id, pr.mail_id, pr.session_id, pr.agent_name, pr.question,
|
||
pr.options, pr.context, COALESCE(pr.kind, '') AS kind, pr.multi_select, pr.result, pr.decided_at, pr.created_at
|
||
FROM permission_requests pr
|
||
JOIN mails m ON m.mail_id = pr.mail_id
|
||
JOIN sessions s ON s.session_id = pr.session_id
|
||
WHERE pr.result IS NULL AND s.status <> 'archived'`
|
||
args := []any{}
|
||
if forUser != "" {
|
||
q += ` AND m.to_name = $1`
|
||
args = append(args, forUser)
|
||
}
|
||
q += ` ORDER BY pr.created_at DESC`
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
reqs := []models.PermissionRequest{}
|
||
for rows.Next() {
|
||
var pr models.PermissionRequest
|
||
var optsJSON []byte
|
||
if err := rows.Scan(&pr.ID, &pr.MailID, &pr.SessionID, &pr.AgentName, &pr.Question,
|
||
&optsJSON, &pr.Context, &pr.Kind, &pr.MultiSelect, &pr.Result, &pr.DecidedAt, &pr.CreatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
json.Unmarshal(optsJSON, &pr.Options)
|
||
reqs = append(reqs, pr)
|
||
}
|
||
return reqs, nil
|
||
}
|
||
|
||
// ListSessionsFor 列出某人类用户参与的会话(owner / 收发 / 抄送);forUser 空表示不限
|
||
func ListSessionsFor(ctx context.Context, forUser string, limit int) ([]models.Session, error) {
|
||
q := `SELECT s.session_id, s.session_alias, s.from_agent, s.subject, s.status,
|
||
s.owner_user_id, s.created_at, s.updated_at,
|
||
COALESCE(s.max_rounds, 0), COALESCE(s.used_rounds, 0),
|
||
COALESCE(NULLIF(s.permission_mode, ''), 'workspace'),
|
||
COALESCE(NULLIF(s.permission_enforcement, ''), 'advisory'),
|
||
(SELECT COUNT(*) FROM mails m WHERE m.session_id = s.session_id)
|
||
FROM sessions s
|
||
WHERE s.status <> 'archived'`
|
||
args := []any{}
|
||
if forUser != "" {
|
||
q += ` AND (s.owner_user_id = (SELECT user_id FROM users WHERE username = $1)
|
||
OR EXISTS (
|
||
SELECT 1 FROM mails mm
|
||
WHERE mm.session_id = s.session_id
|
||
AND (mm.from_name = $1 OR mm.to_name = $1
|
||
OR ` + db.CCHas("mm.cc_list", 1) + `)
|
||
))`
|
||
args = append(args, forUser)
|
||
}
|
||
q += ` ORDER BY s.updated_at DESC`
|
||
if limit > 0 {
|
||
q += fmt.Sprintf(` LIMIT %d`, limit)
|
||
}
|
||
|
||
rows, err := db.DB.QueryContext(ctx, q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
sessions := []models.Session{}
|
||
for rows.Next() {
|
||
var s models.Session
|
||
// 列数必须与上面的 SELECT 一一对应 —— 预算两列曾经只加进了查询而没加进这里,
|
||
// 结果 /me/sessions 整个 500,联系人栅拉不到任何数据。
|
||
if err := rows.Scan(&s.ID, &s.Alias, &s.FromAgent, &s.Subject, &s.Status,
|
||
&s.OwnerUserID, &s.CreatedAt, &s.UpdatedAt,
|
||
&s.MaxRounds, &s.UsedRounds,
|
||
&s.PermissionMode, &s.PermissionEnforcement, &s.MailCount); err != nil {
|
||
return nil, err
|
||
}
|
||
sessions = append(sessions, s)
|
||
}
|
||
return sessions, rows.Err()
|
||
}
|
||
|
||
// CountUnreadInSession 统计某人在某会话内的未读数(含被抄送)
|
||
func CountUnreadInSession(ctx context.Context, name string, sessionID uuid.UUID) (int, error) {
|
||
ccProbe, _ := json.Marshal([]map[string]string{{"name": name}})
|
||
var n int
|
||
err := db.DB.QueryRowContext(ctx, `
|
||
SELECT COUNT(*) FROM mails
|
||
WHERE session_id = $1
|
||
AND status = 'unread'
|
||
AND (to_name = $2 OR cc_list @> $3::jsonb)
|
||
`, sessionID, name, string(ccProbe)).Scan(&n)
|
||
return n, err
|
||
}
|
||
|
||
// SetMailRenameProposal 记录某封邮件里 Agent 提议的新会话别名。
|
||
//
|
||
// 单独一条 UPDATE 而不是塞进 CreateMail 的参数表:提议是可选的旁支信息,
|
||
// 让三个调用点都多传两个几乎总是空串的参数不值当。
|
||
func SetMailRenameProposal(ctx context.Context, mailID uuid.UUID, alias, reason string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE mails SET rename_alias = $1, rename_reason = $2 WHERE mail_id = $3`,
|
||
alias, reason, mailID)
|
||
return err
|
||
}
|
||
|
||
// PendingRenameProposal 返回某会话里**最新一条尚未处理**的改名提议。
|
||
//
|
||
// 「尚未处理」= 提议的别名既不是当前别名(已接受),也不在驳回记录里。
|
||
// 无提议时返回 ("", "", nil)。
|
||
//
|
||
// 只看最新一条:Agent 干活过程中可能多次提议,最后那条才是它现在的结论。
|
||
func PendingRenameProposal(ctx context.Context, sessionID uuid.UUID) (alias, reason string, err error) {
|
||
var cur, dismissed *string
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT session_alias, rename_dismissed FROM sessions WHERE session_id = $1`,
|
||
sessionID).Scan(&cur, &dismissed)
|
||
if err != nil {
|
||
return "", "", err
|
||
}
|
||
|
||
var a, rs *string
|
||
err = db.DB.QueryRowContext(ctx, `
|
||
SELECT rename_alias, rename_reason FROM mails
|
||
WHERE session_id = $1 AND rename_alias IS NOT NULL AND rename_alias <> ''
|
||
ORDER BY created_at DESC, mail_id DESC
|
||
LIMIT 1
|
||
`, sessionID).Scan(&a, &rs)
|
||
if errors.Is(err, sql.ErrNoRows) {
|
||
return "", "", nil
|
||
}
|
||
if err != nil {
|
||
return "", "", err
|
||
}
|
||
if a == nil || *a == "" {
|
||
return "", "", nil
|
||
}
|
||
// 已经改成这个名字了 = 提议已被接受,不必再提示
|
||
if cur != nil && *cur == *a {
|
||
return "", "", nil
|
||
}
|
||
// 用户驳回过这个建议
|
||
if dismissed != nil && *dismissed == *a {
|
||
return "", "", nil
|
||
}
|
||
if rs != nil {
|
||
reason = *rs
|
||
}
|
||
return *a, reason, nil
|
||
}
|
||
|
||
// DismissRenameProposal 记下用户驳回了哪个建议,好让提示条不再反复弹。
|
||
//
|
||
// 只存最后驳回的那一个而不是一张列表:Agent 每次提的名字都不同,
|
||
// 攒一张历史表除了占地方没有别的用处 —— 需要判断的只是「当前这条提议是否被否过」。
|
||
func DismissRenameProposal(ctx context.Context, sessionID uuid.UUID, alias string) error {
|
||
_, err := db.DB.ExecContext(ctx,
|
||
`UPDATE sessions SET rename_dismissed = $1 WHERE session_id = $2`,
|
||
alias, sessionID)
|
||
return err
|
||
}
|
||
|
||
// MarkMailsReadFor 把一批邮件标记为某收件人已读,返回实际影响的行数。
|
||
//
|
||
// **鉴权写进 WHERE 而不是先查后改**:`to_name = $1 OR cc 含 $1` 直接放在
|
||
// UPDATE 条件里,于是「不是发给我的邮件」根本改不动 —— 既省掉一次查询,
|
||
// 也没有「查完到改之间邮件被转走」的时间窗。
|
||
//
|
||
// 已经是 read 的不计入影响行数(`status = 'unread'` 条件),
|
||
// 调用方据此知道这次真正标掉了几封。
|
||
func MarkMailsReadFor(ctx context.Context, recipient string, ids []uuid.UUID) (int, error) {
|
||
if len(ids) == 0 {
|
||
return 0, nil
|
||
}
|
||
|
||
// IN 子句的占位符按方言编号(SQLite 与 PG 都认 $N)。
|
||
// 不用一条条 UPDATE:一次网络往返 + 一次事务,SQLite 单写者下差别明显。
|
||
ph := make([]string, len(ids))
|
||
args := make([]any, 0, len(ids)+1)
|
||
args = append(args, recipient)
|
||
for i, id := range ids {
|
||
ph[i] = fmt.Sprintf("$%d", i+2)
|
||
args = append(args, id)
|
||
}
|
||
|
||
res, err := db.DB.ExecContext(ctx, `
|
||
UPDATE mails SET status = 'read'
|
||
WHERE mail_id IN (`+strings.Join(ph, ",")+`)
|
||
AND status = 'unread'
|
||
AND (to_name = $1 OR `+db.CCHas("cc_list", 1)+`)
|
||
`, args...)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
n, _ := res.RowsAffected()
|
||
return int(n), nil
|
||
}
|
||
|
||
// MarkAllInboxReadFor 把某收件人收件箱里全部未读标为已读,返回影响行数。
|
||
//
|
||
// 排除已归档会话:那些邮件在收件箱里根本看不到,
|
||
// 标掉它们只会让「标记了 N 封」这个数字与用户看到的对不上。
|
||
func MarkAllInboxReadFor(ctx context.Context, recipient string) (int, error) {
|
||
res, err := db.DB.ExecContext(ctx, `
|
||
UPDATE mails SET status = 'read'
|
||
WHERE status = 'unread'
|
||
AND (to_name = $1 OR `+db.CCHas("cc_list", 1)+`)
|
||
AND session_id IN (SELECT session_id FROM sessions WHERE status <> 'archived')
|
||
`, recipient)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
n, _ := res.RowsAffected()
|
||
return int(n), nil
|
||
}
|
||
|
||
// ─── Agent 删除 ───
|
||
|
||
// DeleteAgent 删除 Agent 的全部运行态,保留邮件历史,名字进保留名单。
|
||
//
|
||
// 删除范围:
|
||
// - agent_keys(全部撤销)
|
||
// - agent_platform_sessions(清掉镜像)
|
||
// - models_scope(模型范围配置)
|
||
// - rate_limits(速率限制计数器)
|
||
// - calendar_events:cancelled(不触发、不静默空转)
|
||
// - agents 行本身
|
||
//
|
||
// 保留范围:
|
||
// - mails(历史是审计凭据,不能删)
|
||
// - sessions(与 mails 一起组成线索)
|
||
//
|
||
// 取消邮件:Agent 名在 mails 里的字段(from_name / to_name)不改 ——
|
||
// 那是历史记录的固有属性。未来要"查这封信是谁发的"仍然能查到。
|
||
func DeleteAgent(ctx context.Context, name string) (int, error) {
|
||
tx, err := db.DB.BeginTx(ctx, nil)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
defer tx.Rollback()
|
||
|
||
// 撤销全部密钥
|
||
r1, _ := tx.ExecContext(ctx, `DELETE FROM agent_keys WHERE agent_name = $1`, name)
|
||
keysRevoked, _ := r1.RowsAffected()
|
||
|
||
// 清镜像
|
||
tx.ExecContext(ctx, `DELETE FROM agent_platform_sessions WHERE agent_name = $1`, name)
|
||
|
||
// 清模型范围
|
||
tx.ExecContext(ctx, `DELETE FROM agent_allowed_models WHERE agent_name = $1`, name)
|
||
tx.ExecContext(ctx, `DELETE FROM agent_model_catalog WHERE agent_name = $1`, name)
|
||
|
||
// 清速率限制
|
||
tx.ExecContext(ctx, `DELETE FROM rate_limits WHERE agent_name = $1`, name)
|
||
|
||
// 日历事件置 cancelled:Agent 创建的提醒不该继续触发并投给一个不存在的收件人
|
||
tx.ExecContext(ctx,
|
||
`UPDATE calendar_events SET status = 'cancelled', updated_at = $2
|
||
WHERE created_by = $1 AND status = 'active'`, name, time.Now())
|
||
|
||
// 删 agents 行
|
||
if _, err := tx.ExecContext(ctx, `DELETE FROM agents WHERE agent_name = $1`, name); err != nil {
|
||
return 0, err
|
||
}
|
||
|
||
return int(keysRevoked), tx.Commit()
|
||
}
|
||
|
||
// IsRetiredAgentName 检查一个名字是否已被删除(用于注册时拒绝同名重建)。
|
||
func IsRetiredAgentName(ctx context.Context, name string) (bool, error) {
|
||
// 如果 agents 表里没有这个名字,且没有任何邮件引用它,就当「已退役」。
|
||
// 更严格的做法是建一张 retired_agents 表,但现有数据量下这个查询够了。
|
||
var cnt int
|
||
err := db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM agents WHERE agent_name = $1`, name).Scan(&cnt)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
if cnt > 0 {
|
||
return false, nil // 还活着
|
||
}
|
||
// 检查有没有历史邮件用这个名字
|
||
err = db.DB.QueryRowContext(ctx,
|
||
`SELECT COUNT(*) FROM mails WHERE from_name = $1 OR to_name = $1`, name).Scan(&cnt)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
return cnt > 0, nil
|
||
}
|