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
This commit is contained in:
2026-09-02 07:15:43 +05:00
parent 6d1686b2be
commit 8dcb50a07c
16 changed files with 1511 additions and 104 deletions
+58 -9
View File
@@ -2,6 +2,7 @@ package service
import (
"fmt"
"sort"
"sync"
"time"
@@ -242,7 +243,8 @@ func (e *trafficLossError) Error() string {
)
}
// enforcePeerAccess завершает сессии пиров, которым доступ уже закрыт.
// enforcePeerAccess приводит ЖИВЫЕ СЕССИИ в соответствие с сохранённым
// состоянием.
//
// Политика берётся из peerAccessDenied — той же функции, по которой пира
// пускает или не пускает авторизация. Собственного SQL-условия здесь больше
@@ -250,6 +252,25 @@ func (e *trafficLossError) Error() string {
// расходилось на границах 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 внутри уже
// отсоединённой горутины. Параллельность здесь не нужна: обращений к базе
@@ -259,12 +280,20 @@ func enforcePeerAccess(apiPort int64, trafficStatsSecret string) error {
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
}
@@ -273,25 +302,45 @@ func enforcePeerAccess(apiPort int64, trafficStatsSecret string) error {
if len(authIDs) == 0 {
return nil
}
// Порядок ключей карты в Go случаен; сортировка делает и обращение к
// `/kick`, и журнал воспроизводимыми.
sort.Strings(authIDs)
now := time.Now().UnixMilli()
kick := make([]string, 0, len(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 Hysteria не знает. Раньше здесь стояло
// `*item.AuthId` без проверки — паника на повреждённой строке
// внутри отсоединённой горутины.
authID := authIDOf(peer)
if authID == "" {
continue
}
if peerAccessDenied(peer, now) {
kick = append(kick, authID)
}
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)
}
}
+192
View File
@@ -13,6 +13,7 @@ import (
"hy2xs-admin/dao"
"hy2xs-admin/model/bo"
"hy2xs-admin/model/constant"
"hy2xs-admin/model/dto"
)
// Цикл учёта проверяется против НАСТОЯЩЕГО Traffic Stats API.
@@ -289,6 +290,197 @@ func TestCronIgnoresOfflinePeers(t *testing.T) {
}
}
// --- Сверка живых сессий -----------------------------------------------------
// Живая сессия, которой в базе больше ничего не соответствует, завершается.
//
// Прежний обход шёл по НАЙДЕННЫМ пирам, поэтому authID, которого нет в базе,
// молча выпадал: `dao.ListPeer("auth_id in ?")` просто не возвращала строку.
// Такое состояние возникает после неудавшегося второго шага удаления или
// импорта, заменившего `auth_id`, и восстановить его переподключением нельзя —
// авторизация нового значения не знает. Сессия жила неограниченно долго.
func TestCronKicksSessionWithoutPeerRow(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
online: map[string]int64{"ghost-auth-id": 1},
})
seedPeer(t, "alpha1", "alpha-auth-id")
CronHandleAccount()
if got := stub.kicked(); len(got) != 1 || got[0] != "ghost-auth-id" {
t.Fatalf("сессия без пира в базе не завершена: %v", got)
}
}
// Превышение лимита устройств — свойство живых сессий, а не хранимого
// состояния пира, поэтому peerAccessDenied его не видит и видеть не должен.
// Без этой проверки неудавшийся разрыв при снижении `maxDevices` оставался бы
// навсегда: повторное сохранение формы сравнивает `1 < 1` и разрыва не делает.
func TestCronKicksWhenOnlineExceedsMaxDevices(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
online: map[string]int64{"alpha-auth-id": 3},
})
id := seedPeer(t, "alpha1", "alpha-auth-id")
peerUsage(t, id, map[string]interface{}{"max_devices": int64(1)})
CronHandleAccount()
if got := stub.kicked(); len(got) != 1 || got[0] != "alpha-auth-id" {
t.Fatalf("превышение лимита устройств не отключено: %v", got)
}
}
// Граница: устройств ровно столько, сколько разрешено, — рвать нечего.
func TestCronDoesNotKickAtExactDeviceLimit(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
online: map[string]int64{"alpha-auth-id": 3},
})
// seedPeer создаёт пира с maxDevices = 3.
seedPeer(t, "alpha1", "alpha-auth-id")
CronHandleAccount()
if got := stub.kicked(); len(got) != 0 {
t.Fatalf("пир на границе лимита отключён: %v", got)
}
}
// Повреждённая граница — не «безлимит». На пути авторизации такая строка ведёт
// к отказу, и живая сессия обязана следовать тому же правилу.
func TestCronKicksPeerWithUnusableMaxDevices(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
online: map[string]int64{"alpha-auth-id": 1},
})
id := seedPeer(t, "alpha1", "alpha-auth-id")
peerUsage(t, id, map[string]interface{}{"max_devices": int64(0)})
CronHandleAccount()
if got := stub.kicked(); len(got) != 1 {
t.Fatalf("пир с непригодным лимитом устройств не отключён: %v", got)
}
}
// Отказ базы НЕ является основанием рвать сессии.
//
// «Пира нет» и «прочитать не удалось» — разные ответы, и решение «сессии
// неизвестны, значит лишние» на втором из них отключило бы всех подключённых
// пиров сразу при недоступной SQLite. Проверка существует именно потому, что
// правило «неизвестный authID -> kick» делает это различие решающим.
func TestCronSendsNoKickWhenPeerLookupFails(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
online: map[string]int64{"alpha-auth-id": 1},
})
seedPeer(t, "alpha1", "alpha-auth-id")
apiPort, err := GetHysteria2ApiPort()
if err != nil {
t.Fatalf("порт Traffic Stats API: %v", err)
}
// Порт и секрет читаются из конфига Hysteria, поэтому база после этого уже
// не нужна ни для чего, кроме самой выборки пиров.
if err := dao.CloseSqliteDB(); err != nil {
t.Fatalf("не удалось закрыть базу: %v", err)
}
if err := enforcePeerAccess(apiPort, testTrafficStatsSecret); err == nil {
t.Fatal("отказ базы не сообщён вызывающему")
}
if got := stub.kicked(); len(got) != 0 {
t.Fatalf("отказ базы привёл к разрыву сессий: %v", got)
}
}
// Учёт выданных разрешений чистится по фактической картине подключений.
func TestCronSweepsAdmissionsOfOfflinePeers(t *testing.T) {
newTestDB(t)
startAccountStats(t, &accountStatsStub{online: map[string]int64{}})
seedPeer(t, "alpha1", "alpha-auth-id")
// Разрешение выдано давно и уже протухло, подключения так и не случилось.
if !reserveDeviceSlot("alpha-auth-id", 0, 3, time.Now().Add(-2*pendingAdmissionTTL)) {
t.Fatal("подготовка учёта: разрешение отклонено")
}
if admissionEntries() != 1 {
t.Fatal("подготовка учёта: запись не создана")
}
CronHandleAccount()
if got := admissionEntries(); got != 0 {
t.Fatalf("учёт не убран: записей %d", got)
}
}
// --- Сходимость после неудавшегося разрыва -----------------------------------
// Импорт заменил `auth_id`, а разрыв старой сессии не удался. Повторить его
// операцией импорта невозможно: в базе уже новое значение, и повтор того же
// файла разорвал бы именно его. Сходимость обеспечивает cron.
func TestCronReconcilesSessionAfterFailedImportKick(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
kickStatus: http.StatusInternalServerError,
online: map[string]int64{"old-auth-id": 1},
})
seedPeer(t, "keeper", "old-auth-id")
requireDisconnectError(t, UpsertPeerExport([]bo.PeerExport{importItem("keeper", "new-auth-id")}))
after := snapshotPeers(t)["keeper"]
if after.AuthId == nil || *after.AuthId != "new-auth-id" {
t.Fatalf("импорт не применён: authId=%v", after.AuthId)
}
stub.mu.Lock()
stub.kickStatus = 0
stub.kickedKeys = nil
stub.mu.Unlock()
CronHandleAccount()
if got := stub.kicked(); len(got) != 1 || got[0] != "old-auth-id" {
t.Fatalf("старая сессия не завершена следующим циклом учёта: %v", got)
}
}
// Лимит устройств снижен, разрыв не удался, оператор повторяет сохранение
// формы — и получает успех без разрыва, потому что новое значение уже в базе.
// Единственный механизм схождения здесь — cron.
func TestCronReconcilesSessionAfterFailedMaxDevicesReduction(t *testing.T) {
newTestDB(t)
stub := startAccountStats(t, &accountStatsStub{
kickStatus: http.StatusInternalServerError,
online: map[string]int64{"alpha-auth-id": 3},
})
// seedPeer создаёт пира с maxDevices = 3.
id := seedPeer(t, "alpha1", "alpha-auth-id")
requireDisconnectError(t, UpdatePeer(id, dto.PeerUpdateDto{MaxDevices: int64Ptr(1)}))
// Повтор формы: значение то же самое, разрыва не будет — и это правильно,
// иначе каждое сохранение любой правки рвало бы сессии.
if err := UpdatePeer(id, dto.PeerUpdateDto{MaxDevices: int64Ptr(1)}); err != nil {
t.Fatalf("повторное сохранение формы отказало: %v", err)
}
stub.mu.Lock()
stub.kickStatus = 0
stub.kickedKeys = nil
stub.mu.Unlock()
CronHandleAccount()
if got := stub.kicked(); len(got) != 1 || got[0] != "alpha-auth-id" {
t.Fatalf("превышение лимита не устранено следующим циклом учёта: %v", got)
}
}
// --- Порядок и устройство цикла ----------------------------------------------
// Enforcement принимает решение по счётчикам, поэтому счётчики обязаны быть
+16
View File
@@ -114,6 +114,22 @@ func Hysteria2Auth(conPass string) (int64, string, error) {
// недоступность здесь — не штатное состояние, а аномалия, и пускать
// подключения без единственной проверки, которая ещё не выполнена, значит
// молча снять лимит со всех пиров сразу.
// Чтение `/online` и резервация места — ОДНА последовательность, и она
// выполняется под замком этого пира.
//
// Без замка снимки приходили в резервацию в произвольном порядке, и
// устаревший откатывал учёт назад: разрешение, уже признанное проявившимся,
// возвращалось в «свободное место». Подробный разбор — в начале
// peer_admission.go.
//
// Замок берётся именно здесь, а не раньше: до этой точки известен только
// секрет, а сериализовать нужно подключения ОДНОГО пира, то есть замок
// невозможно взять, пока не прочитан его authId. Всё, что выше, — работа с
// базой и политикой доступа, и разным пирам она не мешает.
unlockAdmission := lockPeerAdmission(*peer.AuthId)
defer unlockAdmission()
onlineUsers, err := hysteria2Online()
if err != nil {
logrus.WithError(err).
+70
View File
@@ -100,3 +100,73 @@ func peerAccessDenied(peer entity.Peer, now int64) bool {
return false
}
// Живая сессия сверяется с сохранённым состоянием ПОВТОРЯЕМО.
//
// Что было. Приведение сессий к состоянию базы выполнялось ровно один раз — в
// той же операции, которая это состояние записала. Порядок «сначала запись,
// потом `/kick`» правильный, и откат при неудаче разрыва делать нельзя: часть
// операции, закрывающая доступ, уже достигнута. Но второй попытки после
// неудачи не существовало вовсе, и два состояния оставались навсегда.
//
// Первое — замена `auth_id` импортом:
//
// импорт old-auth -> new-auth, COMMIT прошёл
// /kick old-auth -> 500
// оператор повторяет тот же импорт
// applyPeerImportEntry читает из базы уже new-auth и рвёт ЕГО
//
// Старый идентификатор после первой же неудачи не хранился нигде, а cron его
// пропускал: `dao.ListPeer("auth_id in ?")` просто не возвращала строку, и
// authID, которого нет в базе, молча выпадал из обхода. Живая QUIC-сессия
// удалённого или переподписанного пира продолжалась сколько угодно долго.
//
// Второе — снижение `maxDevices`:
//
// 5 -> 1, запись прошла, /kick -> 500
// оператор повторяет сохранение формы
// updateRequiresReconcile сравнивает 1 < 1 -> false, разрыва нет
//
// Для `disabled` повторяемость сделана специально (условие смотрит на
// ЗАПРОШЕННОЕ состояние, а не на переход), для квоты и срока её обеспечивает
// cron через peerAccessDenied. Лимит устройств в политику доступа не входит и
// входить не должен — это свойство не пира, а его сессий, — поэтому здесь у
// него не было ни одного механизма схождения.
//
// Оба состояния закрывает один и тот же приём: cron сверяет не «кого из
// известных пиров пора отключить», а КАЖДЫЙ authID, который Hysteria считает
// живым. Отдельная таблица retry, очередь отложенных операций и хранимый
// «список того, что не удалось разорвать» для этого не нужны: `/online` и есть
// список живых сессий, и сверять его достаточно.
// peerSessionNeedsReconcile отвечает, устарела ли живая сессия пира.
//
// Вторым экземпляром политики доступа не является: disabled, quota, expiry и
// ban остаются целиком за peerAccessDenied, и эта функция их не повторяет, а
// вызывает. Своего здесь ровно одно — инвариант живых сессий, которого в
// хранимом состоянии пира нет: число подключённых устройств.
//
// доступ закрыт -> сессия устарела
// maxDevices непригоден -> сессия устарела
// устройств больше, чем разрешено -> сессии устарели
//
// Непригодный `maxDevices` ведёт к разрыву по той же причине, по которой он
// ведёт к отказу в авторизации: повреждённая граница — это не «безлимит».
//
// Число устройств берётся из `/online`, который по официальному контракту
// Traffic Stats API возвращает количество экземпляров клиента Hysteria
// («устройства»), а не число proxy-потоков. То есть сравнение с `maxDevices`
// здесь опирается на upstream-контракт, а не на предположение.
//
// Выбирать «лишнее устройство» не нужно и невозможно: `/kick` оперирует
// идентификатором клиента. После разрыва клиенты переподключаются, и
// admission пропустит ровно столько, сколько разрешено теперь.
func peerSessionNeedsReconcile(peer entity.Peer, onlineDevices int64, now int64) bool {
if peerAccessDenied(peer, now) {
return true
}
if peer.MaxDevices == nil || *peer.MaxDevices < 1 {
return true
}
return onlineDevices > *peer.MaxDevices
}
+151 -1
View File
@@ -42,6 +42,37 @@ import (
// «соединение установлено / не установлено» математически точной системы
// резервирования не построить. Он закрывает конкретный и реальный TOCTOU —
// параллельные HTTP-auth одного процесса, — и делает это fail-closed.
//
// Второй TOCTOU: ПЕРЕУПОРЯДОЧИВАНИЕ снимков `/online`.
//
// Учёта выданных разрешений самого по себе оказалось недостаточно, и это
// отдельный дефект, а не оттенок первого. Сетевой запрос выполнялся вне
// мьютекса, поэтому снимки приходили в критическую секцию в произвольном
// порядке — более старый мог обогнать более новый:
//
// 1. A получает разрешение при `/online = 0`; pending = [A], lastOnline = 0
// 2. B читает `/online = 0` и задерживается на обратном пути
// 3. A действительно подключается, Hysteria показывает `/online = 1`
// 4. C читает `/online = 1` и входит в резервацию ПЕРВЫМ:
// разрешение A признано проявившимся, lastOnline = 1, C получает отказ
// 5. B входит со своим устаревшим `online = 0`
// 6. `online > lastOnline` ложно, после чего lastOnline откатывается в 0
// 7. `0 + pending(0) < 1` -> B получает ALLOW
//
// При `maxDevices = 1` подключений становится два. Это НЕ data race: вся
// работа с памятью защищена мьютексом, поэтому детектор гонок здесь молчит
// принципиально, и поймать дефект может только семантическая проверка.
//
// Лечится это не глобальным мьютексом вокруг сети — он сериализовал бы
// подключения всех пиров через один HTTP-обмен, — а замком на ОДИН authID:
// см. lockPeerAdmission. Конкурируют между собой только авторизации одного и
// того же пира, а их упорядоченность и есть требуемое свойство: снимок,
// прочитанный под замком, не может оказаться старше уже обработанного.
//
// Оба механизма нужны одновременно и закрывают разные половины:
//
// замок по authID — снимки не переупорядочиваются;
// учёт разрешений — снимок не успевает измениться к следующему запросу.
// pendingAdmissionTTL — срок жизни выданного разрешения, которое ещё не
// проявилось в `/online`.
@@ -74,6 +105,67 @@ var deviceAdmissions = struct {
byAuthID map[string]*admissionState
}{byAuthID: map[string]*admissionState{}}
// admissionGate — замок одного authID вместе со счётчиком тех, кому он сейчас
// нужен.
//
// Счётчик существует ради удаления записи. Без него карта замков росла бы по
// одной записи на каждый когда-либо авторизовавшийся authId и не уменьшалась
// бы никогда — то есть та же утечка, от которой в учёте разрешений защищает
// forgetIfIdle, только этажом выше.
type admissionGate struct {
mu sync.Mutex
// waiting — сколько вызывающих держат замок или ждут его. Пока значение
// больше нуля, запись обязана оставаться в карте: удалив её, второй
// вызывающий создал бы НОВЫЙ замок и разошёлся бы с первым.
waiting int
}
var admissionGates = struct {
sync.Mutex
byAuthID map[string]*admissionGate
}{byAuthID: map[string]*admissionGate{}}
// lockPeerAdmission сериализует последовательность «прочитать `/online` ->
// занять место» для ОДНОГО пира и возвращает функцию освобождения.
//
// Замок именно по authID, а не один на процесс, и это существенно. Внутри него
// выполняется сетевой запрос к Traffic Stats API, поэтому общий замок означал
// бы, что все подключения всех пиров выстраиваются в очередь за одним
// HTTP-обменом. Здесь же конкурируют только авторизации одного пира — то есть
// ровно те, для которых порядок и решается.
//
// Время удержания ограничено сверху таймаутом самого обращения к Traffic Stats
// API (proxy.Hysteria2Api ставит контекст на 3 секунды), поэтому «застрявшая»
// Hysteria не превращает замок в бессрочный.
//
// Карта замков и карта учёта разрешений намеренно раздельны: первая описывает,
// кто сейчас проходит авторизацию, вторая — что уже выдано. Совмещение их в
// одной структуре означало бы удерживать мьютекс учёта на время сетевого
// запроса.
func lockPeerAdmission(authID string) func() {
admissionGates.Lock()
gate := admissionGates.byAuthID[authID]
if gate == nil {
gate = &admissionGate{}
admissionGates.byAuthID[authID] = gate
}
gate.waiting++
admissionGates.Unlock()
gate.mu.Lock()
return func() {
gate.mu.Unlock()
admissionGates.Lock()
defer admissionGates.Unlock()
gate.waiting--
if gate.waiting == 0 {
delete(admissionGates.byAuthID, authID)
}
}
}
// reserveDeviceSlot решает, есть ли для нового подключения свободное место, и
// занимает его.
//
@@ -83,6 +175,12 @@ var deviceAdmissions = struct {
// передаёт сюда уже полученное число. Под блокировкой остаются только
// несколько операций с map: держать её на время HTTP-обмена значило бы
// сериализовать все подключения всех пиров через один сетевой запрос.
//
// Упорядоченность снимков обеспечивает не этот мьютекс, а замок по authID:
// вызывающий обязан удерживать lockPeerAdmission от чтения `/online` и до
// возврата отсюда. Без него сюда попадал бы снимок старше уже обработанного, и
// строка `state.lastOnline = online` откатывала бы учёт назад — см. описание
// второго TOCTOU в начале файла.
func reserveDeviceSlot(authID string, online int64, maxDevices int64, now time.Time) bool {
deviceAdmissions.Lock()
defer deviceAdmissions.Unlock()
@@ -159,10 +257,62 @@ func (s *admissionState) forgetIfIdle(authID string) {
}
}
// sweepDeviceAdmissions убирает записи о пирах, за которыми ничего не числится.
//
// Зачем это нужно отдельно от forgetIfIdle. Тот срабатывает ТОЛЬКО на ветке
// отказа: после успешной выдачи разрешения запись остаётся с непустым pending,
// а когда разрешение протухает, снять запись уже некому — следующего обращения
// к этому authId может не быть никогда. Так в карте оставались пиры, удалённые
// из панели, и старые authId, переписанные импортом: за время жизни процесса
// она только росла.
//
// Решение принимается по ФАКТИЧЕСКОЙ картине подключений, а не по хранимому
// lastOnline: последний обновляется только на пути авторизации, поэтому у
// отключившегося пира он остаётся прежним сколь угодно долго.
//
// Удаление записи, у которой нет ни одного действующего разрешения и нет
// подключений, не меняет ни одного будущего решения. Следующая резервация
// начнёт с чистой записи и придёт к тому же ответу: dropMaterialized на пустом
// списке — no-op, а `state.lastOnline` в любом случае перезаписывается
// пришедшим значением до сравнения с лимитом.
//
// Замок authID здесь не берётся намеренно: удерживать его на всём обходе карты
// значило бы останавливать авторизацию каждые 30 секунд. Гонка с параллельной
// авторизацией безопасна — она либо уже добавила разрешение (тогда pending не
// пуст и запись остаётся), либо ещё не дошла до учёта (тогда она создаст
// запись заново, и это ровно та же чистая запись).
func sweepDeviceAdmissions(online map[string]int64, now time.Time) {
deviceAdmissions.Lock()
defer deviceAdmissions.Unlock()
for authID, state := range deviceAdmissions.byAuthID {
state.dropExpired(now)
if len(state.pending) > 0 {
continue
}
if online[authID] > 0 {
continue
}
delete(deviceAdmissions.byAuthID, authID)
}
}
// resetDeviceAdmissions очищает учёт. Существует ради тестов: состояние здесь
// принадлежит процессу, и без сброса тесты видели бы резервации друг друга.
func resetDeviceAdmissions() {
deviceAdmissions.Lock()
defer deviceAdmissions.Unlock()
deviceAdmissions.byAuthID = map[string]*admissionState{}
deviceAdmissions.Unlock()
// Замки сбрасывать НЕЛЬЗЯ: удерживаемый кем-то замок, потерянный из карты,
// перестал бы исключать второго вызывающего. Вместо сброса тесты проверяют,
// что после завершения работы карта пуста сама.
}
// heldPeerAdmissionGates — число живых замков. Существует ради тестов: утечка
// здесь выглядит именно как незакрытая запись в карте.
func heldPeerAdmissionGates() int {
admissionGates.Lock()
defer admissionGates.Unlock()
return len(admissionGates.byAuthID)
}
+402 -74
View File
@@ -2,6 +2,7 @@ package service
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"sync"
@@ -22,37 +23,268 @@ import (
// Последовательный тест этого не поймает никогда — нужен барьер, на котором
// оба запроса гарантированно видят ОДНО И ТО ЖЕ состояние Hysteria.
// --- Барьерный тест против настоящего HTTP -----------------------------------
// --- Параллельные подключения одного пира ------------------------------------
// startBarrierTrafficStats поднимает Traffic Stats API, который задерживает
// первые `hold` обращений к `/online` до тех пор, пока не придут все.
//
// Так воспроизводится ровно то состояние гонки, которое случается на живом
// сервере: оба запроса авторизации прочитали статистику до того, как хоть один
// из них успел превратиться в подключение.
func startBarrierTrafficStats(t *testing.T, online map[string]int64, hold int) {
// authResults прогоняет count одновременных авторизаций и возвращает их исходы.
func authResults(t *testing.T, secret string, count int) []error {
t.Helper()
var mu sync.Mutex
arrived := 0
release := make(chan struct{})
var wg sync.WaitGroup
results := make([]error, count)
for i := range results {
wg.Add(1)
go func(idx int) {
defer wg.Done()
_, _, err := Hysteria2Auth(secret)
results[idx] = err
}(i)
}
wg.Wait()
return results
}
func allowedCount(results []error) int {
allowed := 0
for _, err := range results {
if err == nil {
allowed++
}
}
return allowed
}
// Главная регрессия AUTH-03: при `online = max-1` разрешение обязан получить
// ровно ОДИН из двух одновременных запросов.
//
// Барьера, заставлявшего оба запроса увидеть один снимок, здесь больше нет — и
// это следствие исправления, а не упрощение теста. Авторизации одного authId
// сериализованы замком (см. lockPeerAdmission), поэтому одновременно внутри
// `/online` они оказаться не могут, и барьер на двоих просто не собрался бы.
// Доказываемое свойство от этого не изменилось: `/online` отвечает обоим
// одинаково — именно так и ведёт себя Hysteria, пока клиент ещё устанавливает
// соединение, — и без учёта выданных разрешений оба сравнивали бы `2 < 3`.
func TestHysteria2AuthHoldsDeviceLimitUnderConcurrency(t *testing.T) {
newTestDB(t)
// seedPeer создаёт пира с maxDevices = 3, поэтому online = 2 — это
// последнее свободное место.
startTrafficStats(t, &trafficStatsStub{online: map[string]int64{"alpha-auth-id": 2}})
seedPeer(t, "alpha1", "alpha-auth-id")
if allowed := allowedCount(authResults(t, "alpha1-secret", 2)); allowed != 1 {
t.Fatalf("на последнее свободное место допущено %d подключений из 2", allowed)
}
}
// Свободных мест два — проходят оба: механизм не должен превращаться в отказ
// всем, кроме первого.
func TestHysteria2AuthAdmitsBothWhenTwoSlotsFree(t *testing.T) {
newTestDB(t)
startTrafficStats(t, &trafficStatsStub{online: map[string]int64{"alpha-auth-id": 1}})
seedPeer(t, "alpha1", "alpha-auth-id")
for idx, err := range authResults(t, "alpha1-secret", 2) {
if err != nil {
t.Fatalf("подключение %d отклонено при двух свободных местах: %v", idx, err)
}
}
}
// sequencedTrafficStats — Traffic Stats API, у которого ОДНО заранее названное
// обращение к `/online` удерживается до команды теста.
//
// Именно так воспроизводится переупорядочивание снимков: удерживаемый ответ
// содержит картину на момент ПРИХОДА запроса, а к моменту его доставки картина
// уже другая. Ответ формируется до блокировки намеренно — иначе тест проверял
// бы не устаревший снимок, а свежий.
type sequencedTrafficStats struct {
mu sync.Mutex
online map[string]int64
calls int
holdCall int
held chan struct{}
released chan struct{}
}
func startSequencedTrafficStats(t *testing.T, online map[string]int64, holdCall int) *sequencedTrafficStats {
t.Helper()
stats := &sequencedTrafficStats{
online: online,
holdCall: holdCall,
held: make(chan struct{}),
released: make(chan struct{}),
}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/online":
mu.Lock()
arrived++
last := arrived == hold
mu.Unlock()
stats.mu.Lock()
stats.calls++
call := stats.calls
snapshot := make(map[string]int64, len(stats.online))
for key, value := range stats.online {
snapshot[key] = value
}
stats.mu.Unlock()
if call == stats.holdCall {
close(stats.held)
select {
case <-stats.released:
case <-time.After(5 * time.Second):
// Тест не дал команды — отпускаем, чтобы падение было по
// существу, а не по таймауту всего прогона.
}
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(snapshot)
case "/kick":
w.WriteHeader(http.StatusOK)
default:
w.WriteHeader(http.StatusNotFound)
}
}))
t.Cleanup(server.Close)
pointHysteriaConfigAt(t, server.URL)
if err := dao.UpsertConfigValue(constant.Hysteria2TrafficStatsSecret, testTrafficStatsSecret); err != nil {
t.Fatalf("не удалось записать секрет Traffic Stats API: %v", err)
}
return stats
}
func (s *sequencedTrafficStats) setOnline(online map[string]int64) {
s.mu.Lock()
defer s.mu.Unlock()
s.online = online
}
func (s *sequencedTrafficStats) awaitHeld(t *testing.T) {
t.Helper()
select {
case <-s.held:
case <-time.After(5 * time.Second):
t.Fatal("удерживаемое обращение к /online так и не пришло")
}
}
func (s *sequencedTrafficStats) release() {
close(s.released)
}
// Устаревший снимок `/online` не возвращает уже занятое место.
//
// Регрессия второго TOCTOU. Учёт выданных разрешений сам по себе его не
// закрывал: сетевой запрос выполнялся вне мьютекса, поэтому снимки приходили в
// резервацию в произвольном порядке, и более старый откатывал `lastOnline`
// назад:
//
// A получил разрешение при online = 0
// B прочитал online = 0 и задержался
// A подключился, Hysteria показывает online = 1
// C прочитал online = 1 и первым вошёл в резервацию:
// разрешение A признано проявившимся, lastOnline = 1, C отклонён
// B входит со своим устаревшим 0 -> lastOnline снова 0 -> B ДОПУЩЕН
//
// При maxDevices = 1 подключений становилось два. Детектор гонок здесь
// бесполезен принципиально: вся работа с памятью защищена мьютексом, и гонка
// тут логическая, а не по памяти.
func TestHysteria2AuthRejectsStaleOnlineSnapshot(t *testing.T) {
newTestDB(t)
// Удерживается ВТОРОЕ обращение к `/online` — то самое, которое в разборе
// принадлежит запросу B.
stats := startSequencedTrafficStats(t, map[string]int64{"alpha-auth-id": 0}, 2)
id := seedPeer(t, "alpha1", "alpha-auth-id")
if err := dao.UpdatePeer([]int64{id}, map[string]interface{}{"max_devices": 1}); err != nil {
t.Fatalf("не удалось выставить лимит устройств: %v", err)
}
// A занимает единственное место.
if _, _, err := Hysteria2Auth("alpha1-secret"); err != nil {
t.Fatalf("первое подключение отклонено: %v", err)
}
// B читает `/online = 0` и застревает на обратном пути.
var bErr error
bDone := make(chan struct{})
go func() {
defer close(bDone)
_, _, bErr = Hysteria2Auth("alpha1-secret")
}()
stats.awaitHeld(t)
// A действительно подключился: Hysteria показывает одно устройство.
stats.setOnline(map[string]int64{"alpha-auth-id": 1})
// C авторизуется уже по НОВОМУ снимку. На исправленном коде он ждёт замок и
// до `/online` не доходит, поэтому ожидание ограничено: тест не должен
// зависеть от того, успел C или нет.
var cErr error
cDone := make(chan struct{})
go func() {
defer close(cDone)
_, _, cErr = Hysteria2Auth("alpha1-secret")
}()
select {
case <-cDone:
case <-time.After(300 * time.Millisecond):
}
stats.release()
<-bDone
<-cDone
if bErr == nil {
t.Fatal("устаревший снимок /online вернул уже занятое место: допущено второе устройство при лимите 1")
}
if cErr == nil {
t.Fatal("допущено второе устройство при лимите 1")
}
}
// --- Параллельные подключения разных пиров -----------------------------------
// barrierTrafficStats — Traffic Stats API, который задерживает первые `hold`
// обращений к `/online`, пока не придут все.
//
// Так проверяется, что замок авторизации НЕ является общим на процесс: два
// запроса разных пиров обязаны оказаться внутри `/online` одновременно. С
// глобальным замком барьер не собрался бы никогда.
type barrierTrafficStats struct {
mu sync.Mutex
arrived int
collected bool
release chan struct{}
}
func startBarrierTrafficStats(t *testing.T, online map[string]int64, hold int) *barrierTrafficStats {
t.Helper()
barrier := &barrierTrafficStats{release: make(chan struct{})}
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/online":
barrier.mu.Lock()
barrier.arrived++
last := barrier.arrived == hold
if last {
barrier.collected = true
}
barrier.mu.Unlock()
if last {
close(release)
close(barrier.release)
} else {
select {
case <-release:
case <-barrier.release:
case <-time.After(5 * time.Second):
// Барьер не собрался — отпускаем, чтобы тест упал по
// существу, а не по таймауту всего прогона.
// существу (см. requireCollected), а не по таймауту всего
// прогона.
}
}
@@ -70,71 +302,28 @@ func startBarrierTrafficStats(t *testing.T, online map[string]int64, hold int) {
if err := dao.UpsertConfigValue(constant.Hysteria2TrafficStatsSecret, testTrafficStatsSecret); err != nil {
t.Fatalf("не удалось записать секрет Traffic Stats API: %v", err)
}
return barrier
}
// Главная регрессия AUTH-03: при `online = max-1` разрешение обязан получить
// ровно ОДИН из двух одновременных запросов.
func TestHysteria2AuthHoldsDeviceLimitUnderConcurrency(t *testing.T) {
newTestDB(t)
// seedPeer создаёт пира с maxDevices = 3, поэтому online = 2 — это
// последнее свободное место.
startBarrierTrafficStats(t, map[string]int64{"alpha-auth-id": 2}, 2)
seedPeer(t, "alpha1", "alpha-auth-id")
// requireCollected требует, чтобы барьер действительно собрался.
//
// Без этой проверки тест проходил бы и при глобальном замке: барьер молча
// разошёлся бы по таймауту, а исходы авторизации остались бы прежними.
func (b *barrierTrafficStats) requireCollected(t *testing.T) {
t.Helper()
var wg sync.WaitGroup
results := make([]error, 2)
for i := range results {
wg.Add(1)
go func(idx int) {
defer wg.Done()
_, _, err := Hysteria2Auth("alpha1-secret")
results[idx] = err
}(i)
}
wg.Wait()
allowed := 0
for _, err := range results {
if err == nil {
allowed++
}
}
if allowed != 1 {
t.Fatalf("на последнее свободное место допущено %d подключений из 2", allowed)
}
}
// Свободных мест два — проходят оба: механизм не должен превращаться в
// сериализацию подключений.
func TestHysteria2AuthAdmitsBothWhenTwoSlotsFree(t *testing.T) {
newTestDB(t)
startBarrierTrafficStats(t, map[string]int64{"alpha-auth-id": 1}, 2)
seedPeer(t, "alpha1", "alpha-auth-id")
var wg sync.WaitGroup
results := make([]error, 2)
for i := range results {
wg.Add(1)
go func(idx int) {
defer wg.Done()
_, _, err := Hysteria2Auth("alpha1-secret")
results[idx] = err
}(i)
}
wg.Wait()
for idx, err := range results {
if err != nil {
t.Fatalf("подключение %d отклонено при двух свободных местах: %v", idx, err)
}
b.mu.Lock()
defer b.mu.Unlock()
if !b.collected {
t.Fatalf("запросы разных пиров не оказались в /online одновременно: пришло %d — авторизация сериализована глобально", b.arrived)
}
}
// Резервации принадлежат КОНКРЕТНОМУ пиру: занятое место одного не должно
// закрывать доступ другому.
// закрывать доступ другому, а замок одного не должен задерживать другого.
func TestDeviceAdmissionsAreIsolatedPerPeer(t *testing.T) {
newTestDB(t)
startBarrierTrafficStats(t, map[string]int64{"alpha-auth-id": 2, "bravo-auth-id": 0}, 2)
barrier := startBarrierTrafficStats(t, map[string]int64{"alpha-auth-id": 2, "bravo-auth-id": 0}, 2)
seedPeer(t, "alpha1", "alpha-auth-id")
seedPeer(t, "bravo2", "bravo-auth-id")
@@ -151,6 +340,8 @@ func TestDeviceAdmissionsAreIsolatedPerPeer(t *testing.T) {
}()
wg.Wait()
barrier.requireCollected(t)
if alphaErr != nil {
t.Fatalf("первое подключение пира на последнее место отклонено: %v", alphaErr)
}
@@ -312,6 +503,143 @@ func TestReserveDeviceSlotForgetsIdlePeer(t *testing.T) {
}
}
// --- Замок последовательности «/online -> резервация» ------------------------
// Авторизации ОДНОГО пира исключают друг друга: только так снимок, прочитанный
// вторым, не может оказаться старше уже обработанного.
func TestPeerAdmissionGateSerializesSameAuthID(t *testing.T) {
first := lockPeerAdmission("gate-auth")
entered := make(chan struct{})
done := make(chan struct{})
go func() {
defer close(done)
second := lockPeerAdmission("gate-auth")
close(entered)
second()
}()
select {
case <-entered:
t.Fatal("второй вызывающий вошёл в критическую секцию при удерживаемом замке")
case <-time.After(100 * time.Millisecond):
}
first()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("замок не освободился после возврата функции освобождения")
}
}
// Разные пиры не мешают друг другу: замок именно по authId, а не один на
// процесс. Внутри него выполняется сетевой запрос, поэтому общий замок
// выстроил бы подключения всех пиров в одну очередь.
func TestPeerAdmissionGateDoesNotSerializeDifferentAuthIDs(t *testing.T) {
held := lockPeerAdmission("gate-alpha")
defer held()
done := make(chan struct{})
go func() {
defer close(done)
other := lockPeerAdmission("gate-bravo")
other()
}()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("замок одного пира задержал авторизацию другого")
}
}
// Карта замков не растёт: запись живёт ровно столько, сколько есть желающие её
// взять. Без счётчика ссылок здесь появлялась бы строка на каждый когда-либо
// авторизовавшийся authId — та же утечка, от которой в учёте разрешений
// защищает forgetIfIdle.
func TestPeerAdmissionGateLeavesNoEntriesBehind(t *testing.T) {
if got := heldPeerAdmissionGates(); got != 0 {
t.Fatalf("перед проверкой уже удерживается %d замков", got)
}
for i := 0; i < 50; i++ {
unlock := lockPeerAdmission("gate-transient")
unlock()
}
var wg sync.WaitGroup
for i := 0; i < 20; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
unlock := lockPeerAdmission(fmt.Sprintf("gate-%d", idx%3))
unlock()
}(i)
}
wg.Wait()
if got := heldPeerAdmissionGates(); got != 0 {
t.Fatalf("после освобождения осталось %d замков", got)
}
}
// --- Уборка учёта ------------------------------------------------------------
// Запись о пире, за которым не числится ни подключений, ни действующих
// разрешений, убирается.
//
// forgetIfIdle этого не делал: он срабатывает только на ветке отказа, а после
// успешной выдачи разрешения запись оставалась с непустым pending, и снять её
// было некому. В карте накапливались удалённые пиры и старые authId,
// переписанные импортом.
func TestSweepDeviceAdmissionsForgetsIdlePeers(t *testing.T) {
resetDeviceAdmissions()
t.Cleanup(resetDeviceAdmissions)
now := time.Now()
if !reserveDeviceSlot(admissionAuthID, 0, 3, now) {
t.Fatal("разрешение отклонено при свободном месте")
}
if admissionEntries() != 1 {
t.Fatal("учёт не запомнил выданное разрешение")
}
// Разрешение ещё действует — запись обязана остаться, даже если пира нет в
// `/online`: подключение может проявиться в любой момент.
sweepDeviceAdmissions(map[string]int64{}, now)
if admissionEntries() != 1 {
t.Fatal("уборка сняла действующее разрешение")
}
// Разрешение протухло, подключений нет: помнить нечего.
sweepDeviceAdmissions(map[string]int64{}, now.Add(pendingAdmissionTTL+time.Second))
if got := admissionEntries(); got != 0 {
t.Fatalf("запись о неактивном пире осталась: записей %d", got)
}
}
// Пир, который сейчас на связи, из учёта не убирается: его состояние ещё
// участвует в решениях.
func TestSweepDeviceAdmissionsKeepsOnlinePeers(t *testing.T) {
resetDeviceAdmissions()
t.Cleanup(resetDeviceAdmissions)
now := time.Now()
if !reserveDeviceSlot(admissionAuthID, 1, 3, now) {
t.Fatal("разрешение отклонено при свободном месте")
}
sweepDeviceAdmissions(
map[string]int64{admissionAuthID: 1},
now.Add(pendingAdmissionTTL+time.Second),
)
if got := admissionEntries(); got != 1 {
t.Fatalf("уборка сняла запись подключённого пира: записей %d", got)
}
}
// Учёт выдерживает параллельный доступ и не выдаёт больше мест, чем есть.
// Проверка ловит дефект и без детектора гонок: он виден по числу разрешений.
func TestReserveDeviceSlotIsConcurrencySafe(t *testing.T) {