From c6a24a94b02afacee70b525fafa442d42273204a Mon Sep 17 00:00:00 2001 From: subochev Date: Tue, 15 Sep 2026 01:00:10 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20=D1=84=D0=BE=D0=BD=D0=BE=D0=B2=D1=8B?= =?UTF-8?q?=D0=B9=20=D0=BF=D0=B8=D0=BD=D0=B3=20=D0=B0=D0=BF=D1=81=D1=82?= =?UTF-8?q?=D1=80=D0=B8=D0=BC=D0=BE=D0=B2=20=D0=B2=20backoff-=D0=BE=D1=82?= =?UTF-8?q?=D0=BA=D0=B0=D1=82=D0=B5=20(probe=5Finterval)=20=E2=80=94=20?= =?UTF-8?q?=D0=BC=D0=B3=D0=BD=D0=BE=D0=B2=D0=B5=D0=BD=D0=BD=D0=BE=D0=B5=20?= =?UTF-8?q?=D0=B2=D0=BE=D0=B7=D0=B2=D1=80=D0=B0=D1=89=D0=B5=D0=BD=D0=B8?= =?UTF-8?q?=D0=B5=20=D0=B2=20=D0=BF=D1=83=D0=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CONFIG.md | 69 ++++++++ README.md | 8 + .../kotlin/pw/binom/llmproxy/HealthProber.kt | 133 ++++++++++++++++ .../kotlin/pw/binom/llmproxy/Main.kt | 27 ++++ .../kotlin/pw/binom/llmproxy/Metrics.kt | 30 ++++ .../kotlin/pw/binom/llmproxy/Timeouts.kt | 9 ++ .../pw/binom/llmproxy/ConfigLogicTest.kt | 28 ++++ .../pw/binom/llmproxy/HealthProberTest.kt | 148 ++++++++++++++++++ 8 files changed, 452 insertions(+) create mode 100644 src/commonMain/kotlin/pw/binom/llmproxy/HealthProber.kt create mode 100644 src/commonTest/kotlin/pw/binom/llmproxy/HealthProberTest.kt diff --git a/CONFIG.md b/CONFIG.md index 637aea9..96390d5 100644 --- a/CONFIG.md +++ b/CONFIG.md @@ -85,6 +85,8 @@ providers: allow_fallbacks: false backoff: P1M # опционально; потолок экспоненциального backoff # на весь провайдер (ISO-8601: P1M, PT15M, P1D…) + probe_interval: PT1S # опционально; фоновый пинг апстримов в откате + # (как часто проверять, жив ли провайдер) - id: local-llama url: "http://10.0.0.5:8080/v1" @@ -103,6 +105,8 @@ upstreams: ignore: [deepseek] backoff: PT15M # опционально; собственный backoff этой модели # (приоритет над backoff провайдера) + probe_interval: PT0S # опционально; пинг отключён у этой модели + # (приоритет над probe_interval провайдера) - id: local-qwen provider: local-llama @@ -385,6 +389,68 @@ upstreams: При старте выводится, что настроено: `[llm-proxy] first_byte_timeout: upstreams=routerai-gpt4o=PT0S providers=routerai=PT30S default=PT30S`. +### Фоновый пинг апстримов в откате (`probe_interval`) + +Когда апстрим «уходит в аут» надолго (backoff растёт до потолка), без пинга +он вернётся в ротацию только по истечении интервала отката — это может быть +минуты простоя, хотя модель уже могла ожить. `probe_interval` включает +фоновый «пингер»: пока апстрим в backoff-откате, HealthProber каждые +`probe_interval` шлёт ему минимальный запрос (`2+2=?`, `max_tokens=4`, +`stream=true`) и слушает первый байт ответа. Как только апстрим ответил — +его снимают с отката немедленно (`backoff.recordSuccess`), и он возвращается +в общий пул, не дожидаясь истечения backoff-интервала. + +| Поле | Тип / дефолт | Значение | +|---|---|---| +| `providers[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал пинга для всех моделей провайдера | +| `upstreams[].probe_interval` | ISO-8601-длительность / отсутствует | Интервал для конкретной модели (приоритет над провайдерским) | + +**Дефолт** (если не задан нигде): `PT1S`. + +**Приоритет — модели.** Если `upstreams[].probe_interval` задан, он +перекрывает провайдерский; если нет — берётся провайдерский; если нет и там — +дефолт `PT1S`. + +**Отключение:** `PT0S` на провайдере или апстриме отключает пинг для этого +апстрима — тогда он возвращается в пул только по истечении backoff-интервала. + +Особенности: + +- пингуются **только** апстримы, которые сейчас в backoff-откате; живой + апстрим пинг не получает (проверка состояния — каждые `probe_interval`); +- пинг **не занимает слот конкурентности** (идёт в обход `max_concurrency`), + чтобы не отжимать слот у живых запросов; +- ошибка пинга **не** считается в backoff (иначе каждая попытка удваивала + бы интервал отката); учитывается только успех; +- признак «ответил» — HTTP 2xx и первый байт тела за `first_byte_timeout` + (тот же watchdog, что и у обычных запросов); +- на пинг применяются патчи провайдера и апстрима (как к обычному запросу), + чтобы он валидировал реальный путь; +- при восстановлении — лог: + `[llm-proxy] probe upstream= ответил — снят с backoff, вернулся в пул`, + метрика `llm_proxy_probe_pings_total{result="ok"}`; +- при молчании — лог `[llm-proxy] probe upstream= молчит — остаюсь в backoff`, + метрика `llm_proxy_probe_pings_total{result="fail"}`. + +Значение — ISO-8601-длительность, формат как у `backoff` (`PT1S`, `PT500MS`, +`PT5S`, …). + +```yaml +providers: + - id: routerai + url: "https://routerai.ru/api/v1" + probe_interval: PT1S # пингуем все модели провайдера раз в секунду + +upstreams: + - id: routerai-gpt4o + provider: routerai + model: gpt-4o + probe_interval: PT0S # у модели пинг отключён (ждать истечения backoff) +``` + +При старте выводится, что настроено: +`[llm-proxy] probe_interval: upstreams=routerai-gpt4o=PT0S providers=routerai=PT1S default=PT1S`. + ### Prometheus-метрики (`/metrics`) Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics` @@ -395,6 +461,7 @@ upstreams: | Метрика | Тип | Смысл | |---|---|---| | `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `timeout` — апстрим не прислал первый байт стрима за `first_byte_timeout`, `cancelled` — клиент отвалился посреди стрима | +| `llm_proxy_probe_pings_total{upstream, provider, result}` | counter | фоновые пинги апстримов в backoff-откате (HealthProber). `result`: `ok` — апстрим ответил (снят с отката), `fail` — молчит (остался в откате) | | `llm_proxy_upstream_inflight{upstream, provider}` | gauge | занятые слоты апстрима прямо сейчас (конкурентность) | | `llm_proxy_upstream_fail_streak{upstream, provider}` | gauge | счётчик сбоев подряд (backoff): сколько раз подряд упал | | `llm_proxy_upstream_cooling_seconds{upstream, provider}` | gauge | сколько секунд апстрим ещё в backoff-откате (0 = жив) | @@ -494,6 +561,7 @@ data class ProviderConf( val reasoning_empty_ok: Boolean = false, val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601) val first_byte_timeout: Duration? = null, // окно первого байта стрима (ISO-8601) + val probe_interval: Duration? = null, // интервал фонового пинга (ISO-8601) ) data class UpstreamConf( @@ -505,6 +573,7 @@ data class UpstreamConf( val think_tags: String? = null, // переопределение think-режима модели val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским) val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским) + val probe_interval: Duration? = null, // интервал фонового пинга (приоритет над провайдерским) ) @Serializable diff --git a/README.md b/README.md index 2a3f326..74681fd 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,14 @@ non-stream. в backoff и роутер фейловерит на следующего свободного апстрима. Приоритет: `upstreams[].first_byte_timeout` → `providers[].first_byte_timeout` → `PT30S`. +Пока апстрим «отдыхает» в backoff-откате, его можно фоново пинговать: +`probe_interval` (ISO-8601, дефолт `PT1S`, `PT0S` отключает) — HealthProber +шлёт апстриму минимальный запрос (`2+2=?`, `max_tokens=4`) и, как только тот +ответил, снимает его с отката немедленно, не дожидаясь истечения +backoff-интервала. Пинг не занимает слот конкурентности и не засчитывается +в backoff. Приоритет: `upstreams[].probe_interval` → +`providers[].probe_interval` → `PT1S`. + Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по `model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) — подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/HealthProber.kt b/src/commonMain/kotlin/pw/binom/llmproxy/HealthProber.kt new file mode 100644 index 0000000..7e139b5 --- /dev/null +++ b/src/commonMain/kotlin/pw/binom/llmproxy/HealthProber.kt @@ -0,0 +1,133 @@ +package pw.binom.llmproxy + +import io.ktor.client.HttpClient +import io.ktor.client.call.body +import io.ktor.client.request.headers +import io.ktor.client.request.preparePost +import io.ktor.client.request.setBody +import io.ktor.utils.io.ByteReadChannel +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Job +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import kotlinx.serialization.json.JsonArray +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.JsonPrimitive +import kotlin.time.Duration + +/** + * Фоновый «пингер» апстримов, которые в данный момент в backoff-откате. + * + * Пока апстрим молчит (в откате), HealthProber раз в [interval] шлёт ему + * минимальный chat.completion (`2+2=?`, `max_tokens=4`) и слушает первый байт + * ответа. Как только апстрим ответил — вызывается [BackoffRegistry.recordSuccess], + * и он возвращается в общий пул сразу, не дожидаясь истечения backoff-интервала. + * Это ускоряет восстановление после долгого простоя провайдера. + * + * Пинг НЕ занимает слот конкурентности ([UpstreamCounter]) — он идёт в обход, + * чтобы не отжимать слот у живых запросов. Ошибки пинга НЕ считаются в backoff + * (иначе каждая попытка удваивала бы интервал отката); считается только успех. + */ +class HealthProber( + private val providersById: Map, + private val upstreams: List, + private val http: HttpClient, + private val backoff: BackoffRegistry, + private val metrics: MetricsRegistry, + private val scope: CoroutineScope, +) { + private val jobs = mutableMapOf() + + /** Запустить пинг-корутины для всех апстримов с включённым пингом. */ + fun start() { + for (up in upstreams) { + val interval = effectiveProbeInterval(up, providersById[up.provider]) + if (interval != null) { + jobs[up.id] = scope.launch { probeLoop(up, interval) } + } + } + } + + /** Цикл пинга одного апстрима: пингуем только пока он в откате. */ + internal suspend fun probeLoop(up: UpstreamConf, interval: Duration, pingImpl: suspend (UpstreamConf) -> Boolean = ::ping) { + while (currentCoroutineContext().isActive) { + if (backoff.isCooling(up)) { + val healthy = pingImpl(up) + metrics.recordProbe(up, if (healthy) "ok" else "fail") + if (healthy) { + backoff.recordSuccess(up) + log.info { "[llm-proxy] probe upstream=${up.id} ответил — снят с backoff, вернулся в пул" } + } else { + log.info { "[llm-proxy] probe upstream=${up.id} молчит — остаюсь в backoff" } + } + } + delay(interval) + } + } + + /** Одна пинг-попытка: минимальный stream-запрос, слушаем первый байт. */ + private suspend fun ping(up: UpstreamConf): Boolean { + val provider = providersById[up.provider] ?: return false + val url = provider.url.trimEnd('/') + "/chat/completions" + val providerKey = resolveEnv(provider.key) + val body = buildProbeBody(up, provider) + val firstByteTimeout = effectiveFirstByteTimeout(up, provider) ?: DEFAULT_FIRST_BYTE_TIMEOUT + return try { + http.preparePost(url) { + headers { + if (providerKey.isNotEmpty()) { + set("Authorization", "Bearer $providerKey") + } + set("Content-Type", "application/json") + } + setBody(body.toString()) + }.execute { resp -> + if (resp.status.value !in 200..299) { + return@execute false + } + val ch = resp.body() + val buf = ByteArray(64) + readAvailableWithTimeout(ch, buf, firstByteTimeout.inWholeMilliseconds) > 0 + } + } catch (e: CancellationException) { + throw e + } catch (e: Exception) { + log.info { "[llm-proxy] probe upstream=${up.id} ошибка: ${e.message}" } + false + } + } +} + +/** + * Собрать тело пинг-запроса: минимальный chat.completion, на который любой + * живой апстрим ответит сразу. Применяются патчи провайдера и апстрима (в том + * же порядке, что и для реальных запросов), чтобы пинг валидировал реальный + * путь запроса, но без клиентских данных и модели. + */ +internal fun buildProbeBody(up: UpstreamConf, provider: ProviderConf?): JsonObject { + val base = JsonObject( + mapOf( + "model" to JsonPrimitive(up.model), + "messages" to JsonArray( + listOf( + JsonObject( + mapOf( + "role" to JsonPrimitive("user"), + "content" to JsonPrimitive("2+2=?"), + ), + ), + ), + ), + "max_tokens" to JsonPrimitive(4), + "stream" to JsonPrimitive(true), + ), + ) + var acc = base + listOf(provider?.patch, up.patch).filterNotNull().forEach { patch -> + acc = merge(acc, patch) + } + return acc +} \ No newline at end of file diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt index 466401b..c32df96 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Main.kt @@ -28,6 +28,9 @@ import io.ktor.utils.io.ByteReadChannel import io.ktor.utils.io.ByteWriteChannel import io.ktor.utils.io.writeFully import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.async import kotlinx.coroutines.cancelAndJoin import kotlinx.coroutines.coroutineScope @@ -115,6 +118,13 @@ fun main() { config.providers.filter { it.first_byte_timeout != null }.map { p -> "${p.id}=${p.first_byte_timeout?.toIsoString()}" }.joinToString(", ") + " default=${DEFAULT_FIRST_BYTE_TIMEOUT.toIsoString()}" } + log.info { + "[llm-proxy] probe_interval: upstreams=" + + config.upstreams.filter { it.probe_interval != null }.map { up -> "${up.id}=${up.probe_interval?.toIsoString()}" }.joinToString(", ") + + " providers=" + + config.providers.filter { it.probe_interval != null }.map { p -> "${p.id}=${p.probe_interval?.toIsoString()}" }.joinToString(", ") + + " default=${DEFAULT_PROBE_INTERVAL.toIsoString()}" + } log.info { "[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}" } @@ -124,6 +134,8 @@ fun main() { val http = createHttpClient() val sessions = SessionRegistry() + val prober = HealthProber(providersById, config.upstreams, http, backoff, metrics, CoroutineScope(SupervisorJob() + Dispatchers.Default)) + prober.start() startServer(config.server.host, config.server.port) { proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics) } @@ -527,6 +539,17 @@ internal fun effectiveFirstByteTimeout(up: UpstreamConf, provider: ProviderConf? return if (raw <= Duration.ZERO) null else raw } +/** + * Эффективный интервал фонового пинга апстрима в backoff-откате: значение у + * апстрима, если задано; иначе у провайдера; иначе [DEFAULT_PROBE_INTERVAL]. + * Возвращает null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда + * HealthProber этот апстрим не пингует (ждём естественного истечения бэкоффа). + */ +internal fun effectiveProbeInterval(up: UpstreamConf, provider: ProviderConf?): Duration? { + val raw = up.probe_interval ?: provider?.probe_interval ?: DEFAULT_PROBE_INTERVAL + return if (raw <= Duration.ZERO) null else raw +} + /** * Сигнализирует, что апстрим не прислал ни одного байта тела стрим-ответа за * [timeoutMs] мс. Обрабатывается как обычный сбой (фейловер, backoff, @@ -1083,6 +1106,7 @@ data class ProviderConf( val reasoning_empty_ok: Boolean = false, val backoff: Duration? = null, val first_byte_timeout: Duration? = null, + val probe_interval: Duration? = null, ) data class UpstreamConf( @@ -1094,6 +1118,7 @@ data class UpstreamConf( val think_tags: String? = null, val backoff: Duration? = null, val first_byte_timeout: Duration? = null, + val probe_interval: Duration? = null, ) data class ModelConf( @@ -1153,6 +1178,7 @@ internal fun parseConfig(root: YamlElement): Config { reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false, backoff = m.durationOrNull("backoff"), first_byte_timeout = m.durationOrNull("first_byte_timeout"), + probe_interval = m.durationOrNull("probe_interval"), ) } @@ -1167,6 +1193,7 @@ internal fun parseConfig(root: YamlElement): Config { think_tags = m.strOrNull("think_tags"), backoff = m.durationOrNull("backoff"), first_byte_timeout = m.durationOrNull("first_byte_timeout"), + probe_interval = m.durationOrNull("probe_interval"), ) } diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt index fb27f57..4eaedd0 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Metrics.kt @@ -10,6 +10,8 @@ import kotlinx.coroutines.sync.withLock * Метрики (всё в памяти, сброс при рестарте): * - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик * chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | timeout | cancelled`; + * - `llm_proxy_probe_pings_total{upstream,provider,result}` — счётчик фоновых + * пингов апстримов в backoff-откате (HealthProber); `result`: `ok | fail`; * - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты); * - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев * (счётчик backoff); @@ -19,6 +21,7 @@ import kotlinx.coroutines.sync.withLock class MetricsRegistry { private val lock = Mutex() private val requests = HashMap() + private val probes = HashMap() private data class RequestKey( val model: String, @@ -27,6 +30,12 @@ class MetricsRegistry { val result: String, ) + private data class ProbeKey( + val upstream: String, + val provider: String, + val result: String, + ) + /** Учёт chat-запроса по исходу (значения `result` — в описании класса). */ suspend fun record(model: String, up: UpstreamConf, result: String) { lock.withLock { @@ -35,6 +44,14 @@ class MetricsRegistry { } } + /** Учёт пинга апстрима в backoff-откате (result: `ok`/`fail`). */ + suspend fun recordProbe(up: UpstreamConf, result: String) { + lock.withLock { + val key = ProbeKey(up.id, up.provider, result) + probes[key] = (probes[key] ?: 0L) + 1 + } + } + /** Текстовая выгрузка в формате Prometheus на текущий момент. */ suspend fun render( config: Config, @@ -56,6 +73,19 @@ class MetricsRegistry { ) } } + sb.append("# HELP llm_proxy_probe_pings_total Фоновые пинги апстримов в backoff-откате\n") + sb.append("# TYPE llm_proxy_probe_pings_total counter\n") + lock.withLock { + probes.entries + .sortedWith(compareBy({ it.key.upstream }, { it.key.provider }, { it.key.result })) + .forEach { (k, n) -> + sb.append( + "llm_proxy_probe_pings_total{upstream=\"" + esc(k.upstream) + + "\",provider=\"" + esc(k.provider) + + "\",result=\"" + esc(k.result) + "\"} $n\n", + ) + } + } sb.append("# HELP llm_proxy_upstream_inflight Занято слотов апстрима сейчас\n") sb.append("# TYPE llm_proxy_upstream_inflight gauge\n") for (up in config.upstreams) { diff --git a/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt b/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt index 9c330f9..00621f2 100644 --- a/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt +++ b/src/commonMain/kotlin/pw/binom/llmproxy/Timeouts.kt @@ -15,3 +15,12 @@ const val UPSTREAM_REQUEST_TIMEOUT_MS: Long = 5 * 60_000L * Отключается явным `first_byte_timeout: PT0S` на апстриме или провайдере. */ val DEFAULT_FIRST_BYTE_TIMEOUT: Duration = 30.seconds + +/** + * Default интервал фоновой проверки («пинга») апстрима, который сейчас в + * backoff-откате: пока апстрим молчит, HealthProber шлёт ему минимальный + * запрос (2+2=?, max_tokens=4) каждые [DEFAULT_PROBE_INTERVAL], и как только + * он отвечает — снимает его с отката раньше, чем истечёт интервал бэкоффа. + * Отключается явным `probe_interval: PT0S` на апстриме или провайдере. + */ +val DEFAULT_PROBE_INTERVAL: Duration = 1.seconds diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt index 71979bb..329dc4b 100644 --- a/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt +++ b/src/commonTest/kotlin/pw/binom/llmproxy/ConfigLogicTest.kt @@ -615,4 +615,32 @@ class ConfigLogicTest { assertEquals(0.seconds, cfg.upstreams[0].first_byte_timeout) assertEquals(null, cfg.upstreams[1].first_byte_timeout) } + + @Test + fun parseConfigReadsProbeIntervalOnProviderAndUpstream() { + val yaml = """ + providers: + - id: p1 + url: "https://x.ru/api/v1" + probe_interval: PT2S + - id: p2 + url: "https://y.ru/api/v1" + upstreams: + - id: u1 + provider: p1 + model: real-1 + probe_interval: PT0S + - id: u2 + provider: p2 + model: real-2 + models: + - name: m1 + upstreams: [u1, u2] + """.trimIndent() + val cfg = parseConfig(Yaml.decodeYamlFromString(yaml)) + assertEquals(2.seconds, cfg.providers[0].probe_interval) + assertEquals(null, cfg.providers[1].probe_interval) + assertEquals(0.seconds, cfg.upstreams[0].probe_interval) + assertEquals(null, cfg.upstreams[1].probe_interval) + } } diff --git a/src/commonTest/kotlin/pw/binom/llmproxy/HealthProberTest.kt b/src/commonTest/kotlin/pw/binom/llmproxy/HealthProberTest.kt new file mode 100644 index 0000000..031c237 --- /dev/null +++ b/src/commonTest/kotlin/pw/binom/llmproxy/HealthProberTest.kt @@ -0,0 +1,148 @@ +package pw.binom.llmproxy + +import kotlinx.coroutines.delay +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.TestScope +import kotlinx.coroutines.test.runTest +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.jsonArray +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.jsonPrimitive +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.time.Duration.Companion.milliseconds +import kotlin.time.Duration.Companion.minutes +import kotlin.time.Duration.Companion.seconds + +class HealthProberTest { + + private val json = Json { ignoreUnknownKeys = true } + + @Test + fun effectiveProbeIntervalPrefersUpstreamThenProviderThenDefault() { + val provider = ProviderConf("p", "https://x", "", null, null, probe_interval = 2.seconds) + val upstream = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 5.milliseconds) + val upstreamNoOwn = UpstreamConf("u2", "p", "m") + assertEquals(5.milliseconds, effectiveProbeInterval(upstream, provider)) + assertEquals(2.seconds, effectiveProbeInterval(upstreamNoOwn, provider)) + assertEquals(DEFAULT_PROBE_INTERVAL, effectiveProbeInterval(upstreamNoOwn, null)) + } + + @Test + fun effectiveProbeIntervalZeroDisablesProber() { + val up = UpstreamConf("u", "p", "m", null, null, null, null, null, probe_interval = 0.milliseconds) + assertNull(effectiveProbeInterval(up, null)) + + val providerZero = ProviderConf("p", "https://x", "", null, null, probe_interval = 0.milliseconds) + assertNull(effectiveProbeInterval(UpstreamConf("u2", "p", "m"), providerZero)) + } + + @Test + fun buildProbeBodyIsMinimalChatWithModelReplacedAndPatchesApplied() { + val provider = ProviderConf( + "p", "https://x", "", null, + Json.parseToJsonElement("""{"provider":{"allow_fallbacks":false}}""").jsonObject, + probe_interval = 1.seconds, + ) + val up = UpstreamConf( + "u", "p", "real-1", null, + Json.parseToJsonElement("""{"temperature":0.5}""").jsonObject, + null, null, null, 1.seconds, + ) + + val body = buildProbeBody(up, provider) + + assertEquals("real-1", body["model"]?.jsonPrimitive?.content) + assertEquals("2+2=?", body["messages"]?.jsonArray?.first()?.jsonObject?.get("content")?.jsonPrimitive?.content) + assertEquals("true", body["stream"]?.jsonPrimitive?.content) + assertEquals("4", body["max_tokens"]?.jsonPrimitive?.content) + assertEquals("false", body["provider"]?.jsonObject?.get("allow_fallbacks")?.jsonPrimitive?.content) + assertEquals("0.5", body["temperature"]?.jsonPrimitive?.content) + } + + @Test + fun buildProbeBodyIgnoresModelPatch() { + val up = UpstreamConf("u", "p", "real-1") + val provider = ProviderConf("p", "https://x") + val body = buildProbeBody(up, provider) + assertEquals( + """{"model":"real-1","messages":[{"role":"user","content":"2+2=?"}],"max_tokens":4,"stream":true}""", + body.toString(), + ) + } + + @Test + fun probeLoopRecoversUpstreamWhenPingSucceeds() = runTest { + val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes) + val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds) + val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up)) + backoff.recordFailure(up) + assertEquals(true, backoff.isCooling(up)) + + val metrics = MetricsRegistry() + val prober = HealthProber( + providersById = mapOf("p" to provider), + upstreams = listOf(up), + http = createHttpClient(), + backoff = backoff, + metrics = metrics, + scope = this, + ) + var pings = 0 + val job = launch { + prober.probeLoop(up, 5.milliseconds) { pinged -> + pings++ + assertEquals("u", pinged.id) + true + } + } + runUntil { pings >= 1 } + // успех уже зафиксирован (recordSuccess) — даём циклу несколько + // итераций и убеждаемся, что повторных пингов не было + delay(10.milliseconds) + job.cancel() + + assertEquals(1, pings) + assertEquals(false, backoff.isCooling(up)) + } + + @Test + fun probeLoopKeepsCoolingWhenPingFails() = runTest { + val provider = ProviderConf("p", "https://x", "", null, null, backoff = 5.minutes) + val up = UpstreamConf("u", "p", "m", null, null, null, backoff = 5.minutes, probe_interval = 5.milliseconds) + val backoff = BackoffRegistry(mapOf("p" to provider), mapOf("u" to up)) + backoff.recordFailure(up) + + val metrics = MetricsRegistry() + val prober = HealthProber( + providersById = mapOf("p" to provider), + upstreams = listOf(up), + http = createHttpClient(), + backoff = backoff, + metrics = metrics, + scope = this, + ) + var pings = 0 + val job = launch { + prober.probeLoop(up, 5.milliseconds) { _ -> + pings++ + false + } + } + runUntil { pings >= 2 } + job.cancel() + + assertEquals(true, backoff.isCooling(up)) + } + + private suspend fun TestScope.runUntil(condition: () -> Boolean) { + var i = 0 + while (!condition() && i < 10_000) { + delay(1) + i++ + } + if (!condition()) throw AssertionError("condition not met in time") + } + +} \ No newline at end of file