Files
founder 8dcb50a07c fix(admin): дать отзыву доступа вторую попытку, а лимиту устройств — порядок снимков
Предыдущий проход сделал правильным порядок «сначала долговременная запись,
потом разрыв сессии» и правильно запретил откат при неудаче разрыва. Способа
прийти к согласованному состоянию ПОТОМ он не дал: у двух операций повтор не
работал вовсе.

Импорт, заменивший 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
2026-09-02 07:15:43 +05:00

351 lines
18 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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)
}