package service import ( "encoding/json" "errors" "net/http" "net/http/httptest" "strings" "sync" "testing" "time" "hy2xs-admin/dao" "hy2xs-admin/model/bo" "hy2xs-admin/model/constant" "hy2xs-admin/model/dto" ) // Цикл учёта проверяется против НАСТОЯЩЕГО Traffic Stats API. // // accountStatsStub добавляет к trafficStatsStub то, чего у него нет: ответ // `GET /traffic`. Разделять их не нужно — на живом сервере это один и тот же // API на одном порту, и джоба ходит в оба маршрута подряд. type accountStatsStub struct { mu sync.Mutex traffic map[string]bo.Hysteria2UserTraffic online map[string]int64 trafficStatus int onlineStatus int kickStatus int trafficCalls int onlineCalls int kickCalls int kickedKeys [][]string // trafficCleared запоминает, просила ли админка обнулить счётчики. trafficCleared []bool // usageAtOnline — суммарный расход пиров на момент запроса `/online`. // Именно этим доказывается порядок «сначала учёт, потом enforcement»: // после джобы оба шага уже выполнены и проверять там нечего. usageAtOnline []map[string]int64 } func (s *accountStatsStub) usageSnapshot() map[string]int64 { usage := map[string]int64{} peers, err := dao.ListPeer("1=1") if err != nil { return usage } for _, peer := range peers { if peer.AuthId == nil { continue } var total int64 if peer.DownloadBytes != nil { total += *peer.DownloadBytes } if peer.UploadBytes != nil { total += *peer.UploadBytes } usage[*peer.AuthId] = total } return usage } func startAccountStats(t *testing.T, stub *accountStatsStub) *accountStatsStub { t.Helper() if stub == nil { stub = &accountStatsStub{} } if stub.traffic == nil { stub.traffic = map[string]bo.Hysteria2UserTraffic{} } if stub.online == nil { stub.online = map[string]int64{} } server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { stub.mu.Lock() defer stub.mu.Unlock() switch r.URL.Path { case "/traffic": stub.trafficCalls++ stub.trafficCleared = append(stub.trafficCleared, r.URL.Query().Get("clear") == "1") if stub.trafficStatus != 0 { w.WriteHeader(stub.trafficStatus) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(stub.traffic) case "/online": stub.onlineCalls++ stub.usageAtOnline = append(stub.usageAtOnline, stub.usageSnapshot()) if stub.onlineStatus != 0 { w.WriteHeader(stub.onlineStatus) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(stub.online) case "/kick": stub.kickCalls++ var keys []string if err := json.NewDecoder(r.Body).Decode(&keys); err != nil { w.WriteHeader(http.StatusBadRequest) return } stub.kickedKeys = append(stub.kickedKeys, keys) if stub.kickStatus != 0 { w.WriteHeader(stub.kickStatus) return } 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 stub } func (s *accountStatsStub) kicked() []string { s.mu.Lock() defer s.mu.Unlock() out := []string{} for _, keys := range s.kickedKeys { out = append(out, keys...) } return out } // peerUsage помещает пиру расход и настройки доступа. func peerUsage(t *testing.T, id int64, updates map[string]interface{}) { t.Helper() if err := dao.UpdatePeer([]int64{id}, updates); err != nil { t.Fatalf("подготовка состояния пира: %v", err) } } // --- Границы принудительного отключения -------------------------------------- // Главная регрессия QUOTA-01: cron требовал СТРОГОГО превышения квоты, а // авторизация отказывала уже при равенстве. Пир с исчерпанной квотой не // пускался заново, но его живая сессия не разрывалась никогда — он продолжал // пользоваться доступом, пока не переподключался сам. func TestCronKicksPeerAtExactQuota(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{}{ "quota_bytes": int64(1_000), "download_bytes": int64(600), "upload_bytes": int64(400), }) CronHandleAccount() if got := stub.kicked(); len(got) != 1 || got[0] != "alpha-auth-id" { t.Fatalf("пир с исчерпанной квотой не отключён: %v", got) } } // Нулевая квота — это ноль байтов, а не безлимит. Прежнее условие // `quota_bytes > 0` такую строку не рассматривало вовсе. func TestCronKicksPeerWithZeroQuota(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{}{"quota_bytes": int64(0)}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("пир с нулевой квотой не отключён: %v", got) } } func TestCronDoesNotKickUnlimitedQuota(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{}{ "quota_bytes": int64(-1), "download_bytes": int64(1 << 40), }) CronHandleAccount() if got := stub.kicked(); len(got) != 0 { t.Fatalf("безлимитный пир отключён по квоте: %v", got) } } func TestCronKicksPeerWhenExpiryEqualsNow(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{online: map[string]int64{"alpha-auth-id": 1}}) id := seedPeer(t, "alpha1", "alpha-auth-id") // Срок в недавнем прошлом: «момент наступил» и «момент прошёл» — по // контракту одно и то же, а точное совпадение с now в тесте недостижимо. peerUsage(t, id, map[string]interface{}{"expires_at": time.Now().UnixMilli() - 1}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("пир с истёкшим сроком не отключён: %v", got) } } // Блокировка «до» момента, который уже наступил, закончилась: пира отключать // не за что. func TestCronDoesNotKickAfterBanExpired(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{}{"banned_until": time.Now().UnixMilli() - 1}) CronHandleAccount() if got := stub.kicked(); len(got) != 0 { t.Fatalf("пир с истёкшей блокировкой отключён: %v", got) } } func TestCronKicksBannedPeer(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{}{"banned_until": time.Now().UnixMilli() + 3_600_000}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("заблокированный пир не отключён: %v", got) } } func TestCronKicksDisabledPeer(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{}{"disabled": int64(1)}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("отключённый пир не отключён: %v", got) } } // Действующий пир не трогается, и запрос без единой цели не отправляется вовсе: // раньше POST с пустым массивом уезжал в Hysteria каждые 30 секунд. func TestCronSendsNoKickWithoutTargets(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{online: map[string]int64{"alpha-auth-id": 1}}) seedPeer(t, "alpha1", "alpha-auth-id") CronHandleAccount() stub.mu.Lock() calls := stub.kickCalls stub.mu.Unlock() if calls != 0 { t.Fatalf("вызов /kick без единой цели: %d", calls) } } // Пир, которого Hysteria не считает онлайн, в enforcement не участвует: рвать // у него нечего. func TestCronIgnoresOfflinePeers(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{online: map[string]int64{}}) id := seedPeer(t, "alpha1", "alpha-auth-id") peerUsage(t, id, map[string]interface{}{"disabled": int64(1)}) CronHandleAccount() if got := stub.kicked(); len(got) != 0 { t.Fatalf("офлайн-пир попал в /kick: %v", got) } } // --- Сверка живых сессий ----------------------------------------------------- // Живая сессия, которой в базе больше ничего не соответствует, завершается. // // Прежний обход шёл по НАЙДЕННЫМ пирам, поэтому 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 принимает решение по счётчикам, поэтому счётчики обязаны быть // обновлены ДО него. Раньше обе половины запускались параллельными горутинами, // и превышение квоты замечалось в лучшем случае со следующего тика. func TestCronCollectsTrafficBeforeEnforcing(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{ traffic: map[string]bo.Hysteria2UserTraffic{ "alpha-auth-id": {Rx: 600, Tx: 400}, }, online: map[string]int64{"alpha-auth-id": 1}, }) id := seedPeer(t, "alpha1", "alpha-auth-id") peerUsage(t, id, map[string]interface{}{"quota_bytes": int64(1_000)}) CronHandleAccount() stub.mu.Lock() usage := stub.usageAtOnline stub.mu.Unlock() if len(usage) == 0 { t.Fatal("enforcement не выполнялся") } if got := usage[0]["alpha-auth-id"]; got != 1_000 { t.Fatalf("enforcement увидел расход %d — дельта ещё не была записана", got) } // И следствие: превышение замечено в ТОМ ЖЕ тике, а не в следующем. if kicked := stub.kicked(); len(kicked) != 1 { t.Fatalf("исчерпавший квоту пир не отключён в том же цикле: %v", kicked) } } // Чтение трафика деструктивно по контракту Traffic Stats API: без clear=1 // счётчики Hysteria не обнуляются, и следующий сбор посчитал бы тот же трафик // повторно. func TestCronClearsTrafficCounters(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{ traffic: map[string]bo.Hysteria2UserTraffic{"alpha-auth-id": {Rx: 1, Tx: 1}}, }) seedPeer(t, "alpha1", "alpha-auth-id") CronHandleAccount() stub.mu.Lock() cleared := stub.trafficCleared stub.mu.Unlock() if len(cleared) != 1 || !cleared[0] { t.Fatalf("сбор трафика выполнен без clear=1: %v", cleared) } } // Дельта трафика, которую не удалось приписать пиру, считается потерей и // попадает в исход джобы. // // Чтение `?clear=1` деструктивно по контракту Traffic Stats API: счётчики // Hysteria обнуляются сразу после отправки ответа, поэтому каждая дельта // существует ровно в одном экземпляре. Раньше такой случай делал `continue` и // не оставлял следа вовсе. func TestSaveAccountTrafficCountsLostDeltas(t *testing.T) { newTestDB(t) startAccountStats(t, &accountStatsStub{ traffic: map[string]bo.Hysteria2UserTraffic{ "known-auth-id": {Rx: 10, Tx: 20}, "unknown-auth-id": {Rx: 30, Tx: 40}, }, }) seedPeer(t, "alpha1", "known-auth-id") apiPort, err := GetHysteria2ApiPort() if err != nil { t.Fatalf("порт Traffic Stats API: %v", err) } err = saveAccountTraffic(apiPort, testTrafficStatsSecret) if err == nil { t.Fatal("потеря дельты не сообщена вызывающему") } var loss *trafficLossError if !errors.As(err, &loss) { t.Fatalf("потеря сообщена не как величина: %v", err) } if loss.lost != 1 { t.Fatalf("учтено %d потерь, ожидалась 1", loss.lost) } if !strings.Contains(err.Error(), "1") { t.Errorf("сообщение не называет количество: %q", err.Error()) } // Известный пир при этом обязан получить свою дельту: потеря одной записи // не отменяет остальных. peer := snapshotPeers(t)["alpha1"] if *peer.DownloadBytes != 10 || *peer.UploadBytes != 20 { t.Fatalf("дельта известного пира не записана: %d/%d", *peer.DownloadBytes, *peer.UploadBytes) } } // Мнение systemd на цикл учёта не влияет. // // Прежний гейт `if !Hysteria2IsRunning() { return }` стоял на решении о // применении операции, а util.Exec не отличает «служба неактивна» от // «спросить не удалось»: сломанный systemctl при живой Hysteria молча отключал // и учёт трафика, и принудительное отключение — без единой строки в журнале. func TestCronRunsWhenSystemdSaysStopped(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{online: map[string]int64{"alpha-auth-id": 1}}) withHysteriaServiceState(t, HysteriaServiceInactive) id := seedPeer(t, "alpha1", "alpha-auth-id") peerUsage(t, id, map[string]interface{}{"disabled": int64(1)}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("мнение systemd отключило принудительное отключение: %v", got) } } // Отсутствующее значение секрета Traffic Stats API — отказ джобы, а не паника. // // Раньше здесь стояло `*trafficSecretConfig.Value` без проверки, причём внутри // отсоединённой горутины: разыменование nil роняло бы весь процесс вместе с // обработчиком machine-auth, а не одну джобу. func TestCronSurvivesMissingTrafficSecret(t *testing.T) { newTestDB(t) startAccountStats(t, nil) seedPeer(t, "alpha1", "alpha-auth-id") // Пустое значение ключа неотличимо от его отсутствия: и то и другое // означает «секрета нет». Прежний путь читал `*config.Value` без проверки // и на строке без значения падал с nil-разыменованием. if err := dao.UpsertConfigValue(constant.Hysteria2TrafficStatsSecret, ""); err != nil { t.Fatalf("не удалось стереть секрет: %v", err) } defer func() { if recovered := recover(); recovered != nil { t.Fatalf("отсутствующий секрет уронил джобу учёта: %v", recovered) } }() CronHandleAccount() } // Строка пира без идентификатора не роняет сброс трафика: раньше `*item.Id` // разыменовывался без проверки. func TestCronResetTrafficSurvivesRowWithoutID(t *testing.T) { newTestDB(t) id := seedPeer(t, "alpha1", "alpha-auth-id") peerUsage(t, id, map[string]interface{}{"download_bytes": int64(100), "upload_bytes": int64(200)}) defer func() { if recovered := recover(); recovered != nil { t.Fatalf("сброс трафика упал: %v", recovered) } }() CronResetTraffic() peer := snapshotPeers(t)["alpha1"] if *peer.DownloadBytes != 0 || *peer.UploadBytes != 0 { t.Fatalf("счётчики не сброшены: %d/%d", *peer.DownloadBytes, *peer.UploadBytes) } } // Отказ `/traffic` не отменяет enforcement: политика применяется по уже // известным счётчикам, а не пропускается вместе со сбором. func TestCronEnforcesEvenWhenTrafficCollectionFails(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{ trafficStatus: http.StatusInternalServerError, online: map[string]int64{"alpha-auth-id": 1}, }) id := seedPeer(t, "alpha1", "alpha-auth-id") peerUsage(t, id, map[string]interface{}{"disabled": int64(1)}) CronHandleAccount() if got := stub.kicked(); len(got) != 1 { t.Fatalf("отказ сбора трафика отменил принудительное отключение: %v", got) } } // Второй тик поверх идущего цикла не запускает второй цикл. Проверяется // наблюдаемым следствием: при удерживаемом мьютексе джоба обязана вернуться, // не сходив в Hysteria ни разу. func TestCronHandleAccountSkipsOverlappingTick(t *testing.T) { newTestDB(t) stub := startAccountStats(t, nil) seedPeer(t, "alpha1", "alpha-auth-id") accountJobMutex.Lock() CronHandleAccount() accountJobMutex.Unlock() stub.mu.Lock() calls := stub.trafficCalls + stub.onlineCalls stub.mu.Unlock() if calls != 0 { t.Fatalf("параллельный тик запустил второй цикл учёта: обращений %d", calls) } // А после освобождения обычный тик проходит. CronHandleAccount() stub.mu.Lock() calls = stub.trafficCalls + stub.onlineCalls stub.mu.Unlock() if calls == 0 { t.Fatal("цикл учёта не выполнился после освобождения мьютекса") } } // Джоба СИНХРОННА: планировщик обязан видеть её работу, иначе StopCron // возвращается, releaseResource закрывает SQLite, а недобитые горутины // продолжают писать в закрытое соединение. // // Доказывается тем, что к моменту возврата CronHandleAccount вся работа уже // сделана — при отсоединённых горутинах обращения к Hysteria к этому моменту // ещё не случились бы. func TestCronHandleAccountIsSynchronous(t *testing.T) { newTestDB(t) stub := startAccountStats(t, &accountStatsStub{ traffic: map[string]bo.Hysteria2UserTraffic{"alpha-auth-id": {Rx: 10, Tx: 20}}, online: map[string]int64{"alpha-auth-id": 1}, }) seedPeer(t, "alpha1", "alpha-auth-id") CronHandleAccount() stub.mu.Lock() trafficCalls := stub.trafficCalls onlineCalls := stub.onlineCalls stub.mu.Unlock() if trafficCalls != 1 || onlineCalls != 1 { t.Fatalf("работа не завершена к возврату джобы: /traffic %d, /online %d", trafficCalls, onlineCalls) } // И записанная дельта уже видна: значит цикл дошёл до конца, а не был // передан горутине. peer := snapshotPeers(t)["alpha1"] if *peer.DownloadBytes != 10 || *peer.UploadBytes != 20 { t.Fatalf("дельта не записана к возврату джобы: %d/%d", *peer.DownloadBytes, *peer.UploadBytes) } } // StopCron дожидается запущенной джобы учёта. Раньше внешняя горутина // заканчивалась мгновенно, и планировщику было нечего ждать. func TestStopCronWaitsForAccountJob(t *testing.T) { newTestDB(t) startAccountStats(t, nil) seedPeer(t, "alpha1", "alpha-auth-id") // Джоба удерживается занятым мьютексом: пока он не освобождён, ни один // цикл учёта не идёт, и StopCron обязан вернуться без ожидания. done := make(chan struct{}) go func() { defer close(done) accountJobMutex.Lock() defer accountJobMutex.Unlock() time.Sleep(50 * time.Millisecond) }() if err := InitCron(); err != nil { t.Fatalf("InitCron: %v", err) } StopCron() <-done if count := CronEntryCount(); count != 0 { t.Fatalf("после остановки осталось %d записей", count) } }