diff --git a/README.md b/README.md index 692d198..0acdfd1 100644 --- a/README.md +++ b/README.md @@ -72,7 +72,7 @@ http://服务器地址:8080 - 同步间隔:1–60 分钟 - 历史保留时间:30、60、90、180 或 365 天 - 提前提醒时间:重置前 1–1440 分钟 -- 是否发送重置前提醒和重置后确认 +- 是否发送重置前提醒和重置后确认;重置后确认也会通过额度百分比回落识别并提醒官方活动、临时补发等提前重置 提醒只会针对 Codex app-server 返回的限额窗口发送。发送失败的提醒会在计划时间后的六小时内自动重试。 diff --git a/backend/internal/app/api.go b/backend/internal/app/api.go index a59aeda..e03fb4c 100644 --- a/backend/internal/app/api.go +++ b/backend/internal/app/api.go @@ -449,14 +449,66 @@ func (a *App) syncAccount(ctx context.Context, id int64) error { _, _ = a.store.DB.Exec("INSERT INTO daily_usage(account_id,date,total_tokens,fetched_at) VALUES(?,?,?,?) ON CONFLICT(account_id,date) DO UPDATE SET total_tokens=excluded.total_tokens,fetched_at=excluded.fetched_at", id, x.StartDate, x.Tokens, d.FetchedAt) } } - for _, x := range d.Limits { - _, _ = a.store.DB.Exec("INSERT INTO limit_snapshots(limit_id,window_type,used_percent,duration_mins,resets_at,fetched_at,account_id) VALUES(?,?,?,?,?,?,?)", x.LimitID, x.WindowType, x.UsedPercent, x.WindowDurationMinutes, x.ResetsAt, d.FetchedAt, id) + resetDetected, e := a.storeLimitSnapshots(d) + if e != nil { + return e } rt.dash = d _ = a.store.UpdateAccount(id, d.Account.Email, d.Account.PlanType, d.Account.Connected) + if resetDetected { + go a.processReminders() + } return nil } +const resetDropTolerance = 0.01 + +func (a *App) storeLimitSnapshots(d Dashboard) (bool, error) { + tx, err := a.store.DB.Begin() + if err != nil { + return false, err + } + defer tx.Rollback() + g := a.general() + resetDetected := false + for _, x := range d.Limits { + var previousID, previousFetchedAt, previousResetsAt int64 + var previousUsed float64 + err = tx.QueryRow(`SELECT id,used_percent,resets_at,fetched_at FROM limit_snapshots + WHERE account_id=? AND limit_id=? AND window_type=? ORDER BY fetched_at DESC,id DESC LIMIT 1`, + d.AccountID, x.LimitID, x.WindowType).Scan(&previousID, &previousUsed, &previousResetsAt, &previousFetchedAt) + if err != nil && err != sql.ErrNoRows { + return false, err + } + age := d.FetchedAt - previousFetchedAt + if err == nil && g.NotifyAfter && age >= 0 && age <= int64((6*time.Hour).Seconds()) && previousUsed-x.UsedPercent > resetDropTolerance { + kind := "detected_after" + key := fmt.Sprintf("%d:%s:%s:detected:%d", d.AccountID, x.LimitID, x.WindowType, previousID) + now := time.Unix(d.FetchedAt, 0) + if previousResetsAt <= d.FetchedAt && now.Sub(time.Unix(previousResetsAt, 0)) <= 6*time.Hour { + kind = "after" + key = fmt.Sprintf("%d:%s:%s:%d:after", d.AccountID, x.LimitID, x.WindowType, previousResetsAt) + } + body := fmt.Sprintf("Codex [%s] %s/%s 额度已重置:已用 %.1f%% → %.1f%%,下次重置时间 %s。", + d.DisplayName, x.LimitID, x.WindowType, previousUsed, x.UsedPercent, time.Unix(x.ResetsAt, 0).Format(time.RFC3339)) + _, err = tx.Exec(`INSERT OR IGNORE INTO notifications + (dedupe_key,channel,kind,status,attempts,last_error,scheduled_at,sent_at,body) + VALUES(?,?,?,'pending',0,'',?,NULL,?)`, key, "configured", kind, d.FetchedAt, body) + if err != nil { + return false, err + } + resetDetected = true + } + if _, err = tx.Exec("INSERT INTO limit_snapshots(limit_id,window_type,used_percent,duration_mins,resets_at,fetched_at,account_id) VALUES(?,?,?,?,?,?,?)", x.LimitID, x.WindowType, x.UsedPercent, x.WindowDurationMinutes, x.ResetsAt, d.FetchedAt, d.AccountID); err != nil { + return false, err + } + } + if err = tx.Commit(); err != nil { + return false, err + } + return resetDetected, nil +} + type rawLimit struct { LimitID string `json:"limitId"` LimitName *string `json:"limitName"` diff --git a/backend/internal/app/app.go b/backend/internal/app/app.go index b32078f..783eb6b 100644 --- a/backend/internal/app/app.go +++ b/backend/internal/app/app.go @@ -35,6 +35,7 @@ type App struct { mu sync.RWMutex runtimes map[int64]*accountRuntime loginAttempts sync.Map + reminderMu sync.Mutex } type accountRuntime struct { client codexClient diff --git a/backend/internal/app/notify.go b/backend/internal/app/notify.go index 831bf9c..12d8d0d 100644 --- a/backend/internal/app/notify.go +++ b/backend/internal/app/notify.go @@ -341,6 +341,8 @@ func num(n *int64) string { return fmt.Sprintf("%d", *n) } func (a *App) processReminders() { + a.reminderMu.Lock() + defer a.reminderMu.Unlock() g := a.general() a.mu.RLock() ds := make([]Dashboard, 0, len(a.runtimes)) @@ -368,41 +370,58 @@ func (a *App) processReminders() { continue } key := fmt.Sprintf("%d:%s:%s:%d:%s", d.AccountID, x.LimitID, x.WindowType, x.ResetsAt, kind) - var exists int - if a.store.DB.QueryRow("SELECT 1 FROM notifications WHERE dedupe_key=? AND status='sent'", key).Scan(&exists) == nil { - continue - } body := fmt.Sprintf("Codex [%s] %s/%s 剩余 %.1f%%,重置时间 %s。", d.DisplayName, x.LimitID, x.WindowType, 100-x.UsedPercent, time.Unix(x.ResetsAt, 0).Format(time.RFC3339)) - ok := true - errs := []string{} - if t, e := a.telegramSecret(); e == nil && t.Enabled && t.ChatID != 0 { - if e = tgSend(t, body); e != nil { - ok = false - errs = append(errs, e.Error()) - } - } - if s, e := a.smtpSecret(); e == nil && s.Enabled { - if e = sendSMTP(s, "Codex 用量重置提醒", body); e != nil { - ok = false - errs = append(errs, e.Error()) - } - } - status := "sent" - var sent any = time.Now().Unix() - if !ok { - status = "failed" - sent = nil - } - _, _ = a.store.DB.Exec(`INSERT INTO notifications(dedupe_key,channel,kind,status,attempts,last_error,scheduled_at,sent_at) - VALUES(?,?,?,?,?,?,?,?) - ON CONFLICT(dedupe_key) DO UPDATE SET - status=excluded.status, - attempts=notifications.attempts+1, - last_error=excluded.last_error, - sent_at=excluded.sent_at`, key, "configured", kind, status, 1, strings.Join(errs, "; "), at.Unix(), sent) + _, _ = a.store.DB.Exec(`INSERT OR IGNORE INTO notifications + (dedupe_key,channel,kind,status,attempts,last_error,scheduled_at,sent_at,body) + VALUES(?,?,?,'pending',0,'',?,NULL,?)`, key, "configured", kind, at.Unix(), body) } } } + a.sendPendingReminders(now) +} + +func (a *App) sendPendingReminders(now time.Time) { + type pendingReminder struct { + key string + body string + } + rows, err := a.store.DB.Query(`SELECT dedupe_key,body FROM notifications + WHERE status!='sent' AND scheduled_at<=? AND scheduled_at>=? ORDER BY scheduled_at`, now.Unix(), now.Add(-6*time.Hour).Unix()) + if err != nil { + return + } + pending := []pendingReminder{} + for rows.Next() { + var p pendingReminder + if rows.Scan(&p.key, &p.body) == nil && p.body != "" { + pending = append(pending, p) + } + } + _ = rows.Close() + for _, p := range pending { + ok := true + errs := []string{} + if t, e := a.telegramSecret(); e == nil && t.Enabled && t.ChatID != 0 { + if e = tgSend(t, p.body); e != nil { + ok = false + errs = append(errs, e.Error()) + } + } + if s, e := a.smtpSecret(); e == nil && s.Enabled { + if e = sendSMTP(s, "Codex 用量重置提醒", p.body); e != nil { + ok = false + errs = append(errs, e.Error()) + } + } + status := "sent" + var sent any = now.Unix() + if !ok { + status = "failed" + sent = nil + } + _, _ = a.store.DB.Exec(`UPDATE notifications SET status=?,attempts=attempts+1,last_error=?,sent_at=? WHERE dedupe_key=?`, + status, strings.Join(errs, "; "), sent, p.key) + } } var _ = bufio.ErrInvalidUnreadByte diff --git a/backend/internal/app/reminder_test.go b/backend/internal/app/reminder_test.go new file mode 100644 index 0000000..f4766ac --- /dev/null +++ b/backend/internal/app/reminder_test.go @@ -0,0 +1,112 @@ +package app + +import ( + "strconv" + "strings" + "testing" + "time" + + "codex-helper/internal/store" +) + +func newReminderTestApp(t *testing.T) *App { + t.Helper() + s, err := store.Open(t.TempDir()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = s.DB.Close() }) + return &App{store: s, runtimes: map[int64]*accountRuntime{}} +} + +func reminderDashboard(fetchedAt int64, used float64, resetsAt int64) Dashboard { + return Dashboard{ + AccountID: 1, + DisplayName: "测试账号", + FetchedAt: fetchedAt, + Limits: []LimitBucket{{ + LimitID: "codex", + WindowType: "primary", + UsedPercent: used, + WindowDurationMinutes: 300, + ResetsAt: resetsAt, + }}, + } +} + +func notificationCount(t *testing.T, a *App) int { + t.Helper() + var count int + if err := a.store.DB.QueryRow("SELECT COUNT(*) FROM notifications").Scan(&count); err != nil { + t.Fatal(err) + } + return count +} + +func TestStoreLimitSnapshotsDetectsEarlyReset(t *testing.T) { + a := newReminderTestApp(t) + now := time.Now().Unix() + if detected, err := a.storeLimitSnapshots(reminderDashboard(now, 42, now+3600)); err != nil || detected { + t.Fatalf("initial snapshot: detected=%v err=%v", detected, err) + } + if detected, err := a.storeLimitSnapshots(reminderDashboard(now+60, 3, now+7200)); err != nil || !detected { + t.Fatalf("reset snapshot: detected=%v err=%v", detected, err) + } + var kind, body string + if err := a.store.DB.QueryRow("SELECT kind,body FROM notifications").Scan(&kind, &body); err != nil { + t.Fatal(err) + } + if kind != "detected_after" || !strings.Contains(body, "42.0% → 3.0%") || !strings.Contains(body, "测试账号") { + t.Fatalf("notification kind=%q body=%q", kind, body) + } +} + +func TestStoreLimitSnapshotsUsesScheduledAfterDedupeKey(t *testing.T) { + a := newReminderTestApp(t) + now := time.Now().Unix() + resetAt := now + 30 + _, _ = a.storeLimitSnapshots(reminderDashboard(now, 70, resetAt)) + detected, err := a.storeLimitSnapshots(reminderDashboard(now+60, 0, now+3600)) + if err != nil || !detected { + t.Fatalf("detected=%v err=%v", detected, err) + } + var key, kind string + if err = a.store.DB.QueryRow("SELECT dedupe_key,kind FROM notifications").Scan(&key, &kind); err != nil { + t.Fatal(err) + } + exact := "1:codex:primary:" + strconv.FormatInt(resetAt, 10) + ":after" + if key != exact || kind != "after" { + t.Fatalf("key=%q kind=%q, want %q after", key, kind, exact) + } +} + +func TestStoreLimitSnapshotsIgnoresNonResetChanges(t *testing.T) { + tests := []struct { + name string + oldUsed float64 + newUsed float64 + age time.Duration + notifyAfter bool + }{ + {name: "increase", oldUsed: 10, newUsed: 20, age: time.Minute, notifyAfter: true}, + {name: "tolerance", oldUsed: 10, newUsed: 9.995, age: time.Minute, notifyAfter: true}, + {name: "old snapshot", oldUsed: 50, newUsed: 0, age: 7 * time.Hour, notifyAfter: true}, + {name: "disabled", oldUsed: 50, newUsed: 0, age: time.Minute, notifyAfter: false}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + a := newReminderTestApp(t) + g := defaults() + g.NotifyAfter = tt.notifyAfter + if err := a.store.SetJSON("general", g); err != nil { + t.Fatal(err) + } + now := time.Now().Unix() + _, _ = a.storeLimitSnapshots(reminderDashboard(now, tt.oldUsed, now+3600)) + detected, err := a.storeLimitSnapshots(reminderDashboard(now+int64(tt.age.Seconds()), tt.newUsed, now+7200)) + if err != nil || detected || notificationCount(t, a) != 0 { + t.Fatalf("detected=%v notifications=%d err=%v", detected, notificationCount(t, a), err) + } + }) + } +} diff --git a/backend/internal/store/store.go b/backend/internal/store/store.go index 6b5f482..30a2424 100644 --- a/backend/internal/store/store.go +++ b/backend/internal/store/store.go @@ -39,16 +39,50 @@ CREATE TABLE IF NOT EXISTS sessions (token_hash TEXT PRIMARY KEY, expires_at INT CREATE TABLE IF NOT EXISTS daily_usage (date TEXT PRIMARY KEY, total_tokens INTEGER NOT NULL, fetched_at INTEGER NOT NULL); CREATE TABLE IF NOT EXISTS limit_snapshots (id INTEGER PRIMARY KEY AUTOINCREMENT, limit_id TEXT NOT NULL, window_type TEXT NOT NULL, used_percent REAL NOT NULL, duration_mins INTEGER NOT NULL, resets_at INTEGER NOT NULL, fetched_at INTEGER NOT NULL); CREATE INDEX IF NOT EXISTS idx_limits_time ON limit_snapshots(fetched_at); -CREATE TABLE IF NOT EXISTS notifications (dedupe_key TEXT PRIMARY KEY, channel TEXT NOT NULL, kind TEXT NOT NULL, status TEXT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, last_error TEXT, scheduled_at INTEGER NOT NULL, sent_at INTEGER); +CREATE TABLE IF NOT EXISTS notifications (dedupe_key TEXT PRIMARY KEY, channel TEXT NOT NULL, kind TEXT NOT NULL, status TEXT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, last_error TEXT, scheduled_at INTEGER NOT NULL, sent_at INTEGER, body TEXT NOT NULL DEFAULT ''); CREATE TABLE IF NOT EXISTS telegram_updates (id INTEGER PRIMARY KEY CHECK(id=1), offset INTEGER NOT NULL DEFAULT 0); INSERT OR IGNORE INTO telegram_updates(id,offset) VALUES(1,0); `) if err != nil { return err } + if err = s.ensureNotificationBody(); err != nil { + return err + } return s.migrateAccounts() } +func (s *Store) ensureNotificationBody() error { + rows, err := s.DB.Query("PRAGMA table_info(notifications)") + if err != nil { + return err + } + hasBody := false + for rows.Next() { + var cid, notnull, pk int + var name, typ string + var def any + if err = rows.Scan(&cid, &name, &typ, ¬null, &def, &pk); err != nil { + rows.Close() + return err + } + hasBody = hasBody || name == "body" + } + if err = rows.Close(); err != nil { + return err + } + if hasBody { + return nil + } + // Legacy unsent rows have no recoverable message body. Remove them so the + // reminder scheduler can recreate every still-applicable retry from the + // current dashboard instead of having INSERT OR IGNORE collide with an + // empty-body row. Sent rows remain as dedupe records. + _, err = s.DB.Exec(`ALTER TABLE notifications ADD COLUMN body TEXT NOT NULL DEFAULT ''; + DELETE FROM notifications WHERE status != 'sent';`) + return err +} + func (s *Store) migrateAccounts() error { tx, err := s.DB.Begin() if err != nil { diff --git a/backend/internal/store/store_test.go b/backend/internal/store/store_test.go index d83b262..5459c13 100644 --- a/backend/internal/store/store_test.go +++ b/backend/internal/store/store_test.go @@ -174,6 +174,44 @@ func TestPopulatedLegacyLimitSnapshotsMigrateToDefaultAccount(t *testing.T) { } } +func TestLegacyNotificationsGainBodyColumn(t *testing.T) { + dir := t.TempDir() + db, err := sql.Open("sqlite", filepath.Join(dir, "codex-helper.db")) + if err != nil { + t.Fatal(err) + } + _, err = db.Exec(`CREATE TABLE notifications ( + dedupe_key TEXT PRIMARY KEY, channel TEXT NOT NULL, kind TEXT NOT NULL, + status TEXT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, last_error TEXT, + scheduled_at INTEGER NOT NULL, sent_at INTEGER + ); + INSERT INTO notifications(dedupe_key,channel,kind,status,scheduled_at) + VALUES('pending-key','configured','after','pending',1); + INSERT INTO notifications(dedupe_key,channel,kind,status,scheduled_at,sent_at) + VALUES('sent-key','configured','after','sent',1,2);`) + if err != nil { + t.Fatal(err) + } + _ = db.Close() + + s, err := Open(dir) + if err != nil { + t.Fatal(err) + } + defer s.DB.Close() + var count int + if err = s.DB.QueryRow("SELECT COUNT(*) FROM notifications WHERE dedupe_key='pending-key'").Scan(&count); err != nil || count != 0 { + t.Fatalf("legacy pending notification was not removed: count=%d err=%v", count, err) + } + if err = s.DB.QueryRow("SELECT COUNT(*) FROM notifications WHERE dedupe_key='sent-key'").Scan(&count); err != nil || count != 1 { + t.Fatalf("sent dedupe notification was not preserved: count=%d err=%v", count, err) + } + if _, err = s.DB.Exec(`INSERT INTO notifications + (dedupe_key,channel,kind,status,scheduled_at,body) VALUES('key','configured','after','pending',1,'message')`); err != nil { + t.Fatalf("body column was not added: %v", err) + } +} + func TestBackupIncludesCommittedWALData(t *testing.T) { dir := t.TempDir() s, err := Open(dir)