feat: notify on unexpected quota resets
This commit is contained in:
@@ -72,7 +72,7 @@ http://服务器地址:8080
|
||||
- 同步间隔:1–60 分钟
|
||||
- 历史保留时间:30、60、90、180 或 365 天
|
||||
- 提前提醒时间:重置前 1–1440 分钟
|
||||
- 是否发送重置前提醒和重置后确认
|
||||
- 是否发送重置前提醒和重置后确认;重置后确认也会通过额度百分比回落识别并提醒官方活动、临时补发等提前重置
|
||||
|
||||
提醒只会针对 Codex app-server 返回的限额窗口发送。发送失败的提醒会在计划时间后的六小时内自动重试。
|
||||
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,40 +370,57 @@ 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))
|
||||
_, _ = 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, body); e != nil {
|
||||
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 用量重置提醒", body); e != nil {
|
||||
if e = sendSMTP(s, "Codex 用量重置提醒", p.body); e != nil {
|
||||
ok = false
|
||||
errs = append(errs, e.Error())
|
||||
}
|
||||
}
|
||||
status := "sent"
|
||||
var sent any = time.Now().Unix()
|
||||
var sent any = 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(`UPDATE notifications SET status=?,attempts=attempts+1,last_error=?,sent_at=? WHERE dedupe_key=?`,
|
||||
status, strings.Join(errs, "; "), sent, p.key)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user