From 0ba9b7d7393637c389e764449ecbd3caaa9dfec0 Mon Sep 17 00:00:00 2001 From: subochev Date: Sun, 13 Sep 2026 20:33:36 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D1=8D=D0=BA=D1=81=D0=BF=D0=BE=D0=BD?= =?UTF-8?q?=D0=B5=D0=BD=D1=86=D0=B8=D0=B0=D0=BB=D1=8C=D0=BD=D1=8B=D0=B9=20?= =?UTF-8?q?backoff=20(providers[].backoff=20/=20upstreams[].backoff)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Повторяющиеся сбои апстрима (5xx/429/402 и сетевые ошибки) увязываются экспоненциальным откатом: первая ошибка — отдых 1s, далее 2s, 4s, … до потолка из конфига (ISO-8601: P1M, PT15M…); успешный запрос сбрасывает счётчик. Приоритет: upstreams[].backoff (свой счётчик на модель) над providers[].backoff (общий счётчик на все модели провайдера). - Backoff.kt: BackoffGuard (Mutex/@Volatile, kotlinx-datetime) + BackoffRegistry (скоупы провайдер/модель); - pickFreeUpstream пропускает апстримы в откате (лог «~Ns — пропускаю»); - все апстримы модели в откате → 503 «all upstreams cooling, retry in Ns» + заголовок Retry-After; - некорректный ISO-8601 в backoff — warn + игнор поля; - BackoffTest (8 тестов) + обновлённые тесты pickFreeUpstream; - CONFIG.md/README/TESTING.md. Проверено: jvmTest, compileKotlinLinuxX64, linuxX64Test. --- CONFIG.md | 69 +++++++++++ README.md | 8 +- TESTING.md | 1 + .../kotlin/pw/binom/llmproxy/Backoff.kt | 111 ++++++++++++++++++ .../kotlin/pw/binom/llmproxy/Main.kt | 77 ++++++++++-- .../kotlin/pw/binom/llmproxy/BackoffTest.kt | 109 +++++++++++++++++ .../pw/binom/llmproxy/ConfigLogicTest.kt | 27 +++-- 7 files changed, 378 insertions(+), 24 deletions(-) create mode 100644 src/commonMain/kotlin/pw/binom/llmproxy/Backoff.kt create mode 100644 src/commonTest/kotlin/pw/binom/llmproxy/BackoffTest.kt diff --git a/CONFIG.md b/CONFIG.md index d5568ea..ae1cab2 100644 --- a/CONFIG.md +++ b/CONFIG.md @@ -46,6 +46,10 @@ модели (фейловер); если свободных не осталось — отдаём `503` от прокси. `402` (insufficient balance / исчерпан лимит токенов) — ошибка аккаунта провайдера, но фейловер имеет смысл: следующий апстрим может быть платёжеспособен. +Повторяющиеся сбои уводятся из ротации экспоненциальным **backoff** (поле +`backoff`, ниже): упавший апстрим «отдыхает» 1s, 2s, 4s, … до заданного потолка, +и прокси его не трогает. Если **все** апстримы модели в откате — `503` с +заголовком `Retry-After`. ## Формат файла @@ -79,6 +83,8 @@ providers: patch: # уровень провайдера: ко всем его запросам provider: allow_fallbacks: false + backoff: P1M # опционально; потолок экспоненциального backoff + # на весь провайдер (ISO-8601: P1M, PT15M, P1D…) - id: local-llama url: "http://10.0.0.5:8080/v1" @@ -95,6 +101,8 @@ upstreams: patch: # уровень апстрима provider: ignore: [deepseek] + backoff: PT15M # опционально; собственный backoff этой модели + # (приоритет над backoff провайдера) - id: local-qwen provider: local-llama @@ -258,6 +266,58 @@ providers: # reasoning_empty_ok: true # опционально; дефолт false ``` +### Экспоненциальный backoff (`backoff`) + +Если апстрим регулярно ошибается (`5xx`/`429`/`402` или сетевые ошибки), +прокси не бьёт по нему на каждом запросе, а отправляет в **откат** +(cooldown) — это паттерн экспоненциального бэкоффа / circuit breaker: +чем дольше сервис молчит, тем дольше мы к нему не ходим. Пока апстрим в +откате, роутер пропускает его и берёт следующий по списку модели. + +| Поле | Тип / дефолт | Значение | +|---|---|---| +| `providers[].backoff` | ISO-8601-длительность / отсутствует | Потолок отката **на весь провайдер**: счётчик общий для всех его моделей | +| `upstreams[].backoff` | ISO-8601-длительность / отсутствует | Потолок отката **на конкретную модель**: счётчик индивидуальный | + +**Приоритет — модели.** Если `upstreams[].backoff` задан, у этой модели +собственный счётчик и потолок (провайдерский `backoff` на неё не действует). +Если у модели не задан, но задан у провайдера — счётчик общий на провайдера: +сбой на одной модели охлаждает и все остальные модели этого провайдера. +Если не задан нигде — backoff для этого апстрима выключен (остаются только +конкурентность и фейловер). + +Поведение: + +- **первая** ошибка → отдых 1s; каждая следующая **удваивает** интервал + (2s, 4s, 8s, …) до заданного потолка (cap); +- **успешный** запрос сбрасывает счётчик (откат и удвоение начинаются заново); +- апстрим в откате пропускается при выборе (лог: + `upstream= в backoff-откате (~Ns) — пропускаю`); +- если **все** апстримы модели в откате — прокси отдаёт `503` + (`all upstreams cooling, retry in Ns`) с заголовком `Retry-After: N` + (секунд до выхода первого апстрима из отката). + +Значение — ISO-8601-длительность (формат, в котором сериализуется +`kotlin.time.Duration`): `PT30M` (30 минут), `PT1H15M`, `P1D` (сутки), +`P1M` (месяц). Некорректное значение — предупреждение в лог и поле +просто игнорируется. + +```yaml +providers: + - id: routerai + url: "https://routerai.ru/api/v1" + backoff: P1M # потолок на весь провайдер + +upstreams: + - id: routerai-gpt4o + provider: routerai + model: gpt-4o + backoff: PT15M # у модели своё: потолок 15m, провайдерский P1M не действует +``` + +При старте выводится, что настроено: +`[llm-proxy] backoff: upstreams=routerai-gpt4o=PT15M providers=routerai=P1M`. + ### Пример сборки тела (многослойный `patch`) Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим @@ -342,6 +402,7 @@ data class ProviderConf( val think_tags: String? = null, val reasoning_field: String? = null, val reasoning_empty_ok: Boolean = false, + val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601) ) @Serializable @@ -351,6 +412,8 @@ data class UpstreamConf( val model: String, // реальное имя модели у провайдера val max_concurrency: Int? = null,// опционально; null/0 = безлимит val patch: JsonObject? = null, // к запросам этой апстрим-модели + val think_tags: String? = null, // переопределение think-режима модели + val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским) ) @Serializable @@ -511,6 +574,12 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet() освобождён слот `active[id]=N/limit`. - **Ошибка апстрима** (не `5xx`/`429`/`402`, а сетевая/таймаут): как сейчас — лог с сообщением. +- **Backoff**: при фейловере строка дополняется интервалом отдыха + (`(backoff: отдых PT…)`); при пропуске апстрима в откате — + `upstream= в backoff-откате (~Ns) — пропускаю`; при старте — список + настроенных backoff (`backoff: upstreams=… providers=…`). +- **Все в откате** — `503 all upstreams cooling, retry in Ns` для + `model=<витрина>` с заголовком `Retry-After: N`. Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать существующий парсинг логов (если он есть). diff --git a/README.md b/README.md index 0e553ec..90023a9 100644 --- a/README.md +++ b/README.md @@ -20,10 +20,16 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` ( обоих — безлимит. Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или -апстрима (`off` по умолчанию, `split` — рассуждения из `…` уходят +апстрима (`off` по умолчанию, `split` — рассуждения из `think`-тегов уходят в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в non-stream. +Повторяющиеся сбои апстрима увязываются экспоненциальным backoff: поле +`backoff` (ISO-8601-потолок, напр. `P1M`) задаётся на провайдере (общий +счётчик на его модели) и/или на апстриме (приоритет). Первая ошибка — отдых +1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате +— `503` с `Retry-After`. + | Переменная | Default | Описание | |---|---|---| | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | diff --git a/TESTING.md b/TESTING.md index 0875ab1..128a1e2 100644 --- a/TESTING.md +++ b/TESTING.md @@ -20,6 +20,7 @@ | `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; | | `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. | | `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. | +| `BackoffTest` | Экспоненциальный backoff: удвоение интервала отката до потолка (cap), сброс счётчика при успехе, окно охлаждения (`isCoolingAt`/`remainingAt`), ISO-8601-разбор cap (`PT30M`, `P1D`, `PT1M30S`), скоупы: общий провайдерский счётчик на все его модели, приоритет `upstreams[].backoff` над провайдерским, независимость скоупов, backoff не настроен → без отката. | ## Проверка качества тестов (мутационная приёмка) diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Backoff.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Backoff.kt new file mode 100644 index 0000000..d2c9212 --- /dev/null +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Backoff.kt @@ -0,0 +1,111 @@ +package pw.binom.llmproxy + +import kotlinx.coroutines.sync.Mutex +import kotlinx.datetime.Clock +import kotlinx.datetime.Instant +import kotlin.concurrent.Volatile +import kotlin.time.Duration +import kotlin.time.Duration.Companion.seconds + +/** + * Экспоненциальный бэкофф для одной области отдыха: первая ошибка — отдых + * [base] (1 с), каждая следующая удваивает интервал вплоть до [cap]. + * Успешный запрос сбрасывает счётчик и отдых. + */ +class BackoffGuard(val cap: Duration, val base: Duration = 1.seconds) { + private val lock = Mutex() + @Volatile + private var failures: Int = 0 + @Volatile + private var interval: Duration = Duration.ZERO + @Volatile + private var coolingUntil: Instant = Instant.DISTANT_PAST + + fun isCoolingAt(now: Instant): Boolean = now < coolingUntil + + /** Сколько осталось до выхода из отдыха (0 — не в откате). */ + fun remainingAt(now: Instant): Duration = + if (coolingUntil > now) coolingUntil - now else Duration.ZERO + + /** Учёт ошибки; возвращает интервал, на который область уходит в отдых. */ + suspend fun recordFailure(now: Instant): Duration { + lock.lock() + try { + failures++ + interval = if (failures == 1) base else interval * 2 + if (interval > cap) interval = cap + coolingUntil = now + interval + } finally { + lock.unlock() + } + return interval + } + + /** Успех: сброс счётчика и отдыха. */ + suspend fun recordSuccess() { + lock.lock() + try { + failures = 0 + interval = Duration.ZERO + coolingUntil = Instant.DISTANT_PAST + } finally { + lock.unlock() + } + } +} + +/** + * Реестр областей бэкофф-отдыха. Область апстрима: + * - у модели (записи upstreams) задан `backoff` — свой счётчик и потолок (приоритет); + * - иначе у провайдера задан `backoff` — общий для всех моделей провайдера счётчик; + * - ни там, ни там — бэкофф для этого апстрима выключен. + */ +class BackoffRegistry( + private val providers: Map, + private val upstreams: Map, +) { + private val upstreamGuards = mutableMapOf() + private val providerGuards = mutableMapOf() + private val lock = Mutex() + + private fun guardForLocked(up: UpstreamConf): BackoffGuard? { + up.backoff?.let { cap -> + return upstreamGuards.getOrPut(up.id) { BackoffGuard(cap) } + } + providers[up.provider]?.backoff?.let { cap -> + return providerGuards.getOrPut(up.provider) { BackoffGuard(cap) } + } + return null + } + + private suspend fun locked(block: () -> T?): T? { + lock.lock() + return try { + block() + } finally { + lock.unlock() + } + } + + suspend fun isCooling(up: UpstreamConf): Boolean { + val g = locked { guardForLocked(up) } ?: return false + return g.isCoolingAt(Clock.System.now()) + } + + /** Секунд до выхода апстрима из отдыха (0 — не в откате). */ + suspend fun remainingSeconds(up: UpstreamConf): Long { + val g = locked { guardForLocked(up) } ?: return 0L + return g.remainingAt(Clock.System.now()).inWholeSeconds + } + + /** Учёт ошибки; возвращает интервал отдыха, или null, если бэкофф для апстрима выключен. */ + suspend fun recordFailure(up: UpstreamConf): Duration? { + val g = locked { guardForLocked(up) } ?: return null + return g.recordFailure(Clock.System.now()) + } + + /** Успех: сброс области апстрима. */ + suspend fun recordSuccess(up: UpstreamConf) { + locked { guardForLocked(up) }?.recordSuccess() + } +} diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index c2bd32a..b0d579e 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -46,6 +46,7 @@ import net.mamoe.yamlkt.YamlList import net.mamoe.yamlkt.YamlLiteral import net.mamoe.yamlkt.YamlMap import kotlin.concurrent.Volatile +import kotlin.time.Duration import kotlin.time.Duration.Companion.seconds import kotlin.time.TimeSource @@ -86,11 +87,18 @@ fun main() { val active = config.upstreams.associate { up -> up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider])) } + val backoff = BackoffRegistry(providersById, upstreamsById) log.info { "[llm-proxy] загружено: providers=${config.providers.size}, " + "upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)" } + log.info { + "[llm-proxy] backoff: upstreams=" + + config.upstreams.filter { it.backoff != null }.map { up -> "${up.id}=${up.backoff?.toIsoString()}" }.joinToString(", ") + + " providers=" + + config.providers.filter { it.backoff != null }.map { p -> "${p.id}=${p.backoff?.toIsoString()}" }.joinToString(", ") + } log.info { "[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}" } @@ -101,7 +109,7 @@ fun main() { val http = createHttpClient() val sessions = SessionRegistry() startServer(config.server.host, config.server.port) { - proxyModule(config, providersById, upstreamsById, active, sessions, http) + proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff) } } @@ -114,10 +122,11 @@ fun Application.proxyModule( active: Map, sessions: SessionRegistry, http: HttpClient, + backoff: BackoffRegistry, ) { routing { post("/v1/chat/completions") { - handleChat(call, config, providersById, upstreamsById, active, sessions, http) + handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff) } get("/v1/models") { handleModels(call, config) @@ -133,6 +142,7 @@ private suspend fun handleChat( active: Map, sessions: SessionRegistry, http: HttpClient, + backoff: BackoffRegistry, ) { log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" } val raw = call.receiveText() @@ -165,7 +175,7 @@ private suspend fun handleChat( var anyClaimed = false while (true) { - val up = pickFreeUpstream(pool, active, failed) ?: break + val up = pickFreeUpstream(pool, active, failed, backoff, modelName) ?: break anyClaimed = true val start = TimeSource.Monotonic.markNow() try { @@ -230,7 +240,11 @@ private suspend fun handleChat( }.execute { resp -> upstreamStatus = resp.status.value if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) { - log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" } + val rest = backoff.recordFailure(up) + log.warn { + "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" + + (rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "") + } failed.add(up.id) failover = true return@execute @@ -245,6 +259,7 @@ private suspend fun handleChat( return@execute } responded = true + backoff.recordSuccess(up) if (clientWantsStream) { val ct = resp.headers["Content-Type"] ?: "text/event-stream" val status = HttpStatusCode.fromValue(upstreamStatus) @@ -294,7 +309,11 @@ private suspend fun handleChat( log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" } throw e } catch (e: Exception) { - log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" } + val rest = backoff.recordFailure(up) + log.error { + "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" + + (rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "") + } failed.add(up.id) continue } finally { @@ -305,7 +324,17 @@ private suspend fun handleChat( if (anyClaimed) { call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) } else { - call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) + val retryIn = pool.map { backoff.remainingSeconds(it) }.filter { it > 0 }.minOrNull() + if (retryIn != null) { + call.response.headers.append("Retry-After", retryIn.toString()) + call.respondText( + errorJson("all upstreams cooling, retry in ${retryIn}s"), + ContentType.Application.Json, + HttpStatusCode.ServiceUnavailable, + ) + } else { + call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) + } } } @@ -489,15 +518,28 @@ class UpstreamCounter(private val limit: Int, initial: Int = 0) { /** * Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`, - * у которого свободен слот и который ещё не в `excluded` (не упал ранее). - * Сразу занимает слот (через [tryClaim]). Если свободных нет — возвращает null. + * у которого свободен слот, который ещё не в `excluded` (не упал ранее) и + * не в backoff-откате. Сразу занимает слот (через [tryClaim]). + * Если свободных нет — возвращает null. */ -internal fun pickFreeUpstream( +internal suspend fun pickFreeUpstream( pool: List, active: Map, excluded: Set, -): UpstreamConf? = - pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) } + backoff: BackoffRegistry, + modelName: String = "?", +): UpstreamConf? { + for (up in pool) { + if (up.id in excluded) continue + val rest = backoff.remainingSeconds(up) + if (rest > 0) { + log.info { "[llm-proxy] model=$modelName upstream=${up.id} в backoff-откате (~${rest}c) — пропускаю" } + continue + } + if (tryClaim(up, active)) return up + } + return null +} private suspend fun handleModels(call: ApplicationCall, config: Config) { log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" } @@ -903,6 +945,7 @@ data class ProviderConf( val think_tags: String? = null, val reasoning_field: String? = null, val reasoning_empty_ok: Boolean = false, + val backoff: Duration? = null, ) data class UpstreamConf( @@ -912,6 +955,7 @@ data class UpstreamConf( val max_concurrency: Int? = null, val patch: JsonObject? = null, val think_tags: String? = null, + val backoff: Duration? = null, ) data class ModelConf( @@ -969,6 +1013,7 @@ internal fun parseConfig(root: YamlElement): Config { think_tags = m.strOrNull("think_tags"), reasoning_field = m.strOrNull("reasoning_field"), reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false, + backoff = m.durationOrNull("backoff"), ) } @@ -981,6 +1026,7 @@ internal fun parseConfig(root: YamlElement): Config { max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, think_tags = m.strOrNull("think_tags"), + backoff = m.durationOrNull("backoff"), ) } @@ -1006,5 +1052,14 @@ internal fun Map.str(key: String): String = internal fun Map.strOrNull(key: String): String? = (this[key] as? YamlLiteral)?.content +/** ISO-8601-длительность (PT30M, P1M, PT1H15M…); некорректное значение — warn + null. */ +internal fun Map.durationOrNull(key: String): Duration? = + strOrNull(key)?.let { raw -> + runCatching { Duration.parse(raw) }.getOrElse { + log.warn { "[llm-proxy] config: поле '$key' — некорректная ISO-8601-длительность '$raw', проигнорировано" } + null + } + } + internal fun Map.yamlMapOrNull(key: String): YamlMap? = this[key] as? YamlMap diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/BackoffTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/BackoffTest.kt new file mode 100644 index 0000000..1ce4ab4 --- /dev/null +++ b/src/commonTest/kotlin/pw/binom/llmproxy/BackoffTest.kt @@ -0,0 +1,109 @@ +package pw.binom.llmproxy + +import kotlinx.coroutines.test.runTest +import kotlinx.datetime.Instant +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNull +import kotlin.test.assertTrue +import kotlin.time.Duration +import kotlin.time.Duration.Companion.days +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.minutes +import kotlin.time.Duration.Companion.seconds + +class BackoffTest { + + private val t0 = Instant.fromEpochSeconds(1_700_000_000) + + @Test + fun intervalsDoubleUntilCap() = runTest { + val g = BackoffGuard(cap = 30.seconds) + val seq = (1..7).map { g.recordFailure(t0) } + assertEquals( + listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds, 30.seconds), + seq, + ) + } + + @Test + fun successResetsCounter() = runTest { + val g = BackoffGuard(cap = 30.seconds) + g.recordFailure(t0) + g.recordFailure(t0) + g.recordSuccess() + assertEquals(1.seconds, g.recordFailure(t0)) + assertFalse(g.isCoolingAt(t0 + 1.days)) + } + + @Test + fun coolingWindow() = runTest { + val g = BackoffGuard(cap = 30.seconds) + g.recordFailure(t0) // отдых 1s + assertTrue(g.isCoolingAt(t0 + 500.milliseconds)) + assertFalse(g.isCoolingAt(t0 + 1.seconds)) + assertEquals(500.milliseconds, g.remainingAt(t0 + 500.milliseconds)) + } + + @Test + fun iso8601CapParsing() { + assertEquals(30.minutes, Duration.parse("PT30M")) + assertEquals(1.days, Duration.parse("P1D")) + assertEquals(90.seconds, Duration.parse("PT1M30S")) + } + + @Test + fun providerScopeSharedBetweenItsModels() = runTest { + val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes) + val a = UpstreamConf(id = "a", provider = "p", model = "m1") + val b = UpstreamConf(id = "b", provider = "p", model = "m2") + val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b)) + + assertEquals(1.seconds, reg.recordFailure(a)) + + // общая скоуп-переменная провайдера: сбой на a охладил и b + assertTrue(reg.isCooling(a)) + assertTrue(reg.isCooling(b)) + } + + @Test + fun upstreamBackoffOverridesProvider() = runTest { + val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes) + val a = UpstreamConf(id = "a", provider = "p", model = "m1") + val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.seconds) + val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b)) + + // своё backoff у b: cap 30s (не 5m провайдера): 1s, 2s, 4s, 8s, 16s, 30s (потолок) + val seq = (1..6).map { reg.recordFailure(b) } + assertEquals(listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds), seq) + // скоупы независимы: сбои b не охлаждают a (провайдерский скоуп) + assertFalse(reg.isCooling(a)) + assertTrue(reg.isCooling(b)) + } + + @Test + fun providerFailureDoesNotCoolUpstreamWithOwnBackoff() = runTest { + val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes) + val a = UpstreamConf(id = "a", provider = "p", model = "m1") + val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.minutes) + val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b)) + + // провайдерский скоуп заведён сбоем на a: + reg.recordFailure(a) + assertTrue(reg.isCooling(a)) + // b живёт своим (provider-scope на него не действует): + assertFalse(reg.isCooling(b)) + } + + @Test + fun noBackoffConfiguredMeansNoCooling() = runTest { + val p = ProviderConf(id = "p", url = "https://p") + val a = UpstreamConf(id = "a", provider = "p", model = "m1") + val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a)) + + assertNull(reg.recordFailure(a)) + assertFalse(reg.isCooling(a)) + assertEquals(0L, reg.remainingSeconds(a)) + } +} diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt index a75f9bc..b5a9443 100644 --- a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt @@ -5,6 +5,7 @@ import kotlin.test.assertEquals import kotlin.test.assertFalse import kotlin.test.assertTrue import io.ktor.http.headersOf +import kotlinx.coroutines.test.runTest import kotlinx.serialization.json.Json import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.jsonArray @@ -14,6 +15,8 @@ import net.mamoe.yamlkt.Yaml class ConfigLogicTest { + private val noBackoff = BackoffRegistry(emptyMap(), emptyMap()) + @Test fun mergeDeepMergesNestedObjectsAndReplacesScalars() { val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject @@ -222,58 +225,58 @@ class ConfigLogicTest { } @Test - fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() { + fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() = runTest { val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 2, null), ) - val up = pickFreeUpstream(pool, active, emptySet()) + val up = pickFreeUpstream(pool, active, emptySet(), noBackoff) assertEquals("u1", up?.id) // слот реально занят assertEquals(1, active.getValue("u1").current) } @Test - fun pickFreeUpstreamSkipsExcluded() { + fun pickFreeUpstreamSkipsExcluded() = runTest { val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 2, null), ) - val up = pickFreeUpstream(pool, active, setOf("u1")) + val up = pickFreeUpstream(pool, active, setOf("u1"), noBackoff) assertEquals("u2", up?.id) } @Test - fun pickFreeUpstreamReturnsNullWhenAllBusy() { + fun pickFreeUpstreamReturnsNullWhenAllBusy() = runTest { val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1 val pool = listOf(UpstreamConf("u1", "p", "m", 1, null)) - assertEquals(null, pickFreeUpstream(pool, active, emptySet())) + assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff)) } @Test - fun pickFreeUpstreamImplementsFailoverOrder() { + fun pickFreeUpstreamImplementsFailoverOrder() = runTest { // dead исключён (упал ранее) — выбирается следующий живой u1 val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1)) val pool = listOf( UpstreamConf("dead", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null), ) - val up = pickFreeUpstream(pool, active, setOf("dead")) + val up = pickFreeUpstream(pool, active, setOf("dead"), noBackoff) assertEquals("u1", up?.id) } @Test - fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() { + fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() = runTest { val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1)) val pool = listOf( UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 1, null), ) - assertEquals("u1", pickFreeUpstream(pool, active, emptySet())?.id) - assertEquals("u2", pickFreeUpstream(pool, active, emptySet())?.id) - assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id) + assertEquals("u1", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id) + assertEquals("u2", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id) + assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id) } @Test