6 Commits

Author SHA1 Message Date
subochev c6a24a94b0 feat: фоновый пинг апстримов в backoff-откате (probe_interval) — мгновенное возвращение в пул
Build LLM Proxy / Build and push (release) Successful in 42s
2026-09-15 01:00:10 +03:00
subochev 249b83eb51 feat: watchdog на первый байт стрима (first_byte_timeout) — фейловер + backoff при молчании апстрима 2026-09-15 00:37:51 +03:00
subochev 293964d8f1 fix: reasoning_field на каждом assistant-сообщении (не только с tool_calls)
Build LLM Proxy / Build and push (release) Successful in 41s
Console Go (deepseek в thinking-режиме) в stateful-сессиях требует
reasoning_content на КАЖДОМ assistant-сообщении истории — живые 400
«The reasoning_content in the thinking mode must be passed back to the
API» (19:01, session 890fc602) шли именно на чистых ассистент-турах
без tool_calls, которые старый гейт пропускал. Снял гейт: поле
дописывается на каждом assistant-сообщении (текст из reasoning /
reasoning_details, пустая строка при reasoning_empty_ok) — как это
делает сам opencode («Deepseek requires all assistant messages to have
reasoning on them»). Тесты переписаны под новое поведение.
2026-09-13 22:36:33 +03:00
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
14 changed files with 1643 additions and 54 deletions
+246 -11
View File
@@ -41,9 +41,15 @@
запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос
или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать,
что он перегружен: такой ответ от провайдера — это уже сбой, а не способ
управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`, мы
управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`/`402`, мы
освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка
модели (фейловер); если свободных не осталось — отдаём `503` от прокси.
`402` (insufficient balance / исчерпан лимит токенов) — ошибка аккаунта
провайдера, но фейловер имеет смысл: следующий апстрим может быть платёжеспособен.
Повторяющиеся сбои уводятся из ротации экспоненциальным **backoff** (поле
`backoff`, ниже): упавший апстрим «отдыхает» 1s, 2s, 4s, … до заданного потолка,
и прокси его не трогает. Если **все** апстримы модели в откате — `503` с
заголовком `Retry-After`.
## Формат файла
@@ -77,6 +83,10 @@ providers:
patch: # уровень провайдера: ко всем его запросам
provider:
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"
@@ -93,6 +103,10 @@ upstreams:
patch: # уровень апстрима
provider:
ignore: [deepseek]
backoff: PT15M # опционально; собственный backoff этой модели
# (приоритет над backoff провайдера)
probe_interval: PT0S # опционально; пинг отключён у этой модели
# (приоритет над probe_interval провайдера)
- id: local-qwen
provider: local-llama
@@ -221,12 +235,15 @@ upstreams:
### Нативное поле рассуждений (`reasoning_field`)
Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**)
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в каждом
assistant-сообщении с непустым `tool_calls`. Если его нет — апстрим отвечает
`400` (`The `reasoning_content` in the thinking mode must be passed back to the
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в **каждом**
assistant-сообщении истории (не только в тех, что с `tool_calls`). Если его нет —
апстрим отвечает `400`
(`The `reasoning_content` in the thinking mode must be passed back to the
API.`). Клиенты при этом рассуждения держат в своих форматах: `reasoning`
(строка) и/или `reasoning_details` (массив `{type:"reasoning.text", text, ...}`)
— а нативное поле `reasoning_content` в истории могут и не передавать.
(Точно так же делает и сам opencode: для deepseek он добавляет reasoning-часть
на каждом assistant-сообщении, даже пустую.)
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
@@ -235,15 +252,15 @@ API.`). Клиенты при этом рассуждения держат в с
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в каждом
assistant-сообщении с непустым `tool_calls` — берёт текст из `reasoning`
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в
**каждом** assistant-сообщении — берёт текст из `reasoning`
(если это непустая строка), иначе склеивает `reasoning_details[*].text`
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
непустое поле не перезаписывается. Ничего при этом не убирается и не
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
поле просто дополняется. Тронуто только assistant-сообщение с непустым
`tool_calls`; сообщения без `tool_calls` и `user`/`tool`/`system` не меняются.
поле просто дополняется. Тронуты только assistant-сообщения; `user`/`tool`/
`system` не меняются.
Если текста рассуждений нет вовсе — поле не добавляется, кроме случая
`reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если
`reasoning_field` не задан — тело не меняется вовсе.
@@ -256,6 +273,208 @@ 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=<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`.
### Таймаут первого байта стрима (`first_byte_timeout`)
Для `stream=true` запросов прокси следит, чтобы апстрим начал отдавать
**что-нибудь** (первый байт тела) в течение заданного окна. Если апстрим
молчит дольше — это трактуется как его сбой: слот конкурентности
освобождается, ошибка идёт в backoff-откат, и роутер **фейловерит** на
следующего свободного апстрима модели (при его наличии). Это закрывает
кейс, когда апстрим принял запрос, но «завис» и никогда не начинает
стримить.
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно ожидания первого байта для всех моделей провайдера |
| `upstreams[].first_byte_timeout` | ISO-8601-длительность / отсутствует | Окно для конкретной модели (приоритет над провайдерским) |
**Дефолт** (если не задан нигде): `PT30S`.
**Приоритет — модели.** Если `upstreams[].first_byte_timeout` задан, он
перекрывает провайдерский; если нет — берётся провайдерский; если нет и там —
дефолт `PT30S`.
**Отключение:** `PT0S` на провайдере или апстриме отключает watchdog
для этого апстрима — тогда первый байт ждём без лимита (в рамках общего
`upstream: request_timeout`, см. ниже).
**Применяется только к `stream=true`.** Для `stream=false` (не-стриминговый
ответ) действует только общий лимит `UPSTREAM_REQUEST_TIMEOUT_MS` (`5m`) —
строки SSE там склеиваются в один ответ, и таймаут первого байта к нему
неприменим.
Поведение:
- watchdog висит **на первом чтении** из тела стрим-ответа; как только
пришёл хоть какой-то байт — лимит отключается, дальше стрим читается
свободно (вплоть до общего request_timeout);
- `data: [DONE]` приходит в теле — это уже «первый байт», таймаут не
сработает;
- при срабатывании — лог:
`[llm-proxy] chat model=<m> upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`
(+ `(backoff: отдых …)` если задан backoff), метрика
`llm_proxy_requests_total{result="timeout"}`;
- закрытие апстримом пустого стрима (EOF без единого байта) — **не** таймаут:
это валидный короткий ответ, он проксируется как обычно.
Значение — ISO-8601-длительность, формат как у `backoff` (`PT30S`,
`PT1M30S`, `P1D`, …).
```yaml
providers:
- id: routerai
url: "https://routerai.ru/api/v1"
first_byte_timeout: PT30S # окно для всех моделей провайдера
upstreams:
- id: routerai-gpt4o
provider: routerai
model: gpt-4o
first_byte_timeout: PT0S # у модели watchdog отключён
```
При старте выводится, что настроено:
`[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=<id> ответил — снят с backoff, вернулся в пул`,
метрика `llm_proxy_probe_pings_total{result="ok"}`;
- при молчании — лог `[llm-proxy] probe upstream=<id> молчит — остаюсь в 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`
(`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` — сетевая ошибка/таймаут, `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 = жив) |
Метки: `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|timeout"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -340,15 +559,21 @@ 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)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (ISO-8601)
val probe_interval: Duration? = null, // интервал фонового пинга (ISO-8601)
)
@Serializable
data class UpstreamConf(
val id: String, // наш внутренний id апстрима
val provider: String, // ссылка на ProviderConf.id
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 на модель (приоритет над провайдерским)
val first_byte_timeout: Duration? = null, // окно первого байта стрима (приоритет над провайдерским)
val probe_interval: Duration? = null, // интервал фонового пинга (приоритет над провайдерским)
)
@Serializable
@@ -454,7 +679,7 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
`upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch
== null` пропускаем. Итог — финальное тело запроса;
- шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`;
- при ответе `5xx`/`429` от провайдера — `release(upstream)` и переходим к
- при ответе `5xx`/`429`/`402` от провайдера — `release(upstream)` и переходим к
следующему свободному апстриму из списка (фейловер); исчерпали список —
отдаём `503` (`all upstreams failed`);
- при успехе/отмене клиента — `release(upstream)` в `finally` и выходим.
@@ -507,8 +732,18 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
свободному `upstream=<id>`.
- **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима,
освобождён слот `active[id]=N/limit`.
- **Ошибка апстрима** (не `5xx`/`429`, а сетевая/таймаут): как сейчас — лог с
- **Ошибка апстрима** (не `5xx`/`429`/`402`, а сетевая/таймаут): как сейчас — лог с
сообщением.
- **Backoff**: при фейловере строка дополняется интервалом отдыха
(`(backoff: отдых PT…)`); при пропуске апстрима в откате —
`upstream=<id> в backoff-откате (~Ns) — пропускаю`; при старте — список
настроенных backoff (`backoff: upstreams=… providers=…`).
- **Таймаут первого байта стрима**: при срабатывании watchdog'а —
`upstream=<id> TIMEOUT: первый байт не пришёл за <N>ms → фейловер`;
при старте — список настроенных `first_byte_timeout`
(`first_byte_timeout: upstreams=… providers=… default=PT30S`).
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
`model=<витрина>` с заголовком `Retry-After: N`.
Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать
существующий парсинг логов (если он есть).
+26 -1
View File
@@ -20,10 +20,35 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (
обоих — безлимит.
Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или
апстрима (`off` по умолчанию, `split` — рассуждения из `<think>…</think>` уходят
апстрима (`off` по умолчанию, `split` — рассуждения из `think`-тегов уходят
в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в
non-stream.
Повторяющиеся сбои апстрима увязываются экспоненциальным backoff: поле
`backoff` (ISO-8601-потолок, напр. `P1M`) задаётся на провайдере (общий
счётчик на его модели) и/или на апстриме (приоритет). Первая ошибка — отдых
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
— `503` с `Retry-After`.
Для `stream=true` запросов действует watchdog на первый байт ответа:
`first_byte_timeout` (ISO-8601, дефолт `PT30S`, `PT0S` отключает) — если
апстрим не прислал ни одного байта за окно, слот освобождается, ошибка идёт
в 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. Список
метрик — в CONFIG.md, раздел «Prometheus-метрики».
| Переменная | Default | Описание |
|---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+17 -1
View File
@@ -14,12 +14,14 @@
| --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, неприкосновенность сообщений без `tool_calls` (в т.ч. `[]`) и `user`/`tool`, пустая строка при `emptyOk`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, заполнение поля на **каждом** assistant-сообщении (с `tool_calls`, без, `tool_calls: []` — без текста и `emptyOk=false` не трогаем; `emptyOk=true` — пустая строка), неприкосновенность `user`/`tool`/`system`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
| `NullToleranceTest` | Устойчивость разбора ответов апстрима к JSON-`null` (`choices`, `delta.tool_calls`, `delta.content`) — пропуск вместо исключения; |
| `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; |
| `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`) и конфиг-парсер `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`. |
## Проверка качества тестов (мутационная приёмка)
@@ -81,6 +83,20 @@ Console Go; тело — как у opencode: assistant + `reasoning` + `reasonin
`assistantWithoutToolCallsIsUntouched` и `assistantWithEmptyToolCallsIsUntouched`;
возврат `obj["choices"]?.jsonArray` в `rebuildFromChunks` → падает `rebuildFromChunksToleratesNullChoices`.
**Поправка после боевого 400 (2026-09-13):** живые stateful-сессии Console Go
(сессия, где глюч уже видел thinking-ответы) всё-таки требовали
`reasoning_content` и на assistant-сообщениях **без** `tool_calls` — свежая
сессия такой истории прощала (все формы выше — 200). Опция «на каждом
assistant-сообщении» совпадает с тем, что делает сам opencode («Deepseek
requires all assistant messages to have reasoning on them», пустая
reasoning-часть дописывается тоже). Охрана `tool_calls` в
`applyReasoningField` снята: поле дописывается на каждом assistant-сообщении
(текст — из `reasoning`/`reasoning_details`, при `emptyOk` — пустая строка).
Тесты `assistantWithoutToolCallsIsUntouched`/`assistantWithEmptyToolCallsIsUntouched`
заменены на `assistantWithoutToolCallsAlsoGetsField`,
`assistantWithoutReasoningAndNoEmptyOkIsUntouched`,
`assistantWithoutToolCallsWithEmptyOkFillsEmpty`, `assistantWithEmptyToolCallsGetsField`.
## Известное ограничение
Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
@@ -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()
}
}
@@ -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<String, ProviderConf>,
private val upstreams: List<UpstreamConf>,
private val http: HttpClient,
private val backoff: BackoffRegistry,
private val metrics: MetricsRegistry,
private val scope: CoroutineScope,
) {
private val jobs = mutableMapOf<String, Job>()
/** Запустить пинг-корутины для всех апстримов с включённым пингом. */
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<ByteReadChannel>()
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
}
+275 -24
View File
@@ -28,6 +28,14 @@ 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
import kotlinx.coroutines.selects.onTimeout
import kotlinx.coroutines.selects.select
import kotlinx.coroutines.sync.Mutex
import kotlinx.datetime.Clock
import kotlinx.serialization.Serializable
@@ -46,6 +54,10 @@ 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.days
import kotlin.time.Duration.Companion.hours
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds
import kotlin.time.TimeSource
@@ -86,11 +98,33 @@ fun main() {
val active = config.upstreams.associate { up ->
up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider]))
}
val backoff = BackoffRegistry(providersById, upstreamsById)
val metrics = MetricsRegistry()
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] first_byte_timeout: upstreams=" +
config.upstreams.filter { it.first_byte_timeout != null }.map { up -> "${up.id}=${up.first_byte_timeout?.toIsoString()}" }.joinToString(", ") +
" providers=" +
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}"
}
@@ -100,8 +134,10 @@ 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)
proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
}
}
@@ -114,10 +150,19 @@ fun Application.proxyModule(
active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) {
routing {
get("/metrics") {
call.respondText(
metrics.render(config, active, backoff),
ContentType.parse("text/plain; version=0.0.4"),
HttpStatusCode.OK,
)
}
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") {
handleModels(call, config)
@@ -133,6 +178,8 @@ private suspend fun handleChat(
active: Map<String, UpstreamCounter>,
sessions: SessionRegistry,
http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) {
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val raw = call.receiveText()
@@ -165,7 +212,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 {
@@ -229,8 +276,13 @@ private suspend fun handleChat(
setBody(forwarded.toString())
}.execute { resp ->
upstreamStatus = resp.status.value
if (upstreamStatus >= 500 || upstreamStatus == 429) {
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" }
if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) {
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)
failover = true
return@execute
@@ -238,26 +290,30 @@ private suspend fun handleChat(
if (upstreamStatus in 400..499) {
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
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"
call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus))
responded = true
return@execute
}
responded = true
backoff.recordSuccess(up)
if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus)
val thinkMode = effectiveThinkTags(up, provider)
val firstByteTimeout = effectiveFirstByteTimeout(up, provider)
call.respondBytesWriter(ContentType.parse(ct), status) {
val ch = resp.body<ByteReadChannel>()
if (thinkMode == "off") {
// Флажок не выставлен — сырой байтовый passthrough + выравнивание конца потока.
streamRawWithDoneContract(ch)
streamRawWithDoneContract(ch, firstByteTimeout)
} else {
streamSseWithThinkTags(ch, thinkMode)
streamSseWithThinkTags(ch, thinkMode, firstByteTimeout)
}
}
metrics.record(modelName, up, "ok")
log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true"
@@ -281,6 +337,7 @@ private suspend fun handleChat(
}
}.getOrDefault(ContentType.parse(ct))
call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus))
metrics.record(modelName, up, "ok")
log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false"
@@ -290,11 +347,26 @@ private suspend fun handleChat(
if (failover) continue
if (responded) return
} catch (e: FirstByteTimeoutException) {
val rest = backoff.recordFailure(up)
metrics.record(modelName, up, "timeout")
log.warn {
"[llm-proxy] chat model=$modelName upstream=${up.id} TIMEOUT: первый байт не пришёл за ${e.timeoutMs}ms → фейловер" +
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
}
failed.add(up.id)
continue
} catch (e: CancellationException) {
metrics.record(modelName, up, "cancelled")
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)
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)
continue
} finally {
@@ -305,7 +377,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)
}
}
}
@@ -445,6 +527,38 @@ internal fun effectiveThinkTags(up: UpstreamConf, provider: ProviderConf?): Stri
}
}
/**
* Эффективный first-byte-timeout для стрим-ответа: значение у апстрима, если
* задано; иначе у провайдера; иначе [DEFAULT_FIRST_BYTE_TIMEOUT]. Возвращает
* null, если явно отключено (`Duration.ZERO` — `PT0S`): тогда стрим
* ждёт первого байта без лимита (в пределах общего
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
*/
internal fun effectiveFirstByteTimeout(up: UpstreamConf, provider: ProviderConf?): Duration? {
val raw = up.first_byte_timeout ?: provider?.first_byte_timeout ?: DEFAULT_FIRST_BYTE_TIMEOUT
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,
* метрика `timeout`), но отделён от [CancellationException] (отмена клиентом)
* и от generic `Exception` (сетевые ошибки).
*/
class FirstByteTimeoutException(val timeoutMs: Long) :
RuntimeException("upstream did not send first byte within ${timeoutMs}ms")
/** Освободить слот апстрима (в finally по завершении проксирования). */
internal fun release(up: UpstreamConf, active: Map<String, UpstreamCounter>) {
active.getValue(up.id).release()
@@ -489,15 +603,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<UpstreamConf>,
active: Map<String, UpstreamCounter>,
excluded: Set<String>,
): 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)}" }
@@ -649,11 +776,14 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
/**
* Достроить нативное поле рассуждений для апстримов, которые его требуют
* (Console Go / deepseek в thinking-режиме): если у провайдера объявлено
* `reasoningField`, то в каждом assistant-сообщении с непустым `tool_calls`
* добавляем это поле, ЕСЛИ его там ещё нет. Текст берём из `reasoning`
* `reasoningField`, то в каждом assistant-сообщении добавляем это поле,
* ЕСЛИ его там ещё нет (не только в тех, что с tool_calls — deepseek
* требует `reasoning_content` на КАЖДОМ ассистент-сообщении в thinking-режиме,
* см. «must be passed back to the API»). Текст берём из `reasoning`
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
* пишем пустую строку, только если `emptyOk`.
* пишем пустую строку, только если `emptyOk` — апстрим требует самого
* НАЛИЧИЯ поля, даже пустого (так делает и сам opencode).
* Тело возвращается без изменений (тот же объект), если менять нечего.
*/
internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject {
@@ -664,8 +794,6 @@ internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Bool
val msg = el as? JsonObject ?: return@map el
val role = (msg["role"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (role != "assistant") return@map el
val toolCalls = msg["tool_calls"] as? JsonArray
if (toolCalls == null || toolCalls.isEmpty()) return@map el
val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (existing != null && existing.isNotEmpty()) return@map el
val text = reasoningTextOf(msg)
@@ -765,12 +893,26 @@ internal class StreamEndDetector(private val windowSize: Int = STREAM_SCAN_WINDO
* Сырой байтовый passthrough стрима (`think_tags: off`) с выравниванием контракта
* конца потока: байты уходят клиенту как есть, а если апстрим закрыл поток, не
* прислав `data: [DONE]`, но `finish_reason` в потоке был — дописываем маркер.
* Если задан [firstByteTimeout] (non-null), первый байт тела должен прийти
* за это время, иначе бросаем [FirstByteTimeoutException]; после первого
* байта watchdog выключается (ждём ответа вплоть до общего
* [UPSTREAM_REQUEST_TIMEOUT_MS]).
*/
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteReadChannel) {
internal suspend fun ByteWriteChannel.streamRawWithDoneContract(
source: ByteReadChannel,
firstByteTimeout: Duration? = null,
) {
val detector = StreamEndDetector()
val buf = ByteArray(8192)
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
var isFirstRead = timeoutMs > 0L
while (true) {
val n = source.readAvailable(buf)
val n = if (isFirstRead) {
isFirstRead = false
readAvailableWithTimeout(source, buf, timeoutMs)
} else {
source.readAvailable(buf)
}
if (n == -1) break
if (n > 0) {
writeFully(buf, 0, n)
@@ -791,14 +933,30 @@ internal suspend fun ByteWriteChannel.streamRawWithDoneContract(source: ByteRead
* записывается как `data: <json>\n\n` (событие-граница SSE), а при ошибке парса —
* исходная строка. Каждую строку сразу `flush()`, чтобы стрим не «залипал» в
* буфере. В конце потока накопленные хвосты сплиттеров сбрасываются финиш-чанком.
* Если задан [firstByteTimeout] (non-null), первая строка должна прийти за это
* время, иначе [FirstByteTimeoutException] (watchdog только на первой строке).
*/
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadChannel, thinkMode: String) {
internal suspend fun ByteWriteChannel.streamSseWithThinkTags(
source: ByteReadChannel,
thinkMode: String,
firstByteTimeout: Duration? = null,
) {
val splitters = mutableMapOf<Int, ThinkTagSplitter>()
val addReasoning = thinkMode == "split"
var sawDone = false
var sawFinishReason = false
val timeoutMs = firstByteTimeout?.inWholeMilliseconds ?: 0L
var isFirstRead = timeoutMs > 0L
suspend fun readNext(): String? = if (isFirstRead) {
isFirstRead = false
readLineWithTimeout(source, timeoutMs)
} else {
source.readLine(LineEnding.Lenient)
}
while (true) {
val line = source.readLine(LineEnding.Lenient) ?: break
val line = readNext() ?: break
when {
line.startsWith("data:") -> {
val payload = line.removePrefix("data:").trim()
@@ -851,6 +1009,49 @@ internal suspend fun ByteWriteChannel.streamSseWithThinkTags(source: ByteReadCha
}
}
/**
* Прочитать из [source] в [buf] с жёстким лимитом [timeoutMs]: если за это
* время ни одного байта не пришло, бросить [FirstByteTimeoutException].
* `select` отменяет «проигравшую» ветку; при таймауте явно отменяем и
* дожидаемся чтение (`cancelAndJoin`), чтобы соединение с апстримом
* корректно закрылось. Сентинел [Int.MIN_VALUE] невозможен как результат
* [ByteReadChannel.readAvailable] (это `-1` для EOF или `>= 0` для байт).
*/
internal suspend fun readAvailableWithTimeout(source: ByteReadChannel, buf: ByteArray, timeoutMs: Long): Int =
coroutineScope {
val readJob = async { source.readAvailable(buf) }
val result = select<Int> {
readJob.onAwait { it }
onTimeout(timeoutMs) { Int.MIN_VALUE }
}
if (result == Int.MIN_VALUE) {
readJob.cancelAndJoin()
throw FirstByteTimeoutException(timeoutMs)
}
result
}
/**
* Прочитать строку из [source] с лимитом [timeoutMs]: если за это время ни
* одной строки/байта не пришло, бросить [FirstByteTimeoutException].
* [ByteReadChannel.readLine] возвращает `null` и при EOF, и (через `onTimeout`)
* при таймауте; различаем по состоянию [readJob]: если ветка чтения не
* завершилась к моменту возврата из `select` — таймаут (select отменил её).
*/
internal suspend fun readLineWithTimeout(source: ByteReadChannel, timeoutMs: Long): String? =
coroutineScope {
val readJob = async { source.readLine(LineEnding.Lenient) }
val result = select<String?> {
readJob.onAwait { it }
onTimeout(timeoutMs) { null }
}
if (result == null && !readJob.isCompleted) {
readJob.cancelAndJoin()
throw FirstByteTimeoutException(timeoutMs)
}
result
}
/**
* Пересобрать SSE-чанк: по каждому choice (ключ `index`, дефолт 0) взять
* `delta.content` (только если это JSON-строка; массив частей не трогаем) и
@@ -903,6 +1104,9 @@ data class ProviderConf(
val think_tags: String? = null,
val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class UpstreamConf(
@@ -912,6 +1116,9 @@ data class UpstreamConf(
val max_concurrency: Int? = null,
val patch: JsonObject? = null,
val think_tags: String? = null,
val backoff: Duration? = null,
val first_byte_timeout: Duration? = null,
val probe_interval: Duration? = null,
)
data class ModelConf(
@@ -969,6 +1176,9 @@ 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"),
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
probe_interval = m.durationOrNull("probe_interval"),
)
}
@@ -981,6 +1191,9 @@ 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"),
first_byte_timeout = m.durationOrNull("first_byte_timeout"),
probe_interval = m.durationOrNull("probe_interval"),
)
}
@@ -1006,5 +1219,43 @@ internal fun Map<String, YamlElement>.str(key: String): String =
internal fun Map<String, YamlElement>.strOrNull(key: String): String? =
(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? =
this[key] as? YamlMap
@@ -0,0 +1,130 @@
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 | 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);
* - `llm_proxy_upstream_cooling_seconds{upstream,provider}` — сколько секунд
* апстрим ещё в backoff-откате (0 = жив).
*/
class MetricsRegistry {
private val lock = Mutex()
private val requests = HashMap<RequestKey, Long>()
private val probes = HashMap<ProbeKey, Long>()
private data class RequestKey(
val model: String,
val upstream: String,
val provider: String,
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 {
val key = RequestKey(model, up.id, up.provider, result)
requests[key] = (requests[key] ?: 0L) + 1
}
}
/** Учёт пинга апстрима в 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,
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_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) {
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"
}
@@ -1,4 +1,26 @@
package pw.binom.llmproxy
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
/** Предел времени на весь запрос к апстриму, мс: от отправки до приёма всего ответа (включая стриминг). */
const val UPSTREAM_REQUEST_TIMEOUT_MS: Long = 5 * 60_000L
/**
* Default «watchdog» на первый байт тела стрим-ответа: если апстрим не прислал
* ни одного байта за это время, прокси считает это сбоем апстрима и
* фейловерит (освобождает слот конкурентности, считает ошибку в backoff и
* идёт к следующему апстриму модели). Применяется ТОЛЬКО к stream=true
* запросам; для non-stream работает общий [UPSTREAM_REQUEST_TIMEOUT_MS].
* Отключается явным `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
@@ -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))
}
}
@@ -4,7 +4,9 @@ import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
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 +16,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 +226,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
@@ -583,4 +587,60 @@ class ConfigLogicTest {
assertEquals(null, cfg.providers[0].reasoning_field)
assertEquals(false, cfg.providers[0].reasoning_empty_ok)
}
@Test
fun parseConfigReadsFirstByteTimeoutOnProviderAndUpstream() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
first_byte_timeout: PT45S
- id: p2
url: "https://y.ru/api/v1"
upstreams:
- id: u1
provider: p1
model: real-1
first_byte_timeout: PT0S
- id: u2
provider: p2
model: real-2
models:
- name: m1
upstreams: [u1, u2]
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals(45.seconds, cfg.providers[0].first_byte_timeout)
assertEquals(null, cfg.providers[1].first_byte_timeout)
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)
}
}
@@ -0,0 +1,218 @@
package pw.binom.llmproxy
import io.ktor.utils.io.ByteChannel
import io.ktor.utils.io.close
import io.ktor.utils.io.readAvailable
import io.ktor.utils.io.writeFully
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.milliseconds
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.advanceTimeBy
import kotlinx.coroutines.test.runCurrent
import kotlinx.coroutines.test.runTest
/**
* Тесты watchdog'а на первый байт тела стрим-ответа.
*
* Сценарии:
* - Апстрим не прислал ничего за N мс → [FirstByteTimeoutException] на первой
* попытке чтения (после первой попытки watchdog отключается).
* - Апстрим закрыл канал до таймаута, не прислав данных → возврат EOF (-1 /
* null), НЕ [FirstByteTimeoutException].
* - Апстрим прислал первый байт вовремя → читаем его, watchdog отключён.
* - Когда таймаут не задан (`null`) → читаем без лимита (как и было).
*/
class FirstByteTimeoutTest {
/** Окно для watchdog'а в тестах: виртуальное время `runTest` его мгновенно прокручивает. */
private val watchWindow = 100.milliseconds
@Test
fun readAvailableTimesOutWhenNoBytesAndNotClosed() = runTest {
val src = ByteChannel(autoFlush = true)
val buf = ByteArray(64)
val done = CompletableDeferred<Throwable?>()
launch {
try {
readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
assertEquals(watchWindow.inWholeMilliseconds, e.timeoutMs)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
}
@Test
fun readAvailableReturnsEofImmediatelyWhenChannelClosedEmpty() = runTest {
// Канал сразу закрыт без данных: -1 (EOF), никакого таймаута.
val src = ByteChannel(autoFlush = true)
src.close(null)
val buf = ByteArray(64)
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
assertEquals(-1, n, "закрытый пустой канал → EOF, не таймаут")
}
@Test
fun readAvailableReturnsBytesWhenDataArrivesInTime() = runTest {
val src = ByteChannel(autoFlush = true)
val payload = "hello".encodeToByteArray()
src.writeFully(payload)
val buf = ByteArray(64)
val n = readAvailableWithTimeout(src, buf, watchWindow.inWholeMilliseconds)
assertEquals(payload.size, n)
assertEquals(payload.decodeToString(), buf.decodeToString(0, n))
}
@Test
fun readLineTimesOutWhenNoLinesAndNotClosed() = runTest {
val src = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
}
@Test
fun readLineReturnsNullWhenChannelClosedEmpty() = runTest {
// Закрытый пустой канал → null (EOF), не таймаут. Это критично: иначе бы
// мы фейловерили легитимные короткие ответы, где апстрим закрыл стрим
// до первой строки.
val src = ByteChannel(autoFlush = true)
src.close(null)
val line = readLineWithTimeout(src, watchWindow.inWholeMilliseconds)
assertEquals(null, line)
}
@Test
fun rawStreamTimesOutWhenUpstreamSilent() = runTest {
// Сырой passthrough: апстрим молчит вечно → на первой попытке чтения
// вылетает FirstByteTimeoutException, в out ничего не пишется.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
out.streamRawWithDoneContract(src, watchWindow)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
// Подтверждаем, что в out ничего не утекло.
out.close(null)
val buf = ByteArray(64)
val n = out.readAvailable(buf)
assertEquals(-1, n, "до таймаута клиенту ничего не должно уйти")
}
@Test
fun thinkStreamTimesOutWhenUpstreamSilent() = runTest {
// SSE-вариант: первая строка не приходит → FirstByteTimeoutException,
// в out ничего не пишется.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val done = CompletableDeferred<Throwable?>()
launch {
try {
out.streamSseWithThinkTags(src, "split", watchWindow)
done.complete(AssertionError("expected FirstByteTimeoutException"))
} catch (e: FirstByteTimeoutException) {
done.complete(null)
} catch (e: Throwable) {
done.complete(e)
}
}
advanceTimeBy(watchWindow.inWholeMilliseconds + 50)
runCurrent()
val err = done.await()
if (err != null) throw err
out.close(null)
val buf = ByteArray(64)
val n = out.readAvailable(buf)
assertEquals(-1, n)
}
@Test
fun thinkStreamFirstLineInTimeThenReadsUnlimited() = runTest {
// Первый байт/строка пришли вовремя → watchdog выключен, остальной
// поток читается обычным порядком. Проверяем, что НЕ бросило
// FirstByteTimeoutException и данные дошли до out.
val src = ByteChannel(autoFlush = true)
val out = ByteChannel(autoFlush = true)
val payload = "data: {\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hi\"}}]}\n\n"
src.writeFully(payload.encodeToByteArray())
src.close(null)
out.streamSseWithThinkTags(src, "split", watchWindow)
out.close(null)
val sb = StringBuilder()
val buf = ByteArray(256)
while (true) {
val n = out.readAvailable(buf)
if (n == -1) break
if (n > 0) sb.append(buf.decodeToString(0, n))
}
assertTrue(sb.toString().contains("data:"), "данные дошли до out: $sb")
}
@Test
fun effectiveFirstByteTimeoutPrefersUpstreamThenProviderThenDefault() {
// upstream > provider > default
val provider = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 20.milliseconds)
val upstream = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 5.milliseconds)
val upstreamNoOwn = UpstreamConf("u", "p", "m")
assertEquals(5.milliseconds, effectiveFirstByteTimeout(upstream, provider))
assertEquals(20.milliseconds, effectiveFirstByteTimeout(upstreamNoOwn, provider))
assertEquals(DEFAULT_FIRST_BYTE_TIMEOUT, effectiveFirstByteTimeout(upstreamNoOwn, null))
}
@Test
fun effectiveFirstByteTimeoutZeroDisablesWatchdog() {
// 0s / PT0S → null: passthrough без watchdog'а.
val up = UpstreamConf("u", "p", "m", null, null, null, null, first_byte_timeout = 0.milliseconds)
assertEquals(null, effectiveFirstByteTimeout(up, null))
val providerZero = ProviderConf("p", "https://x", "", null, null, first_byte_timeout = 0.milliseconds)
val upNoOwn = UpstreamConf("u", "p", "m")
assertEquals(null, effectiveFirstByteTimeout(upNoOwn, providerZero))
}
}
@@ -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")
}
}
@@ -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))
}
}
@@ -74,20 +74,42 @@ class ReasoningFieldTest {
}
@Test
fun assistantWithoutToolCallsIsUntouched() {
// Без tool_calls assistant-сообщение не трогаем (тело идентично).
fun assistantWithoutToolCallsAlsoGetsField() {
// DeepSeek требует reasoning_content на КАЖДОМ assistant-сообщении
// (thinking-режим), а не только на тех, что с tool_calls.
val m = """{"role":"assistant","reasoning":"R","content":"hi"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("R", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun assistantWithoutReasoningAndNoEmptyOkIsUntouched() {
// assistant-сообщение без рассуждений и без tool_calls:
// текста нет, emptyOk=false — поле не добавляем.
val m = """{"role":"assistant","content":"hi"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
}
@Test
fun assistantWithEmptyToolCallsIsUntouched() {
// Пустой массив tool_calls = не трогаем (тело идентично).
fun assistantWithoutToolCallsWithEmptyOkFillsEmpty() {
// Без tool_calls и без текста, но emptyOk=true — дописываем пустую строку,
// т.к. апстриму важно само НАЛИЧИЕ поля.
val m = """{"role":"assistant","content":"hi"}"""
val out = apply(m, "reasoning_content", true)
val m0 = outMsgs(out)[0]
assertEquals("", str(m0, "reasoning_content"))
assertTrue(m0.containsKey("reasoning_content"))
}
@Test
fun assistantWithEmptyToolCallsGetsField() {
// Пустой массив tool_calls больше не исключает сообщение из обработки:
// текст из reasoning дописывается в нативное поле.
val m = """{"role":"assistant","tool_calls":[],"reasoning":"R"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertEquals("R", str(outMsgs(out)[0], "reasoning_content"))
}
@Test