package service import ( "errors" "fmt" "strings" "sync" "time" "github.com/robfig/cron/v3" "github.com/sirupsen/logrus" "hy2xs-admin/dao" "hy2xs-admin/model/constant" ) // Планировщик принадлежит ПРОЦЕССУ, а не HTTP-серверу. // // Что было. InitCron жил в middleware, вызывался из runServer и на каждом // вызове создавал новый cron.New(). Ссылка на планировщик никуда не // сохранялась, а releaseResource() закрывал только SQLite и Hysteria — // cron.Stop() не вызывался нигде. При этом смена RESET_TRAFFIC_CRON делала // StopServer(), из-за чего runServer возвращался, а точка входа в cmd.go // крутила его в `for {}` и запускала всё заново. // // Итог: каждая смена расписания добавляла ЦЕЛЫЙ дублирующий набор джоб, а // старое расписание сброса трафика продолжало работать. После двух правок // «@monthly → @weekly → @monthly» на процессе висели три планировщика, // утроенный CollectMetricsSnapshot и три разных расписания сброса // одновременно. Плюс окно, в котором джобы старого планировщика били в уже // закрытое соединение SQLite: releaseResource() отрабатывал раньше, чем // следующий runServer успевал открыть базу. // // Правильная граница: фиксированные джобы регистрируются один раз за жизнь // процесса, а расписание сброса трафика перепланируется на месте по своему // EntryID. HTTP-сервер к смене настройки отношения не имеет вообще. // cronShutdownTimeout ограничивает ожидание уже запущенных джоб при остановке. // systemd по умолчанию даёт юниту 90 секунд, так что запас есть. const cronShutdownTimeout = 10 * time.Second var ( cronMu sync.Mutex cronScheduler *cron.Cron resetTrafficEntryID cron.EntryID // cronOwnedJobs учитывает джобы, запущенные ВНЕ расписания. // // Стартовая уборка статистики намеренно выполняется в фоне: на большой базе // она заметно долгая, и держать на ней запуск сервиса незачем. Но // cron.Stop() ждёт только то, что запустил сам планировщик, поэтому // необслуженная горутина переживала закрытие SQLite — ровно та же болезнь, // от которой лечится весь этот файл, только меньшего масштаба. cronOwnedJobs sync.WaitGroup ) // ValidateResetTrafficCron проверяет выражение ТЕМ ЖЕ парсером, которым его // потом будет разбирать runtime. // // cron.New() без опций собирает parser из Minute|Hour|Dom|Month|Dow|Descriptor, // и ровно его же использует cron.ParseStandard. Поэтому «валидно на входе API» // и «планируется в рантайме» здесь не могут разойтись — а разойтись они могли // бы, если бы валидация была написана собственной регуляркой. // // Пустая строка — легальное значение и означает «автоматический сброс // выключен»: в панели поле clearable, и оператор имеет право его очистить. func ValidateResetTrafficCron(expression string) error { trimmed := strings.TrimSpace(expression) if trimmed == "" { return nil } if _, err := cron.ParseStandard(trimmed); err != nil { return fmt.Errorf("%s: невалидное cron-выражение %q: %v", constant.ResetTrafficCron, expression, err) } return nil } // InitCron поднимает единственный планировщик процесса. func InitCron() error { cronMu.Lock() defer cronMu.Unlock() if cronScheduler != nil { return errors.New("cron scheduler is already running") } c := cron.New(cron.WithLocation(time.Now().Location())) fixedJobs := []struct { name string spec string job func() }{ {"CronHandleAccount", "@every 30s", CronHandleAccount}, {"CollectMetricsSnapshot", "@every 10s", CollectMetricsSnapshot}, {"CleanupStatsRetention", "@every 1h", CleanupStatsRetention}, } for _, fixed := range fixedJobs { if _, err := c.AddFunc(fixed.spec, fixed.job); err != nil { logrus.Errorf("cron add func %s err: %v", fixed.name, err) return fmt.Errorf("cron add func %s err", fixed.name) } } expression, err := storedResetTrafficCron() if err != nil { return err } if expression != "" { id, addErr := c.AddFunc(expression, CronResetTraffic) if addErr != nil { // Старт НЕ прерывается, и это осознанное решение. // // Панель отдаёт не только операторский UI: на ней же висит // /internal/hysteria/auth, куда Hysteria ходит при каждом // подключении пира. Отказ старта из-за испорченной строки // расписания положил бы подключения пользователей — цена // несопоставима с отключённым плановым сбросом счётчиков. // // Молчаливой деградации при этом нет: запись через API теперь // валидируется тем же парсером, поэтому попасть сюда можно только // правкой базы в обход продукта, и об этом пишется ERROR. logrus.Errorf( "cron: сохранённое %s=%q невалидно (%v); плановый сброс трафика выключен до исправления настройки", constant.ResetTrafficCron, expression, addErr, ) } else { resetTrafficEntryID = id } } cronOwnedJobs.Add(1) go func() { defer cronOwnedJobs.Done() CleanupStatsRetention() }() c.Start() cronScheduler = c return nil } // RescheduleResetTraffic переносит джобу сброса трафика на новое расписание. // // HTTP-сервер здесь не участвует. Раньше единственным способом применить новое // расписание был перезапуск процесса через StopServer(), и именно он и плодил // планировщики. func RescheduleResetTraffic(expression string) error { if err := ValidateResetTrafficCron(expression); err != nil { return err } cronMu.Lock() defer cronMu.Unlock() if cronScheduler == nil { return errors.New("cron scheduler is not running") } // Старая запись снимается всегда, даже если новая не будет добавлена: // пустое выражение означает «сброс выключен», а не «оставить как было». if resetTrafficEntryID != 0 { cronScheduler.Remove(resetTrafficEntryID) resetTrafficEntryID = 0 } trimmed := strings.TrimSpace(expression) if trimmed == "" { return nil } // Ошибка здесь уже невозможна: выражение прошло тот же парсер выше. // Проверка остаётся, чтобы расхождение двух парсеров не превратилось в // молча пропавшую джобу. id, err := cronScheduler.AddFunc(trimmed, CronResetTraffic) if err != nil { return fmt.Errorf("%s: не удалось запланировать %q: %v", constant.ResetTrafficCron, expression, err) } resetTrafficEntryID = id return nil } // StopCron останавливает планировщик и дожидается уже запущенных джоб. // // Без этого джобы продолжали работать после закрытия SQLite: каждая из них // ходит в базу, и остановка в обратном порядке (сначала планировщик, потом // соединение) — часть контракта завершения процесса. func StopCron() { cronMu.Lock() scheduler := cronScheduler cronScheduler = nil resetTrafficEntryID = 0 cronMu.Unlock() if scheduler == nil { return } // Дожидаемся обеих групп: и джоб, запущенных планировщиком, и фоновых, // которые он не видит. Незавершённая джоба означает работу с базой, которую // releaseResource закроет сразу после возврата отсюда. drained := make(chan struct{}) go func() { <-scheduler.Stop().Done() cronOwnedJobs.Wait() close(drained) }() select { case <-drained: case <-time.After(cronShutdownTimeout): logrus.Warnf("cron: запущенные джобы не завершились за %s, продолжаем остановку", cronShutdownTimeout) } } // ResetTrafficScheduled сообщает, запланирован ли сейчас сброс трафика. // Существует ради тестов: иначе проверить, что старая запись действительно // снята, а не просто добавлена рядом, можно было бы только по времени. func ResetTrafficScheduled() bool { cronMu.Lock() defer cronMu.Unlock() return resetTrafficEntryID != 0 } // CronEntryCount возвращает число активных записей планировщика. // Тоже ради тестов: накопление джоб — это именно рост этого числа. func CronEntryCount() int { cronMu.Lock() defer cronMu.Unlock() if cronScheduler == nil { return 0 } return len(cronScheduler.Entries()) } func storedResetTrafficCron() (string, error) { config, err := dao.GetConfig("key = ?", constant.ResetTrafficCron) if err != nil { return "", err } if config.Value == nil { return "", nil } return strings.TrimSpace(*config.Value), nil }