7 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
subochev 8154485b3b feat: нативное поле рассуждений (reasoning_field) — правка 400 от Console Go
Build LLM Proxy / Build and push (release) Successful in 39s
Диагноз: апстрим deepseek-v4.1-flash (провайдер opencode, Console Go) в thinking-режиме
требует reasoning_content в assistant-сообщениях с tool_calls, а клиент opencode
присылает рассуждения как reasoning + reasoning_details. Отсюда 400
'The reasoning_content in the thinking mode must be passed back to the API'.

- providers[].reasoning_field / reasoning_empty_ok: прокси аддитивно достраивает
  нативное поле в assistant-сообщениях с непустым tool_calls (текст из reasoning
  или reasoning_details[].text, тип reasoning.text); существующее непустое поле
  не перезаписывается, ничего не переименовывается, прочие сообщения не трогаются;
- устойчивость разбора ответов к JSON-null (choices/delta/content/tool_calls) —
  было 290 фейловеров с локального апстрима на платные из-за нашего же исключения;
- лог: класс исключения в сообщении об ошибке; для 4xx логируется тело ответа
  апстрима (читается безопасно: 4xx — не стрим);
- тесты: ReasoningFieldTest (13), NullToleranceTest (6), ConfigLogicTest (+2) — 100 всего;
- CONFIG.md, TESTING.md (фактические замеры A/B против Console Go).
2026-09-12 23:46:42 +03:00
15 changed files with 2071 additions and 64 deletions
+283 -7
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
@@ -218,6 +232,249 @@ upstreams:
рассуждения». Если `content` не строка (мультимодальный массив частей) — ответ
не трогаем. Без флажка (`off`) ответ идёт байт-в-байт как раньше.
### Нативное поле рассуждений (`reasoning_field`)
Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**)
в 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-сообщении, даже пустую.)
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в
**каждом** assistant-сообщении — берёт текст из `reasoning`
(если это непустая строка), иначе склеивает `reasoning_details[*].text`
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
непустое поле не перезаписывается. Ничего при этом не убирается и не
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
поле просто дополняется. Тронуты только assistant-сообщения; `user`/`tool`/
`system` не меняются.
Если текста рассуждений нет вовсе — поле не добавляется, кроме случая
`reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если
`reasoning_field` не задан — тело не меняется вовсе.
```yaml
providers:
- id: opencode
url: "https://opencode.ai/zen/go/v1"
reasoning_field: reasoning_content # Console Go / deepseek в thinking-режиме
# 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` (из примера выше), маршрут уходит на апстрим
@@ -296,18 +553,27 @@ data class ProviderConf(
val id: String,
val url: String,
val key: String = "",
val max_concurrency: Int? = null, // лимит по умолчанию для апстримов провайдера
val patch: JsonObject? = null, // ко всем запросам провайдера
val session_header: String? = null, // заголовок-сессия, считается из истории
val max_concurrency: Int? = null,
val patch: JsonObject? = null,
val session_header: String? = null,
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
@@ -413,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` и выходим.
@@ -466,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-конфигу |
+56 -1
View File
@@ -13,11 +13,15 @@
| Файл | Что проверяет |
| --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига, приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `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`. |
## Проверка качества тестов (мутационная приёмка)
@@ -44,6 +48,57 @@
Правило: боевой код нельзя подгонять под тест; если тест не проходит, неверен тест.
## Поле рассуждений для Console Go (релиз 11) — замеры приёмки
Причина правки: апстрим `deepseek-v4.1-flash` у провайдера `opencode` (Console Go) в thinking-режиме
требует `reasoning_content` в assistant-сообщениях с `tool_calls`, а клиент opencode присылает
рассуждения как `reasoning` + `reasoning_details` — отсюда `400 The reasoning_content in the
thinking mode must be passed back to the API`.
**Границы требования (замер прямыми запросами к Console Go, одинаковое тело):**
| assistant-сообщение | HTTP |
| --- | --- |
| с `tool_calls`, без reasoning вовсе | 400 |
| с `tool_calls`, `reasoning` + `reasoning_details` (форма opencode) | 400 |
| с `tool_calls`, только `reasoning_details` | 400 |
| с `tool_calls`, `reasoning_content` + `reasoning` + `reasoning_details` (аддитивно) | 200 |
| с `tool_calls`, `reasoning_content: ""` | 200 |
| без `tool_calls`, без reasoning | 200 |
**Живая приёмка (локальный инстанс на 8101, конфиг-копия боевого, модель с единственным апстримом
Console Go; тело — как у opencode: assistant + `reasoning` + `reasoning_details` + `tool_calls`):**
| Конфиг | `stream=false` | `stream=true` |
| --- | --- | --- |
| без правки (`reasoning_field` не задан) | **400** — та самая ошибка про `reasoning_content` | **400** |
| с правкой (`reasoning_field: reasoning_content`) | **200**, модель продолжила диалог после tool-результата | **200** |
**Регрессия на боевых цепочках (тот же инстанс, то же тело с `tool_calls`):** `codding-big` → 200
(ушло на `minimax-m3`, поле не добавляется — флаг объявлен только у провайдера `opencode`),
`codding` → 200 (`qwen-3.8`, local), `assistant` → 200 (Console Go с правкой).
**Мутационная приёмка:** 4 мутации из ТЗ отработал кодер, две проверены вручную
(`cleanAllTests jvmTest`): снятие охраны `tool_calls` в `applyReasoningField` → падают
`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. Реальный апстрим проверяется только после деплоя.
Правка поля рассуждений — исключение: она проверена живым инстансом против настоящего Console Go (таблица выше) до релиза.
@@ -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
}
+359 -43
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 {
@@ -176,7 +223,8 @@ private suspend fun handleChat(
continue
}
val patched = buildBody(bodyJson, provider, up, modelConf)
val patched0 = buildBody(bodyJson, provider, up, modelConf)
val patched = applyReasoningField(patched0, provider.reasoning_field, provider.reasoning_empty_ok)
val forwarded = if (clientWantsStream) {
patched
} else {
@@ -228,26 +276,44 @@ 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
}
if (upstreamStatus in 400..499) {
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
val errorBody = runCatching { resp.body<String>() }.getOrDefault("")
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"
@@ -271,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"
@@ -280,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.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 {
@@ -295,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)
}
}
}
@@ -435,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()
@@ -479,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)}" }
@@ -538,27 +675,27 @@ internal fun rebuildFromChunks(sse: String): String {
val data = line.removePrefix("data:").trim()
if (data.isEmpty() || data == "[DONE]") return@forEach
val obj = runCatching { json.parseToJsonElement(data).jsonObject }.getOrNull() ?: return@forEach
if (id.isEmpty()) id = obj["id"]?.jsonPrimitive?.content ?: ""
if (created == null) created = obj["created"]?.jsonPrimitive?.content?.toLongOrNull()
if (model.isEmpty()) model = obj["model"]?.jsonPrimitive?.content ?: ""
if (systemFingerprint == null) systemFingerprint = obj["system_fingerprint"]?.jsonPrimitive?.content
if (serviceTier == null) serviceTier = obj["service_tier"]?.jsonPrimitive?.content
if (id.isEmpty()) id = (obj["id"] as? JsonPrimitive)?.content ?: ""
if (created == null) created = (obj["created"] as? JsonPrimitive)?.content?.toLongOrNull()
if (model.isEmpty()) model = (obj["model"] as? JsonPrimitive)?.content ?: ""
if (systemFingerprint == null) systemFingerprint = (obj["system_fingerprint"] as? JsonPrimitive)?.content
if (serviceTier == null) serviceTier = (obj["service_tier"] as? JsonPrimitive)?.content
if (provider == null) provider = obj["provider"]
(obj["error"] as? JsonObject)?.let { error = it }
(obj["usage"] as? JsonObject)?.let { usage = it }
val chArr = obj["choices"]?.jsonArray ?: return@forEach
val chArr = obj["choices"] as? JsonArray ?: return@forEach
for (ch in chArr) {
val c = ch.jsonObject
val idx = c["index"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0
val c = ch as? JsonObject ?: continue
val idx = (c["index"] as? JsonPrimitive)?.content?.toIntOrNull() ?: 0
val mc = choices.getOrPut(idx) { MutableChoice() }
val delta = c["delta"]?.jsonObject
val delta = c["delta"] as? JsonObject
if (delta != null) {
if (mc.role == null) mc.role = delta["role"]?.jsonPrimitive?.content
delta["content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.content.append(it) }
delta["reasoning_content"]?.jsonPrimitive?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) }
delta["tool_calls"]?.jsonArray?.forEach { tc -> (tc as? JsonObject)?.let { mc.toolCalls.add(it) } }
if (mc.role == null) mc.role = (delta["role"] as? JsonPrimitive)?.content
(delta["content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.content.append(it) }
(delta["reasoning_content"] as? JsonPrimitive)?.content?.takeIf { it != "null" }?.let { mc.reasoning.append(it) }
(delta["tool_calls"] as? JsonArray)?.forEach { tc -> (tc as? JsonObject)?.let { mc.toolCalls.add(it) } }
}
c["finish_reason"]?.jsonPrimitive?.content?.takeIf { it.isNotEmpty() && it != "null" }?.let { mc.finishReason = it }
(c["finish_reason"] as? JsonPrimitive)?.content?.takeIf { it.isNotEmpty() && it != "null" }?.let { mc.finishReason = it }
c["logprobs"]?.let { mc.logprobs = it }
}
}
@@ -610,12 +747,12 @@ internal fun rebuildFromChunks(sse: String): String {
* изменился → null (отдать исходную строку как есть).
*/
internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String? {
val choices = obj["choices"]?.jsonArray ?: return null
val choices = obj["choices"] as? JsonArray ?: return null
val addReasoning = thinkMode == "split"
var changed = false
val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject
val message = choice["message"]?.jsonObject ?: return@map choiceEl
val choice = choiceEl as? JsonObject ?: return@map choiceEl
val message = choice["message"] as? JsonObject ?: return@map choiceEl
val contentStr = (message["content"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl
val splitter = ThinkTagSplitter(thinkMode)
@@ -636,6 +773,56 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
return JsonObject(obj.toMutableMap().apply { this["choices"] = JsonArray(newChoices) }).toString()
}
/**
* Достроить нативное поле рассуждений для апстримов, которые его требуют
* (Console Go / deepseek в thinking-режиме): если у провайдера объявлено
* `reasoningField`, то в каждом assistant-сообщении добавляем это поле,
* ЕСЛИ его там ещё нет (не только в тех, что с tool_calls — deepseek
* требует `reasoning_content` на КАЖДОМ ассистент-сообщении в thinking-режиме,
* см. «must be passed back to the API»). Текст берём из `reasoning`
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
* пишем пустую строку, только если `emptyOk` — апстрим требует самого
* НАЛИЧИЯ поля, даже пустого (так делает и сам opencode).
* Тело возвращается без изменений (тот же объект), если менять нечего.
*/
internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject {
if (field == null) return body
val messages = body["messages"] as? JsonArray ?: return body
var changed = false
val newMessages = messages.map { el ->
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 existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (existing != null && existing.isNotEmpty()) return@map el
val text = reasoningTextOf(msg)
if (text.isEmpty() && !emptyOk) return@map el
changed = true
JsonObject(msg.toMutableMap().apply { this[field] = JsonPrimitive(text) })
}
if (!changed) return body
return JsonObject(body.toMutableMap().apply { this["messages"] = JsonArray(newMessages) })
}
/**
* Текст рассуждений сообщения: `reasoning` (если непустая строка), иначе
* склейка `reasoning_details[*].text` через "\n" — только элементы, у которых
* `type` отсутствует или равен "reasoning.text".
*/
internal fun reasoningTextOf(msg: JsonObject): String {
val reasoning = (msg["reasoning"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (reasoning != null && reasoning.isNotEmpty()) return reasoning
val details = msg["reasoning_details"] as? JsonArray ?: return ""
val parts = details.mapNotNull { el ->
val d = el as? JsonObject ?: return@mapNotNull null
val type = (d["type"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (type != null && type != "reasoning.text") return@mapNotNull null
(d["text"] as? JsonPrimitive)?.takeIf { it.isString }?.content
}
return parts.joinToString("\n")
}
/** Терминальный маркер SSE, который клиенты (в т.ч. Bifrost) ждут как конец потока. */
private const val SSE_DONE_MARKER = "data: [DONE]\n\n"
@@ -665,9 +852,11 @@ internal fun hasNonNullFinishReasonText(text: String): Boolean {
/** Структурная проверка чанка: в `choices[*].finish_reason` есть непустая строка. */
internal fun hasFinishReason(obj: JsonObject): Boolean =
obj["choices"]?.jsonArray?.any { el ->
val fr = (el.jsonObject["finish_reason"] as? JsonPrimitive)?.takeIf { it.isString }?.content
!fr.isNullOrEmpty()
(obj["choices"] as? JsonArray)?.any { el ->
(el as? JsonObject)?.let { choice ->
val fr = (choice["finish_reason"] as? JsonPrimitive)?.takeIf { it.isString }?.content
!fr.isNullOrEmpty()
} ?: false
} ?: false
/**
@@ -704,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)
@@ -730,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()
@@ -790,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-строка; массив частей не трогаем) и
@@ -805,14 +1067,14 @@ internal fun transformThinkChunk(
thinkMode: String,
addReasoning: Boolean,
): String? {
val choices = obj["choices"]?.jsonArray ?: return null
val choices = obj["choices"] as? JsonArray ?: return null
var changed = false
val newChoices = choices.map { choiceEl ->
val choice = choiceEl.jsonObject
val delta = choice["delta"]?.jsonObject
val choice = choiceEl as? JsonObject ?: return@map choiceEl
val delta = choice["delta"] as? JsonObject
val contentStr = (delta?.get("content") as? JsonPrimitive)?.takeIf { it.isString }?.content
if (contentStr == null) return@map choiceEl
val idx = choice["index"]?.jsonPrimitive?.content?.toIntOrNull() ?: 0
val idx = (choice["index"] as? JsonPrimitive)?.content?.toIntOrNull() ?: 0
val splitter = splitters.getOrPut(idx) { ThinkTagSplitter(thinkMode) }
val (newContent, reasoning) = splitter.feed(contentStr)
changed = true
@@ -840,6 +1102,11 @@ data class ProviderConf(
val patch: JsonObject? = null,
val session_header: String? = null,
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(
@@ -849,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(
@@ -904,6 +1174,11 @@ internal fun parseConfig(root: YamlElement): Config {
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
session_header = m.strOrNull("session_header"),
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"),
)
}
@@ -916,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"),
)
}
@@ -941,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
@@ -546,4 +550,97 @@ class ConfigLogicTest {
assertEquals("off", effectiveThinkTags(upNull, prov(null)))
assertEquals("off", effectiveThinkTags(upNull, null))
}
@Test
fun parseConfigReadsReasoningFieldAndEmptyOkOnProvider() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
reasoning_field: reasoning_content
reasoning_empty_ok: true
- id: p2
url: "https://y.ru/api/v1"
reasoning_empty_ok: false
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
assertEquals("reasoning_content", cfg.providers[0].reasoning_field)
assertEquals(true, cfg.providers[0].reasoning_empty_ok)
assertEquals(null, cfg.providers[1].reasoning_field)
assertEquals(false, cfg.providers[1].reasoning_empty_ok)
}
@Test
fun parseConfigDefaultsReasoningFieldsWhenAbsent() {
val yaml = """
providers:
- id: p1
url: "https://x.ru/api/v1"
models:
- name: m1
upstreams: []
""".trimIndent()
val cfg = parseConfig(Yaml.decodeYamlFromString(yaml))
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))
}
}
@@ -0,0 +1,71 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.jsonArray
import kotlinx.serialization.json.jsonObject
import kotlinx.serialization.json.jsonPrimitive
class NullToleranceTest {
private val json = Json { ignoreUnknownKeys = true }
@Test
fun rebuildFromChunksToleratesNullChoices() {
// Чанк с choices:null не должен ронять сборку; значимые поля сохраняются.
val sse = """
data: {"id":"c1","created":123,"model":"m","choices":null}
data: {"id":"c1","created":123,"model":"m","choices":[{"index":0,"delta":{"content":"Hi"}}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
assertEquals("c1", out["id"]?.jsonPrimitive?.content)
assertEquals("m", out["model"]?.jsonPrimitive?.content)
assertEquals(123, out["created"]?.jsonPrimitive?.content?.toLong())
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("Hi", choice?.get("message")?.jsonObject?.get("content")?.jsonPrimitive?.content)
}
@Test
fun rebuildFromChunksToleratesNullDeltaToolCalls() {
val sse = """
data: {"id":"c1","model":"m","choices":[{"index":0,"delta":{"role":"assistant","tool_calls":null},"finish_reason":"stop"}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("assistant", choice?.get("message")?.jsonObject?.get("role")?.jsonPrimitive?.content)
}
@Test
fun rebuildFromChunksToleratesNullContent() {
val sse = """
data: {"id":"c1","model":"m","choices":[{"index":0,"delta":{"role":"assistant","content":null},"finish_reason":"stop"}]}
data: [DONE]
""".trimIndent()
val out = json.parseToJsonElement(rebuildFromChunks(sse)).jsonObject
val choice = out["choices"]?.jsonArray?.get(0)?.jsonObject
assertEquals("assistant", choice?.get("message")?.jsonObject?.get("role")?.jsonPrimitive?.content)
}
@Test
fun hasFinishReasonToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertFalse(hasFinishReason(obj))
}
@Test
fun transformThinkChunkToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertNull(transformThinkChunk(obj, mutableMapOf(), "split", true))
}
@Test
fun transformThinkMessageToleratesNullChoices() {
val obj = json.parseToJsonElement("""{"choices":null}""").jsonObject
assertNull(transformThinkMessage(obj, "split"))
}
}
@@ -0,0 +1,187 @@
package pw.binom.llmproxy
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.JsonPrimitive
import kotlinx.serialization.json.jsonObject
class ReasoningFieldTest {
private val json = Json { ignoreUnknownKeys = true }
// Минимальный непустой tool_call для assistant-сообщения.
private val TC = """{"id":"t1","type":"function","function":{"name":"f","arguments":"{}"}}"""
private fun str(obj: JsonObject, key: String): String? = (obj[key] as? JsonPrimitive)?.content
private fun outMsgs(out: JsonObject): List<JsonObject> =
(out["messages"] as? JsonArray)?.mapNotNull { it as? JsonObject } ?: emptyList()
// Собрать тело `{"messages":[...]}` и прогнать через applyReasoningField.
private fun apply(messagesJson: String, field: String?, emptyOk: Boolean): JsonObject {
val body = json.parseToJsonElement("""{"messages":[$messagesJson]}""").jsonObject
return applyReasoningField(body, field, emptyOk)
}
@Test
fun documentedOpencodeShapeGetsNativeReasoningContent() {
// Форма клиента opencode: reasoning (строка) + reasoning_details (массив),
// нативного поля нет — добавляем его, текст берём из рассуждений.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R","reasoning_details":[{"type":"reasoning.text","text":"R","index":0}]}"""
val out = apply(m, "reasoning_content", false)
val m0 = outMsgs(out)[0]
assertEquals("R", str(m0, "reasoning_content"))
assertEquals("R", str(m0, "reasoning"))
assertTrue((m0["tool_calls"] as? JsonArray)?.isNotEmpty() == true)
assertTrue((m0["reasoning_details"] as? JsonArray)?.isNotEmpty() == true)
}
@Test
fun reasoningDetailsTextIsCopiedIntoReasoningContent() {
// Только reasoning_details (без строки reasoning) — текст всё равно достаётся.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.text","text":"D"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("D", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun reasoningDetailsWithForeignTypeAreIgnored() {
// Чужой тип (reasoning.encrypted) игнорируется, берётся только reasoning.text.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.encrypted","text":"X"},{"type":"reasoning.text","text":"T"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("T", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun multipleReasoningDetailsAreJoinedInOrder() {
// Два reasoning.text склеиваются через "\n" в порядке массива.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning_details":[{"type":"reasoning.text","text":"a"},{"type":"reasoning.text","text":"b"}]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("a\nb", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun existingNonEmptyReasoningContentIsNotOverwritten() {
// Непустое нативное поле не перезаписываем текстом из рассуждений.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R","reasoning_content":"нативное"}"""
val out = apply(m, "reasoning_content", false)
assertEquals("нативное", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
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 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("R", str(outMsgs(out)[0], "reasoning_content"))
}
@Test
fun userAndToolMessagesAreUntouched() {
// user/tool/system с теми же полями — не assistant, не трогаем.
val msgs = """
{"role":"user","tool_calls":[$TC],"reasoning":"R"},
{"role":"tool","tool_calls":[$TC],"reasoning":"R"},
{"role":"system","tool_calls":[$TC],"reasoning":"R"}
""".trimIndent().replace("\n", " ")
val out = apply(msgs, "reasoning_content", false)
val roles = outMsgs(out).map { str(it, "role") }
assertEquals(listOf("user", "tool", "system"), roles)
outMsgs(out).forEach { assertFalse(it.containsKey("reasoning_content")) }
}
@Test
fun emptyTextWithEmptyOkFillsEmptyString() {
// Рассуждений нет вовсе, но emptyOk=true — пишем пустую строку.
val m = """{"role":"assistant","tool_calls":[$TC]}"""
val out = apply(m, "reasoning_content", true)
val m0 = outMsgs(out)[0]
assertEquals("", str(m0, "reasoning_content"))
assertTrue(m0.containsKey("reasoning_content"))
}
@Test
fun emptyTextWithoutEmptyOkLeavesMessageUnchanged() {
// Рассуждений нет и emptyOk=false — сообщение не трогаем.
val m = """{"role":"assistant","tool_calls":[$TC]}"""
val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
}
@Test
fun reasoningFieldNullLeavesBodyUntouched() {
// Поле не объявлено у провайдера (null) — тело не трогаем.
val m = """{"role":"assistant","tool_calls":[$TC],"reasoning":"R"}"""
val out = apply(m, null, false)
assertEquals("""{"messages":[$m]}""", out.toString())
}
@Test
fun messagesAbsentOrNotArrayLeavesBodyUntouched() {
// Нет `messages` — не трогаем.
val noMessages = json.parseToJsonElement("""{"model":"m1"}""").jsonObject
assertEquals(noMessages.toString(), applyReasoningField(noMessages, "reasoning_content", false).toString())
// `messages: null` (JsonNull) — тоже не трогаем.
val nullMessages = json.parseToJsonElement("""{"messages":null}""").jsonObject
val out = applyReasoningField(nullMessages, "reasoning_content", false)
assertEquals(nullMessages.toString(), out.toString())
}
@Test
fun messageOrderAndOtherFieldsArePreserved() {
// Меняется только 2-е (assistant с tool_calls), порядок и остальные поля на месте.
val msgs = """
{"role":"user","content":"hi"},
{"role":"assistant","tool_calls":[$TC],"reasoning":"R","content":""},
{"role":"tool","tool_call_id":"t1","content":"ok"},
{"role":"user","content":"again"}
""".trimIndent().replace("\n", " ")
val out = apply(msgs, "reasoning_content", false)
val list = outMsgs(out)
assertEquals(listOf("user", "assistant", "tool", "user"), list.map { str(it, "role") })
assertEquals("R", str(list[1], "reasoning_content"))
assertTrue((list[1]["tool_calls"] as? JsonArray)?.isNotEmpty() == true)
assertTrue(list[1].containsKey("reasoning"))
assertEquals("hi", str(list[0], "content"))
assertEquals("ok", str(list[2], "content"))
assertEquals("again", str(list[3], "content"))
}
}