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
674 lines
27 KiB
Go
674 lines
27 KiB
Go
package service
|
||
|
||
import (
|
||
"encoding/json"
|
||
"fmt"
|
||
"net/http"
|
||
"net/http/httptest"
|
||
"sync"
|
||
"testing"
|
||
"time"
|
||
|
||
"hy2xs-admin/dao"
|
||
"hy2xs-admin/model/constant"
|
||
)
|
||
|
||
// Лимит устройств проверяется на ПАРАЛЛЕЛЬНЫХ запросах авторизации.
|
||
//
|
||
// Прежняя проверка сравнивала ответ `/online` с maxDevices и сразу отвечала
|
||
// «allow»: между чтением и ответом место ничем не удерживалось, поэтому два
|
||
// одновременных подключения при `online = max-1` получали разрешение оба и
|
||
// объявленный лимит превышался.
|
||
//
|
||
// Последовательный тест этого не поймает никогда — нужен барьер, на котором
|
||
// оба запроса гарантированно видят ОДНО И ТО ЖЕ состояние Hysteria.
|
||
|
||
// --- Параллельные подключения одного пира ------------------------------------
|
||
|
||
// authResults прогоняет count одновременных авторизаций и возвращает их исходы.
|
||
func authResults(t *testing.T, secret string, count int) []error {
|
||
t.Helper()
|
||
|
||
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":
|
||
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(barrier.release)
|
||
} else {
|
||
select {
|
||
case <-barrier.release:
|
||
case <-time.After(5 * time.Second):
|
||
// Барьер не собрался — отпускаем, чтобы тест упал по
|
||
// существу (см. requireCollected), а не по таймауту всего
|
||
// прогона.
|
||
}
|
||
}
|
||
|
||
w.Header().Set("Content-Type", "application/json")
|
||
_ = json.NewEncoder(w).Encode(online)
|
||
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 barrier
|
||
}
|
||
|
||
// requireCollected требует, чтобы барьер действительно собрался.
|
||
//
|
||
// Без этой проверки тест проходил бы и при глобальном замке: барьер молча
|
||
// разошёлся бы по таймауту, а исходы авторизации остались бы прежними.
|
||
func (b *barrierTrafficStats) requireCollected(t *testing.T) {
|
||
t.Helper()
|
||
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
if !b.collected {
|
||
t.Fatalf("запросы разных пиров не оказались в /online одновременно: пришло %d — авторизация сериализована глобально", b.arrived)
|
||
}
|
||
}
|
||
|
||
// Резервации принадлежат КОНКРЕТНОМУ пиру: занятое место одного не должно
|
||
// закрывать доступ другому, а замок одного не должен задерживать другого.
|
||
func TestDeviceAdmissionsAreIsolatedPerPeer(t *testing.T) {
|
||
newTestDB(t)
|
||
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")
|
||
|
||
var wg sync.WaitGroup
|
||
var alphaErr, bravoErr error
|
||
wg.Add(2)
|
||
go func() {
|
||
defer wg.Done()
|
||
_, _, alphaErr = Hysteria2Auth("alpha1-secret")
|
||
}()
|
||
go func() {
|
||
defer wg.Done()
|
||
_, _, bravoErr = Hysteria2Auth("bravo2-secret")
|
||
}()
|
||
wg.Wait()
|
||
|
||
barrier.requireCollected(t)
|
||
|
||
if alphaErr != nil {
|
||
t.Fatalf("первое подключение пира на последнее место отклонено: %v", alphaErr)
|
||
}
|
||
if bravoErr != nil {
|
||
t.Fatalf("резервация чужого пира закрыла доступ: %v", bravoErr)
|
||
}
|
||
}
|
||
|
||
// Последовательно тот же лимит тоже держится: второй запрос видит место,
|
||
// занятое первым, хотя `/online` ещё показывает прежнее число.
|
||
func TestHysteria2AuthCountsPendingAdmissionSequentially(t *testing.T) {
|
||
newTestDB(t)
|
||
// Число НЕ меняется между запросами — именно так и ведёт себя Hysteria,
|
||
// пока клиент ещё устанавливает соединение.
|
||
startTrafficStats(t, &trafficStatsStub{online: map[string]int64{"alpha-auth-id": 1}})
|
||
seedPeer(t, "alpha1", "alpha-auth-id")
|
||
|
||
// max = 3, online = 1 -> свободно два места.
|
||
if _, _, err := Hysteria2Auth("alpha1-secret"); err != nil {
|
||
t.Fatalf("первое подключение отклонено: %v", err)
|
||
}
|
||
if _, _, err := Hysteria2Auth("alpha1-secret"); err != nil {
|
||
t.Fatalf("второе подключение отклонено: %v", err)
|
||
}
|
||
// Третье превысило бы лимит: 1 онлайн + 2 выданных разрешения.
|
||
if _, _, err := Hysteria2Auth("alpha1-secret"); err == nil {
|
||
t.Fatal("подключение сверх лимита принято: выданные разрешения не учтены")
|
||
}
|
||
}
|
||
|
||
// --- Единица учёта ------------------------------------------------------------
|
||
|
||
const admissionAuthID = "auth-under-test"
|
||
|
||
func TestReserveDeviceSlotAllowsUpToLimit(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
for i := 1; i <= 3; i++ {
|
||
if !reserveDeviceSlot(admissionAuthID, 0, 3, now) {
|
||
t.Fatalf("разрешение %d из 3 отклонено", i)
|
||
}
|
||
}
|
||
if reserveDeviceSlot(admissionAuthID, 0, 3, now) {
|
||
t.Fatal("выдано четвёртое разрешение при лимите 3")
|
||
}
|
||
}
|
||
|
||
// Рост числа онлайн-устройств означает, что выданные разрешения превратились в
|
||
// подключения. Не сняв их, админка посчитала бы одно устройство дважды, и
|
||
// лимит стал бы вдвое строже объявленного.
|
||
func TestReserveDeviceSlotAbsorbsMaterializedAdmissions(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
if !reserveDeviceSlot(admissionAuthID, 0, 2, now) {
|
||
t.Fatal("первое разрешение отклонено")
|
||
}
|
||
// Клиент подключился: Hysteria теперь видит одно устройство.
|
||
if !reserveDeviceSlot(admissionAuthID, 1, 2, now) {
|
||
t.Fatal("проявившееся разрешение посчитано дважды")
|
||
}
|
||
// Теперь занято: одно подключение плюс одно выданное разрешение.
|
||
if reserveDeviceSlot(admissionAuthID, 1, 2, now) {
|
||
t.Fatal("выдано разрешение сверх лимита")
|
||
}
|
||
}
|
||
|
||
// Разрешение, за которым не последовало подключения, освобождает место само.
|
||
func TestReserveDeviceSlotExpiresPendingAdmission(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
if !reserveDeviceSlot(admissionAuthID, 0, 1, now) {
|
||
t.Fatal("первое разрешение отклонено")
|
||
}
|
||
if reserveDeviceSlot(admissionAuthID, 0, 1, now) {
|
||
t.Fatal("выдано разрешение сверх лимита 1")
|
||
}
|
||
|
||
// Клиент так и не подключился.
|
||
later := now.Add(pendingAdmissionTTL + time.Second)
|
||
if !reserveDeviceSlot(admissionAuthID, 0, 1, later) {
|
||
t.Fatal("протухшее разрешение не освободило место")
|
||
}
|
||
}
|
||
|
||
// Место не занимается отказом: иначе серия отклонённых попыток удерживала бы
|
||
// слоты на всё время TTL.
|
||
func TestReserveDeviceSlotDoesNotConsumeSlotOnRefusal(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
for i := 0; i < 5; i++ {
|
||
if reserveDeviceSlot(admissionAuthID, 3, 3, now) {
|
||
t.Fatal("выдано разрешение при исчерпанном лимите")
|
||
}
|
||
}
|
||
// Одно устройство отключилось — место обязано быть свободно немедленно.
|
||
if !reserveDeviceSlot(admissionAuthID, 2, 3, now) {
|
||
t.Fatal("отклонённые попытки заняли места")
|
||
}
|
||
}
|
||
|
||
// admissionEntries — размер учёта.
|
||
func admissionEntries() int {
|
||
deviceAdmissions.Lock()
|
||
defer deviceAdmissions.Unlock()
|
||
return len(deviceAdmissions.byAuthID)
|
||
}
|
||
|
||
// Учёт не растёт от повторных обращений: одна запись на пира, сколько бы
|
||
// попыток он ни сделал.
|
||
//
|
||
// Границы роста здесь две, и обе существенны. Верхняя — число пиров: сюда
|
||
// попадают только authId, прошедшие поиск по secret_digest и всю политику
|
||
// доступа, поэтому произвольный ключ извне добавить нельзя. Нижняя — запись
|
||
// исчезает, как только помнить о пире нечего (см. следующий тест).
|
||
func TestReserveDeviceSlotKeepsOneEntryPerPeer(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
for i := 0; i < 100; i++ {
|
||
reserveDeviceSlot("transient-auth", 1, 1, now)
|
||
}
|
||
|
||
if got := admissionEntries(); got != 1 {
|
||
t.Fatalf("повторные попытки одного пира дали %d записей", got)
|
||
}
|
||
}
|
||
|
||
// Запись исчезает, когда о пире нечего помнить: он не онлайн и выданных
|
||
// разрешений за ним нет. Без этого карта накапливала бы по строке на каждый
|
||
// когда-либо авторизовавшийся authId, включая давно удалённых пиров.
|
||
func TestReserveDeviceSlotForgetsIdlePeer(t *testing.T) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
now := time.Now()
|
||
if !reserveDeviceSlot(admissionAuthID, 1, 3, now) {
|
||
t.Fatal("разрешение отклонено при свободном месте")
|
||
}
|
||
if admissionEntries() != 1 {
|
||
t.Fatal("учёт не запомнил выданное разрешение")
|
||
}
|
||
|
||
// Пир отключился целиком, а выданное разрешение протухло: помнить нечего.
|
||
later := now.Add(pendingAdmissionTTL + time.Second)
|
||
// Лимит 0 — разрешение не выдаётся, поэтому запись остаться не должна.
|
||
reserveDeviceSlot(admissionAuthID, 0, 0, later)
|
||
|
||
if got := admissionEntries(); got != 0 {
|
||
t.Fatalf("запись о неактивном пире осталась: записей %d", got)
|
||
}
|
||
}
|
||
|
||
// --- Замок последовательности «/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) {
|
||
resetDeviceAdmissions()
|
||
t.Cleanup(resetDeviceAdmissions)
|
||
|
||
const workers = 50
|
||
const limit = int64(7)
|
||
|
||
now := time.Now()
|
||
var wg sync.WaitGroup
|
||
var mu sync.Mutex
|
||
allowed := 0
|
||
|
||
for i := 0; i < workers; i++ {
|
||
wg.Add(1)
|
||
go func() {
|
||
defer wg.Done()
|
||
if reserveDeviceSlot(admissionAuthID, 0, limit, now) {
|
||
mu.Lock()
|
||
allowed++
|
||
mu.Unlock()
|
||
}
|
||
}()
|
||
}
|
||
wg.Wait()
|
||
|
||
if int64(allowed) != limit {
|
||
t.Fatalf("выдано %d разрешений при лимите %d", allowed, limit)
|
||
}
|
||
}
|