8dcb50a07c
Предыдущий проход сделал правильным порядок «сначала долговременная запись, потом разрыв сессии» и правильно запретил откат при неудаче разрыва. Способа прийти к согласованному состоянию ПОТОМ он не дал: у двух операций повтор не работал вовсе. Импорт, заменивший auth_id: после неудавшегося /kick старое значение не хранится нигде, повтор того же файла читает из базы уже новое и рвёт его, а cron пропускал незнакомый authID молча — dao.ListPeer просто не возвращала строку. Живая сессия оставалась навсегда. Снижение maxDevices: повтор формы даёт 1 < 1 -> false, разрыва больше нет. Лимит устройств в политику доступа не входит и входить не должен — это свойство сессий, — поэтому механизма схождения у него не было. enforcePeerAccess стал сверкой живых сессий: обход идёт по каждому authID из /online. Нет строки в базе -> kick; peerAccessDenied -> kick; непригодный maxDevices -> kick; устройств больше разрешённого -> kick. Отказ базы при этом не рвёт ничего. Ни таблицы отложенных операций, ни очереди retry: список живых сессий уже есть, и это /online. Отдельно закрыт второй TOCTOU лимита устройств. Учёт выданных разрешений закрыл сравнение двух одинаковых снимков, но сетевой запрос выполнялся вне блокировки, поэтому снимки приходили в резервацию в произвольном порядке и устаревший откатывал lastOnline назад, возвращая уже занятое место. Это не data race — память защищена мьютексом, и -race здесь молчит принципиально. Последовательность «прочитать /online -> занять место» выполняется под замком по authId; глобальный замок не годится, внутри идёт сетевой запрос. Учёт разрешений больше не растёт бесконечно: запись снималась только на ветке отказа, поэтому в карте копились удалённые пиры и переписанные импортом идентификаторы. Уборка идёт по фактической картине подключений. Гейты приёмки доращены под все три инварианта и проверены в обе стороны. Go 1.26.7 -> 1.26.8. Документация приведена в соответствие в двух местах, где описывала снятую архитектуру. Разбор: docs/acceptance/2026-09-02-v1.0.0-rc3-preflight-findings.md
351 lines
18 KiB
Go
351 lines
18 KiB
Go
package service
|
||
|
||
import (
|
||
"fmt"
|
||
"sort"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/sirupsen/logrus"
|
||
"gorm.io/gorm"
|
||
"hy2xs-admin/dao"
|
||
"hy2xs-admin/model/entity"
|
||
"hy2xs-admin/proxy"
|
||
"hy2xs-admin/util"
|
||
)
|
||
|
||
// Джоба учёта принадлежит планировщику, а не собственным горутинам.
|
||
//
|
||
// Что было:
|
||
//
|
||
// CronHandleAccount()
|
||
// -> go func()
|
||
// -> go saveAccountTraffic()
|
||
// -> go kickAccount()
|
||
//
|
||
// Три уровня отсоединённых горутин. Для cron.Cron джоба заканчивалась почти
|
||
// мгновенно — сразу после запуска внешней, — поэтому StopCron(), который
|
||
// честно ждёт `scheduler.Stop().Done()`, не ждал НИЧЕГО из настоящей работы.
|
||
// Завершение процесса выглядело так: планировщик отчитался «джоб не осталось»,
|
||
// releaseResource() закрыл SQLite, а внутренние горутины продолжали писать
|
||
// трафик и рвать сессии в уже закрытое соединение. Это ровно та болезнь, от
|
||
// которой лечится cron_scheduler.go, только протащенная внутрь одной джобы.
|
||
//
|
||
// Второе следствие того же устройства было тише и хуже. Обе внутренние
|
||
// горутины запускались ПАРАЛЛЕЛЬНО, поэтому принудительное отключение читало
|
||
// счётчики трафика ДО того, как в них попадала только что снятая дельта. При
|
||
// тридцатисекундном тике это значит, что превышение квоты замечалось в лучшем
|
||
// случае со следующего цикла, а на границе — не замечалось вовсе.
|
||
//
|
||
// Теперь джоба синхронна, порядок внутри неё строгий, а взаимное исключение
|
||
// даёт один мьютекс на весь цикл: сбор трафика и enforcement больше не могут
|
||
// ни разъехаться во времени, ни наложиться сами на себя.
|
||
|
||
// accountJobMutex сериализует цикл учёта.
|
||
//
|
||
// Заменяет пару trafficMutex + kickMutex. Раздельные мьютексы защищали каждую
|
||
// половину от самой себя, но не защищали пару от расщепления: при затянувшемся
|
||
// сборе трафика следующий тик мог запустить enforcement поверх предыдущего
|
||
// сбора. Одного мьютекса на весь цикл достаточно и, в отличие от двух, он
|
||
// выражает действительный инвариант — «в любой момент времени выполняется не
|
||
// более одного цикла учёта».
|
||
var accountJobMutex sync.Mutex
|
||
|
||
// CronHandleAccount — один синхронный цикл учёта: собрать трафик, затем
|
||
// применить политику доступа.
|
||
//
|
||
// Состояние службы по systemd здесь НЕ спрашивается. Прежний гейт
|
||
//
|
||
// if !Hysteria2IsRunning() { return }
|
||
//
|
||
// стоял на решении о применении операции, а Hysteria2IsRunning для этого
|
||
// непригоден по собственному объявлению: util.Exec схлопывает «systemctl
|
||
// вернул 3, служба неактивна» и «запустить systemctl не удалось» в одну
|
||
// ошибку. То есть сломанный systemctl при живой Hysteria молча отключал и учёт
|
||
// трафика, и принудительное отключение — без единой строки в журнале.
|
||
//
|
||
// Нужные системы спрашиваются напрямую: `/traffic`, `/online`, `/kick`. Если
|
||
// Hysteria действительно не работает, вызов вернёт ошибку, и она будет
|
||
// записана. Если сломан systemctl, а Hysteria жива, учёт продолжит работать.
|
||
func CronHandleAccount() {
|
||
// Пропуск тика при уже идущем цикле — не отказ: следующий тик через 30
|
||
// секунд, а очередь из накопившихся циклов ничего бы не дала.
|
||
if !accountJobMutex.TryLock() {
|
||
return
|
||
}
|
||
defer accountJobMutex.Unlock()
|
||
|
||
apiPort, err := GetHysteria2ApiPort()
|
||
if err != nil {
|
||
logrus.WithError(err).Error("cron: не удалось определить порт Traffic Stats API; цикл учёта пропущен")
|
||
return
|
||
}
|
||
|
||
// Секрет берётся общей функцией, которая отличает «ключа нет» от пустого
|
||
// значения. Раньше здесь стояло `*trafficSecretConfig.Value` без единой
|
||
// проверки: строка в таблице `config` без значения роняла бы процесс
|
||
// паникой на разыменовании nil — причём внутри отсоединённой горутины, где
|
||
// её некому перехватить, то есть падал бы весь сервис вместе с
|
||
// обработчиком machine-auth.
|
||
secret, err := hysteria2TrafficSecret()
|
||
if err != nil {
|
||
logrus.WithError(err).Error("cron: секрет Traffic Stats API недоступен; цикл учёта пропущен")
|
||
return
|
||
}
|
||
|
||
// Порядок обязателен: enforcement принимает решение по счётчикам, поэтому
|
||
// счётчики должны быть уже обновлены.
|
||
if err := saveAccountTraffic(apiPort, secret); err != nil {
|
||
logrus.WithError(err).Error("cron: сбор трафика завершился с ошибкой")
|
||
}
|
||
|
||
if err := enforcePeerAccess(apiPort, secret); err != nil {
|
||
logrus.WithError(err).Error("cron: принудительное отключение завершилось с ошибкой")
|
||
}
|
||
}
|
||
|
||
// CronResetTraffic обнуляет счётчики трафика всех пиров по расписанию.
|
||
func CronResetTraffic() {
|
||
peers, err := dao.ListPeer("1=1")
|
||
if err != nil {
|
||
logrus.WithError(err).Error("cron: не удалось прочитать пиров для сброса трафика")
|
||
return
|
||
}
|
||
ids := make([]int64, 0, len(peers))
|
||
for _, item := range peers {
|
||
// Строка без идентификатора — повреждённые данные. Раньше здесь
|
||
// стояло `*item.Id` без проверки, то есть такая строка роняла джобу
|
||
// паникой, а вместе с ней и процесс.
|
||
if item.Id == nil {
|
||
logrus.Error("cron: строка пира без идентификатора пропущена при сбросе трафика")
|
||
continue
|
||
}
|
||
ids = append(ids, *item.Id)
|
||
}
|
||
if len(ids) == 0 {
|
||
return
|
||
}
|
||
for _, chunk := range util.SplitArr(ids, 100) {
|
||
if err := dao.UpdatePeer(chunk, map[string]interface{}{"download_bytes": 0, "upload_bytes": 0}); err != nil {
|
||
logrus.WithError(err).Error("cron: сброс трафика части пиров не выполнен")
|
||
continue
|
||
}
|
||
}
|
||
}
|
||
|
||
// saveAccountTraffic переносит накопленный Hysteria трафик в базу.
|
||
//
|
||
// Чтение ДЕСТРУКТИВНОЕ: `?clear=1` обнуляет счётчики Hysteria сразу после
|
||
// того, как ответ отправлен (официальный контракт Traffic Stats API). Значит
|
||
// каждая дельта существует ровно в одном экземпляре, и потерянная здесь
|
||
// потеряна навсегда.
|
||
//
|
||
// Полностью закрыть это окно можно только сменой модели учёта — недеструктивным
|
||
// `GET /traffic` с долговременными checkpoint'ами верхних счётчиков и
|
||
// вычислением дельты на стороне админки. Это отдельная подсистема с обработкой
|
||
// перезапуска и сброса счётчиков Hysteria, и в текущем проходе она намеренно
|
||
// не вводится: квота здесь — операционная граница доступа, а не биллинговый
|
||
// учёт с финансово значимым каждым байтом.
|
||
//
|
||
// Чего это НЕ оправдывает — молчания. Раньше отказ записи внутри цикла делал
|
||
// `continue`, и дельта конкретного пира исчезала, не оставив следа в исходе
|
||
// джобы. Теперь каждая потеря считается и попадает в возвращаемую ошибку.
|
||
func saveAccountTraffic(apiPort int64, trafficStatsSecret string) error {
|
||
users, err := proxy.NewHysteria2Api(apiPort).ListUsers(true, trafficStatsSecret)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if len(users) == 0 {
|
||
return nil
|
||
}
|
||
|
||
nowMs := time.Now().UnixMilli()
|
||
hourStart := nowMs - (nowMs % int64(time.Hour/time.Millisecond))
|
||
|
||
lost := 0
|
||
for key, traffic := range users {
|
||
rxBytes := traffic.Rx
|
||
txBytes := traffic.Tx
|
||
if rxBytes == 0 && txBytes == 0 {
|
||
continue
|
||
}
|
||
|
||
peer, peerErr := dao.GetPeer("auth_id = ?", key)
|
||
if peerErr != nil {
|
||
// Пир, которого админка не знает: удалён между сбором и записью
|
||
// либо создан в обход панели. Дельта уже обнулена в Hysteria и
|
||
// приписывать её некому.
|
||
logrus.WithError(peerErr).
|
||
WithField("authId", key).
|
||
Warn("cron: трафик получен для неизвестного пира и не записан")
|
||
lost++
|
||
continue
|
||
}
|
||
if peer.Id == nil {
|
||
logrus.WithField("authId", key).
|
||
Error("cron: строка пира без идентификатора; трафик не записан")
|
||
lost++
|
||
continue
|
||
}
|
||
|
||
authId := key
|
||
if peer.AuthId != nil && *peer.AuthId != "" {
|
||
authId = *peer.AuthId
|
||
}
|
||
|
||
sample := entity.TrafficSample{
|
||
PeerId: peer.Id,
|
||
AuthId: &authId,
|
||
RxBytes: &rxBytes,
|
||
TxBytes: &txBytes,
|
||
SampledAt: &nowMs,
|
||
}
|
||
if err := dao.SaveTrafficSample(sample); err != nil {
|
||
logrus.WithError(err).
|
||
WithField("peerId", *peer.Id).
|
||
Error("cron: не удалось сохранить отсчёт трафика")
|
||
// Отсчёт — история для графиков; счётчики пира важнее, и попытка
|
||
// их обновить продолжается.
|
||
}
|
||
|
||
if err := dao.UpdatePeer([]int64{*peer.Id}, map[string]interface{}{
|
||
"download_bytes": gorm.Expr("download_bytes + ?", rxBytes),
|
||
"upload_bytes": gorm.Expr("upload_bytes + ?", txBytes),
|
||
}); err != nil {
|
||
logrus.WithError(err).
|
||
WithField("peerId", *peer.Id).
|
||
Error("cron: счётчики пира не обновлены; дельта Hysteria уже обнулена и потеряна")
|
||
lost++
|
||
continue
|
||
}
|
||
|
||
_ = dao.UpsertTrafficAggregateHourly(*peer.Id, hourStart, rxBytes, txBytes)
|
||
}
|
||
|
||
if lost > 0 {
|
||
return &trafficLossError{lost: lost}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// trafficLossError сообщает, сколько дельт не удалось записать.
|
||
//
|
||
// Отдельный тип, а не fmt.Errorf, потому что количество здесь — величина, а не
|
||
// украшение фразы: чтение `?clear=1` деструктивно, поэтому «потеряно 1 из 200»
|
||
// и «потеряно 200 из 200» — разные события, и различать их должен уметь не
|
||
// только человек, читающий журнал.
|
||
type trafficLossError struct{ lost int }
|
||
|
||
func (e *trafficLossError) Error() string {
|
||
return fmt.Sprintf(
|
||
"дельт трафика не записано и потеряно безвозвратно: %d",
|
||
e.lost,
|
||
)
|
||
}
|
||
|
||
// enforcePeerAccess приводит ЖИВЫЕ СЕССИИ в соответствие с сохранённым
|
||
// состоянием.
|
||
//
|
||
// Политика берётся из peerAccessDenied — той же функции, по которой пира
|
||
// пускает или не пускает авторизация. Собственного SQL-условия здесь больше
|
||
// нет, и это главное свойство: пока правило было записано в двух местах, оно
|
||
// расходилось на границах quota, expiry и ban, и исчерпавший квоту пир не
|
||
// пускался заново, но и не отключался никогда.
|
||
//
|
||
// Обход идёт по КАЖДОМУ authID, который Hysteria считает живым, а не по
|
||
// найденным в базе пирам. Прежняя реализация читала
|
||
//
|
||
// peers, err := dao.ListPeer("auth_id in ?", chunk)
|
||
// for _, peer := range peers { ... }
|
||
//
|
||
// и потому не видела сессий, которым в базе больше ничего не соответствует.
|
||
// Это не теоретический случай: `auth_id` перезаписывает импорт, а строку
|
||
// целиком убирает удаление. Обе операции рвут старую сессию сами, но их второй
|
||
// шаг может не удаться — и тогда единственным местом, где о ней ещё известно,
|
||
// остаётся сам `/online`. Пропуская незнакомый идентификатор молча, cron
|
||
// оставлял такую сессию жить неограниченно долго. Подробности — в
|
||
// peer_session-разделе peer_access.go.
|
||
//
|
||
// Отказ базы НЕ приводит к разрыву. «Пира нет» и «прочитать не удалось» —
|
||
// разные ответы, и второй не даёт права рвать ничьи сессии: недоступная SQLite
|
||
// иначе означала бы отключение всех подключённых пиров сразу. Ошибка чтения
|
||
// прекращает цикл до единого обращения к `/kick`.
|
||
//
|
||
// Обход последовательный. Прежняя реализация раскладывала online-пиров на
|
||
// чанки по 10 и запускала по горутине на чанк с sync.WaitGroup внутри уже
|
||
// отсоединённой горутины. Параллельность здесь не нужна: обращений к базе
|
||
// столько же, а `/kick` всё равно один на весь набор.
|
||
func enforcePeerAccess(apiPort int64, trafficStatsSecret string) error {
|
||
online, err := proxy.NewHysteria2Api(apiPort).OnlineUsers(trafficStatsSecret)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// Учёт выданных разрешений чистится по фактической картине подключений, и
|
||
// это единственное место продукта, где она известна целиком. Делается это
|
||
// до любых решений: уборка ни на что не влияет и ничего не рвёт.
|
||
sweepDeviceAdmissions(online, time.Now())
|
||
|
||
if len(online) == 0 {
|
||
return nil
|
||
}
|
||
|
||
authIDs := make([]string, 0, len(online))
|
||
for authID := range online {
|
||
// Пустой ключ ничему не соответствует: рвать по нему нечего, и в
|
||
// dedup disconnectAuthIDs он всё равно не попал бы.
|
||
if authID == "" {
|
||
continue
|
||
}
|
||
authIDs = append(authIDs, authID)
|
||
}
|
||
if len(authIDs) == 0 {
|
||
return nil
|
||
}
|
||
// Порядок ключей карты в Go случаен; сортировка делает и обращение к
|
||
// `/kick`, и журнал воспроизводимыми.
|
||
sort.Strings(authIDs)
|
||
|
||
known := make(map[string]entity.Peer, len(authIDs))
|
||
for _, chunk := range util.SplitArr(authIDs, 100) {
|
||
peers, err := dao.ListPeer("auth_id in ?", chunk)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
for _, peer := range peers {
|
||
// Строка без authId Hysteria не знает. Раньше здесь стояло
|
||
// `*item.AuthId` без проверки — паника на повреждённой строке
|
||
// внутри отсоединённой горутины.
|
||
authID := authIDOf(peer)
|
||
if authID == "" {
|
||
continue
|
||
}
|
||
known[authID] = peer
|
||
}
|
||
}
|
||
|
||
now := time.Now().UnixMilli()
|
||
kick := make([]string, 0, len(authIDs))
|
||
for _, authID := range authIDs {
|
||
peer, found := known[authID]
|
||
if !found {
|
||
// Сессия, которой в базе больше ничего не соответствует: пир удалён
|
||
// либо его идентификатор заменён импортом, а разрыв в тот момент не
|
||
// удался. Восстановить такое состояние переподключением нельзя —
|
||
// авторизация нового значения не знает, — поэтому единственный
|
||
// правильный исход тот же, что и у первой попытки.
|
||
logrus.WithField("authId", authID).
|
||
Warn("cron: живая сессия без пира в базе; сессия завершается")
|
||
kick = append(kick, authID)
|
||
continue
|
||
}
|
||
if peerSessionNeedsReconcile(peer, online[authID], now) {
|
||
kick = append(kick, authID)
|
||
}
|
||
}
|
||
|
||
// Пустой набор до `/kick` не доходит: раньше запрос с пустым массивом в
|
||
// теле уезжал в Hysteria каждые 30 секунд.
|
||
return disconnectAuthIDs(kick)
|
||
}
|