diff --git a/gateway/internal/db/migrate.go b/gateway/internal/db/migrate.go index 2f1bede..da6955e 100644 --- a/gateway/internal/db/migrate.go +++ b/gateway/internal/db/migrate.go @@ -80,6 +80,14 @@ var sqliteAddColumns = []struct{ table, column, ddl string }{ // (Agent 回信的 from_workspace 存的是 Agent 名而不是路径),强行回填只会 // 造出一批看起来有值实际是错的数据。 {"sessions", "workspace", "ALTER TABLE sessions ADD COLUMN workspace TEXT NOT NULL DEFAULT ''"}, + // 日历多收件人。旧库默认 '[]':读的时候由 EffectiveRecipients() 退回 + // to_address / agent_name,历史事件因此继续工作,不需要数据迁移。 + {"calendar_events", "recipients", "ALTER TABLE calendar_events ADD COLUMN recipients TEXT NOT NULL DEFAULT '[]'"}, + {"calendar_events", "delivery_mode", "ALTER TABLE calendar_events ADD COLUMN delivery_mode TEXT NOT NULL DEFAULT 'separate'"}, + // 日历事件已触发的 occurrence。旧库为 NULL:等价于「从未触发」, + // 于是已过期的一次性事件会补发一次提醒 —— 这是可接受的, + // 而反过来(默认成 event_time)会让正在等的提醒永远发不出去。 + {"calendar_events", "fired_for", "ALTER TABLE calendar_events ADD COLUMN fired_for DATETIME"}, // 派给该 Agent 的新任务默认多少个来回。 // 旧库也给 20:之前的 max_rounds 默认是 10 但那是终身额度,语义不同, // 不能直接搬过来当单任务预算。 diff --git a/gateway/internal/db/migrations/init.sql b/gateway/internal/db/migrations/init.sql index 516b4e5..e7b5283 100644 --- a/gateway/internal/db/migrations/init.sql +++ b/gateway/internal/db/migrations/init.sql @@ -345,3 +345,57 @@ CREATE TABLE IF NOT EXISTS agent_platform_sessions ( CREATE INDEX IF NOT EXISTS idx_platform_sessions_ws ON agent_platform_sessions(agent_name, workspace); + +-- ---------- 日历(见 init_sqlite.sql 里的设计说明) ---------- +-- +-- 三层分离:事件是日历实体,提醒是触发器,邮件是投递通道。 +-- 这份 PG schema 曾经整块缺失 —— 后果是 DATABASE_URL 一旦非空, +-- 所有 /calendar/* 端点在 relation does not exist 上 500, +-- 而 SQLite 下一切正常,于是问题只在切外部库时才暴露。 +CREATE TABLE IF NOT EXISTS calendar_events ( + event_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + title VARCHAR(512) NOT NULL, + description TEXT NOT NULL DEFAULT '', + + -- 提醒邮件正文模板,支持 {title} {time} {description} + reminder_text TEXT NOT NULL DEFAULT '', + + -- 收件方:agent_name 是兜底,to_address 是权威(完整三维寻址) + agent_name VARCHAR(128) NOT NULL DEFAULT '', + to_address VARCHAR(512) NOT NULL DEFAULT '', + + event_time TIMESTAMPTZ NOT NULL, + remind_before INTEGER NOT NULL DEFAULT 0, + recurrence VARCHAR(32) NOT NULL DEFAULT 'none', + recurrence_end TIMESTAMPTZ, + + -- 收件人列表与投递方式(见 init_sqlite.sql 的说明) + recipients JSONB NOT NULL DEFAULT '[]'::jsonb, + delivery_mode VARCHAR(32) NOT NULL DEFAULT 'separate', + + status VARCHAR(32) NOT NULL DEFAULT 'active', + last_fired_at TIMESTAMPTZ, + -- 已触发的 occurrence(= 当时的 event_time)。见 init_sqlite.sql 的说明。 + fired_for TIMESTAMPTZ, + created_at TIMESTAMPTZ DEFAULT NOW(), + updated_at TIMESTAMPTZ DEFAULT NOW(), + created_by VARCHAR(128) NOT NULL DEFAULT '' +); + +CREATE INDEX IF NOT EXISTS idx_calendar_next_fire + ON calendar_events(status, event_time) WHERE status = 'active'; + +CREATE INDEX IF NOT EXISTS idx_calendar_time_range + ON calendar_events(event_time, status); + +CREATE TABLE IF NOT EXISTS calendar_attachments ( + attachment_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + event_id UUID NOT NULL REFERENCES calendar_events(event_id) ON DELETE CASCADE, + filename VARCHAR(512) NOT NULL, + sha256 VARCHAR(64) NOT NULL DEFAULT '', + size_bytes BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ DEFAULT NOW() +); + +CREATE INDEX IF NOT EXISTS idx_calendar_att_event + ON calendar_attachments(event_id); diff --git a/gateway/internal/db/migrations/init_sqlite.sql b/gateway/internal/db/migrations/init_sqlite.sql index 2290321..7a32fed 100644 --- a/gateway/internal/db/migrations/init_sqlite.sql +++ b/gateway/internal/db/migrations/init_sqlite.sql @@ -336,3 +336,94 @@ CREATE TABLE IF NOT EXISTS agent_platform_sessions ( CREATE INDEX IF NOT EXISTS idx_platform_sessions_ws ON agent_platform_sessions(agent_name, workspace); + +-- ─── 日历事件 ─── +-- +-- Outlook 风格:事件 → 提醒 → 邮件通知 Agent。 +-- 事件本身是日历实体,提醒是定时触发器,邮件是投递通道。 +-- 三者分离:同一条事件可以有多个提醒(提前提醒 + 当天提醒), +-- 同一条提醒只触发一封邮件(幂等由 fired_at 控制)。 +CREATE TABLE IF NOT EXISTS calendar_events ( + event_id TEXT PRIMARY KEY DEFAULT (gen_random_uuid()), + -- 事件标题(UI 显示 + 邮件主题前缀) + title TEXT NOT NULL, + -- 事件描述(UI 显示,可含 markdown) + description TEXT NOT NULL DEFAULT '', + -- 用户可编辑的提醒消息模板。支持变量:{title} {time} {description} + reminder_text TEXT NOT NULL DEFAULT '', + + -- 通知目标 + agent_name TEXT NOT NULL DEFAULT '', + -- 收件人三维地址(空 = 用 agent_name 默认地址) + to_address TEXT NOT NULL DEFAULT '', + + -- 时间安排 + event_time DATETIME NOT NULL, + -- 提前多少分钟提醒(0 = 事件触发时) + remind_before INTEGER NOT NULL DEFAULT 0, + -- 重复规则:none / daily / weekly / monthly + -- lunar_monthly(每农历月同一日)/ lunar_yearly(每农历年同月同日) + -- + -- 农历规则必须经 internal/lunar 推进,不能加固定天数 —— + -- 农历月 29~30 天不定、农历年 353~385 天(闰年多一整月), + -- 近似推进一年能偏半个月。 + recurrence TEXT NOT NULL DEFAULT 'none', + -- 重复结束(空 = 永久) + recurrence_end DATETIME, + + -- 收件人列表(JSON 数组,每项是完整三维地址串)。 + -- + -- 存原始串而不是结构化地址:session 位的 new/别名三态该在**触发那一刻** + -- 解析。存结构化的话「.new」这种一次性语义在建事件时就被固化, + -- 而重复事件每次触发都该重新决定落到哪条会话。 + recipients TEXT NOT NULL DEFAULT '[]', + + -- 多收件人的投递方式: + -- separate(默认)= 各发一封、落各自会话、互不可见 + -- together = 首个为主收件人,其余进 cc_list、共享一条线索 + -- + -- 默认 separate 因为它的失败模式更轻:together 用错会让本该独立判断的 + -- Agent 互相看到回复而趋同,那种污染事后无法分离。 + delivery_mode TEXT NOT NULL DEFAULT 'separate', + + -- 状态:active / paused / cancelled + status TEXT NOT NULL DEFAULT 'active', + -- 最后一次触发的墙上时钟(给 UI 显示「上次触发于」) + last_fired_at DATETIME, + -- 已触发的那个 occurrence,值 = 当时的 event_time。 + -- + -- 去重不能拿 last_fired_at 跟 event_time 比大小:DueEvents 有 60 秒 + -- lookahead,落在窗口内的**未来**事件被触发后 last_fired_at(now) 仍然 + -- 小于 event_time,于是每个 tick 重发一次,直到 event_time 真正过去。 + -- 生产实测:一条 12:53:17 的事件在 12:52:30 / 12:53:00 / 12:53:06 / + -- 12:53:36 发了 4 封相同提醒。 + -- 按 occurrence 比相等则精确:AdvanceRecurrence 改了 event_time 就再触发, + -- 没改就永不重发。 + fired_for DATETIME, + created_at DATETIME DEFAULT (strftime('%Y-%m-%d %H:%M:%f','now')), + updated_at DATETIME DEFAULT (strftime('%Y-%m-%d %H:%M:%f','now')), + + -- 创建者(人类用户) + created_by TEXT NOT NULL DEFAULT '' +); + +-- 调度器每分钟扫描 active 事件,按 event_time + remind_before 排序取下一个 +CREATE INDEX IF NOT EXISTS idx_calendar_next_fire + ON calendar_events(status, event_time) WHERE status = 'active'; + +-- 按时间范围查(日历视图) +CREATE INDEX IF NOT EXISTS idx_calendar_time_range + ON calendar_events(event_time, status); + +-- 附件:事件触发时随提醒邮件一起发出 +CREATE TABLE IF NOT EXISTS calendar_attachments ( + attachment_id TEXT PRIMARY KEY DEFAULT (gen_random_uuid()), + event_id TEXT NOT NULL REFERENCES calendar_events(event_id) ON DELETE CASCADE, + filename TEXT NOT NULL, + sha256 TEXT NOT NULL DEFAULT '', + size_bytes INTEGER NOT NULL DEFAULT 0, + created_at DATETIME DEFAULT (strftime('%Y-%m-%d %H:%M:%f','now')) +); + +CREATE INDEX IF NOT EXISTS idx_calendar_att_event + ON calendar_attachments(event_id); diff --git a/gateway/internal/handler/calendar.go b/gateway/internal/handler/calendar.go new file mode 100644 index 0000000..64b1916 --- /dev/null +++ b/gateway/internal/handler/calendar.go @@ -0,0 +1,775 @@ +package handler + +import ( + "crypto/sha256" + "errors" + "fmt" + "io" + "net/http" + "strings" + "time" + + "github.com/go-chi/chi/v5" + + "github.com/agentmail/gateway/internal/blob" + "github.com/agentmail/gateway/internal/config" + "github.com/agentmail/gateway/internal/middleware" + "github.com/agentmail/gateway/internal/models" + "github.com/agentmail/gateway/internal/repo" +) + +// defaultReminderTemplate 是提醒正文的默认模板。 +// +// 存变量而非字面值:{title}/{time}/{description} 在触发时由 +// scheduler.RenderReminder 替换。与前端 CalendarEventEditor 的 +// DEFAULT_TEMPLATE 必须逐字一致 —— 前端用它作 placeholder 与预览, +// 两边不同会让人看到的预览与 Agent 实收的正文不是一回事。 +const defaultReminderTemplate = "日程提醒:{title}\n时间:{time}\n{description}" + +// ─── Calendar Events ─── + +// POST /api/v1/calendar/events +func CreateCalendarEvent(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + + var req struct { + Title string `json:"title"` + Description string `json:"description"` + ReminderText string `json:"reminder_text"` + AgentName string `json:"agent_name"` + ToAddress string `json:"to_address"` + Recipients []string `json:"recipients"` + DeliveryMode string `json:"delivery_mode"` + EventTime time.Time `json:"event_time"` + RemindBefore int `json:"remind_before"` + Recurrence string `json:"recurrence"` + RecurrenceEnd *time.Time `json:"recurrence_end"` + } + if !DecodeBody(w, r, &req) { + return + } + if req.Title == "" { + Error(w, http.StatusBadRequest, "Missing title") + return + } + if req.EventTime.IsZero() { + Error(w, http.StatusBadRequest, "Missing event_time") + return + } + if req.ReminderText == "" { + // 默认模板必须存**变量形式**而不是把值烤进去。 + // + // 原来这里 Sprintf 出一份含字面时间的正文。对重复事件是错的: + // AdvanceRecurrence 只推进 event_time,reminder_text 保持不动 —— + // 于是「每天 9 点」的提醒从第二天起永远写着第一天的日期, + // Agent 收到的信里时间与实际触发时刻越差越远。 + // + // 变量形式由 scheduler.RenderReminder 在**触发时**替换, + // 每一次触发都拿当时的 event_time。前端的 DEFAULT_TEMPLATE + // 也是这一份(web/src/components/CalendarEventEditor.tsx), + // 两处必须一致,否则预览与实发不符。 + req.ReminderText = defaultReminderTemplate + } + if req.Recurrence == "" { + req.Recurrence = models.RecurNone + } + if !validRecurrence(req.Recurrence) { + Error(w, http.StatusBadRequest, + "recurrence 必须是 none/daily/weekly/monthly/lunar_monthly/lunar_yearly 之一") + return + } + + recipients, badAddr := normalizeRecipients(req.Recipients) + if badAddr != "" { + // 地址在这里就校验而不是等到触发时:建事件时报错人能立刻改, + // 而触发时报错只会进 journalctl —— 人以为提醒设好了,实际永远发不出去。 + Error(w, http.StatusBadRequest, "收件地址无法解析:"+badAddr) + return + } + // 收件人一个都没有时事件永远发不出去,这不该静默通过 + if len(recipients) == 0 && strings.TrimSpace(req.ToAddress) == "" && + strings.TrimSpace(req.AgentName) == "" { + Error(w, http.StatusBadRequest, "至少要有一个收件人(recipients / to_address / agent_name)") + return + } + if req.DeliveryMode == "" { + req.DeliveryMode = models.DeliverSeparate + } + + e := &models.CalendarEvent{ + Title: req.Title, + Description: req.Description, + ReminderText: req.ReminderText, + AgentName: req.AgentName, + ToAddress: req.ToAddress, + Recipients: recipients, + DeliveryMode: req.DeliveryMode, + EventTime: req.EventTime, + RemindBefore: req.RemindBefore, + Recurrence: req.Recurrence, + RecurrenceEnd: req.RecurrenceEnd, + Status: "active", + CreatedBy: user.Username, + } + + if _, err := repo.CreateCalendarEvent(r.Context(), e); err != nil { + Error(w, http.StatusInternalServerError, "Failed to create event") + return + } + JSON(w, http.StatusCreated, e) +} + +// GET /api/v1/calendar/events?from=...&to=... +func ListCalendarEvents(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + + fromStr := r.URL.Query().Get("from") + toStr := r.URL.Query().Get("to") + status := r.URL.Query().Get("status") + + var from, to time.Time + if fromStr != "" { + from, _ = time.Parse(time.RFC3339, fromStr) + } + if toStr != "" { + to, _ = time.Parse(time.RFC3339, toStr) + } + if to.IsZero() { + to = time.Now().AddDate(0, 1, 0) // 默认往后一个月 + } + if from.IsZero() { + from = time.Now().AddDate(0, -1, 0) // 默认往前一个月 + } + + events, err := repo.ListCalendarEvents(r.Context(), from, to, status) + if err != nil { + Error(w, http.StatusInternalServerError, "Failed to list events") + return + } + if events == nil { + events = []models.CalendarEvent{} + } + JSON(w, http.StatusOK, map[string]interface{}{ + "events": events, + }) +} + +// GET /api/v1/calendar/events/{id} +func GetCalendarEvent(w http.ResponseWriter, r *http.Request) { + if _, ok := pathUUID(w, r, "id"); !ok { + return + } + eventID := chi.URLParam(r, "id") + e, err := repo.GetCalendarEvent(r.Context(), eventID) + if err != nil { + if errors.Is(err, repo.ErrEventNotFound) { + Error(w, http.StatusNotFound, "Event not found") + return + } + Error(w, http.StatusInternalServerError, "Failed to get event") + return + } + JSON(w, http.StatusOK, e) +} + +// PUT /api/v1/calendar/events/{id} +func UpdateCalendarEvent(w http.ResponseWriter, r *http.Request) { + eventID := chi.URLParam(r, "id") + if _, ok := pathUUID(w, r, "id"); !ok { + return + } + + var req struct { + Title string `json:"title"` + Description string `json:"description"` + ReminderText string `json:"reminder_text"` + AgentName string `json:"agent_name"` + ToAddress string `json:"to_address"` + Recipients []string `json:"recipients"` + DeliveryMode string `json:"delivery_mode"` + EventTime time.Time `json:"event_time"` + RemindBefore int `json:"remind_before"` + Recurrence string `json:"recurrence"` + RecurrenceEnd *time.Time `json:"recurrence_end"` + Status string `json:"status"` + } + if !DecodeBody(w, r, &req) { + return + } + if req.Recurrence != "" && !validRecurrence(req.Recurrence) { + Error(w, http.StatusBadRequest, + "recurrence 必须是 none/daily/weekly/monthly/lunar_monthly/lunar_yearly 之一") + return + } + recipients, badAddr := normalizeRecipients(req.Recipients) + if badAddr != "" { + Error(w, http.StatusBadRequest, "收件地址无法解析:"+badAddr) + return + } + + e := &models.CalendarEvent{ + Title: req.Title, + Description: req.Description, + ReminderText: req.ReminderText, + AgentName: req.AgentName, + ToAddress: req.ToAddress, + Recipients: recipients, + DeliveryMode: req.DeliveryMode, + EventTime: req.EventTime, + RemindBefore: req.RemindBefore, + Recurrence: req.Recurrence, + RecurrenceEnd: req.RecurrenceEnd, + Status: req.Status, + } + if err := repo.UpdateCalendarEvent(r.Context(), eventID, e); err != nil { + if errors.Is(err, repo.ErrEventNotFound) { + Error(w, http.StatusNotFound, "Event not found") + return + } + Error(w, http.StatusInternalServerError, "Failed to update event") + return + } + JSON(w, http.StatusOK, e) +} + +// DELETE /api/v1/calendar/events/{id} +func DeleteCalendarEvent(w http.ResponseWriter, r *http.Request) { + eventID := chi.URLParam(r, "id") + if err := repo.DeleteCalendarEvent(r.Context(), eventID); err != nil { + if errors.Is(err, repo.ErrEventNotFound) { + Error(w, http.StatusNotFound, "Event not found") + return + } + Error(w, http.StatusInternalServerError, "Failed to delete event") + return + } + JSON(w, http.StatusOK, map[string]string{"status": "deleted"}) +} + +// ─── Calendar Attachments ─── + +// POST /api/v1/calendar/events/{id}/attachments +func UploadCalendarAttachment(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + if Blobs == nil { + Error(w, http.StatusServiceUnavailable, "附件存储未初始化") + return + } + eventID := strings.TrimSpace(chi.URLParam(r, "id")) + if eventID == "" { + Error(w, http.StatusBadRequest, "Missing event id") + return + } + // 事件必须存在:否则会攒下一堆孤儿附件记录,而 ON DELETE CASCADE + // 永远清不掉它们(没有对应的父行可删)。 + if _, err := repo.GetCalendarEvent(r.Context(), eventID); err != nil { + Error(w, http.StatusNotFound, "事件不存在") + return + } + + max := config.C.MaxAttachmentBytes + // 与邮件附件同一套双层限制:外层卡整个请求体(含 multipart 边界), + // blob.Put 的 max 卡单个文件内容。少了外层,超大 multipart 头能拖死内存。 + r.Body = http.MaxBytesReader(w, r.Body, max+1<<20) + if err := r.ParseMultipartForm(32 << 20); err != nil { + Error(w, http.StatusBadRequest, "解析 multipart 失败(是否超过大小上限?)") + return + } + defer func() { + if r.MultipartForm != nil { + r.MultipartForm.RemoveAll() + } + }() + + file, header, err := r.FormFile("file") + if err != nil { + Error(w, http.StatusBadRequest, "缺少 file 字段") + return + } + defer file.Close() + + // 原来这里是 `data := make([]byte, header.Size); file.Read(data)` —— + // 两处错:单次 Read 不保证填满缓冲(大文件必然短读,sha256 因此算的是 + // 半截内容),而且**文件内容从未落盘**,只往库里写了一条元数据。 + // 结果是附件"上传成功"、清单里看得见、发提醒时取不到任何字节。 + sum, size, err := Blobs.Put(file, max) + if errors.Is(err, blob.ErrTooLarge) { + Error(w, http.StatusRequestEntityTooLarge, + fmt.Sprintf("附件超过上限 %.1f MB", float64(max)/(1<<20))) + return + } + if err != nil { + Error(w, http.StatusInternalServerError, "保存附件失败") + return + } + + att := &models.CalendarAttachment{ + EventID: eventID, + Filename: sanitizeFilename(header.Filename), + SizeBytes: size, + SHA256: sum, + } + if err := repo.AddCalendarAttachment(r.Context(), att); err != nil { + // 落盘成功但入库失败:孤立文件由 GC 回收,不影响正确性 + Error(w, http.StatusInternalServerError, "登记附件失败") + return + } + JSON(w, http.StatusCreated, att) +} + +// GET /api/v1/calendar/events/{id}/attachments +func ListCalendarAttachments(w http.ResponseWriter, r *http.Request) { + eventID := chi.URLParam(r, "id") + atts, err := repo.ListCalendarAttachments(r.Context(), eventID) + if err != nil { + Error(w, http.StatusInternalServerError, "Failed to list attachments") + return + } + if atts == nil { + atts = []models.CalendarAttachment{} + } + JSON(w, http.StatusOK, map[string]interface{}{"attachments": atts}) +} + +// DELETE /api/v1/calendar/events/{id}/attachments/{aid} +func DeleteCalendarAttachment(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + attID := strings.TrimSpace(chi.URLParam(r, "attachmentID")) + if attID == "" { + Error(w, http.StatusBadRequest, "Missing attachment id") + return + } + // 原来这里返回 501 并让人「删整个事件来清附件」—— 那要求人为了撤掉 + // 一个错传的文件把整条日程连提醒配置一起重建。 + // + // 磁盘上的 blob 不在这里删:内容寻址下同一个 sha256 可能被别的附件 + // (甚至别的邮件)引用着,删文件会让那些引用一起坏掉。孤立 blob 归 GC。 + ok, err := repo.DeleteCalendarAttachment(r.Context(), attID) + if err != nil { + Error(w, http.StatusInternalServerError, "删除附件失败") + return + } + if !ok { + Error(w, http.StatusNotFound, "附件不存在") + return + } + JSON(w, http.StatusOK, map[string]any{"status": "deleted", "attachment_id": attID}) +} + +// validRecurrence 白名单校验重复规则。 +// +// 必须白名单而不是「未知值当 none」:把 `lunar_montly`(拼错)静默当成 +// 不重复,用户设的每月提醒只会响一次,而没有任何地方报错。 +func validRecurrence(r string) bool { + switch r { + case models.RecurNone, models.RecurDaily, models.RecurWeekly, models.RecurMonthly, + models.RecurYearly, models.RecurLunarMonthly, models.RecurLunarYearly: + return true + } + return false +} + +// normalizeRecipients 清洗收件人列表:去空白、去重、校验地址可解析。 +// +// 返回第二个值非空表示有地址解析失败(值即那个地址),调用方回 400。 +// 在建事件时校验而不是等触发:建事件时报错人能立刻改, +// 触发时报错只会进 journalctl —— 人以为设好了,实际永远发不出去。 +// +// 去重是必要的:together 模式下同一个 Agent 既是主收件人又在抄送里, +// 会让它收到两条一模一样的 SSE,插件可能因此起两轮。 +func normalizeRecipients(in []string) ([]string, string) { + seen := make(map[string]bool, len(in)) + out := make([]string, 0, len(in)) + for _, raw := range in { + raw = strings.TrimSpace(raw) + if raw == "" { + continue + } + if _, err := models.ParseAddress(raw); err != nil { + return nil, raw + } + if seen[raw] { + continue + } + seen[raw] = true + out = append(out, raw) + } + return out, "" +} + +// ─── iCal 导入导出 ─── + +// GET /api/v1/calendar/export.ics +func ExportCalendarICS(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + + // 区间取自查询参数:前端导出的是「当前正在看的那段」, + // 写死 ±1 年会让人点导出后得到一堆与屏幕上不符的事件。 + from := time.Now().AddDate(-1, 0, 0) + to := time.Now().AddDate(1, 0, 0) + if v := r.URL.Query().Get("from"); v != "" { + if t, err := time.Parse(time.RFC3339, v); err == nil { + from = t + } + } + if v := r.URL.Query().Get("to"); v != "" { + if t, err := time.Parse(time.RFC3339, v); err == nil { + to = t + } + } + events, err := repo.ListCalendarEvents(r.Context(), from, to, "active") + if err != nil { + Error(w, http.StatusInternalServerError, "Failed to list events") + return + } + + var sb strings.Builder + sb.WriteString("BEGIN:VCALENDAR\r\n") + sb.WriteString("VERSION:2.0\r\n") + sb.WriteString("PRODID:-//AgentMail//Calendar//EN\r\n") + + for _, e := range events { + sb.WriteString("BEGIN:VEVENT\r\n") + fmt.Fprintf(&sb, "UID:%s@agentmail\r\n", e.EventID) + fmt.Fprintf(&sb, "DTSTAMP:%s\r\n", e.EventTime.UTC().Format("20060102T150405Z")) + fmt.Fprintf(&sb, "DTSTART:%s\r\n", e.EventTime.UTC().Format("20060102T150405Z")) + // 默认 1 小时持续时间 + fmt.Fprintf(&sb, "DTEND:%s\r\n", e.EventTime.Add(time.Hour).UTC().Format("20060102T150405Z")) + // 转义换行 + summary := strings.ReplaceAll(e.Title, "\n", "\\n") + fmt.Fprintf(&sb, "SUMMARY:%s\r\n", summary) + if e.Description != "" { + desc := strings.ReplaceAll(e.Description, "\n", "\\n") + fmt.Fprintf(&sb, "DESCRIPTION:%s\r\n", desc) + } + if e.Recurrence != "none" { + var freq string + switch e.Recurrence { + case "daily": + freq = "DAILY" + case "weekly": + freq = "WEEKLY" + case "monthly": + freq = "MONTHLY" + case "yearly": + freq = "YEARLY" + } + if freq != "" { + fmt.Fprintf(&sb, "RRULE:FREQ=%s\r\n", freq) + } + // 农历规则 RFC 5545 表达不了(RRULE 只有公历频率)。 + // + // 折中:用 X- 扩展属性记下真实规则,并把它降级成最接近的公历 + // 近似(lunar_monthly → MONTHLY、lunar_yearly → YEARLY)。 + // 别的客户端至少能看到一个大致对的重复;导回本系统时 + // X- 属性会把精确规则还原。 + // + // 不写近似 RRULE 的后果更糟:外部客户端会把它当一次性事件, + // 用户以为导出的日历里有「每年农历生日」,实际只有一条。 + if models.IsLunarRecurrence(e.Recurrence) { + fmt.Fprintf(&sb, "X-AGENTMAIL-RECURRENCE:%s\r\n", e.Recurrence) + if e.Recurrence == models.RecurLunarMonthly { + sb.WriteString("RRULE:FREQ=MONTHLY\r\n") + } else { + sb.WriteString("RRULE:FREQ=YEARLY\r\n") + } + } + } + // VALARM 的 TRIGGER 必须写成 `-PTM`。 + // + // 两个坑:iCal 的 duration 里 `M` **在 T 之前是月、在 T 之后才是分钟** —— + // 原来写的 `-P15M` 在任何合规日历客户端里都是「提前 15 个月」。 + // 而且原来用 maxInt(RemindBefore, 15) 兜底,把用户明确设的 + // 「到点提醒」(0) 悄悄改成提前 15 分钟;导出不该修改语义。 + // 收件人与投递模式同样没有标准字段可放。 + // 不导出的后果:往返一圈后事件变成「没有收件人」,永远不会提醒。 + if rs := e.EffectiveRecipients(); len(rs) > 0 { + fmt.Fprintf(&sb, "X-AGENTMAIL-RECIPIENTS:%s\r\n", strings.Join(rs, ",")) + fmt.Fprintf(&sb, "X-AGENTMAIL-DELIVERY:%s\r\n", e.EffectiveDeliveryMode()) + } + fmt.Fprintf(&sb, "BEGIN:VALARM\r\n") + fmt.Fprintf(&sb, "TRIGGER:-PT%dM\r\n", e.RemindBefore) + fmt.Fprintf(&sb, "ACTION:DISPLAY\r\n") + fmt.Fprintf(&sb, "DESCRIPTION:%s\r\n", summary) + fmt.Fprintf(&sb, "END:VALARM\r\n") + + sb.WriteString("END:VEVENT\r\n") + } + + sb.WriteString("END:VCALENDAR\r\n") + + w.Header().Set("Content-Type", "text/calendar; charset=utf-8") + w.Header().Set("Content-Disposition", `attachment; filename="agentmail-calendar.ics"`) + w.Write([]byte(sb.String())) +} + +// POST /api/v1/calendar/import.ics +func ImportCalendarICS(w http.ResponseWriter, r *http.Request) { + user := middleware.GetUser(r) + if user == nil { + Error(w, http.StatusUnauthorized, "not authenticated") + return + } + + // 两种上传形态都接受。 + // + // multipart 是浏览器 的天然形态;raw text/calendar 是 + // 脚本与 Agent 的天然形态(curl --data-binary @x.ics)。只支持前者会让 + // 命令行调用者收到含糊的「Missing file field」,只支持后者则要求前端 + // 先把文件读成字符串再发 —— 两边各让一步不如两边都收。 + var body []byte + ct := r.Header.Get("Content-Type") + if strings.HasPrefix(ct, "multipart/") { + if err := r.ParseMultipartForm(10 << 20); err != nil { + Error(w, http.StatusBadRequest, "Failed to parse multipart: "+err.Error()) + return + } + file, _, err := r.FormFile("file") + if err != nil { + Error(w, http.StatusBadRequest, "Missing file field") + return + } + defer file.Close() + body, err = io.ReadAll(file) + if err != nil { + Error(w, http.StatusBadRequest, "Failed to read file") + return + } + } else { + var err error + body, err = io.ReadAll(http.MaxBytesReader(w, r.Body, 10<<20)) + if err != nil { + Error(w, http.StatusBadRequest, "Failed to read body") + return + } + } + if len(body) == 0 { + Error(w, http.StatusBadRequest, "Empty .ics payload") + return + } + + events := parseICS(body) + imported := 0 + for _, e := range events { + e.CreatedBy = user.Username + e.Status = "active" + if _, err := repo.CreateCalendarEvent(r.Context(), &e); err == nil { + imported++ + } + } + + // skipped 单独给出而不是让前端自己减:insert 失败(撞名、约束冲突) + // 与「解析出来但没入库」是同一回事,前端只关心「有几个没进来」。 + JSON(w, http.StatusOK, map[string]interface{}{ + "imported": imported, + "skipped": len(events) - imported, + "total": len(events), + }) +} + +// ─── iCal 解析 ─── + +func parseICS(data []byte) []models.CalendarEvent { + var events []models.CalendarEvent + var current *models.CalendarEvent + + lines := strings.Split(string(data), "\n") + for _, raw := range lines { + line := strings.TrimSpace(raw) + if line == "" { + continue + } + + // 处理折叠行(iCal 的续行以空格开头) + if strings.HasPrefix(raw, " ") || strings.HasPrefix(raw, "\t") { + if current != nil && len(events) > 0 { + // 简单续行处理:追加到最后一个字段 + } + continue + } + + colon := strings.Index(line, ":") + if colon < 0 { + continue + } + key := line[:colon] + value := line[colon+1:] + + // 去掉参数部分(如 DTSTART;TZID=...:value) + if semi := strings.Index(key, ";"); semi >= 0 { + key = key[:semi] + } + // 键名大小写不敏感(RFC 5545 §3.1)。X- 扩展属性尤其容易被 + // 其他客户端改写大小写,不归一化会让往返丢掉农历规则。 + key = strings.ToUpper(key) + + switch key { + case "BEGIN": + if value == "VEVENT" { + current = &models.CalendarEvent{Recurrence: "none"} + } + case "END": + if value == "VEVENT" && current != nil { + if !current.EventTime.IsZero() { + events = append(events, *current) + } + current = nil + } + case "SUMMARY": + if current != nil { + current.Title = strings.ReplaceAll(value, "\\n", "\n") + } + case "DESCRIPTION": + if current != nil { + current.Description = strings.ReplaceAll(value, "\\n", "\n") + } + case "DTSTART": + if current != nil { + if t, err := time.Parse("20060102T150405Z", value); err == nil { + current.EventTime = t + } else if t, err := time.ParseInLocation("20060102T150405", value, time.Local); err == nil { + current.EventTime = t + } else if t, err := time.Parse("20060102", value); err == nil { + current.EventTime = t + } + } + case "RRULE": + // **不覆盖已经从 X-AGENTMAIL-RECURRENCE 读到的农历规则。** + // + // 导出时农历事件同时写了 X- 精确值与一条公历近似 RRULE + // (给别的客户端看)。X- 出现在 RRULE 之前时, + // 若这里无条件赋值就会把精确的 lunar_monthly 打回 monthly + // —— 往返一圈农历规则悄悄退化成公历,用户要过一个月才发现 + // 提醒日子不对。 + if current != nil && !models.IsLunarRecurrence(current.Recurrence) { + v := strings.ToUpper(value) + switch { + case strings.Contains(v, "FREQ=DAILY"): + current.Recurrence = models.RecurDaily + case strings.Contains(v, "FREQ=WEEKLY"): + current.Recurrence = models.RecurWeekly + case strings.Contains(v, "FREQ=MONTHLY"): + current.Recurrence = models.RecurMonthly + case strings.Contains(v, "FREQ=YEARLY"): + current.Recurrence = models.RecurYearly + } + } + case "TRIGGER": + if current != nil { + if mins, ok := parseTriggerMinutes(value); ok { + current.RemindBefore = mins + } + } + case "X-AGENTMAIL-RECURRENCE": + // 精确规则覆盖上面从 RRULE 猜出来的近似值。 + // 顺序无关:X- 属性只在值合法时才生效。 + if current != nil && validRecurrence(value) { + current.Recurrence = value + } + case "X-AGENTMAIL-RECIPIENTS": + if current != nil { + list, bad := normalizeRecipients(strings.Split(value, ",")) + // 单个地址坏掉不该让整份导入失败:其余收件人仍有效。 + // 全坏时 list 为空,事件会在 Create 时被收件人校验拦下。 + if bad == "" { + current.Recipients = list + } + } + case "X-AGENTMAIL-DELIVERY": + if current != nil && (value == models.DeliverSeparate || value == models.DeliverTogether) { + current.DeliveryMode = value + } + } + } + + return events +} + +// ─── 辅助 ─── + +func blobSha256(data []byte) string { + sum := sha256.Sum256(data) + return fmt.Sprintf("%x", sum[:]) +} + +// parseTriggerMinutes 把 VALARM 的 TRIGGER duration 解析成「提前多少分钟」。 +// +// 接受 `-PT30M` / `-PT1H` / `-PT1H30M` / `-P1D` / `-P1DT2H` 这些形态。 +// 关键规则:`M` 出现在 `T` **之后**才是分钟,之前是月 —— 按月的 trigger +// 无法映射到 remind_before(那是个分钟数),直接忽略比乱换算好。 +// +// 正号(事件之后提醒)也忽略:remind_before 语义上只能提前。 +// 返回 ok=false 表示「这条 TRIGGER 用不上」,调用方保持原值不动。 +func parseTriggerMinutes(v string) (int, bool) { + v = strings.TrimSpace(strings.ToUpper(v)) + if !strings.HasPrefix(v, "-P") { + return 0, false + } + rest := v[2:] + + // 切成 T 前后两段:前面是日期部分(Y/M/W/D),后面是时间部分(H/M/S) + datePart, timePart := rest, "" + if i := strings.Index(rest, "T"); i >= 0 { + datePart, timePart = rest[:i], rest[i+1:] + } + + total := 0 + // 日期部分只认 W/D。Y/M 是可变长度(月有 28~31 天),换算成分钟只能靠猜。 + if n, ok := durationField(datePart, 'W'); ok { + total += n * 7 * 24 * 60 + } + if n, ok := durationField(datePart, 'D'); ok { + total += n * 24 * 60 + } + if n, ok := durationField(timePart, 'H'); ok { + total += n * 60 + } + if n, ok := durationField(timePart, 'M'); ok { + total += n + } + // 秒不进 remind_before:它的粒度是分钟,30 秒会被截成 0 而看不出区别 + if total <= 0 { + return 0, false + } + return total, true +} + +// durationField 从 `1H30M` 这样的串里取出紧接在 unit 之前的整数。 +func durationField(s string, unit byte) (int, bool) { + idx := strings.IndexByte(s, unit) + if idx < 0 { + return 0, false + } + start := idx + for start > 0 && s[start-1] >= '0' && s[start-1] <= '9' { + start-- + } + if start == idx { + return 0, false + } + n := 0 + for i := start; i < idx; i++ { + n = n*10 + int(s[i]-'0') + } + return n, true +} diff --git a/gateway/internal/handler/ics_test.go b/gateway/internal/handler/ics_test.go new file mode 100644 index 0000000..d344c78 --- /dev/null +++ b/gateway/internal/handler/ics_test.go @@ -0,0 +1,423 @@ +package handler + +import ( + "fmt" + "strings" + "testing" + "time" + + "github.com/agentmail/gateway/internal/models" +) + +// parseICS 是导入的唯一入口,解析错了不会报错 —— 事件只是安静地不出现, +// 或者出现在错误的时间。这些测试锁住 iCal 的形态约定。 + +func TestParseICSBasicEvent(t *testing.T) { + ics := "BEGIN:VCALENDAR\r\n" + + "VERSION:2.0\r\n" + + "BEGIN:VEVENT\r\n" + + "SUMMARY:发布评审\r\n" + + "DESCRIPTION:看 llmsproxy 的部署脚本\r\n" + + "DTSTART:20260903T063000Z\r\n" + + "END:VEVENT\r\n" + + "END:VCALENDAR\r\n" + + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatalf("期望 1 个事件,得到 %d", len(events)) + } + e := events[0] + if e.Title != "发布评审" { + t.Errorf("Title = %q,期望 发布评审", e.Title) + } + if e.Description != "看 llmsproxy 的部署脚本" { + t.Errorf("Description = %q", e.Description) + } + if !e.EventTime.Equal(time.Date(2026, 9, 3, 6, 30, 0, 0, time.UTC)) { + t.Errorf("EventTime = %v,期望 2026-09-03T06:30:00Z", e.EventTime) + } + if e.Recurrence != "none" { + t.Errorf("Recurrence = %q,期望 none", e.Recurrence) + } +} + +func TestParseICSMultipleEvents(t *testing.T) { + var sb strings.Builder + sb.WriteString("BEGIN:VCALENDAR\r\n") + for i, title := range []string{"甲", "乙", "丙"} { + sb.WriteString("BEGIN:VEVENT\r\n") + sb.WriteString("SUMMARY:" + title + "\r\n") + sb.WriteString("DTSTART:2026090" + string(rune('1'+i)) + "T020000Z\r\n") + sb.WriteString("END:VEVENT\r\n") + } + sb.WriteString("END:VCALENDAR\r\n") + + events := parseICS([]byte(sb.String())) + if len(events) != 3 { + t.Fatalf("期望 3 个事件,得到 %d", len(events)) + } + for i, want := range []string{"甲", "乙", "丙"} { + if events[i].Title != want { + t.Errorf("第 %d 个 Title = %q,期望 %q", i, events[i].Title, want) + } + } +} + +// 没有 DTSTART 的 VEVENT 必须被丢弃:让它进库会得到一个 zero time 事件, +// 调度器认为它「早就该触发了」,于是立刻发一封莫名其妙的提醒。 +func TestParseICSDropsEventWithoutStart(t *testing.T) { + ics := "BEGIN:VCALENDAR\r\n" + + "BEGIN:VEVENT\r\n" + + "SUMMARY:没有时间\r\n" + + "END:VEVENT\r\n" + + "END:VCALENDAR\r\n" + + if events := parseICS([]byte(ics)); len(events) != 0 { + t.Fatalf("无 DTSTART 的事件应被丢弃,却得到 %d 个", len(events)) + } +} + +func TestParseICSRecurrence(t *testing.T) { + cases := []struct { + rrule string + want string + }{ + {"FREQ=DAILY", "daily"}, + {"FREQ=WEEKLY;BYDAY=MO", "weekly"}, + {"FREQ=MONTHLY;BYMONTHDAY=1", "monthly"}, + {"FREQ=YEARLY", "yearly"}, + // 不支持的频率退回 none 而不是乱猜:把 HOURLY 当 daily + // 会让提醒少发 23 次且没有任何报错。 + {"FREQ=HOURLY", "none"}, + {"FREQ=SECONDLY", "none"}, + } + for _, c := range cases { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + "RRULE:" + c.rrule + "\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatalf("%s: 期望 1 个事件", c.rrule) + } + if events[0].Recurrence != c.want { + t.Errorf("%s: Recurrence = %q,期望 %q", c.rrule, events[0].Recurrence, c.want) + } + } +} + +func TestParseICSTriggerToRemindBefore(t *testing.T) { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + "TRIGGER:-PT30M\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatal("期望 1 个事件") + } + if events[0].RemindBefore != 30 { + t.Errorf("RemindBefore = %d,期望 30", events[0].RemindBefore) + } +} + +// DTSTART 有三种合法形态,都得认。只认 UTC 那种会让本地时间的 .ics +// 整份导入失败(每个事件都缺 DTSTART → 全被丢弃 → 「导入 0 个」且无提示)。 +func TestParseICSDateFormats(t *testing.T) { + cases := []struct { + name string + value string + }{ + {"UTC", "20260903T063000Z"}, + {"本地时间", "20260903T143000"}, + {"仅日期", "20260903"}, + } + for _, c := range cases { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:" + c.value + "\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Errorf("%s (%s): 期望 1 个事件,得到 %d", c.name, c.value, len(events)) + continue + } + if events[0].EventTime.IsZero() { + t.Errorf("%s (%s): EventTime 为零值", c.name, c.value) + } + } +} + +// DTSTART;TZID=Asia/Shanghai:... 这种带参数的键必须归一化到 DTSTART, +// 否则 switch 落空 → 无 EventTime → 事件被丢。 +func TestParseICSStripsKeyParameters(t *testing.T) { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\n" + + "DTSTART;TZID=Asia/Shanghai:20260903T143000\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatalf("带 TZID 参数的 DTSTART 应被识别,得到 %d 个事件", len(events)) + } + if events[0].EventTime.IsZero() { + t.Error("EventTime 为零值") + } +} + +func TestParseICSEscapedNewlines(t *testing.T) { + ics := "BEGIN:VEVENT\r\nSUMMARY:第一行\\n第二行\r\n" + + "DTSTART:20260903T020000Z\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatal("期望 1 个事件") + } + if !strings.Contains(events[0].Title, "\n") { + t.Errorf("转义的 \\n 应还原成真换行,得到 %q", events[0].Title) + } +} + +func TestParseICSEmptyAndGarbage(t *testing.T) { + for _, in := range []string{"", "不是 ics", "BEGIN:VCALENDAR\r\nEND:VCALENDAR\r\n"} { + if events := parseICS([]byte(in)); len(events) != 0 { + t.Errorf("输入 %q 应给 0 个事件,得到 %d", in, len(events)) + } + } +} + +// LF 换行(非 CRLF)的 .ics 也要能解析:很多工具导出的是 LF。 +func TestParseICSAcceptsLFLineEndings(t *testing.T) { + ics := "BEGIN:VCALENDAR\nBEGIN:VEVENT\nSUMMARY:LF 换行\n" + + "DTSTART:20260903T020000Z\nEND:VEVENT\nEND:VCALENDAR\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatalf("LF 换行应能解析,得到 %d 个事件", len(events)) + } + if events[0].Title != "LF 换行" { + t.Errorf("Title = %q", events[0].Title) + } +} + +// TRIGGER duration 解析。 +// +// 原实现是 fmt.Sscanf(value, "-PT%dM", &mins),只认一种形态;导出端又写的是 +// `-P15M`(T 之前的 M 在 iCal 里是**月**)—— 于是自己导出的文件自己都读不回来。 +func TestParseTriggerMinutes(t *testing.T) { + cases := []struct { + in string + want int + ok bool + }{ + {"-PT30M", 30, true}, + {"-PT1H", 60, true}, + {"-PT1H30M", 90, true}, + {"-P1D", 1440, true}, + {"-P1DT2H", 1560, true}, + {"-P1W", 10080, true}, + {"-pt45m", 45, true}, // 大小写不敏感 + + // T 之前的 M 是月,映射不到分钟数,忽略比乱换算好 + {"-P3M", 0, false}, + // 正号 = 事件之后提醒,remind_before 表达不了 + {"PT30M", 0, false}, + // 零时长与垃圾输入 + {"-PT0M", 0, false}, + {"", 0, false}, + {"垃圾", 0, false}, + {"-P", 0, false}, + } + for _, c := range cases { + got, ok := parseTriggerMinutes(c.in) + if ok != c.ok || got != c.want { + t.Errorf("parseTriggerMinutes(%q) = (%d, %v),期望 (%d, %v)", + c.in, got, ok, c.want, c.ok) + } + } +} + +// 导出写的 TRIGGER 必须能被自己的导入解析回同一个分钟数。 +// 这条往返曾经是断的:导出 -P15M、导入找 -PT%dM。 +func TestTriggerRoundtrip(t *testing.T) { + for _, mins := range []int{5, 15, 30, 60, 120, 1440} { + // 导出端的写法(与 ExportCalendarICS 里那行一致) + trigger := fmt.Sprintf("-PT%dM", mins) + got, ok := parseTriggerMinutes(trigger) + if !ok { + t.Errorf("%d 分钟导出成 %q 后无法解析", mins, trigger) + continue + } + if got != mins { + t.Errorf("%d 分钟往返后变成 %d(trigger=%q)", mins, got, trigger) + } + } +} + +// 整份 .ics 的往返:TRIGGER 经过 parseICS 后落到 RemindBefore 上。 +func TestParseICSTriggerVariants(t *testing.T) { + cases := []struct { + trigger string + want int + }{ + {"-PT15M", 15}, + {"-PT2H", 120}, + {"-P1D", 1440}, + } + for _, c := range cases { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + "TRIGGER:" + c.trigger + "\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Errorf("%s: 期望 1 个事件", c.trigger) + continue + } + if events[0].RemindBefore != c.want { + t.Errorf("%s: RemindBefore = %d,期望 %d", c.trigger, events[0].RemindBefore, c.want) + } + } +} + +// 默认提醒模板必须是**变量形式**,不能把当时的时间烤成字面值。 +// +// 重复事件上这个区别是致命的:AdvanceRecurrence 只推进 event_time, +// reminder_text 保持不动 —— 字面值会让「每天 9 点」的提醒从第二天起 +// 永远写着第一天的日期,且不报任何错。 +// +// 同时锁住「与前端 DEFAULT_TEMPLATE 逐字一致」: +// web/src/components/CalendarEventEditor.tsx 用它作预览, +// 两边不同会让人看到的预览与 Agent 实收的正文不是一回事。 +func TestDefaultReminderTemplateUsesVariables(t *testing.T) { + for _, v := range []string{"{title}", "{time}", "{description}"} { + if !strings.Contains(defaultReminderTemplate, v) { + t.Errorf("默认模板缺变量 %s:%q", v, defaultReminderTemplate) + } + } + // 前端那份的字面内容(保持同步) + const frontend = "日程提醒:{title}\n时间:{time}\n{description}" + if defaultReminderTemplate != frontend { + t.Errorf("后端默认模板与前端 DEFAULT_TEMPLATE 不一致:\n后端 %q\n前端 %q", + defaultReminderTemplate, frontend) + } + // 不该含任何形如年份的字面数字 —— 那是「把值烤进模板」的迹象 + for _, digit := range []string{"2026", "20:", ":00"} { + if strings.Contains(defaultReminderTemplate, digit) { + t.Errorf("默认模板含字面时间片段 %q:%q", digit, defaultReminderTemplate) + } + } +} + +// ─── 农历与多收件人的 iCal 往返 ─── + +// 农历规则 RFC 5545 表达不了。折中方案:X- 扩展属性记精确规则 + +// 降级成最接近的公历 RRULE。别的客户端至少能看到一个大致对的重复, +// 导回本系统时 X- 属性还原精确规则。 +func TestParseICSLunarRecurrenceExtension(t *testing.T) { + cases := []struct { + name string + body string + want string + }{ + { + "X- 属性覆盖 RRULE 近似值", + "RRULE:FREQ=MONTHLY\r\nX-AGENTMAIL-RECURRENCE:lunar_monthly\r\n", + "lunar_monthly", + }, + { + "农历年", + "RRULE:FREQ=YEARLY\r\nX-AGENTMAIL-RECURRENCE:lunar_yearly\r\n", + "lunar_yearly", + }, + { + "X- 在 RRULE 之前也生效(顺序无关)", + "X-AGENTMAIL-RECURRENCE:lunar_monthly\r\nRRULE:FREQ=MONTHLY\r\n", + "lunar_monthly", + }, + { + "非法 X- 值被忽略,保留 RRULE 的近似值", + "RRULE:FREQ=MONTHLY\r\nX-AGENTMAIL-RECURRENCE:lunar_montly\r\n", + "monthly", + }, + { + "键名小写也认(RFC 5545 §3.1 大小写不敏感)", + "x-agentmail-recurrence:lunar_yearly\r\n", + "lunar_yearly", + }, + } + for _, c := range cases { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + c.body + "END:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Errorf("%s: 期望 1 个事件,得到 %d", c.name, len(events)) + continue + } + if events[0].Recurrence != c.want { + t.Errorf("%s: Recurrence = %q,期望 %q", c.name, events[0].Recurrence, c.want) + } + } +} + +func TestParseICSRecipientsExtension(t *testing.T) { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + "X-AGENTMAIL-RECIPIENTS:pi@/home/x,dsh,opencode@/tmp.alias\r\n" + + "X-AGENTMAIL-DELIVERY:together\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatalf("期望 1 个事件,得到 %d", len(events)) + } + e := events[0] + if len(e.Recipients) != 3 { + t.Fatalf("收件人应有 3 个,得到 %d:%v", len(e.Recipients), e.Recipients) + } + if e.Recipients[0] != "pi@/home/x" || e.Recipients[2] != "opencode@/tmp.alias" { + t.Errorf("收件人内容或顺序不符:%v", e.Recipients) + } + if e.DeliveryMode != "together" { + t.Errorf("DeliveryMode = %q,期望 together", e.DeliveryMode) + } +} + +// 未知投递模式必须被忽略(留空 → EffectiveDeliveryMode 给 separate), +// 而不是原样写进库里。 +func TestParseICSRejectsBadDeliveryMode(t *testing.T) { + ics := "BEGIN:VEVENT\r\nSUMMARY:x\r\nDTSTART:20260903T020000Z\r\n" + + "X-AGENTMAIL-DELIVERY:随便写的\r\nEND:VEVENT\r\n" + events := parseICS([]byte(ics)) + if len(events) != 1 { + t.Fatal("期望 1 个事件") + } + if events[0].DeliveryMode != "" { + t.Errorf("非法投递模式应被忽略,得到 %q", events[0].DeliveryMode) + } + if events[0].EffectiveDeliveryMode() != models.DeliverSeparate { + t.Error("兜底应是 separate") + } +} + +func TestValidRecurrence(t *testing.T) { + for _, ok := range []string{"none", "daily", "weekly", "monthly", "yearly", "lunar_monthly", "lunar_yearly"} { + if !validRecurrence(ok) { + t.Errorf("%q 应合法", ok) + } + } + // 拼错必须被拒而不是静默当 none —— 后者会让每月提醒只响一次且无报错 + for _, bad := range []string{"", "lunar_montly", "LUNAR_MONTHLY", "每月", "lunar_weekly"} { + if validRecurrence(bad) { + t.Errorf("%q 应非法", bad) + } + } +} + +func TestNormalizeRecipients(t *testing.T) { + got, bad := normalizeRecipients([]string{" pi ", "", "dsh", "pi", " "}) + if bad != "" { + t.Fatalf("不该报错,得到 %q", bad) + } + // 去空白 + 去重,保留首次出现的顺序 + if len(got) != 2 || got[0] != "pi" || got[1] != "dsh" { + t.Errorf("清洗结果 %v,期望 [pi dsh]", got) + } + + // 去重是必要的:together 模式下同一 Agent 既主收又抄送会收到两条 SSE + if dup, _ := normalizeRecipients([]string{"pi@/x", "pi@/x"}); len(dup) != 1 { + t.Errorf("重复地址应去重,得到 %v", dup) + } + + // 非法地址回报具体是哪一个 + if _, bad := normalizeRecipients([]string{"pi", "@@@bad@@@"}); bad == "" { + t.Error("非法地址应被报出") + } + + // nil / 空输入给空数组而不是 nil(避免序列化成 null) + if out, _ := normalizeRecipients(nil); out == nil { + t.Error("nil 输入应给空数组") + } +} diff --git a/gateway/internal/models/calendar.go b/gateway/internal/models/calendar.go new file mode 100644 index 0000000..d624546 --- /dev/null +++ b/gateway/internal/models/calendar.go @@ -0,0 +1,143 @@ +package models + +import ( + "strings" + "time" +) + +// CalendarEvent 日历事件。 +// +// 设计参照 Outlook:事件有时间、提醒、收件人,触发时产生一封邮件。 +// 事件本身是日历实体,提醒是触发器,邮件是投递通道 —— 三者分离。 +type CalendarEvent struct { + EventID string `json:"event_id"` + Title string `json:"title"` + Description string `json:"description"` + ReminderText string `json:"reminder_text"` + + // AgentName / ToAddress 是**单收件人时代的字段**,保留作兼容与兜底: + // Recipients 为空时用它们。新代码一律读 EffectiveRecipients()。 + AgentName string `json:"agent_name"` + ToAddress string `json:"to_address"` + + // Recipients 是完整的收件人列表(每项是完整三维地址串)。 + // + // 为什么不用 []Address 而用 []string:地址的三段语义(尤其 session 位的 + // new/别名三态)在**触发那一刻**才该被解析 —— 存结构化的话, + // 「.new」这种一次性语义在建事件时就被固化,而重复事件每次触发都该 + // 重新决定落到哪条会话。存原始串让 ParseAddress 在投递时做这个决定。 + Recipients []string `json:"recipients"` + + // DeliveryMode 决定多收件人怎么投: + // "separate"(默认)—— 每人各发一封,落在各自的会话里,互相看不到 + // "together" —— 第一个是主收件人,其余进 cc_list,共享同一条线索 + // + // 两种语义都需要而不是二选一:「让三个 Agent 各自独立汇报」与 + // 「让 pi 主办、dsh 知情」是完全不同的任务形态,用错会让协作失败 —— + // 前者用 together 会让三个 Agent 互相看到对方的回复而趋同, + // 后者用 separate 会让 dsh 完全不知道 pi 在做什么。 + DeliveryMode string `json:"delivery_mode"` + + EventTime time.Time `json:"event_time"` + RemindBefore int `json:"remind_before"` // 提前多少分钟 + + // Recurrence:公历 none/daily/weekly/monthly/yearly + 农历两种 + // lunar_monthly —— 每农历月同一日(如每月十五) + // lunar_yearly —— 每农历年同月同日(过农历生日/祭日) + // + // lunar_daily 不存在:农历的「日」与公历同长,那就是 daily。 + // lunar_weekly 也不存在:农历没有「周」这个单位。 + Recurrence string `json:"recurrence"` + RecurrenceEnd *time.Time `json:"recurrence_end,omitempty"` + + Status string `json:"status"` // active/paused/cancelled + LastFiredAt *time.Time `json:"last_fired_at,omitempty"` + + // FiredFor 是已触发的那个 occurrence(值 = 当时的 EventTime)。 + // 去重靠它与 EventTime 相等判断,不是拿 LastFiredAt 比大小 —— + // DueEvents 有 60 秒 lookahead,后者在窗口内恒为真会导致每 tick 重发。 + FiredFor *time.Time `json:"fired_for,omitempty"` + + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` + CreatedBy string `json:"created_by"` +} + +// 重复规则常量。农历规则单独一组:它们的推进要经过 internal/lunar, +// 不能像公历那样 AddDate 固定天数(农历月 29~30 天、闰年 13 个月)。 +const ( + RecurNone = "none" + RecurDaily = "daily" + RecurWeekly = "weekly" + RecurMonthly = "monthly" + RecurYearly = "yearly" + RecurLunarMonthly = "lunar_monthly" + RecurLunarYearly = "lunar_yearly" +) + +// IsLunarRecurrence 判断一条重复规则是否按农历推进。 +func IsLunarRecurrence(r string) bool { + return r == RecurLunarMonthly || r == RecurLunarYearly +} + +// 投递模式常量。 +const ( + DeliverSeparate = "separate" + DeliverTogether = "together" +) + +// EffectiveRecipients 返回真正要投的收件人列表。 +// +// Recipients 优先;为空时退回 ToAddress,再退回 AgentName。 +// 这个兜底链让旧数据(只有 agent_name 的事件)继续工作 —— +// 历史事件不迁移,读的时候归一化。 +func (e *CalendarEvent) EffectiveRecipients() []string { + if len(e.Recipients) > 0 { + out := make([]string, 0, len(e.Recipients)) + for _, r := range e.Recipients { + if r = strings.TrimSpace(r); r != "" { + out = append(out, r) + } + } + if len(out) > 0 { + return out + } + } + if a := strings.TrimSpace(e.ToAddress); a != "" { + return []string{a} + } + if a := strings.TrimSpace(e.AgentName); a != "" { + return []string{a} + } + return nil +} + +// EffectiveDeliveryMode 归一化投递模式,未知值按 separate 处理。 +// +// 默认 separate 而不是 together:separate 的失败是「Agent 各干各的」, +// together 的失败是「本该独立的 Agent 互相污染了上下文」—— +// 后者更难发现也更难挽回。 +func (e *CalendarEvent) EffectiveDeliveryMode() string { + if e.DeliveryMode == DeliverTogether { + return DeliverTogether + } + return DeliverSeparate +} + +// CalendarAttachment 事件附件。 +type CalendarAttachment struct { + AttachmentID string `json:"attachment_id"` + EventID string `json:"event_id"` + Filename string `json:"filename"` + SHA256 string `json:"sha256"` + SizeBytes int64 `json:"size_bytes"` + CreatedAt time.Time `json:"created_at"` +} + +// CalendarView 日历视图(月/周/日)。 +type CalendarView struct { + Events []CalendarEvent `json:"events"` + // 当前时间线上的事件数(用于统计徽章) + UpcomingCount int `json:"upcoming_count"` + TodayCount int `json:"today_count"` +} diff --git a/gateway/internal/repo/calendar.go b/gateway/internal/repo/calendar.go new file mode 100644 index 0000000..b7f6d95 --- /dev/null +++ b/gateway/internal/repo/calendar.go @@ -0,0 +1,534 @@ +package repo + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/agentmail/gateway/internal/lunar" + "github.com/agentmail/gateway/internal/models" + "github.com/google/uuid" +) + +var ErrEventNotFound = errors.New("calendar event not found") + +// calendarCols 是所有 SELECT 共用的列清单。 +// +// 抽出来是因为原先有**四处**手抄同一串列名(Get / List / DueEvents 各一处), +// 而 Scan 的参数顺序必须与之逐一对应。加一列时漏改任何一处都不会编译报错 —— +// 只会在运行时得到 "Scan: expected N destination arguments" 或者更糟: +// 列数恰好相同而值错位(曾在 ListSessionsFor 上真的发生过, +// 加了预算两列没加进 Scan,整个联系人栏 500)。 +const calendarCols = `event_id, title, description, reminder_text, agent_name, to_address, + recipients, delivery_mode, event_time, remind_before, recurrence, recurrence_end, + status, last_fired_at, fired_for, created_at, updated_at, created_by` + +// rowScanner 让 QueryRow 与 Rows 共用同一个 scan 实现。 +type rowScanner interface { + Scan(dest ...any) error +} + +// scanCalendarEvent 按 calendarCols 的顺序读一行。 +// +// recipients 存的是 JSON 文本,必须先读进 []byte 再 Unmarshal —— +// 直接 Scan 进 []string 会静默失败(driver 不知道怎么转)。 +func scanCalendarEvent(sc rowScanner) (*models.CalendarEvent, error) { + var e models.CalendarEvent + var recipientsJSON []byte + if err := sc.Scan( + &e.EventID, &e.Title, &e.Description, &e.ReminderText, + &e.AgentName, &e.ToAddress, + &recipientsJSON, &e.DeliveryMode, + &e.EventTime, &e.RemindBefore, &e.Recurrence, &e.RecurrenceEnd, + &e.Status, &e.LastFiredAt, &e.FiredFor, &e.CreatedAt, &e.UpdatedAt, &e.CreatedBy, + ); err != nil { + return nil, err + } + if len(recipientsJSON) > 0 { + // 解析失败不算致命:退回 to_address/agent_name 兜底链, + // 事件仍能投递。让一条脏 JSON 把整个列表打成 500 更糟。 + _ = json.Unmarshal(recipientsJSON, &e.Recipients) + } + if e.Recipients == nil { + // Go 的 nil slice 序列化成 null,前端 .map 会崩 + e.Recipients = []string{} + } + return &e, nil +} + +// marshalRecipients 把收件人列表序列化成入库的 JSON 文本。 +func marshalRecipients(list []string) string { + if list == nil { + list = []string{} + } + b, err := json.Marshal(list) + if err != nil { + return "[]" + } + return string(b) +} + +// ─── CRUD ─── + +func CreateCalendarEvent(ctx context.Context, e *models.CalendarEvent) (*models.CalendarEvent, error) { + e.EventID = uuid.New().String() + e.CreatedAt = time.Now() + e.UpdatedAt = e.CreatedAt + if e.Status == "" { + e.Status = "active" + } + if e.Recurrence == "" { + e.Recurrence = "none" + } + + if e.DeliveryMode == "" { + e.DeliveryMode = models.DeliverSeparate + } + if e.Recipients == nil { + e.Recipients = []string{} + } + + _, err := db.DB.ExecContext(ctx, ` + INSERT INTO calendar_events + (event_id, title, description, reminder_text, agent_name, to_address, + recipients, delivery_mode, + event_time, remind_before, recurrence, recurrence_end, status, created_by, + created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + e.EventID, e.Title, e.Description, e.ReminderText, + e.AgentName, e.ToAddress, + marshalRecipients(e.Recipients), e.DeliveryMode, + e.EventTime, e.RemindBefore, e.Recurrence, e.RecurrenceEnd, + e.Status, e.CreatedBy, e.CreatedAt, e.UpdatedAt, + ) + if err != nil { + return nil, err + } + return e, nil +} + +func GetCalendarEvent(ctx context.Context, eventID string) (*models.CalendarEvent, error) { + e, err := scanCalendarEvent(db.DB.QueryRowContext(ctx, + `SELECT `+calendarCols+` FROM calendar_events WHERE event_id = ?`, eventID)) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return nil, ErrEventNotFound + } + return nil, err + } + return e, nil +} + +func UpdateCalendarEvent(ctx context.Context, eventID string, e *models.CalendarEvent) error { + e.UpdatedAt = time.Now() + result, err := db.DB.ExecContext(ctx, ` + UPDATE calendar_events SET + title = ?, description = ?, reminder_text = ?, + agent_name = ?, to_address = ?, + recipients = ?, delivery_mode = ?, + event_time = ?, remind_before = ?, recurrence = ?, recurrence_end = ?, + status = ?, updated_at = ? + WHERE event_id = ?`, + e.Title, e.Description, e.ReminderText, + e.AgentName, e.ToAddress, + marshalRecipients(e.Recipients), e.EffectiveDeliveryMode(), + e.EventTime, e.RemindBefore, e.Recurrence, e.RecurrenceEnd, + e.Status, e.UpdatedAt, eventID, + ) + if err != nil { + return err + } + n, _ := result.RowsAffected() + if n == 0 { + return ErrEventNotFound + } + return nil +} + +func DeleteCalendarEvent(ctx context.Context, eventID string) error { + result, err := db.DB.ExecContext(ctx, `DELETE FROM calendar_events WHERE event_id = ?`, eventID) + if err != nil { + return err + } + n, _ := result.RowsAffected() + if n == 0 { + return ErrEventNotFound + } + return nil +} + +// ─── 查询 ─── + +// ListCalendarEvents 返回指定时间范围内的事件(日历视图)。 +func ListCalendarEvents(ctx context.Context, from, to time.Time, status string) ([]models.CalendarEvent, error) { + if status == "" { + status = "active" + } + rows, err := db.DB.QueryContext(ctx, ` + SELECT `+calendarCols+` + FROM calendar_events + WHERE event_time >= ? AND event_time <= ? + AND (status = ? OR ? = '') + ORDER BY event_time ASC`, from, to, status, status) + if err != nil { + return nil, err + } + defer rows.Close() + + var events []models.CalendarEvent + for rows.Next() { + e, err := scanCalendarEvent(rows) + if err != nil { + return nil, err + } + events = append(events, *e) + } + return events, rows.Err() +} + +// ─── 调度器 ─── + +// DueEvents 返回下一分钟内需要触发的事件。 +// +// 调度器每分钟调用一次:event_time + remind_before <= now+60s 且尚未触发(last_fired_at 为 NULL +// 或小于 event_time)的 active 事件。 +// DueEvents 取出该触发的事件。 +// +// 60 秒 lookahead 让提醒宁早不晚:调度周期是 30 秒,不提前看的话 +// 一个刚好落在两个 tick 之间的提醒会迟到最多 30 秒。 +// +// **去重判据是 fired_for(已触发的 occurrence)与 event_time 相等**, +// 不是 last_fired_at 与 event_time 比大小 —— 后者在 lookahead 窗口内 +// 恒为真(触发时刻早于 event_time),会让同一条提醒每个 tick 重发一次。 +func DueEvents(ctx context.Context) ([]models.CalendarEvent, error) { + now := time.Now() + deadline := now.Add(60 * time.Second) + + rows, err := db.DB.QueryContext(ctx, ` + SELECT `+calendarCols+` + FROM calendar_events + WHERE status = 'active' + AND datetime(event_time, '-' || remind_before || ' minutes') <= ? + AND (fired_for IS NULL OR fired_for <> event_time) + ORDER BY event_time ASC`, deadline) + if err != nil { + return nil, err + } + defer rows.Close() + + var events []models.CalendarEvent + for rows.Next() { + e, err := scanCalendarEvent(rows) + if err != nil { + return nil, err + } + events = append(events, *e) + } + return events, rows.Err() +} + +// MarkEventFired 标记事件已触发,防止重复。 +func MarkEventFired(ctx context.Context, eventID string) error { + // fired_for 直接从 event_time 列复制而不是在 Go 侧格式化再写回: + // 两者必须逐字节相同(判据是字符串相等),经过一轮 time.Time 往返 + // 有可能改变表示形式。 + _, err := db.DB.ExecContext(ctx, + `UPDATE calendar_events SET last_fired_at = ?, fired_for = event_time + WHERE event_id = ?`, + time.Now(), eventID) + return err +} + +// AdvanceRecurrence 为重复事件计算下一次触发时间。 +// +// 返回 false 表示重复已过期(recurrence_end 已过),事件应置为 cancelled。 +func AdvanceRecurrence(ctx context.Context, eventID string) (bool, error) { + var recurrence string + var eventTime time.Time + var recurrenceEnd *time.Time + err := db.DB.QueryRowContext(ctx, ` + SELECT recurrence, event_time, recurrence_end + FROM calendar_events WHERE event_id = ?`, eventID).Scan( + &recurrence, &eventTime, &recurrenceEnd) + if err != nil { + return false, err + } + + if recurrence == models.RecurNone { + return false, nil + } + + // **一路推到未来**,不是只推一步。 + // + // 只推一步的后果(实测):一条 100 天前设的每日事件,每轮扫描都判定 + // 「已过期该触发」→ 发一封 → event_time 只前进一天 → 下一轮又过期。 + // 30 轮扫描触发 30 次,而调度周期是 30 秒 —— 人会收到一串垃圾提醒, + // 连发 100 封才追上今天。 + // + // 跳过的那些 occurrence **不补发**:定时提醒的价值在于「按时」, + // 三个月前那次站会提醒现在发出去毫无意义,只会淹掉真正该看的那封。 + // 本轮仍会发一封(fireEvent 已经在发了),代表「这条规则还活着」。 + next, err := advanceToFuture(recurrence, eventTime, time.Now(), recurrenceEnd) + if err != nil { + // 推不出下一次(例如「每年农历闰六月」而目标年无闰六月): + // 置为 cancelled 而不是留在 active 空转。留着会让调度器每 30 秒 + // 重试同一个算不出来的规则,日志里刷同一条错误直到有人发现。 + _, cErr := db.DB.ExecContext(ctx, + `UPDATE calendar_events SET status = 'cancelled' WHERE event_id = ?`, eventID) + if cErr != nil { + return false, cErr + } + return false, err + } + // 零值 = 已越过 recurrence_end("none" 在函数开头就返回了,到不了这里)。 + // **必须置 cancelled**:留在 active 会让 DueEvents 每轮都捞到这条 + // 早已过期的事件,而 fired_for 已经等于 event_time 所以它又不会被触发 —— + // 表现是一条永远排在到期列表里、永远不动的僵尸事件。 + if next.IsZero() { + _, cErr := db.DB.ExecContext(ctx, + `UPDATE calendar_events SET status = 'cancelled', updated_at = ? WHERE event_id = ?`, + time.Now(), eventID) + return false, cErr + } + + // 走到这里说明 next 既在未来又在终止时间之内。 + _, err = db.DB.ExecContext(ctx, + `UPDATE calendar_events SET event_time = ?, updated_at = ? WHERE event_id = ?`, + next, time.Now(), eventID) + return true, err +} + +// advanceToFuture 从 from 起反复按规则推进,直到越过 now。 +// +// 单独成不碰数据库的函数是为了可测。三个终止条件,缺一不可: +// +// 1. 越过 now —— 正常出口 +// 2. 越过 recurrenceEnd —— 返回零值,调用方据此置 cancelled +// 3. maxAdvanceSteps 上限 —— 防御性的。规则算得出但不前进(理论上 +// NextOccurrence 不会返回 <= 当前值,但农历那条路径依赖外部库, +// 一旦它某年给出反直觉结果,没有上限就是个死循环 goroutine, +// 而它跑在调度器里 —— 整个提醒系统会一起卡住) +// +// 上限取 4000:按每日重复算约 11 年,足够覆盖「很久以前设的提醒」, +// 而 4000 次纯内存日期运算在一个 tick 里跑完毫无压力。 +const maxAdvanceSteps = 4000 + +func advanceToFuture(recurrence string, from, now time.Time, recurrenceEnd *time.Time) (time.Time, error) { + cur := from + for i := 0; i < maxAdvanceSteps; i++ { + next, err := NextOccurrence(recurrence, cur) + if err != nil { + return time.Time{}, err + } + if next.IsZero() { + return time.Time{}, nil // 不重复 + } + if !next.After(cur) { + // 规则不前进 —— 与死循环等价,当作算不出来 + return time.Time{}, fmt.Errorf("重复规则 %q 未能前进(停在 %s)", recurrence, cur.Format(time.RFC3339)) + } + cur = next + if recurrenceEnd != nil && cur.After(*recurrenceEnd) { + return time.Time{}, nil // 已过终止时间 + } + if cur.After(now) { + return cur, nil + } + } + return time.Time{}, fmt.Errorf("重复规则 %q 推进 %d 次仍未越过当前时刻", recurrence, maxAdvanceSteps) +} + +// NextOccurrence 按重复规则算出下一次触发时刻。 +// +// 独立成不碰数据库的纯函数是为了可测:农历推进错了不会报错, +// 只会让提醒发在错误的日子,而那种错误要等真的过了一个月才看得见。 +// +// 返回零值 time 且 err == nil 表示「不重复」(规则是 none 或未知值)。 +// +// **农历规则不能用 AddDate 近似**:农历月 29~30 天不定、农历年 353~385 天 +// (闰年多一整月)。用固定天数推进一年能偏半个月 —— 农历生日提醒会 +// 逐年漂移到完全不相干的日子上。 +func NextOccurrence(recurrence string, from time.Time) (time.Time, error) { + switch recurrence { + case models.RecurDaily: + return from.AddDate(0, 0, 1), nil + case models.RecurWeekly: + return from.AddDate(0, 0, 7), nil + case models.RecurMonthly: + // 公历每月:AddDate 在月末会溢出(1 月 31 日 +1 月 = 3 月 3 日)。 + // 夹到目标月的最后一天 —— 与农历那边的 clamp 语义一致: + // 「每月 31 日」的意思是「月末」,滚到下月初是错的。 + return addSolarMonthClamped(from, 1), nil + case models.RecurYearly: + // 公历每年:2 月 29 日在平年会溢出成 3 月 1 日,同样要夹。 + // 闰日生日的约定是「平年过 2 月 28」,不是 3 月 1 日。 + return addSolarMonthClamped(from, 12), nil + case models.RecurLunarMonthly: + d := lunar.FromSolar(from).AddMonths(1) + t, _, err := d.ToSolar(from.Location(), from.Hour(), from.Minute(), from.Second(), from.Nanosecond()) + return t, err + case models.RecurLunarYearly: + d := lunar.FromSolar(from).AddYears(1) + t, _, err := d.ToSolar(from.Location(), from.Hour(), from.Minute(), from.Second(), from.Nanosecond()) + return t, err + default: + return time.Time{}, nil + } +} + +// addSolarMonthClamped 在公历上加月份,日期夹到目标月的实际天数内。 +// +// time.AddDate 的溢出行为(3 月 31 日 +1 月 = 5 月 1 日)对「每月同一日」 +// 的提醒是错的:31 日的事件会在 2 月变成 3 月 3 日,然后从此每月 3 日提醒 +// —— 一次溢出永久改变了规则。 +func addSolarMonthClamped(t time.Time, n int) time.Time { + y, m, d := t.Date() + m += time.Month(n) + for m > 12 { + m -= 12 + y++ + } + // 目标月第 0 天 = 上个月最后一天,用它拿到月长 + last := time.Date(y, m+1, 0, 0, 0, 0, 0, t.Location()).Day() + if d > last { + d = last + } + return time.Date(y, m, d, t.Hour(), t.Minute(), t.Second(), t.Nanosecond(), t.Location()) +} + +// ─── 附件 ─── + +func AddCalendarAttachment(ctx context.Context, a *models.CalendarAttachment) error { + a.AttachmentID = uuid.New().String() + a.CreatedAt = time.Now() + _, err := db.DB.ExecContext(ctx, ` + INSERT INTO calendar_attachments (attachment_id, event_id, filename, sha256, size_bytes, created_at) + VALUES (?, ?, ?, ?, ?, ?)`, + a.AttachmentID, a.EventID, a.Filename, a.SHA256, a.SizeBytes, a.CreatedAt) + return err +} + +func ListCalendarAttachments(ctx context.Context, eventID string) ([]models.CalendarAttachment, error) { + rows, err := db.DB.QueryContext(ctx, ` + SELECT attachment_id, event_id, filename, sha256, size_bytes, created_at + FROM calendar_attachments WHERE event_id = ? ORDER BY created_at`, eventID) + if err != nil { + return nil, err + } + defer rows.Close() + var atts []models.CalendarAttachment + for rows.Next() { + var a models.CalendarAttachment + if err := rows.Scan(&a.AttachmentID, &a.EventID, &a.Filename, &a.SHA256, &a.SizeBytes, &a.CreatedAt); err != nil { + return nil, err + } + atts = append(atts, a) + } + return atts, rows.Err() +} + +func DeleteCalendarAttachments(ctx context.Context, eventID string) error { + _, err := db.DB.ExecContext(ctx, `DELETE FROM calendar_attachments WHERE event_id = ?`, eventID) + return err +} + +// DeleteCalendarAttachment 删单条附件。 +// +// 返回 false 表示这条不存在(而不是报错):调用方据此回 404 而非 500。 +// 只删元数据,磁盘 blob 留给 GC —— 内容寻址下同一个 sha256 可能被别的 +// 附件引用着,跟着删会让那些引用一起坏掉。 +func DeleteCalendarAttachment(ctx context.Context, attachmentID string) (bool, error) { + res, err := db.DB.ExecContext(ctx, + `DELETE FROM calendar_attachments WHERE attachment_id = ?`, attachmentID) + if err != nil { + return false, err + } + n, err := res.RowsAffected() + if err != nil { + return false, err + } + return n > 0, nil +} + +// AttachCalendarFilesToMail 把事件的附件复制成邮件附件。 +// +// 提醒邮件是新建的,附件必须重新挂一份指向同一 sha256 的元数据 —— +// 内容寻址下这不拷磁盘文件,只是多一条记录。 +// +// 缺了这一步的后果:人在事件上传了附件、UI 里看得见、提醒也按时发出, +// 但 Agent 收到的那封信里附件清单是空的 —— 事件附件与邮件附件是两张表, +// 不复制就永远只存在于日历侧。这是「日历附件只记元数据未接投递」的另一半。 +// +// uploader 记为 calendarSender("calendar"):附件随提醒邮件重新分发, +// 其可见范围由该邮件的参与方决定,而不是沿用事件创建者。 +func AttachCalendarFilesToMail(ctx context.Context, eventID string, mailID uuid.UUID, uploader string) (int, error) { + atts, err := ListCalendarAttachments(ctx, eventID) + if err != nil { + return 0, err + } + n := 0 + for _, a := range atts { + // sha256 为空说明这条记录没有真实内容(历史脏数据),跳过而不是 + // 挂一个下载必然 404 的附件 + if a.SHA256 == "" { + continue + } + if _, err := db.DB.ExecContext(ctx, ` + INSERT INTO attachments (mail_id, uploader, filename, content_type, size_bytes, sha256) + VALUES ($1, $2, $3, $4, $5, $6) + `, mailID, uploader, a.Filename, "application/octet-stream", a.SizeBytes, a.SHA256); err != nil { + return n, err + } + n++ + } + return n, nil +} + +// ListCalendarEventsCreatedBy 只返回某个创建者建的事件。 +// +// Agent 侧列表用这个而不是 ListCalendarEvents:Agent 不该看到别人(人类或 +// 其他 Agent)的日程 —— 那里可能有它无权知道的会议、地址、附件名。 +// +// 注意**不是**「发给我的事件」:`recipients` 里有我但我没建的,同样不返回。 +// 理由是那些事件的编辑权不属于我,列出来只会让模型试图改它然后拿到 403。 +// 想知道「谁给我设了提醒」,那条信息在提醒邮件本身里。 +func ListCalendarEventsCreatedBy(ctx context.Context, creator string, from, to time.Time, status string) ([]models.CalendarEvent, error) { + rows, err := db.DB.QueryContext(ctx, ` + SELECT `+calendarCols+` + FROM calendar_events + WHERE created_by = ? + AND event_time >= ? AND event_time <= ? + AND (status = ? OR ? = '') + ORDER BY event_time ASC`, creator, from, to, status, status) + if err != nil { + return nil, err + } + defer rows.Close() + + events := []models.CalendarEvent{} + for rows.Next() { + e, err := scanCalendarEvent(rows) + if err != nil { + return nil, err + } + events = append(events, *e) + } + return events, rows.Err() +} + +// CountActiveEventsBy 数某个创建者当前有多少条生效中的事件。 +// +// 给 Agent 侧的总量上限用。速率限制只压住「短时间内暴建」, +// 压不住「每小时建 19 条、连建一周」—— 而日历事件是长效的, +// 攒下来的每一条都会持续产生提醒邮件。 +func CountActiveEventsBy(ctx context.Context, creator string) (int, error) { + var n int + err := db.DB.QueryRowContext(ctx, + `SELECT COUNT(*) FROM calendar_events WHERE created_by = ? AND status = 'active'`, + creator).Scan(&n) + return n, err +} diff --git a/gateway/internal/repo/calendar_test.go b/gateway/internal/repo/calendar_test.go new file mode 100644 index 0000000..8755e61 --- /dev/null +++ b/gateway/internal/repo/calendar_test.go @@ -0,0 +1,972 @@ +package repo + +import ( + "context" + "testing" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/agentmail/gateway/internal/lunar" + "github.com/agentmail/gateway/internal/models" + "github.com/google/uuid" +) + +func seedEvent(t *testing.T, e *models.CalendarEvent) *models.CalendarEvent { + t.Helper() + if e.Title == "" { + e.Title = "测试事件" + } + if e.EventTime.IsZero() { + e.EventTime = time.Now().Add(time.Hour) + } + out, err := CreateCalendarEvent(context.Background(), e) + if err != nil { + t.Fatalf("建事件: %v", err) + } + return out +} + +func TestCalendarEventCRUD(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + at := time.Now().Add(2 * time.Hour).Truncate(time.Second) + e := seedEvent(t, &models.CalendarEvent{ + Title: "每日站会", + Description: "同步进展", + ReminderText: "日程提醒:{title}", + AgentName: "dsh", + ToAddress: "dsh@/home", + EventTime: at, + RemindBefore: 15, + Recurrence: "daily", + CreatedBy: "jianf", + }) + + if e.EventID == "" { + t.Fatal("建完事件必须有 event_id") + } + if e.Status != "active" { + t.Errorf("新事件默认应为 active,得到 %q", e.Status) + } + + got, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读事件: %v", err) + } + if got.Title != "每日站会" || got.RemindBefore != 15 || got.Recurrence != "daily" { + t.Errorf("读回的字段不符:%+v", got) + } + if !got.EventTime.Equal(at) { + t.Errorf("event_time 读回错位:写 %v 读 %v", at, got.EventTime) + } + + got.Title = "改名后的站会" + got.Status = "paused" + if err := UpdateCalendarEvent(ctx, e.EventID, got); err != nil { + t.Fatalf("改事件: %v", err) + } + again, _ := GetCalendarEvent(ctx, e.EventID) + if again.Title != "改名后的站会" || again.Status != "paused" { + t.Errorf("改后没生效:%+v", again) + } + + if err := DeleteCalendarEvent(ctx, e.EventID); err != nil { + t.Fatalf("删事件: %v", err) + } + if _, err := GetCalendarEvent(ctx, e.EventID); err != ErrEventNotFound { + t.Errorf("删掉后应报 ErrEventNotFound,得到 %v", err) + } +} + +func TestCalendarNotFoundIsTyped(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + // 不存在的 id 要给出可判定的错误,而不是 sql.ErrNoRows —— + // handler 靠它区分 404 与 500。 + if _, err := GetCalendarEvent(ctx, "00000000-0000-0000-0000-000000000000"); err != ErrEventNotFound { + t.Errorf("Get 应报 ErrEventNotFound,得到 %v", err) + } + if err := DeleteCalendarEvent(ctx, "00000000-0000-0000-0000-000000000000"); err != ErrEventNotFound { + t.Errorf("Delete 应报 ErrEventNotFound,得到 %v", err) + } + if err := UpdateCalendarEvent(ctx, "00000000-0000-0000-0000-000000000000", + &models.CalendarEvent{Title: "x", EventTime: time.Now()}); err != ErrEventNotFound { + t.Errorf("Update 应报 ErrEventNotFound,得到 %v", err) + } +} + +func TestDueEventsOnlyReturnsRipe(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + now := time.Now() + + // 已经该响的(事件时间在过去) + ripe := seedEvent(t, &models.CalendarEvent{Title: "该响了", EventTime: now.Add(-time.Minute)}) + // 提前 30 分钟提醒、事件在 20 分钟后 —— 提醒点已过 + early := seedEvent(t, &models.CalendarEvent{ + Title: "提前提醒已到", EventTime: now.Add(20 * time.Minute), RemindBefore: 30, + }) + // 还早(1 小时后,无提前提醒) + future := seedEvent(t, &models.CalendarEvent{Title: "还早", EventTime: now.Add(time.Hour)}) + // 已暂停的不该响 + paused := seedEvent(t, &models.CalendarEvent{Title: "暂停的", EventTime: now.Add(-time.Minute)}) + p, _ := GetCalendarEvent(ctx, paused.EventID) + p.Status = "paused" + if err := UpdateCalendarEvent(ctx, paused.EventID, p); err != nil { + t.Fatalf("暂停: %v", err) + } + + due, err := DueEvents(ctx) + if err != nil { + t.Fatalf("DueEvents: %v", err) + } + + got := map[string]bool{} + for _, e := range due { + got[e.EventID] = true + } + if !got[ripe.EventID] { + t.Error("到期事件没被取出") + } + if !got[early.EventID] { + t.Error("remind_before 已过的事件没被取出") + } + if got[future.EventID] { + t.Error("未到期事件被取出了") + } + if got[paused.EventID] { + t.Error("已暂停的事件被取出了") + } +} + +func TestMarkEventFiredStopsRefiring(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + // 幂等的关键:标记后同一条不该再出现在 DueEvents 里, + // 否则调度器每 30 秒把同一封提醒重发一遍。 + e := seedEvent(t, &models.CalendarEvent{Title: "只该响一次", EventTime: time.Now().Add(-time.Minute)}) + + due, _ := DueEvents(ctx) + if len(due) != 1 { + t.Fatalf("标记前应有 1 条到期,得到 %d", len(due)) + } + + if err := MarkEventFired(ctx, e.EventID); err != nil { + t.Fatalf("标记: %v", err) + } + + due, _ = DueEvents(ctx) + for _, d := range due { + if d.EventID == e.EventID { + t.Error("已标记触发的事件仍出现在 DueEvents 里") + } + } +} + +func TestAdvanceRecurrence(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + base := time.Now().Add(-time.Minute).Truncate(time.Second) + + t.Run("一次性事件不推进", func(t *testing.T) { + e := seedEvent(t, &models.CalendarEvent{Title: "一次性", EventTime: base, Recurrence: "none"}) + advanced, err := AdvanceRecurrence(ctx, e.EventID) + if err != nil { + t.Fatalf("推进: %v", err) + } + if advanced { + t.Error("recurrence=none 不该推进") + } + }) + + for _, tc := range []struct { + rule string + want time.Time + }{ + {"daily", base.AddDate(0, 0, 1)}, + {"weekly", base.AddDate(0, 0, 7)}, + {"monthly", base.AddDate(0, 1, 0)}, + } { + t.Run(tc.rule+" 推进一个周期", func(t *testing.T) { + e := seedEvent(t, &models.CalendarEvent{ + Title: tc.rule, EventTime: base, Recurrence: tc.rule, + }) + advanced, err := AdvanceRecurrence(ctx, e.EventID) + if err != nil { + t.Fatalf("推进: %v", err) + } + if !advanced { + t.Fatal("应该推进") + } + got, _ := GetCalendarEvent(ctx, e.EventID) + if !got.EventTime.Equal(tc.want) { + t.Errorf("下次时间应为 %v,得到 %v", tc.want, got.EventTime) + } + // 推进后 event_time 已在未来,且 last_fired_at 仍为旧值 → + // 必须重新出现在 DueEvents 里等待下一轮(否则重复事件只响一次)。 + if got.Status != "active" { + t.Errorf("推进后应仍为 active,得到 %q", got.Status) + } + }) + } + + t.Run("超过 recurrence_end 则取消", func(t *testing.T) { + end := base.Add(12 * time.Hour) // 下一次(+1 天)会越过它 + e := seedEvent(t, &models.CalendarEvent{ + Title: "快结束了", EventTime: base, Recurrence: "daily", RecurrenceEnd: &end, + }) + advanced, err := AdvanceRecurrence(ctx, e.EventID) + if err != nil { + t.Fatalf("推进: %v", err) + } + if advanced { + t.Error("越过 recurrence_end 时不该报告推进成功") + } + got, _ := GetCalendarEvent(ctx, e.EventID) + if got.Status != "cancelled" { + t.Errorf("越过结束时间应置为 cancelled,得到 %q", got.Status) + } + }) +} + +func TestListCalendarEventsRange(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + now := time.Now() + inRange := seedEvent(t, &models.CalendarEvent{Title: "范围内", EventTime: now.Add(time.Hour)}) + seedEvent(t, &models.CalendarEvent{Title: "太远", EventTime: now.AddDate(0, 3, 0)}) + + events, err := ListCalendarEvents(ctx, now, now.Add(24*time.Hour), "active") + if err != nil { + t.Fatalf("列事件: %v", err) + } + if len(events) != 1 || events[0].EventID != inRange.EventID { + t.Errorf("时间范围过滤不对,得到 %d 条", len(events)) + } +} + +func TestCalendarAttachments(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{Title: "带附件"}) + + if err := AddCalendarAttachment(ctx, &models.CalendarAttachment{ + EventID: e.EventID, Filename: "报表.xlsx", SHA256: "abc", SizeBytes: 2048, + }); err != nil { + t.Fatalf("加附件: %v", err) + } + + atts, err := ListCalendarAttachments(ctx, e.EventID) + if err != nil { + t.Fatalf("列附件: %v", err) + } + if len(atts) != 1 || atts[0].Filename != "报表.xlsx" { + t.Fatalf("附件读回不符:%+v", atts) + } + if atts[0].AttachmentID == "" { + t.Error("附件必须有 attachment_id —— 没有它模型无法在 send_mail 里引用") + } + + // 事件没有附件时返回空而不是报错 + other := seedEvent(t, &models.CalendarEvent{Title: "没附件"}) + if atts, err := ListCalendarAttachments(ctx, other.EventID); err != nil || len(atts) != 0 { + t.Errorf("无附件事件应返回空列表,得到 %d 条 err=%v", len(atts), err) + } +} + +func TestDeleteCalendarAttachment(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{Title: "要删附件"}) + for _, name := range []string{"甲.pdf", "乙.pdf"} { + if err := AddCalendarAttachment(ctx, &models.CalendarAttachment{ + EventID: e.EventID, Filename: name, SHA256: "sum-" + name, SizeBytes: 10, + }); err != nil { + t.Fatalf("加附件 %s: %v", name, err) + } + } + + atts, _ := ListCalendarAttachments(ctx, e.EventID) + if len(atts) != 2 { + t.Fatalf("准备阶段应有 2 个附件,得到 %d", len(atts)) + } + + ok, err := DeleteCalendarAttachment(ctx, atts[0].AttachmentID) + if err != nil { + t.Fatalf("删附件: %v", err) + } + if !ok { + t.Error("删掉存在的附件应返回 true") + } + + left, _ := ListCalendarAttachments(ctx, e.EventID) + if len(left) != 1 { + t.Fatalf("删一个后应剩 1 个,得到 %d", len(left)) + } + if left[0].AttachmentID == atts[0].AttachmentID { + t.Error("删错了对象") + } + + // 不存在的 id 返回 false 而不是报错 —— 调用方据此回 404 而非 500 + ok, err = DeleteCalendarAttachment(ctx, "00000000-0000-0000-0000-000000000000") + if err != nil { + t.Errorf("删不存在的附件不该报错,得到 %v", err) + } + if ok { + t.Error("删不存在的附件应返回 false") + } +} + +// 事件附件必须能复制成邮件附件。 +// +// 少了这一步,附件只存在于日历侧:UI 里看得见、提醒按时发出、 +// 而 Agent 收到的那封信附件清单是空的 —— 两张表互不相通。 +func TestAttachCalendarFilesToMail(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{Title: "带附件的提醒"}) + if err := AddCalendarAttachment(ctx, &models.CalendarAttachment{ + EventID: e.EventID, Filename: "周报.md", SHA256: "deadbeef", SizeBytes: 512, + }); err != nil { + t.Fatalf("加附件: %v", err) + } + // sha256 为空的脏数据必须被跳过:挂上去只会得到一个下载必然 404 的附件 + if err := AddCalendarAttachment(ctx, &models.CalendarAttachment{ + EventID: e.EventID, Filename: "没内容.bin", SHA256: "", SizeBytes: 0, + }); err != nil { + t.Fatalf("加空附件: %v", err) + } + + seedAgentForAttach(t, "pi") + sessionID := seedSessionForAttach(t, "pi") + mailID, err := CreateMail(ctx, sessionID, nil, "calendar", "", "pi", "", "日程提醒:带附件的提醒", "正文", nil) + if err != nil { + t.Fatalf("建邮件: %v", err) + } + + n, err := AttachCalendarFilesToMail(ctx, e.EventID, mailID, "calendar") + if err != nil { + t.Fatalf("挂附件: %v", err) + } + if n != 1 { + t.Fatalf("应只挂 1 个(空 sha256 那条跳过),得到 %d", n) + } + + mailAtts, err := ListAttachmentsFor(ctx, mailID) + if err != nil { + t.Fatalf("列邮件附件: %v", err) + } + if len(mailAtts) != 1 { + t.Fatalf("邮件上应有 1 个附件,得到 %d", len(mailAtts)) + } + if mailAtts[0].Filename != "周报.md" || mailAtts[0].SHA256 != "deadbeef" { + t.Errorf("附件内容不符:%+v", mailAtts[0]) + } + // 内容寻址:复制不产生新的 sha256,指向同一份磁盘文件 + if mailAtts[0].Uploader != "calendar" { + t.Errorf("uploader 应是 calendar,得到 %q", mailAtts[0].Uploader) + } + + // 没有附件的事件挂 0 个且不报错 + empty := seedEvent(t, &models.CalendarEvent{Title: "无附件"}) + if n, err := AttachCalendarFilesToMail(ctx, empty.EventID, mailID, "calendar"); err != nil || n != 0 { + t.Errorf("无附件事件应挂 0 个,得到 %d err=%v", n, err) + } +} + +func seedAgentForAttach(t *testing.T, name string) { + t.Helper() + if _, err := db.DB.ExecContext(context.Background(), + `INSERT INTO agents (agent_name, secret, platform) VALUES ($1, 'x', 'test')`, + name); err != nil { + t.Fatalf("seed agent: %v", err) + } +} + +func seedSessionForAttach(t *testing.T, agentName string) uuid.UUID { + t.Helper() + id := uuid.New() + // 列名是 from_agent 而不是 agent_name(后者是 agents 表的主键名) + if _, err := db.DB.ExecContext(context.Background(), + `INSERT INTO sessions (session_id, from_agent, subject, session_alias, status) + VALUES ($1, $2, '日程提醒', 'cal-test', 'active')`, + id, agentName); err != nil { + t.Fatalf("seed session: %v", err) + } + return id +} + +// 落在 60 秒 lookahead 窗口内的**未来**事件,标记后不得再次到期。 +// +// 生产实测的重发现场:一条 event_time=12:53:17 的事件在 +// 12:52:30 / 12:53:00 / 12:53:06 / 12:53:36 各发了一封相同提醒。 +// 根因是去重判据写成 `last_fired_at < event_time` —— 触发时刻(now) +// 本来就早于 event_time,条件恒真,于是每个 tick 重发一次, +// 直到 event_time 真正过去才自己停下。 +// +// 改成按 occurrence 相等(fired_for = 当时的 event_time)才精确: +// AdvanceRecurrence 改了 event_time 就该再触发,没改就永不重发。 +func TestFiredEventInLookaheadWindowDoesNotRefire(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + // 47 秒后 —— 在 lookahead 窗口内,所以第一次扫描就会入选 + e := seedEvent(t, &models.CalendarEvent{ + Title: "窗口内的未来事件", EventTime: time.Now().Add(47 * time.Second), + }) + + due, err := DueEvents(ctx) + if err != nil { + t.Fatalf("首次扫描: %v", err) + } + if len(due) != 1 { + t.Fatalf("lookahead 应让它提前入选,得到 %d 条", len(due)) + } + + if err := MarkEventFired(ctx, e.EventID); err != nil { + t.Fatalf("标记: %v", err) + } + + // 模拟后续几个 tick + for i := 0; i < 3; i++ { + due, err = DueEvents(ctx) + if err != nil { + t.Fatalf("第 %d 次重扫: %v", i+2, err) + } + for _, d := range due { + if d.EventID == e.EventID { + t.Fatalf("第 %d 次扫描仍判定到期 —— 提醒会被重发", i+2) + } + } + } +} + +// 重复事件推进 event_time 之后必须重新到期: +// 按 occurrence 去重的另一半,漏了它就变成「每个重复事件只响一次」。 +func TestRecurringEventRefiresAfterAdvance(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{ + Title: "每天都要响", + EventTime: time.Now().Add(-time.Minute), + Recurrence: "daily", + }) + + if err := MarkEventFired(ctx, e.EventID); err != nil { + t.Fatalf("标记: %v", err) + } + due, _ := DueEvents(ctx) + for _, d := range due { + if d.EventID == e.EventID { + t.Fatal("标记后不该立刻再次到期") + } + } + + // 推进到下一次(+1 天)后,把时间挪到过去模拟「第二天到了」 + if _, err := AdvanceRecurrence(ctx, e.EventID); err != nil { + t.Fatalf("推进重复: %v", err) + } + if _, err := db.DB.ExecContext(ctx, + `UPDATE calendar_events SET event_time = ? WHERE event_id = ?`, + time.Now().Add(-30*time.Second), e.EventID); err != nil { + t.Fatalf("模拟次日: %v", err) + } + + due, _ = DueEvents(ctx) + found := false + for _, d := range due { + if d.EventID == e.EventID { + found = true + } + } + if !found { + t.Error("event_time 推进后应重新到期,否则重复事件只响一次") + } +} + +// ─── 重复规则推进(NextOccurrence 是纯函数,不碰数据库)─── + +func TestNextOccurrenceSolar(t *testing.T) { + base := time.Date(2026, 9, 3, 9, 30, 0, 0, time.Local) + + cases := []struct { + rule string + want string + }{ + {models.RecurDaily, "2026-09-04"}, + {models.RecurWeekly, "2026-09-10"}, + {models.RecurMonthly, "2026-10-03"}, + } + for _, c := range cases { + got, err := NextOccurrence(c.rule, base) + if err != nil { + t.Errorf("%s: %v", c.rule, err) + continue + } + if got.Format("2006-01-02") != c.want { + t.Errorf("%s: 得到 %s,期望 %s", c.rule, got.Format("2006-01-02"), c.want) + } + // 时钟必须原样保留 + if got.Hour() != 9 || got.Minute() != 30 { + t.Errorf("%s: 时钟被改动 %v", c.rule, got) + } + } + + // none 与未知值都返回零值 + nil error + for _, r := range []string{models.RecurNone, "", "每隔一个蓝月亮"} { + got, err := NextOccurrence(r, base) + if err != nil || !got.IsZero() { + t.Errorf("%q 应返回零值无错,得到 %v err=%v", r, got, err) + } + } +} + +// time.AddDate 的溢出对「每月同一日」是错的:3 月 31 日 +1 月 = 5 月 1 日。 +// 一次溢出会永久改变规则 —— 31 日的事件在 2 月变成 3 月 3 日, +// 然后从此每月 3 日提醒。 +func TestNextOccurrenceMonthlyClampsMonthEnd(t *testing.T) { + cases := []struct { + from string + want string + why string + }{ + {"2026-01-31", "2026-02-28", "1月31日 +1月 → 2月末(2026 非闰年)"}, + {"2026-03-31", "2026-04-30", "3月31日 +1月 → 4月30日"}, + {"2026-05-31", "2026-06-30", "5月31日 +1月 → 6月30日"}, + {"2028-01-31", "2028-02-29", "闰年 2 月有 29 天"}, + {"2026-01-15", "2026-02-15", "月中日期不受影响"}, + } + for _, c := range cases { + from, _ := time.ParseInLocation("2006-01-02", c.from, time.Local) + got, err := NextOccurrence(models.RecurMonthly, from) + if err != nil { + t.Errorf("%s: %v", c.why, err) + continue + } + if got.Format("2006-01-02") != c.want { + t.Errorf("%s: 得到 %s,期望 %s", c.why, got.Format("2006-01-02"), c.want) + } + } +} + +// 农历月推进:公历间隔在 29~30 天之间浮动,不是固定值。 +// 这正是不能用 AddDate 的原因。 +func TestNextOccurrenceLunarMonthly(t *testing.T) { + // 2026-09-03 = 农历七月廿二 + cur := time.Date(2026, 9, 3, 9, 0, 0, 0, time.Local) + gaps := map[int]bool{} + for i := 0; i < 6; i++ { + next, err := NextOccurrence(models.RecurLunarMonthly, cur) + if err != nil { + t.Fatalf("第 %d 次推进: %v", i+1, err) + } + if !next.After(cur) { + t.Fatalf("第 %d 次推进没有前进:%v → %v", i+1, cur, next) + } + gap := int(next.Sub(cur).Hours() / 24) + gaps[gap] = true + // 农历同一日:连续推进后农历「日」应保持 + if d := lunar.FromSolar(next); d.Day != 22 { + t.Errorf("第 %d 次推进后农历日变成 %d(期望 22):%s", i+1, d.Day, d.String()) + } + cur = next + } + // 间隔必须出现过多种值,证明不是固定天数 + if len(gaps) < 2 { + t.Errorf("六次农历月推进的公历间隔只有 %v —— 疑似退化成固定天数", gaps) + } + for g := range gaps { + if g < 28 || g > 31 { + t.Errorf("农历月间隔 %d 天不合理", g) + } + } +} + +// 农历年推进:公历日期每年漂移。用公历 yearly 会固定在同一天, +// 与「过农历生日/祭日」的期望不符 —— 这是农历规则存在的理由。 +func TestNextOccurrenceLunarYearly(t *testing.T) { + cur := time.Date(2026, 9, 3, 9, 0, 0, 0, time.Local) + seen := map[string]bool{} + for i := 0; i < 5; i++ { + next, err := NextOccurrence(models.RecurLunarYearly, cur) + if err != nil { + t.Fatalf("第 %d 次: %v", i+1, err) + } + if !next.After(cur) { + t.Fatalf("第 %d 次没有前进:%v → %v", i+1, cur, next) + } + // 农历月日应保持 + d := lunar.FromSolar(next) + if d.Month != 7 || d.Day != 22 { + t.Errorf("第 %d 次推进后农历变成 %d-%d(期望 7-22)", i+1, d.Month, d.Day) + } + seen[next.Format("01-02")] = true + cur = next + } + if len(seen) < 3 { + t.Errorf("五年公历月日只有 %d 种 —— 农历年重复应漂移", len(seen)) + } +} + +// 农历规则经过数据库这一轮也要正确(AdvanceRecurrence 里调 NextOccurrence)。 +func TestAdvanceRecurrenceLunar(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + start := time.Date(2026, 9, 3, 9, 0, 0, 0, time.Local) + e := seedEvent(t, &models.CalendarEvent{ + Title: "农历每月十五(这里用廿二)", + EventTime: start, + Recurrence: models.RecurLunarMonthly, + }) + + advanced, err := AdvanceRecurrence(ctx, e.EventID) + if err != nil { + t.Fatalf("推进: %v", err) + } + if !advanced { + t.Fatal("农历重复应能推进") + } + + after, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读回: %v", err) + } + if !after.EventTime.After(start) { + t.Errorf("event_time 未前进:%v", after.EventTime) + } + // 农历日保持 + if d := lunar.FromSolar(after.EventTime); d.Day != 22 { + t.Errorf("农历日变成 %d,期望 22(%s)", d.Day, d.String()) + } + // 公历间隔应在一个农历月内 + gap := int(after.EventTime.Sub(start).Hours() / 24) + if gap < 28 || gap > 31 { + t.Errorf("间隔 %d 天不像一个农历月", gap) + } +} + +// ─── 多收件人 ─── + +func TestRecipientsRoundtrip(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{ + Title: "三个 Agent 各自汇报", + Recipients: []string{"pi@/home/program/agentmail", "dsh", "opencode@/tmp"}, + DeliveryMode: models.DeliverSeparate, + }) + + got, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读回: %v", err) + } + if len(got.Recipients) != 3 { + t.Fatalf("收件人应有 3 个,得到 %d:%v", len(got.Recipients), got.Recipients) + } + if got.Recipients[0] != "pi@/home/program/agentmail" { + t.Errorf("顺序或内容不符:%v", got.Recipients) + } + if got.EffectiveDeliveryMode() != models.DeliverSeparate { + t.Errorf("投递模式 = %q", got.EffectiveDeliveryMode()) + } +} + +// 空收件人列表必须序列化成 [](而不是 null):Go 的 nil slice 会变 null, +// 前端 .map 直接崩。 +func TestRecipientsNeverNull(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{Title: "没写收件人"}) + got, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读回: %v", err) + } + if got.Recipients == nil { + t.Error("Recipients 为 nil —— 会序列化成 null 让前端崩") + } + if len(got.Recipients) != 0 { + t.Errorf("应是空数组,得到 %v", got.Recipients) + } +} + +// 旧数据(只有 agent_name / to_address)必须继续工作 —— 历史事件不迁移。 +func TestEffectiveRecipientsFallbackChain(t *testing.T) { + cases := []struct { + name string + e models.CalendarEvent + want []string + }{ + { + "Recipients 优先", + models.CalendarEvent{Recipients: []string{"a", "b"}, ToAddress: "c", AgentName: "d"}, + []string{"a", "b"}, + }, + { + "退回 to_address", + models.CalendarEvent{ToAddress: "pi@/tmp.alias", AgentName: "pi"}, + []string{"pi@/tmp.alias"}, + }, + { + "再退回 agent_name", + models.CalendarEvent{AgentName: "dsh"}, + []string{"dsh"}, + }, + { + "全空给 nil", + models.CalendarEvent{}, + nil, + }, + { + "Recipients 里全是空白时继续退回", + models.CalendarEvent{Recipients: []string{"", " "}, AgentName: "pi"}, + []string{"pi"}, + }, + } + for _, c := range cases { + got := c.e.EffectiveRecipients() + if len(got) != len(c.want) { + t.Errorf("%s: 得到 %v,期望 %v", c.name, got, c.want) + continue + } + for i := range got { + if got[i] != c.want[i] { + t.Errorf("%s: 第 %d 项 %q,期望 %q", c.name, i, got[i], c.want[i]) + } + } + } +} + +// 未知投递模式按 separate 处理:它的失败模式更轻。 +// together 用错会让本该独立判断的 Agent 互相看到回复而趋同,事后无法分离。 +func TestEffectiveDeliveryModeDefaultsToSeparate(t *testing.T) { + for _, in := range []string{"", "separate", "垃圾值", "SEPARATE"} { + e := models.CalendarEvent{DeliveryMode: in} + if got := e.EffectiveDeliveryMode(); got != models.DeliverSeparate { + t.Errorf("DeliveryMode=%q → %q,期望 separate", in, got) + } + } + e := models.CalendarEvent{DeliveryMode: models.DeliverTogether} + if e.EffectiveDeliveryMode() != models.DeliverTogether { + t.Error("together 应被保留") + } +} + +// 公历每年:2 月 29 日在平年必须夹到 2 月 28,不能溢出成 3 月 1 日。 +// 闰日生日的约定是「平年过 2 月 28」。 +func TestNextOccurrenceYearlyClampsLeapDay(t *testing.T) { + cases := []struct { + from string + want string + why string + }{ + {"2028-02-29", "2029-02-28", "闰日 +1 年 → 平年 2 月 28"}, + {"2026-03-15", "2027-03-15", "普通日期不受影响"}, + {"2027-02-28", "2028-02-28", "平年 2/28 → 闰年仍是 2/28(不跳到 29)"}, + } + for _, c := range cases { + from, _ := time.ParseInLocation("2006-01-02", c.from, time.Local) + got, err := NextOccurrence(models.RecurYearly, from) + if err != nil { + t.Errorf("%s: %v", c.why, err) + continue + } + if got.Format("2006-01-02") != c.want { + t.Errorf("%s: 得到 %s,期望 %s", c.why, got.Format("2006-01-02"), c.want) + } + } +} + +// 过期的重复事件必须一次推到未来,不能每轮补发一封。 +// +// 实测的 bug:AdvanceRecurrence 只推进一步 —— 一条 100 天前设的每日事件, +// 每轮扫描都判定「已过期该触发」→ 发一封 → event_time 只前进一天 → +// 下一轮又过期。30 轮扫描触发 30 次,而调度周期是 30 秒, +// 人会收到一串垃圾提醒,连发 100 封才追上今天。 +func TestStaleRecurringEventDoesNotFlood(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{ + Title: "很久以前设的每日提醒", + EventTime: time.Now().AddDate(0, 0, -100), + Recurrence: models.RecurDaily, + }) + + fires := 0 + // 模拟调度器连续跑 30 轮(生产上就是 15 分钟) + for i := 0; i < 30; i++ { + due, err := DueEvents(ctx) + if err != nil { + t.Fatalf("第 %d 轮扫描: %v", i+1, err) + } + hit := false + for _, d := range due { + if d.EventID == e.EventID { + hit = true + } + } + if !hit { + break + } + fires++ + if err := MarkEventFired(ctx, e.EventID); err != nil { + t.Fatalf("标记: %v", err) + } + if _, err := AdvanceRecurrence(ctx, e.EventID); err != nil { + t.Fatalf("推进: %v", err) + } + } + + if fires != 1 { + t.Errorf("过期的每日重复事件触发了 %d 次,应只触发 1 次", fires) + } + after, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读回: %v", err) + } + if !after.EventTime.After(time.Now()) { + t.Errorf("推进后 event_time 仍在过去:%v", after.EventTime) + } + // 只跳到「刚过现在」的那一次,不是跳到很远的将来 + if after.EventTime.After(time.Now().AddDate(0, 0, 2)) { + t.Errorf("推得太远了:%v", after.EventTime) + } +} + +// 农历规则的过期事件同样不能刷屏。 +func TestStaleLunarRecurringDoesNotFlood(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + e := seedEvent(t, &models.CalendarEvent{ + Title: "去年设的农历每月提醒", + EventTime: time.Now().AddDate(-1, 0, 0), + Recurrence: models.RecurLunarMonthly, + }) + + if _, err := AdvanceRecurrence(ctx, e.EventID); err != nil { + t.Fatalf("推进: %v", err) + } + after, _ := GetCalendarEvent(ctx, e.EventID) + if !after.EventTime.After(time.Now()) { + t.Errorf("一年前的农历事件推进后仍在过去:%v", after.EventTime) + } + // 农历日必须保持 + if d := lunar.FromSolar(after.EventTime); d.Day != lunar.FromSolar(e.EventTime).Day { + t.Errorf("农历日从 %d 变成 %d", lunar.FromSolar(e.EventTime).Day, d.Day) + } +} + +// 越过 recurrence_end 时必须置 cancelled 而不是留在 active。 +// +// 留着的表现是一条僵尸事件:DueEvents 每轮都捞到它(event_time 在过去), +// 但 fired_for 已等于 event_time 所以又不触发 —— 永远排在到期列表里不动。 +func TestAdvanceCancelsAfterRecurrenceEnd(t *testing.T) { + setupTestDB(t) + ctx := context.Background() + + end := time.Now().AddDate(0, 0, -1) // 昨天就该停 + e := seedEvent(t, &models.CalendarEvent{ + Title: "已到期的每日重复", + EventTime: time.Now().AddDate(0, 0, -5), + Recurrence: models.RecurDaily, + RecurrenceEnd: &end, + }) + + advanced, err := AdvanceRecurrence(ctx, e.EventID) + if err != nil { + t.Fatalf("推进: %v", err) + } + if advanced { + t.Error("已过 recurrence_end 不该报告推进成功") + } + after, err := GetCalendarEvent(ctx, e.EventID) + if err != nil { + t.Fatalf("读回: %v", err) + } + if after.Status != "cancelled" { + t.Errorf("状态应是 cancelled,得到 %q —— 留在 active 会变僵尸事件", after.Status) + } + // 且不该再出现在到期列表里 + due, _ := DueEvents(ctx) + for _, d := range due { + if d.EventID == e.EventID { + t.Error("已 cancelled 的事件仍出现在 DueEvents") + } + } +} + +// advanceToFuture 是纯函数,单独测三个终止条件。 +func TestAdvanceToFuture(t *testing.T) { + now := time.Date(2026, 9, 3, 12, 0, 0, 0, time.Local) + + t.Run("跨过 now 就停", func(t *testing.T) { + from := now.AddDate(0, 0, -100) + got, err := advanceToFuture(models.RecurDaily, from, now, nil) + if err != nil { + t.Fatalf("推进: %v", err) + } + if !got.After(now) { + t.Errorf("结果 %v 不在 now 之后", got) + } + // 恰好是越过 now 的第一次,不是更远 + if got.After(now.AddDate(0, 0, 1)) { + t.Errorf("推过头了:%v", got) + } + }) + + t.Run("越过 recurrenceEnd 给零值", func(t *testing.T) { + end := now.AddDate(0, 0, -1) + got, err := advanceToFuture(models.RecurDaily, now.AddDate(0, 0, -5), now, &end) + if err != nil { + t.Fatalf("不该报错:%v", err) + } + if !got.IsZero() { + t.Errorf("应给零值,得到 %v", got) + } + }) + + t.Run("不重复给零值无错", func(t *testing.T) { + got, err := advanceToFuture(models.RecurNone, now, now, nil) + if err != nil || !got.IsZero() { + t.Errorf("得到 %v err=%v", got, err) + } + }) + + t.Run("未来的事件原地推一步", func(t *testing.T) { + from := now.AddDate(0, 0, 5) + got, err := advanceToFuture(models.RecurDaily, from, now, nil) + if err != nil { + t.Fatalf("推进: %v", err) + } + // from 已在未来,第一次推进就该返回 + if !got.Equal(from.AddDate(0, 0, 1)) { + t.Errorf("得到 %v,期望 %v", got, from.AddDate(0, 0, 1)) + } + }) + + // 上限是防御性的:农历路径依赖外部库,一旦某年给出反直觉结果, + // 没有上限就是个死循环 goroutine,而它跑在调度器里 —— 整个提醒系统一起卡住 + t.Run("十年前的每日事件也能在上限内追上", func(t *testing.T) { + got, err := advanceToFuture(models.RecurDaily, now.AddDate(-10, 0, 0), now, nil) + if err != nil { + t.Fatalf("十年(约 3650 步)应在 %d 上限内:%v", maxAdvanceSteps, err) + } + if !got.After(now) { + t.Errorf("结果 %v 不在 now 之后", got) + } + }) +} diff --git a/gateway/internal/scheduler/calendar.go b/gateway/internal/scheduler/calendar.go new file mode 100644 index 0000000..9d760aa --- /dev/null +++ b/gateway/internal/scheduler/calendar.go @@ -0,0 +1,384 @@ +package scheduler + +import ( + "context" + "errors" + "fmt" + "log" + "strings" + "sync" + "time" + + "github.com/agentmail/gateway/internal/db" + "github.com/agentmail/gateway/internal/models" + "github.com/agentmail/gateway/internal/repo" + "github.com/agentmail/gateway/internal/sse" + "github.com/google/uuid" +) + +// CalendarScheduler 定时扫描日历事件,触发到期提醒。 +// +// 设计参照 Outlook 的 Exchange 提醒器: +// - 每 30 秒扫描一次(精度到分钟够用,不需要秒级) +// - 到期事件 → 生成一封提醒邮件 → 投递 +// - 重复事件自动推进到下一次 +// - 幂等:last_fired_at 保证同一分钟不触发两次 +// +// 为什么用进程内 goroutine 而不是 cron:调度器与 Gateway 同生命周期, +// 不需要外部依赖,也不需要第二套「谁在跑」的信任模型。 +var ( + schedMu sync.Mutex + schedStop chan struct{} + schedWg sync.WaitGroup +) + +// Start 启动日历调度器。重复调用安全(先停旧的)。 +func Start() { + Stop() + + schedMu.Lock() + defer schedMu.Unlock() + + schedStop = make(chan struct{}) + stop := schedStop + schedWg.Add(1) + + go func() { + defer schedWg.Done() + runLoop(stop) + }() + + log.Printf("[scheduler] 日历调度器已启动(30s 周期)") +} + +// Stop 停止调度器,等待当前扫描完成。 +func Stop() { + schedMu.Lock() + defer schedMu.Unlock() + + if schedStop != nil { + close(schedStop) + schedWg.Wait() + schedStop = nil + log.Printf("[scheduler] 日历调度器已停止") + } +} + +func runLoop(stop chan struct{}) { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + // 启动时先扫一次:进程重启期间到期的事件不该被跳过 + scanAndFire() + + for { + select { + case <-stop: + return + case <-ticker.C: + scanAndFire() + } + } +} + +// scanAndFire 扫一轮到期事件。 +// +// **全程不得 panic 出去**:它跑在后台 goroutine 里,而 Go 的 goroutine panic +// 会直接结束整个进程 —— 一条写坏的提醒把邮件网关整个带倒是荒谬的代价。 +func scanAndFire() { + defer func() { + if rec := recover(); rec != nil { + log.Printf("[scheduler] 扫描崩溃(已捕获,下一轮重试): %v", rec) + } + }() + + // DB 未就绪就什么都不做。 + // + // 调度器的启动时机比它看起来脆弱:main 里它排在 migrate 之后, + // 但 Start() 是个导出函数,测试与将来的调用方都可能在建库前调到它。 + // 而 nil *sql.DB 上调 QueryContext 是空指针解引用,不是一个 error。 + if db.DB == nil { + return + } + + ctx := context.Background() + events, err := repo.DueEvents(ctx) + if err != nil { + log.Printf("[scheduler] 查询到期事件失败: %v", err) + return + } + for _, e := range events { + fireEvent(ctx, e) + } +} + +// RenderReminder 把提醒模板里的变量替换成事件的实际内容。 +// +// 单独提出来是为了可测:模板渲染错了不会立刻报错, +// 只会让 Agent 收到一封写着 `{title}` 的邮件。 +// +// **{time} 必须转本地时区再格式化。** +// DSN 带 `_timezone=UTC`,从库里读回的 EventTime 是 UTC;直接 Format +// 会把人在 +0800 输入的 14:30 写成 06:30,而前端预览用的是本地时间 +// —— 于是预览显示 14:30、Agent 收到 06:30,两边差 8 小时且两边都不报错。 +// 日历是给人看的,人说「下午两点半」指的就是自己时区的那个时刻。 +func RenderReminder(tmpl string, e models.CalendarEvent) string { + body := tmpl + body = strings.ReplaceAll(body, "{title}", e.Title) + body = strings.ReplaceAll(body, "{time}", e.EventTime.Local().Format("2006-01-02 15:04")) + body = strings.ReplaceAll(body, "{description}", e.Description) + return body +} + +// fireEvent 触发一条事件:渲染提醒文本、投递邮件、推进重复。 +// +// 顺序很重要:**先标记已触发再发信**。 +// 反过来的话,发信成功但标记失败会让下一轮再发一遍 —— +// 提醒邮件重复比漏发更糟(Agent 会把同一件事做两次)。 +func fireEvent(ctx context.Context, e models.CalendarEvent) { + // 单条事件崩溃不得让同一轮里剩下的提醒全部落空。 + defer func() { + if rec := recover(); rec != nil { + log.Printf("[scheduler] 事件 %s 触发崩溃(已捕获): %v", short(e.EventID), rec) + } + }() + + // 收件人:Recipients 优先,为空时退回 to_address / agent_name(旧数据) + recipients := e.EffectiveRecipients() + if len(recipients) == 0 { + log.Printf("[scheduler] 事件 %s 没有收件人,跳过(title=%q)", short(e.EventID), e.Title) + // 没有收件人的事件永远发不出去,标记已触发免得每 30 秒重试一次 + _ = repo.MarkEventFired(ctx, e.EventID) + return + } + + if err := repo.MarkEventFired(ctx, e.EventID); err != nil { + log.Printf("[scheduler] 标记事件 %s 已触发失败,本轮不发信: %v", short(e.EventID), err) + return + } + + body := RenderReminder(e.ReminderText, e) + subject := "日程提醒:" + e.Title + + switch e.EffectiveDeliveryMode() { + case models.DeliverTogether: + // 一起发:首个是主收件人,其余进 cc_list —— 所有人共享同一条线索, + // 能看到彼此的回复。适合「pi 主办、dsh 知情」这种有主次的协作。 + primary, cc := recipients[0], recipients[1:] + if err := SendCalendarMail(ctx, e.EventID, primary, subject, body, e.CreatedBy, cc...); err != nil { + log.Printf("[scheduler] 投递提醒失败(%s → %s +%d抄送): %v", + e.Title, primary, len(cc), err) + } else { + log.Printf("[scheduler] 已触发日历提醒: %s → %s(抄送 %d 人,同一线索)", + e.Title, primary, len(cc)) + } + + default: + // 各发一封:每人落在自己的会话里,互相看不到。 + // 适合「让三个 Agent 各自独立汇报」—— 用 together 会让他们互相 + // 看到回复而趋同,那种上下文污染事后无法分离。 + // + // 一个失败不影响其余:三个 Agent 里有一个离线时, + // 另外两个仍该收到提醒。 + ok, failed := 0, 0 + for _, addr := range recipients { + if err := SendCalendarMail(ctx, e.EventID, addr, subject, body, e.CreatedBy); err != nil { + failed++ + log.Printf("[scheduler] 投递提醒失败(%s → %s): %v", e.Title, addr, err) + continue + } + ok++ + } + if failed == 0 { + log.Printf("[scheduler] 已触发日历提醒: %s → %d 个收件人(各自独立会话)", + e.Title, ok) + } else { + log.Printf("[scheduler] 日历提醒 %s:%d 成功 / %d 失败", e.Title, ok, failed) + } + } + + // 重复事件推进到下一次;一次性事件到此结束 + if e.Recurrence != models.RecurNone { + advanced, err := repo.AdvanceRecurrence(ctx, e.EventID) + if err != nil { + log.Printf("[scheduler] 推进重复事件 %s 失败: %v", short(e.EventID), err) + } else if !advanced { + log.Printf("[scheduler] 重复事件 %s 已过 recurrence_end 或无法推进,已置为 cancelled", + short(e.EventID)) + } + } +} + +func short(id string) string { + if len(id) > 8 { + return id[:8] + } + return id +} + +// SendCalendarMail 把一条提醒投成邮件。 +// +// 直接走 repo + sse 而不是自己发一个 HTTP 请求到 /mail/send: +// 调度器是进程内 goroutine,绕出去再进来只是多一次鉴权与序列化。 +// +// **日历提醒不扣会话预算**:预算的语义是「这件事值得模型自主发多少封信」, +// 而提醒是人预先设定的定时任务,不是模型的自主行为。让它扣预算会出现 +// 「每天 9 点的日报提醒把当天的预算吃掉一格」这种反直觉结果。 +// +// 收件地址支持完整三维寻址: +// - `agent` → 该 Agent 的默认会话(长期提醒应该用这个, +// 所有提醒落在同一条线索上,模型看得到历史) +// - `agent@/path` → 指定工作目录的默认会话 +// - `agent@/path.alias` → 指定已存在的会话(不存在则报错,不静默新建) +// - `agent@/path.new` → 每次提醒开一条新会话(适合互不相关的一次性任务) +func SendCalendarMail(ctx context.Context, eventID, toAddr, subject, body, createdBy string, ccAddrs ...string) error { + addr, err := models.ParseAddress(toAddr) + if err != nil { + return fmt.Errorf("收件地址 %q 无法解析: %w", toAddr, err) + } + + // 抄送方逐个解析。单个解析失败只跳过它,不让整封信发不出去 —— + // 主收件人能收到提醒比「抄送名单必须完整」重要。 + ccList := make([]models.Address, 0, len(ccAddrs)) + for _, raw := range ccAddrs { + raw = strings.TrimSpace(raw) + if raw == "" { + continue + } + ca, cErr := models.ParseAddress(raw) + if cErr != nil { + log.Printf("[scheduler] 抄送地址 %q 无法解析,已跳过: %v", raw, cErr) + continue + } + // 抄送方的 session 位不参与会话定位(那是主收件人的事), + // 但必须保留在地址里:cc_list 的 name/path 决定谁能看到这条线索。 + ccList = append(ccList, ca) + } + + sessionID, err := resolveCalendarSession(ctx, addr, subject) + if err != nil { + return err + } + + // 发件人固定为 calendar:它不是任何一个 Agent,也不是人。 + // 用创建者的名字会让 Agent 以为人在实时找它,而人此刻可能在睡觉 —— + // 模型据此判断「要不要马上追问」,来源写错会让它问一个不在线的人。 + mailID, err := repo.CreateMail(ctx, sessionID, nil, + calendarSender, "", addr.Name, addr.Path, subject, body, ccList) + if err != nil { + return err + } + + // 事件附件复制成邮件附件。失败不阻断投递:提醒本身(正文)比附件重要得多, + // 少一个附件也比整条提醒发不出去好 —— 后者会让人以为定时任务坏了。 + if eventID != "" { + if n, aErr := repo.AttachCalendarFilesToMail(ctx, eventID, mailID, calendarSender); aErr != nil { + log.Printf("[scheduler] 事件 %s 的附件挂载失败(提醒仍已发出): %v", short(eventID), aErr) + } else if n > 0 { + log.Printf("[scheduler] 事件 %s 随提醒带了 %d 个附件", short(eventID), n) + } + } + + alias := repo.SessionAliasOf(ctx, sessionID) + sse.Default.SendToRecipient(addr.Name, "new_mail", map[string]interface{}{ + "mail_id": mailID.String(), + "session_id": sessionID.String(), + "from_name": calendarSender, + "subject": subject, + "mail_type": "normal", + "role": "to", + "to_workspace": addr.Path, + "session_alias": alias, + // 回信地址给 calendar 是发不出去的(它不是收件方), + // 给会话自己的地址才让模型能把结果回报到同一条线索上。 + "reply_address": models.FormatAddress(addr.Name, addr.Path, alias), + "self_address": models.FormatAddress(addr.Name, addr.Path, alias), + // 让插件与 UI 能区分「这封是定时提醒」而不是有人在找它 + "origin": "calendar", + }) + sse.Default.SendToRecipient(addr.Name, "session_update", map[string]interface{}{ + "session_id": sessionID.String(), + "status": "active", + }) + + // 抄送方也要收到 SSE。 + // + // 漏了这一步的后果很隐蔽:邮件的 cc_list 里有他们、他们**查**收件箱 + // 能看到这封信,但没有任何事件推给他们 —— 于是插件不会唤起会话, + // Agent 直到下一次补拉(重启时)才发现。对「知情方」而言等于没通知。 + for _, c := range ccList { + alias := repo.SessionAliasOf(ctx, sessionID) + sse.Default.SendToRecipient(c.Name, "new_mail", map[string]interface{}{ + "mail_id": mailID.String(), + "session_id": sessionID.String(), + "from_name": calendarSender, + "subject": subject, + "mail_type": "normal", + "role": "cc", + "to_workspace": c.Path, + "session_alias": alias, + "reply_address": models.FormatAddress(addr.Name, addr.Path, alias), + "self_address": models.FormatAddress(c.Name, c.Path, alias), + "origin": "calendar", + }) + sse.Default.SendToRecipient(c.Name, "session_update", map[string]interface{}{ + "session_id": sessionID.String(), + "status": "active", + }) + } + + // 创建者也该看到提醒发出去了 —— 否则「我设的提醒到底触发了没有」 + // 只能去翻 journalctl。 + if createdBy != "" && createdBy != addr.Name { + sse.Default.SendToRecipient(createdBy, "session_update", map[string]interface{}{ + "session_id": sessionID.String(), + "status": "active", + }) + } + return nil +} + +// calendarSender 是提醒邮件的发件人名。 +// +// 刻意不是人类用户名也不是 Agent 名:这样收件方一眼能看出这封信来自定时任务, +// 而 `from_name` 又不会与命名空间里任何真实账号冲突。 +const calendarSender = "calendar" + +// resolveCalendarSession 按地址的 session 位定位会话。 +// +// 与 handler.resolveTarget 同一套三态语义,但**不受新建会话速率限制**: +// 那条限制是防 Agent 暴开线索的,而日历事件的数量由人在界面上决定。 +func resolveCalendarSession(ctx context.Context, addr models.Address, subject string) (uuid.UUID, error) { + switch addr.Mode() { + case models.SessionNew: + id, err := repo.CreateSession(ctx, nil, calendarSender, subject, addr.Path) + if err != nil { + return uuid.Nil, err + } + // 与发信路径一致:`.new` 建完必须立刻有别名,否则这条会话 + // 除了回复那一封之外再也无法寻址(未命名会话查不到也补全不出来)。 + _, _ = repo.EnsureSessionAlias(ctx, id, repo.AutoAliasFor(addr.Name, subject)) + return id, nil + + case models.SessionNamed: + id, err := repo.FindNamedSessionFor(ctx, addr.Name, addr.Path, addr.Session) + if errors.Is(err, repo.ErrSessionNotFound) { + // 不静默新建:别名指向不存在的会话时报错,与人发信时的语义一致。 + // 静默新建会让「提醒发到哪去了」变成一个查不清的问题。 + return uuid.Nil, fmt.Errorf("会话 %q 不存在于 %s@%s", addr.Session, addr.Name, addr.Path) + } + if err != nil { + return uuid.Nil, err + } + repo.TouchSession(ctx, id) + return id, nil + + default: // SessionDefault + id, err := repo.FindOrCreateDefaultSession(ctx, addr.Name, addr.Path, calendarSender, subject) + if err != nil { + return uuid.Nil, err + } + _, _ = repo.EnsureSessionAlias(ctx, id, repo.AutoAliasFor(addr.Name, subject)) + return id, nil + } +} diff --git a/gateway/internal/scheduler/calendar_test.go b/gateway/internal/scheduler/calendar_test.go new file mode 100644 index 0000000..e01db20 --- /dev/null +++ b/gateway/internal/scheduler/calendar_test.go @@ -0,0 +1,111 @@ +package scheduler + +import ( + "strings" + "testing" + "time" + + "github.com/agentmail/gateway/internal/models" +) + +func TestRenderReminder(t *testing.T) { + at := time.Date(2026, 9, 4, 9, 30, 0, 0, time.Local) + e := models.CalendarEvent{ + Title: "每日站会", + Description: "同步昨天进展与今天计划", + EventTime: at, + } + + t.Run("三个变量都替换", func(t *testing.T) { + got := RenderReminder("日程提醒:{title}\n时间:{time}\n{description}", e) + for _, want := range []string{"每日站会", "2026-09-04 09:30", "同步昨天进展与今天计划"} { + if !strings.Contains(got, want) { + t.Errorf("渲染结果缺 %q:\n%s", want, got) + } + } + if strings.Contains(got, "{") { + t.Errorf("仍有未替换的占位符:\n%s", got) + } + }) + + t.Run("同一变量出现多次全部替换", func(t *testing.T) { + // ReplaceAll 而非 Replace:模板里写两遍 {title} 时 + // 只替换第一处会让 Agent 收到一封半成品邮件。 + got := RenderReminder("{title} —— 请开始 {title}", e) + if strings.Contains(got, "{title}") { + t.Errorf("第二处 {title} 未替换:%s", got) + } + }) + + t.Run("空模板不产生占位符残留", func(t *testing.T) { + if got := RenderReminder("", e); got != "" { + t.Errorf("空模板应渲染成空串,得到 %q", got) + } + }) + + t.Run("没有变量的模板原样返回", func(t *testing.T) { + const plain = "该跑测试了" + if got := RenderReminder(plain, e); got != plain { + t.Errorf("纯文本模板应原样返回,得到 %q", got) + } + }) + + t.Run("描述为空时不留下空行以外的痕迹", func(t *testing.T) { + e2 := e + e2.Description = "" + got := RenderReminder("{title}|{description}|", e2) + if got != "每日站会||" { + t.Errorf("空描述应替换成空串,得到 %q", got) + } + }) + + // {time} 必须是本地时间。DSN 带 _timezone=UTC,从库里读回的 EventTime + // 是 UTC;不转本地就会把人在 +0800 输入的 14:30 写成 06:30, + // 而前端预览用的是本地时间 —— 两边差 8 小时且都不报错。 + t.Run("time 用本地时区而非 UTC", func(t *testing.T) { + // 刻意构造一个 UTC 时刻(模拟从库里 Scan 出来的样子) + utcEvent := models.CalendarEvent{ + Title: "跨时区检查", + EventTime: time.Date(2026, 9, 3, 6, 30, 0, 0, time.UTC), + } + got := RenderReminder("{time}", utcEvent) + want := time.Date(2026, 9, 3, 6, 30, 0, 0, time.UTC).Local().Format("2006-01-02 15:04") + if got != want { + t.Errorf("{time} = %q,期望本地时间 %q", got, want) + } + }) + + // 同一时刻无论以哪个时区的 Location 传进来,渲染结果必须一致 —— + // 它代表的是「墙上时钟的那一刻」,与 Location 的表示方式无关。 + t.Run("同一时刻不同 Location 渲染一致", func(t *testing.T) { + base := time.Date(2026, 9, 3, 6, 30, 0, 0, time.UTC) + a := RenderReminder("{time}", models.CalendarEvent{EventTime: base}) + b := RenderReminder("{time}", models.CalendarEvent{EventTime: base.Local()}) + if a != b { + t.Errorf("UTC 与 Local 表示同一时刻却渲染出不同结果:%q vs %q", a, b) + } + }) +} + +func TestShort(t *testing.T) { + // 日志里截前 8 位;短 id(测试里可能出现)不能 panic + if got := short("0123456789abcdef"); got != "01234567" { + t.Errorf("长 id 应截断成 8 位,得到 %q", got) + } + if got := short("abc"); got != "abc" { + t.Errorf("短 id 应原样返回,得到 %q", got) + } + if got := short(""); got != "" { + t.Errorf("空串应原样返回,得到 %q", got) + } +} + +func TestStartStopIdempotent(t *testing.T) { + // Stop 在没启动时被调(defer 里必然发生)不该 panic; + // Start 两次也不该泄漏 goroutine(第二次先停旧的)。 + Stop() + Start() + Start() + Stop() + Stop() +}