6d1686b2be
Второй разбор того же слоя, уже по состоянию после 162759c. Тема: границы между
частями access-control. Прошлый проход починил одну операцию отзыва доступа и
оставил остальные; правило доступа при этом продолжало существовать в двух
экземплярах. Проведены три границы: состояние пира -> решение о доступе,
сохранённое изменение -> живая сессия, планировщик -> принадлежащая ему работа.
Правило доступа. Оно было записано двумя разными SQL-условиями: одним в выборке
Hysteria2Auth, другим в выборке cron. Второе не является отрицанием первого, и
расхождение приходилось ровно на границы — quota=0, usage=quota, now=expiresAt,
now=bannedUntil: авторизация отказывала, cron сессию не рвал. Условие cron
требовало СТРОГОГО превышения квоты, а счётчики растут порциями по ответу
Traffic Stats API, поэтому точное равенство — обычный исход очередного сбора.
Пир с исчерпанной квотой не пускался заново, но его живая сессия не разрывалась
никогда. Политика вынесена в peerAccessDenied; авторизация ищет пира только по
secret_digest, cron применяет ту же функцию. quota=-1 — единственный безлимит,
quota=0 — ноль байтов, bannedUntil=now — блокировка уже закончилась. Строка без
решающего поля трактуется как повреждённая и ведёт к отказу.
Операции, оставлявшие живую сессию. DeletePeer состоял из одного dao.DeletePeer:
строка исчезала вместе с auth_id, то есть вместе с единственным, чем эту сессию
можно было завершить, — состояние становилось невосстановимым. Разрыв при
изменении выполнялся только при disabled=1, поэтому мимо проходили смена
секрета, урезание квоты ниже израсходованного, перенос срока в прошлое и
снижение maxDevices. Импорт переписывает auth_id, секрет, квоту, срок и disabled
целиком и не трогал сессий вовсе. Все операции идут теперь через один
reconcileLiveSessions, а он — через disconnectAuthIDs, единственный вход к /kick:
он принимает готовые идентификаторы, дедуплицирует их, разбивает на части и не
обращается к базе. Импорт собирает старые auth_id ВНУТРИ транзакции (после
commit их в базе уже нет) и рвёт ПОСЛЕ commit (до него клиент успел бы
переподключиться к ещё не изменённому пиру). Правило асимметрично намеренно:
ограничение применяется немедленно, послабление — нет.
Цикл учёта. CronHandleAccount запускала горутину, которая запускала ещё две, —
для планировщика джоба заканчивалась почти мгновенно, поэтому StopCron не ждал
настоящей работы: releaseResource закрывал SQLite, а горутины продолжали в неё
писать. Параллельность обеих половин означала ещё и то, что enforcement читал
счётчики до записи снятой дельты. Джоба стала синхронной, под одним мьютексом на
весь цикл, порядок строгий. Закрыты три nil-разыменования — trafficSecretConfig,
item.AuthId и item.Id, — каждое из которых роняло процесс целиком вместе с
обработчиком machine-auth. Гейт Hysteria2IsRunning убран: util.Exec не отличает
«служба неактивна» от «спросить не удалось», и сломанный systemctl при живой
Hysteria молча отключал и учёт, и enforcement. Потеря дельты при отказе SQLite
больше не молчит: чтение /traffic?clear=1 деструктивно, и каждая потеря
считается. Checkpoint accounting в 1.0.0 намеренно не вводится — квота здесь
операционный предел доступа, а не учёт с финансово значимым каждым байтом.
Лимит устройств. Между чтением /online и ответом allow место ничем не
удерживалось: при online=max-1 два одновременных запроса получали разрешение
оба. Мьютекс вокруг /online этого не чинит — ответив allow, админка не создаёт
подключение, и следующий запрос продолжает видеть прежнее число. Появился
process-local учёт выданных, но ещё не проявившихся разрешений: решение по сумме
«подключено плюс зарезервировано», рост online снимает соответствующее их число,
протухшие снимаются по внутреннему TTL. Сеть опрашивается вне блокировки.
Гейты. Проверка «авторизация не возвращает успех из ветки ошибки» была записана
регуляркой err != nil \{[\s\S]*?return \*peer\.Id, а ленивый [\s\S]*? свободно
пересекает границы блоков: она даёт совпадение на коде из HEAD, то есть гейт
нельзя было удовлетворить, не сломав продукт. Тело ветки теперь выделяется по
балансу фигурных скобок, и логика проверена в обе стороны. go test -race стал
обязательным шагом сборки: состояние трекера разрешений и мьютекс цикла учёта
принадлежат процессу, и их корректность не наблюдаема ни в go test, ни в go vet;
пропуск при недоступном компиляторе не предусмотрен.
Панель. importPeerApi не объявлял skipErrorToast, а handleImport не имел ни try,
ни catch: после появления частичного результата отказ уходил бы необработанным
отклонением промиса, список не обновлялся бы при уже изменённой базе, а общий
перехватчик показал бы предупреждение красной ошибкой. Формулировка
peer_disconnect_failed во всех трёх местах сделана operation-neutral: через этот
код отчитываются восемь операций, а для удалённого пира прежняя фраза «новые
подключения пира запрещены» просто бессмысленна.
302 lines
15 KiB
Go
302 lines
15 KiB
Go
package service
|
||
|
||
import (
|
||
"fmt"
|
||
"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, и исчерпавший квоту пир не
|
||
// пускался заново, но и не отключался никогда.
|
||
//
|
||
// Обход последовательный. Прежняя реализация раскладывала online-пиров на
|
||
// чанки по 10 и запускала по горутине на чанк с sync.WaitGroup внутри уже
|
||
// отсоединённой горутины. Параллельность здесь не нужна: обращений к базе
|
||
// столько же, а `/kick` всё равно один на весь набор.
|
||
func enforcePeerAccess(apiPort int64, trafficStatsSecret string) error {
|
||
online, err := proxy.NewHysteria2Api(apiPort).OnlineUsers(trafficStatsSecret)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if len(online) == 0 {
|
||
return nil
|
||
}
|
||
|
||
authIDs := make([]string, 0, len(online))
|
||
for authID := range online {
|
||
if authID == "" {
|
||
continue
|
||
}
|
||
authIDs = append(authIDs, authID)
|
||
}
|
||
if len(authIDs) == 0 {
|
||
return nil
|
||
}
|
||
|
||
now := time.Now().UnixMilli()
|
||
kick := make([]string, 0, 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
|
||
}
|
||
if peerAccessDenied(peer, now) {
|
||
kick = append(kick, authID)
|
||
}
|
||
}
|
||
}
|
||
|
||
// Пустой набор до `/kick` не доходит: раньше запрос с пустым массивом в
|
||
// теле уезжал в Hysteria каждые 30 секунд.
|
||
return disconnectAuthIDs(kick)
|
||
}
|