3 Commits
11 ... 12

Author SHA1 Message Date
subochev f70a9fdc9c feat: Prometheus-метрики GET /metrics + ленивый ISO-8601-парсер для backoff
Build LLM Proxy / Build and push (release) Successful in 40s
- GET /metrics (Prometheus text 0.0.4): llm_proxy_requests_total{model,upstream,provider,result}
  (ok/4xx/429/402/5xx/net_err/cancelled), llm_proxy_upstream_inflight,
  llm_proxy_upstream_fail_streak, llm_proxy_upstream_cooling_seconds
- учёт исходов в handleChat (каждая фейловер-попытка — отдельно)
- parseIsoDuration: PnM≈n×30d, PnY≈n×365d (kotlin.time Duration.parse не принимает Y/M)
- remainingSeconds округляет остаток вверх (открытое окно ≥1s)
- доки (CONFIG.md/README/TESTING) + тесты (MetricsTest, BackoffTest)
2026-09-13 20:59:28 +03:00
subochev 0ba9b7d739 feat: экспоненциальный backoff (providers[].backoff / upstreams[].backoff)
Повторяющиеся сбои апстрима (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.
2026-09-13 20:33:36 +03:00
subochev adaffba885 feat: 402 (insufficient balance / лимит токенов) — фейловер на следующий апстрим
402 от провайдера — ошибка аккаунта (баланс/лимит токенов), а не запроса:
следующий апстрим в списке модели может быть платёжеспособен. Теперь при
5xx/429/402 освобождаем слот и пробуем следующего свободного; остальные 4xx
продолжаем отдавать клиенту как есть. Обновлён CONFIG.md (маршрутизация,
лог-секции). Ранее этот правка терялась (3155065 не сохранился) — восстанавливаю.
2026-09-13 20:33:32 +03:00
9 changed files with 672 additions and 29 deletions
+98 -3
View File
@@ -41,9 +41,15 @@
запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос
или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать, или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать,
что он перегружен: такой ответ от провайдера — это уже сбой, а не способ что он перегружен: такой ответ от провайдера — это уже сбой, а не способ
управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`, мы управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`/`402`, мы
освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка
модели (фейловер); если свободных не осталось — отдаём `503` от прокси. модели (фейловер); если свободных не осталось — отдаём `503` от прокси.
`402` (insufficient balance / исчерпан лимит токенов) — ошибка аккаунта
провайдера, но фейловер имеет смысл: следующий апстрим может быть платёжеспособен.
Повторяющиеся сбои уводятся из ротации экспоненциальным **backoff** (поле
`backoff`, ниже): упавший апстрим «отдыхает» 1s, 2s, 4s, … до заданного потолка,
и прокси его не трогает. Если **все** апстримы модели в откате — `503` с
заголовком `Retry-After`.
## Формат файла ## Формат файла
@@ -77,6 +83,8 @@ providers:
patch: # уровень провайдера: ко всем его запросам patch: # уровень провайдера: ко всем его запросам
provider: provider:
allow_fallbacks: false allow_fallbacks: false
backoff: P1M # опционально; потолок экспоненциального backoff
# на весь провайдер (ISO-8601: P1M, PT15M, P1D…)
- id: local-llama - id: local-llama
url: "http://10.0.0.5:8080/v1" url: "http://10.0.0.5:8080/v1"
@@ -93,6 +101,8 @@ upstreams:
patch: # уровень апстрима patch: # уровень апстрима
provider: provider:
ignore: [deepseek] ignore: [deepseek]
backoff: PT15M # опционально; собственный backoff этой модели
# (приоритет над backoff провайдера)
- id: local-qwen - id: local-qwen
provider: local-llama provider: local-llama
@@ -256,6 +266,82 @@ providers:
# reasoning_empty_ok: true # опционально; дефолт false # 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=<id> в backoff-откате (~Ns) — пропускаю`);
- если **все** апстримы модели в откате — прокси отдаёт `503`
(`all upstreams cooling, retry in Ns`) с заголовком `Retry-After: N`
(секунд до выхода первого апстрима из отката).
Значение — ISO-8601-длительность: `PT30M` (30 минут), `PT1H15M`, `P1D` (сутки),
`P1M` (месяц), `P1Y` (год). Годы/месяцы укорачиваются приближённо
(1 год ≈ 365d, 1 месяц ≈ 30d) — kotlin.time `Duration.parse` не принимает
Y/M (у них нет фиксированной длины), конфиг-парсер прокси расширяет формат.
Некорректное значение — предупреждение в лог и поле просто игнорируется.
```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`.
### Prometheus-метрики (`/metrics`)
Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics`
(`text/plain; version=0.0.4`), без авторизации (внутренний контур).
Достаточно включить скрейп в Prometheus/VictoriaMetrics — и в Grafana можно
вести дашборды использования по провайдерам/моделям и алерты на деградацию.
| Метрика | Тип | Смысл |
|---|---|---|
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `cancelled` — клиент отвалился посреди стрима |
| `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 = жив) |
Метки: `model` — витринное имя модели, `upstream`/`provider` — внутренние id из конфига.
Метрики live в памяти: при рестарте прокси сбрасываются (серить их будет Prometheus).
Примеры для Grafana:
- оборот по виртуальным моделям: `sum(rate(llm_proxy_requests_total[5m])) by (model)`;
- оборот по провайдерам: `sum(rate(llm_proxy_requests_total[5m])) by (provider)`;
- «провайдер умер»: `llm_proxy_upstream_cooling_seconds > 0` дольше N минут — алерт;
- доля ошибок провайдера: `sum(rate(llm_proxy_requests_total{result=~"4xx|429|402|5xx|net_err"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
### Пример сборки тела (многослойный `patch`) ### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -340,6 +426,7 @@ data class ProviderConf(
val think_tags: String? = null, val think_tags: String? = null,
val reasoning_field: String? = null, val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false, val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
) )
@Serializable @Serializable
@@ -349,6 +436,8 @@ data class UpstreamConf(
val model: String, // реальное имя модели у провайдера val model: String, // реальное имя модели у провайдера
val max_concurrency: Int? = null,// опционально; null/0 = безлимит val max_concurrency: Int? = null,// опционально; null/0 = безлимит
val patch: JsonObject? = null, // к запросам этой апстрим-модели val patch: JsonObject? = null, // к запросам этой апстрим-модели
val think_tags: String? = null, // переопределение think-режима модели
val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским)
) )
@Serializable @Serializable
@@ -454,7 +543,7 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
`upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch `upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch
== null` пропускаем. Итог — финальное тело запроса; == null` пропускаем. Итог — финальное тело запроса;
- шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`; - шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`;
- при ответе `5xx`/`429` от провайдера — `release(upstream)` и переходим к - при ответе `5xx`/`429`/`402` от провайдера — `release(upstream)` и переходим к
следующему свободному апстриму из списка (фейловер); исчерпали список — следующему свободному апстриму из списка (фейловер); исчерпали список —
отдаём `503` (`all upstreams failed`); отдаём `503` (`all upstreams failed`);
- при успехе/отмене клиента — `release(upstream)` в `finally` и выходим. - при успехе/отмене клиента — `release(upstream)` в `finally` и выходим.
@@ -507,8 +596,14 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
свободному `upstream=<id>`. свободному `upstream=<id>`.
- **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима, - **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима,
освобождён слот `active[id]=N/limit`. освобождён слот `active[id]=N/limit`.
- **Ошибка апстрима** (не `5xx`/`429`, а сетевая/таймаут): как сейчас — лог с - **Ошибка апстрима** (не `5xx`/`429`/`402`, а сетевая/таймаут): как сейчас — лог с
сообщением. сообщением.
- **Backoff**: при фейловере строка дополняется интервалом отдыха
(`(backoff: отдых PT…)`); при пропуске апстрима в откате —
`upstream=<id> в backoff-откате (~Ns) — пропускаю`; при старте — список
настроенных backoff (`backoff: upstreams=… providers=…`).
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
`model=<витрина>` с заголовком `Retry-After: N`.
Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать
существующий парсинг логов (если он есть). существующий парсинг логов (если он есть).
+12 -1
View File
@@ -20,10 +20,21 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (
обоих — безлимит. обоих — безлимит.
Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или
апстрима (`off` по умолчанию, `split` — рассуждения из `<think>…</think>` уходят апстрима (`off` по умолчанию, `split` — рассуждения из `think`-тегов уходят
в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в
non-stream. non-stream.
Повторяющиеся сбои апстрима увязываются экспоненциальным backoff: поле
`backoff` (ISO-8601-потолок, напр. `P1M`) задаётся на провайдере (общий
счётчик на его модели) и/или на апстриме (приоритет). Первая ошибка — отдых
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
— `503` с `Retry-After`.
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
метрик — в CONFIG.md, раздел «Prometheus-метрики».
| Переменная | Default | Описание | | Переменная | Default | Описание |
|---|---|---| |---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+2
View File
@@ -20,6 +20,8 @@
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; | | `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
| `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. | | `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. |
| `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. | | `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. |
| `BackoffTest` | Экспоненциальный backoff: удвоение интервала отката до потолка (cap), сброс счётчика при успехе, окно охлаждения (`isCoolingAt`/`remainingAt`), ISO-8601-разбор cap (`PT30M`, `P1D`, `PT1M30S`) и конфиг-парсер `parseIsoDuration` (ленивые Y/M: `P1M`≈30d, `P1Y`≈365d; мусор — ошибка), скоупы: общий провайдерский счётчик на все его модели, приоритет `upstreams[].backoff` над провайдерским, независимость скоупов, backoff не настроен → без отката. |
| `MetricsTest` | `/metrics`: counters по исходам запросов (`ok/4xx/429/402/5xx/net_err/cancelled`), гейджи `inflight`/`fail_streak`/`cooling_seconds` (общий провайдерский скоуп виден в метриках), экранирование меток, маппинг HTTP-статуса → `result`. |
## Проверка качества тестов (мутационная приёмка) ## Проверка качества тестов (мутационная приёмка)
@@ -0,0 +1,124 @@
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()
}
}
/** Текущая серия сбоев подряд (0 — после успеха или без сбоев). */
fun failStreak(): Int = failures
}
/**
* Реестр областей бэкофф-отдыха. Область апстрима:
* - у модели (записи upstreams) задан `backoff` — свой счётчик и потолок (приоритет);
* - иначе у провайдера задан `backoff` — общий для всех моделей провайдера счётчик;
* - ни там, ни там — бэкофф для этого апстрима выключен.
*/
class BackoffRegistry(
private val providers: Map<String, ProviderConf>,
private val upstreams: Map<String, UpstreamConf>,
) {
private val upstreamGuards = mutableMapOf<String, BackoffGuard>()
private val providerGuards = mutableMapOf<String, BackoffGuard>()
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 <T> 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
val remaining = g.remainingAt(Clock.System.now())
if (remaining <= Duration.ZERO) return 0L
val s = remaining.inWholeSeconds
val isExact = remaining.inWholeNanoseconds == s * 1_000_000_000L
return if (isExact) s else s + 1
}
/** Учёт ошибки; возвращает интервал отдыха, или 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()
}
/** Текущая серия сбоев апстрима (0 — нет сбоев или backoff не настроен). */
suspend fun failStreakOf(up: UpstreamConf): Int {
val g = locked { guardForLocked(up) } ?: return 0
return g.failStreak()
}
}
+115 -12
View File
@@ -46,6 +46,10 @@ import net.mamoe.yamlkt.YamlList
import net.mamoe.yamlkt.YamlLiteral import net.mamoe.yamlkt.YamlLiteral
import net.mamoe.yamlkt.YamlMap import net.mamoe.yamlkt.YamlMap
import kotlin.concurrent.Volatile import kotlin.concurrent.Volatile
import kotlin.time.Duration
import kotlin.time.Duration.Companion.days
import kotlin.time.Duration.Companion.hours
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds import kotlin.time.Duration.Companion.seconds
import kotlin.time.TimeSource import kotlin.time.TimeSource
@@ -86,11 +90,19 @@ fun main() {
val active = config.upstreams.associate { up -> val active = config.upstreams.associate { up ->
up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider])) up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider]))
} }
val backoff = BackoffRegistry(providersById, upstreamsById)
val metrics = MetricsRegistry()
log.info { log.info {
"[llm-proxy] загружено: providers=${config.providers.size}, " + "[llm-proxy] загружено: providers=${config.providers.size}, " +
"upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)" "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 { log.info {
"[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}" "[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}"
} }
@@ -101,7 +113,7 @@ fun main() {
val http = createHttpClient() val http = createHttpClient()
val sessions = SessionRegistry() val sessions = SessionRegistry()
startServer(config.server.host, config.server.port) { startServer(config.server.host, config.server.port) {
proxyModule(config, providersById, upstreamsById, active, sessions, http) proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
} }
} }
@@ -114,10 +126,19 @@ fun Application.proxyModule(
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry, sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) { ) {
routing { routing {
get("/metrics") {
call.respondText(
metrics.render(config, active, backoff),
ContentType.parse("text/plain; version=0.0.4"),
HttpStatusCode.OK,
)
}
post("/v1/chat/completions") { post("/v1/chat/completions") {
handleChat(call, config, providersById, upstreamsById, active, sessions, http) handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
} }
get("/v1/models") { get("/v1/models") {
handleModels(call, config) handleModels(call, config)
@@ -133,6 +154,8 @@ private suspend fun handleChat(
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry, sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) { ) {
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" } log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val raw = call.receiveText() val raw = call.receiveText()
@@ -165,7 +188,7 @@ private suspend fun handleChat(
var anyClaimed = false var anyClaimed = false
while (true) { while (true) {
val up = pickFreeUpstream(pool, active, failed) ?: break val up = pickFreeUpstream(pool, active, failed, backoff, modelName) ?: break
anyClaimed = true anyClaimed = true
val start = TimeSource.Monotonic.markNow() val start = TimeSource.Monotonic.markNow()
try { try {
@@ -229,8 +252,13 @@ private suspend fun handleChat(
setBody(forwarded.toString()) setBody(forwarded.toString())
}.execute { resp -> }.execute { resp ->
upstreamStatus = resp.status.value upstreamStatus = resp.status.value
if (upstreamStatus >= 500 || upstreamStatus == 429) { if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) {
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" } val rest = backoff.recordFailure(up)
metrics.record(modelName, up, statusClass(upstreamStatus))
log.warn {
"[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" +
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
}
failed.add(up.id) failed.add(up.id)
failover = true failover = true
return@execute return@execute
@@ -238,13 +266,15 @@ private suspend fun handleChat(
if (upstreamStatus in 400..499) { if (upstreamStatus in 400..499) {
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть. // 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
val errorBody = runCatching { resp.body<String>() }.getOrDefault("") val errorBody = runCatching { resp.body<String>() }.getOrDefault("")
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" } log.warn { "[llm-proxy] chat model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" }
metrics.record(modelName, up, statusClass(upstreamStatus))
val ct = resp.headers["Content-Type"] ?: "application/json" val ct = resp.headers["Content-Type"] ?: "application/json"
call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus)) call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus))
responded = true responded = true
return@execute return@execute
} }
responded = true responded = true
backoff.recordSuccess(up)
if (clientWantsStream) { if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream" val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus) val status = HttpStatusCode.fromValue(upstreamStatus)
@@ -258,6 +288,7 @@ private suspend fun handleChat(
streamSseWithThinkTags(ch, thinkMode) streamSseWithThinkTags(ch, thinkMode)
} }
} }
metrics.record(modelName, up, "ok")
log.info { log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " + "[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true" "status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true"
@@ -281,6 +312,7 @@ private suspend fun handleChat(
} }
}.getOrDefault(ContentType.parse(ct)) }.getOrDefault(ContentType.parse(ct))
call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus)) call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus))
metrics.record(modelName, up, "ok")
log.info { log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " + "[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false" "status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false"
@@ -291,10 +323,16 @@ private suspend fun handleChat(
if (failover) continue if (failover) continue
if (responded) return if (responded) return
} catch (e: CancellationException) { } catch (e: CancellationException) {
metrics.record(modelName, up, "cancelled")
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" } log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
throw e throw e
} catch (e: Exception) { } 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)
metrics.record(modelName, up, "net_err")
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) failed.add(up.id)
continue continue
} finally { } finally {
@@ -304,10 +342,20 @@ private suspend fun handleChat(
if (anyClaimed) { if (anyClaimed) {
call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
} else {
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 { } else {
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
} }
} }
}
private fun errorJson(msg: String): String = private fun errorJson(msg: String): String =
"""{"error":{"message":"$msg"}}""" """{"error":{"message":"$msg"}}"""
@@ -489,15 +537,28 @@ class UpstreamCounter(private val limit: Int, initial: Int = 0) {
/** /**
* Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`, * Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`,
* у которого свободен слот и который ещё не в `excluded` (не упал ранее). * у которого свободен слот, который ещё не в `excluded` (не упал ранее) и
* Сразу занимает слот (через [tryClaim]). Если свободных нет — возвращает null. * не в backoff-откате. Сразу занимает слот (через [tryClaim]).
* Если свободных нет — возвращает null.
*/ */
internal fun pickFreeUpstream( internal suspend fun pickFreeUpstream(
pool: List<UpstreamConf>, pool: List<UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
excluded: Set<String>, excluded: Set<String>,
): UpstreamConf? = backoff: BackoffRegistry,
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) } 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) { 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)}" } log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
@@ -903,6 +964,7 @@ data class ProviderConf(
val think_tags: String? = null, val think_tags: String? = null,
val reasoning_field: String? = null, val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false, val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
) )
data class UpstreamConf( data class UpstreamConf(
@@ -912,6 +974,7 @@ data class UpstreamConf(
val max_concurrency: Int? = null, val max_concurrency: Int? = null,
val patch: JsonObject? = null, val patch: JsonObject? = null,
val think_tags: String? = null, val think_tags: String? = null,
val backoff: Duration? = null,
) )
data class ModelConf( data class ModelConf(
@@ -969,6 +1032,7 @@ internal fun parseConfig(root: YamlElement): Config {
think_tags = m.strOrNull("think_tags"), think_tags = m.strOrNull("think_tags"),
reasoning_field = m.strOrNull("reasoning_field"), reasoning_field = m.strOrNull("reasoning_field"),
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false, reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
backoff = m.durationOrNull("backoff"),
) )
} }
@@ -981,6 +1045,7 @@ internal fun parseConfig(root: YamlElement): Config {
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
think_tags = m.strOrNull("think_tags"), think_tags = m.strOrNull("think_tags"),
backoff = m.durationOrNull("backoff"),
) )
} }
@@ -1006,5 +1071,43 @@ internal fun Map<String, YamlElement>.str(key: String): String =
internal fun Map<String, YamlElement>.strOrNull(key: String): String? = internal fun Map<String, YamlElement>.strOrNull(key: String): String? =
(this[key] as? YamlLiteral)?.content (this[key] as? YamlLiteral)?.content
/**
* ISO-8601-длительность (PT30M, P1D, PT1H15M…) с расширением для годов/месяцев,
* которые kotlin.time не принимает (у них нет фиксированной длины):
* `PnY` → n×365d, `PnM` → n×30d (приближение, см. CONFIG.md).
* Некорректное значение — warn + null.
*/
internal fun Map<String, YamlElement>.durationOrNull(key: String): Duration? =
strOrNull(key)?.let { raw ->
runCatching { parseIsoDuration(raw) }.getOrElse {
log.warn { "[llm-proxy] config: поле '$key' — некорректная ISO-8601-длительность '$raw', проигнорировано" }
null
}
}
/**
* ISO-8601-длительность для конфига: как [Duration.parse], плюс годы/месяцы
* (`P1M`, `P2Y`, `P1Y3M`) — приближённо 1 год = 365d, 1 месяц = 30d.
* (kotlin.time [Duration.parse] принимает только D/W/H/M(ин)/S.)
*/
internal fun parseIsoDuration(raw: String): Duration {
val m = Regex(
"^P" +
"(?:(\\d+)Y)?" +
"(?:(\\d+)M)?" +
"(?:(\\d+)W)?" +
"(?:(\\d+)D)?" +
"(?:T(?:(\\d+)H)?(?:(\\d+)M)?(?:(\\d+)S)?)?",
RegexOption.IGNORE_CASE,
).matchEntire(raw.trim()) ?: error("не ISO-8601-длительность: $raw")
if ((1..7).all { m.groupValues[it].isEmpty() }) error("не ISO-8601-длительность: $raw")
fun g(i: Int): Long = m.groupValues[i].toLongOrNull() ?: 0L
val days = g(1) * 365 + g(2) * 30 + g(3) * 7 + g(4)
val hours = g(5)
val minutes = g(6)
val seconds = g(7)
return days.days + hours.hours + minutes.minutes + seconds.seconds
}
internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? = internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? =
this[key] as? YamlMap this[key] as? YamlMap
@@ -0,0 +1,100 @@
package pw.binom.llmproxy
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
/**
* Pull-метрики для Prometheus (формат экспозиции `text/plain; version=0.0.4`),
* отдаются на `GET /metrics`.
*
* Метрики (всё в памяти, сброс при рестарте):
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | cancelled`;
* - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
* (счётчик backoff);
* - `llm_proxy_upstream_cooling_seconds{upstream,provider}` — сколько секунд
* апстрим ещё в backoff-откате (0 = жив).
*/
class MetricsRegistry {
private val lock = Mutex()
private val requests = HashMap<RequestKey, Long>()
private data class RequestKey(
val model: String,
val upstream: String,
val provider: String,
val result: String,
)
/** Учёт chat-запроса по исходу (значения `result` — в описании класса). */
suspend fun record(model: String, up: UpstreamConf, result: String) {
lock.withLock {
val key = RequestKey(model, up.id, up.provider, result)
requests[key] = (requests[key] ?: 0L) + 1
}
}
/** Текстовая выгрузка в формате Prometheus на текущий момент. */
suspend fun render(
config: Config,
active: Map<String, UpstreamCounter>,
backoff: BackoffRegistry,
): String {
val sb = StringBuilder()
sb.append("# HELP llm_proxy_requests_total Chat-запросы llm-proxy по исходу\n")
sb.append("# TYPE llm_proxy_requests_total counter\n")
lock.withLock {
requests.entries
.sortedWith(compareBy({ it.key.model }, { it.key.upstream }, { it.key.provider }, { it.key.result }))
.forEach { (k, n) ->
sb.append(
"llm_proxy_requests_total{model=\"" + esc(k.model) +
"\",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) {
val n = active[up.id]?.current ?: 0
sb.append(
"llm_proxy_upstream_inflight{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} $n\n",
)
}
sb.append("# HELP llm_proxy_upstream_fail_streak Серия сбоев апстрима подряд (счётчик backoff)\n")
sb.append("# TYPE llm_proxy_upstream_fail_streak gauge\n")
for (up in config.upstreams) {
sb.append(
"llm_proxy_upstream_fail_streak{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} " +
backoff.failStreakOf(up) + "\n",
)
}
sb.append("# HELP llm_proxy_upstream_cooling_seconds Секунд до выхода апстрима из backoff-отката (0 = не в откате)\n")
sb.append("# TYPE llm_proxy_upstream_cooling_seconds gauge\n")
for (up in config.upstreams) {
sb.append(
"llm_proxy_upstream_cooling_seconds{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} " +
backoff.remainingSeconds(up) + "\n",
)
}
return sb.toString()
}
private fun esc(s: String): String =
s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n")
}
/** HTTP-статус апстрима → класс исхода для метрик. */
internal fun statusClass(status: Int): String = when {
status in 200..399 -> "ok"
status == 429 -> "429"
status == 402 -> "402"
status in 400..499 -> "4xx"
else -> "5xx"
}
@@ -0,0 +1,131 @@
package pw.binom.llmproxy
import kotlinx.coroutines.test.runTest
import kotlinx.datetime.Instant
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFails
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.hours
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 configDurationParsingWithYearsAndMonths() {
// kotlin.time не принимает Y/M — конфиг-парсер расширяет: 1Y≈365d, 1M≈30d
assertEquals(30.days, parseIsoDuration("P1M"))
assertEquals(365.days, parseIsoDuration("P1Y"))
assertEquals(455.days, parseIsoDuration("P1Y3M"))
assertEquals(7.days, parseIsoDuration("P1W"))
assertEquals(2.days + 3.hours + 15.minutes, parseIsoDuration("P2DT3H15M"))
// базовый ISO-8601 без Y/M — как kotlin.time
assertEquals(30.minutes, parseIsoDuration("PT30M"))
}
@Test
fun configDurationParsingRejectsGarbage() {
assertFails { parseIsoDuration("PT30") }
assertFails { parseIsoDuration("1d") }
assertFails { parseIsoDuration("P") }
assertFails { parseIsoDuration("") }
}
@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))
}
}
@@ -5,6 +5,7 @@ import kotlin.test.assertEquals
import kotlin.test.assertFalse import kotlin.test.assertFalse
import kotlin.test.assertTrue import kotlin.test.assertTrue
import io.ktor.http.headersOf import io.ktor.http.headersOf
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonArray
@@ -14,6 +15,8 @@ import net.mamoe.yamlkt.Yaml
class ConfigLogicTest { class ConfigLogicTest {
private val noBackoff = BackoffRegistry(emptyMap(), emptyMap())
@Test @Test
fun mergeDeepMergesNestedObjectsAndReplacesScalars() { fun mergeDeepMergesNestedObjectsAndReplacesScalars() {
val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject
@@ -222,58 +225,58 @@ class ConfigLogicTest {
} }
@Test @Test
fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() { fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 2, 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("u1", up?.id)
// слот реально занят // слот реально занят
assertEquals(1, active.getValue("u1").current) assertEquals(1, active.getValue("u1").current)
} }
@Test @Test
fun pickFreeUpstreamSkipsExcluded() { fun pickFreeUpstreamSkipsExcluded() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 2, 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) assertEquals("u2", up?.id)
} }
@Test @Test
fun pickFreeUpstreamReturnsNullWhenAllBusy() { fun pickFreeUpstreamReturnsNullWhenAllBusy() = runTest {
val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1 val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1
val pool = listOf(UpstreamConf("u1", "p", "m", 1, null)) val pool = listOf(UpstreamConf("u1", "p", "m", 1, null))
assertEquals(null, pickFreeUpstream(pool, active, emptySet())) assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff))
} }
@Test @Test
fun pickFreeUpstreamImplementsFailoverOrder() { fun pickFreeUpstreamImplementsFailoverOrder() = runTest {
// dead исключён (упал ранее) — выбирается следующий живой u1 // dead исключён (упал ранее) — выбирается следующий живой u1
val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1)) val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1))
val pool = listOf( val pool = listOf(
UpstreamConf("dead", "p", "m", 1, null), UpstreamConf("dead", "p", "m", 1, null),
UpstreamConf("u1", "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) assertEquals("u1", up?.id)
} }
@Test @Test
fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() { fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 1, null),
) )
assertEquals("u1", pickFreeUpstream(pool, active, emptySet())?.id) assertEquals("u1", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
assertEquals("u2", pickFreeUpstream(pool, active, emptySet())?.id) assertEquals("u2", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id) assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
} }
@Test @Test
@@ -0,0 +1,74 @@
package pw.binom.llmproxy
import kotlinx.coroutines.test.runTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
class MetricsTest {
private val p = ProviderConf(id = "p", url = "https://p", backoff = 30.seconds)
private val a = UpstreamConf(id = "a", provider = "p", model = "m1")
private val b = UpstreamConf(id = "b", provider = "p", model = "m2")
private val cfg = Config(
server = ServerConf(),
providers = listOf(p),
upstreams = listOf(a, b),
models = emptyList(),
)
@Test
fun requestCounters() = runTest {
val m = MetricsRegistry()
m.record("my-gpt", a, "ok")
m.record("my-gpt", a, "ok")
m.record("my-gpt", a, "429")
m.record("my-gpt", b, "5xx")
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"ok\"} 2\n"), "ok=2: $text")
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"429\"} 1\n"), "429=1: $text")
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"b\",provider=\"p\",result=\"5xx\"} 1\n"), "5xx=1: $text")
assertTrue(text.contains("# TYPE llm_proxy_requests_total counter"), "TYPE: $text")
}
@Test
fun inflightGauge() = runTest {
val m = MetricsRegistry()
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4, 3), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"a\",provider=\"p\"} 3\n"), text)
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"b\",provider=\"p\"} 0\n"), text)
}
@Test
fun backoffGauges() = runTest {
val backoff = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
assertEquals(1.seconds, backoff.recordFailure(a))
val text = MetricsRegistry().render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), backoff)
// общий скоуп провайдера: сбой у a виден и у b
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"a\",provider=\"p\"} 1\n"), text)
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"b\",provider=\"p\"} 1\n"), text)
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"a\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"b\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
}
@Test
fun labelEscaping() = runTest {
val m = MetricsRegistry()
m.record("we\"ird\nmodel", a, "ok")
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("model=\"we\\\"ird\\nmodel\""), text)
}
@Test
fun statusClasses() {
assertEquals("ok", statusClass(200))
assertEquals("ok", statusClass(302))
assertEquals("4xx", statusClass(400))
assertEquals("4xx", statusClass(499))
assertEquals("429", statusClass(429))
assertEquals("402", statusClass(402))
assertEquals("5xx", statusClass(500))
assertEquals("5xx", statusClass(599))
}
}