4 Commits
11 ... 13

Author SHA1 Message Date
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
10 changed files with 730 additions and 47 deletions
+108 -10
View File
@@ -41,9 +41,15 @@
запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос запросов сейчас идёт на каждый апстрим, и по этим счётчикам решаем, брать запрос
или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать, или нет. Мы **не полагаемся** на `503`/`429` от самого апстрима, чтобы узнать,
что он перегружен: такой ответ от провайдера — это уже сбой, а не способ что он перегружен: такой ответ от провайдера — это уже сбой, а не способ
управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`, мы управления нагрузкой. Если выбранный апстрим всё же вернул `5xx`/`429`/`402`, мы
освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка освобождаем его слот и, при наличии, пробуем **следующий свободный** из списка
модели (фейловер); если свободных не осталось — отдаём `503` от прокси. модели (фейловер); если свободных не осталось — отдаём `503` от прокси.
`402` (insufficient balance / исчерпан лимит токенов) — ошибка аккаунта
провайдера, но фейловер имеет смысл: следующий апстрим может быть платёжеспособен.
Повторяющиеся сбои уводятся из ротации экспоненциальным **backoff** (поле
`backoff`, ниже): упавший апстрим «отдыхает» 1s, 2s, 4s, … до заданного потолка,
и прокси его не трогает. Если **все** апстримы модели в откате — `503` с
заголовком `Retry-After`.
## Формат файла ## Формат файла
@@ -77,6 +83,8 @@ providers:
patch: # уровень провайдера: ко всем его запросам patch: # уровень провайдера: ко всем его запросам
provider: provider:
allow_fallbacks: false allow_fallbacks: false
backoff: P1M # опционально; потолок экспоненциального backoff
# на весь провайдер (ISO-8601: P1M, PT15M, P1D…)
- id: local-llama - id: local-llama
url: "http://10.0.0.5:8080/v1" url: "http://10.0.0.5:8080/v1"
@@ -93,6 +101,8 @@ upstreams:
patch: # уровень апстрима patch: # уровень апстрима
provider: provider:
ignore: [deepseek] ignore: [deepseek]
backoff: PT15M # опционально; собственный backoff этой модели
# (приоритет над backoff провайдера)
- id: local-qwen - id: local-qwen
provider: local-llama provider: local-llama
@@ -221,12 +231,15 @@ upstreams:
### Нативное поле рассуждений (`reasoning_field`) ### Нативное поле рассуждений (`reasoning_field`)
Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**) Некоторые шлюзы-апстримы (например, **Console Go / deepseek в thinking-режиме**)
в thinking-режиме **требуют** вернуть им нативное поле рассуждений в каждом в thinking-режиме **требуют** вернуть им нативное поле рассуждений в **каждом**
assistant-сообщении с непустым `tool_calls`. Если его нет — апстрим отвечает assistant-сообщении истории (не только в тех, что с `tool_calls`). Если его нет —
`400` (`The `reasoning_content` in the thinking mode must be passed back to the апстрим отвечает `400`
(`The `reasoning_content` in the thinking mode must be passed back to the
API.`). Клиенты при этом рассуждения держат в своих форматах: `reasoning` API.`). Клиенты при этом рассуждения держат в своих форматах: `reasoning`
(строка) и/или `reasoning_details` (массив `{type:"reasoning.text", text, ...}`) (строка) и/или `reasoning_details` (массив `{type:"reasoning.text", text, ...}`)
— а нативное поле `reasoning_content` в истории могут и не передавать. — а нативное поле `reasoning_content` в истории могут и не передавать.
(Точно так же делает и сам opencode: для deepseek он добавляет reasoning-часть
на каждом assistant-сообщении, даже пустую.)
Поля (только на уровне **провайдера**, это свойство шлюза, а не модели): Поля (только на уровне **провайдера**, это свойство шлюза, а не модели):
@@ -235,15 +248,15 @@ API.`). Клиенты при этом рассуждения держат в с
| `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) | | `providers[].reasoning_field` | строка / отсутствует | Имя нативного поля рассуждений у шлюза-апстрима (пример: `reasoning_content`) |
| `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе | | `providers[].reasoning_empty_ok` | булево / `false` | Писать ли пустую строку, если текста рассуждений нет вовсе |
Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в каждом Зачем: при отправке запроса прокси **аддитивно** достраивает это поле в
assistant-сообщении с непустым `tool_calls` — берёт текст из `reasoning` **каждом** assistant-сообщении — берёт текст из `reasoning`
(если это непустая строка), иначе склеивает `reasoning_details[*].text` (если это непустая строка), иначе склеивает `reasoning_details[*].text`
(только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и (только элементы без `type` или с `type == "reasoning.text"`, через `"\n"`) и
записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее записывает в поле `reasoning_field`, **если его там ещё нет**. Существующее
непустое поле не перезаписывается. Ничего при этом не убирается и не непустое поле не перезаписывается. Ничего при этом не убирается и не
переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте, переименовывается — клиентский `reasoning`/`reasoning_details` остаются на месте,
поле просто дополняется. Тронуто только assistant-сообщение с непустым поле просто дополняется. Тронуты только assistant-сообщения; `user`/`tool`/
`tool_calls`; сообщения без `tool_calls` и `user`/`tool`/`system` не меняются. `system` не меняются.
Если текста рассуждений нет вовсе — поле не добавляется, кроме случая Если текста рассуждений нет вовсе — поле не добавляется, кроме случая
`reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если `reasoning_empty_ok: true` (тогда пишется пустая строка `""`). Если
`reasoning_field` не задан — тело не меняется вовсе. `reasoning_field` не задан — тело не меняется вовсе.
@@ -256,6 +269,82 @@ providers:
# reasoning_empty_ok: true # опционально; дефолт false # reasoning_empty_ok: true # опционально; дефолт false
``` ```
### Экспоненциальный backoff (`backoff`)
Если апстрим регулярно ошибается (`5xx`/`429`/`402` или сетевые ошибки),
прокси не бьёт по нему на каждом запросе, а отправляет в **откат**
(cooldown) — это паттерн экспоненциального бэкоффа / circuit breaker:
чем дольше сервис молчит, тем дольше мы к нему не ходим. Пока апстрим в
откате, роутер пропускает его и берёт следующий по списку модели.
| Поле | Тип / дефолт | Значение |
|---|---|---|
| `providers[].backoff` | ISO-8601-длительность / отсутствует | Потолок отката **на весь провайдер**: счётчик общий для всех его моделей |
| `upstreams[].backoff` | ISO-8601-длительность / отсутствует | Потолок отката **на конкретную модель**: счётчик индивидуальный |
**Приоритет — модели.** Если `upstreams[].backoff` задан, у этой модели
собственный счётчик и потолок (провайдерский `backoff` на неё не действует).
Если у модели не задан, но задан у провайдера — счётчик общий на провайдера:
сбой на одной модели охлаждает и все остальные модели этого провайдера.
Если не задан нигде — backoff для этого апстрима выключен (остаются только
конкурентность и фейловер).
Поведение:
- **первая** ошибка → отдых 1s; каждая следующая **удваивает** интервал
(2s, 4s, 8s, …) до заданного потолка (cap);
- **успешный** запрос сбрасывает счётчик (откат и удвоение начинаются заново);
- апстрим в откате пропускается при выборе (лог:
`upstream=<id> в backoff-откате (~Ns) — пропускаю`);
- если **все** апстримы модели в откате — прокси отдаёт `503`
(`all upstreams cooling, retry in Ns`) с заголовком `Retry-After: N`
(секунд до выхода первого апстрима из отката).
Значение — ISO-8601-длительность: `PT30M` (30 минут), `PT1H15M`, `P1D` (сутки),
`P1M` (месяц), `P1Y` (год). Годы/месяцы укорачиваются приближённо
(1 год ≈ 365d, 1 месяц ≈ 30d) — kotlin.time `Duration.parse` не принимает
Y/M (у них нет фиксированной длины), конфиг-парсер прокси расширяет формат.
Некорректное значение — предупреждение в лог и поле просто игнорируется.
```yaml
providers:
- id: routerai
url: "https://routerai.ru/api/v1"
backoff: P1M # потолок на весь провайдер
upstreams:
- id: routerai-gpt4o
provider: routerai
model: gpt-4o
backoff: PT15M # у модели своё: потолок 15m, провайдерский P1M не действует
```
При старте выводится, что настроено:
`[llm-proxy] backoff: upstreams=routerai-gpt4o=PT15M providers=routerai=P1M`.
### Prometheus-метрики (`/metrics`)
Прокси отдаёт pull-метрики в Prometheus text-формате по `GET /metrics`
(`text/plain; version=0.0.4`), без авторизации (внутренний контур).
Достаточно включить скрейп в Prometheus/VictoriaMetrics — и в Grafana можно
вести дашборды использования по провайдерам/моделям и алерты на деградацию.
| Метрика | Тип | Смысл |
|---|---|---|
| `llm_proxy_requests_total{model, upstream, provider, result}` | counter | chat-запросы по исходу. `result`: `ok` — успех, `4xx` — ошибка запроса, `429`/`402`/`5xx` — исход с апстрима (каждая фейловер-попытка учитывается отдельно), `net_err` — сетевая ошибка/таймаут, `cancelled` — клиент отвалился посреди стрима |
| `llm_proxy_upstream_inflight{upstream, provider}` | gauge | занятые слоты апстрима прямо сейчас (конкурентность) |
| `llm_proxy_upstream_fail_streak{upstream, provider}` | gauge | счётчик сбоев подряд (backoff): сколько раз подряд упал |
| `llm_proxy_upstream_cooling_seconds{upstream, provider}` | gauge | сколько секунд апстрим ещё в backoff-откате (0 = жив) |
Метки: `model` — витринное имя модели, `upstream`/`provider` — внутренние id из конфига.
Метрики live в памяти: при рестарте прокси сбрасываются (серить их будет Prometheus).
Примеры для Grafana:
- оборот по виртуальным моделям: `sum(rate(llm_proxy_requests_total[5m])) by (model)`;
- оборот по провайдерам: `sum(rate(llm_proxy_requests_total[5m])) by (provider)`;
- «провайдер умер»: `llm_proxy_upstream_cooling_seconds > 0` дольше N минут — алерт;
- доля ошибок провайдера: `sum(rate(llm_proxy_requests_total{result=~"4xx|429|402|5xx|net_err"}[10m])) by (provider) / sum(rate(llm_proxy_requests_total[10m])) by (provider)`.
### Пример сборки тела (многослойный `patch`) ### Пример сборки тела (многослойный `patch`)
Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим Берём модель `my-gpt` (из примера выше), маршрут уходит на апстрим
@@ -340,6 +429,7 @@ data class ProviderConf(
val think_tags: String? = null, val think_tags: String? = null,
val reasoning_field: String? = null, val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false, val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null, // потолок backoff на провайдера (ISO-8601)
) )
@Serializable @Serializable
@@ -349,6 +439,8 @@ data class UpstreamConf(
val model: String, // реальное имя модели у провайдера val model: String, // реальное имя модели у провайдера
val max_concurrency: Int? = null,// опционально; null/0 = безлимит val max_concurrency: Int? = null,// опционально; null/0 = безлимит
val patch: JsonObject? = null, // к запросам этой апстрим-модели val patch: JsonObject? = null, // к запросам этой апстрим-модели
val think_tags: String? = null, // переопределение think-режима модели
val backoff: Duration? = null, // потолок backoff на модель (приоритет над провайдерским)
) )
@Serializable @Serializable
@@ -454,7 +546,7 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
`upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch `upstream.patch` → `model.patch`), каждый через `merge(...)`; слои с `patch
== null` пропускаем. Итог — финальное тело запроса; == null` пропускаем. Итог — финальное тело запроса;
- шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`; - шлём на `provider.url + /chat/completions` с `Authorization: Bearer provider.key`;
- при ответе `5xx`/`429` от провайдера — `release(upstream)` и переходим к - при ответе `5xx`/`429`/`402` от провайдера — `release(upstream)` и переходим к
следующему свободному апстриму из списка (фейловер); исчерпали список — следующему свободному апстриму из списка (фейловер); исчерпали список —
отдаём `503` (`all upstreams failed`); отдаём `503` (`all upstreams failed`);
- при успехе/отмене клиента — `release(upstream)` в `finally` и выходим. - при успехе/отмене клиента — `release(upstream)` в `finally` и выходим.
@@ -507,8 +599,14 @@ fun release(u: UpstreamConf) = active.getValue(u.id).decrementAndGet()
свободному `upstream=<id>`. свободному `upstream=<id>`.
- **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима, - **Завершение** (успех/ошибка/отмена клиента): длительность, статус апстрима,
освобождён слот `active[id]=N/limit`. освобождён слот `active[id]=N/limit`.
- **Ошибка апстрима** (не `5xx`/`429`, а сетевая/таймаут): как сейчас — лог с - **Ошибка апстрима** (не `5xx`/`429`/`402`, а сетевая/таймаут): как сейчас — лог с
сообщением. сообщением.
- **Backoff**: при фейловере строка дополняется интервалом отдыха
(`(backoff: отдых PT…)`); при пропуске апстрима в откате —
`upstream=<id> в backoff-откате (~Ns) — пропускаю`; при старте — список
настроенных backoff (`backoff: upstreams=… providers=…`).
- **Все в откате** — `503 all upstreams cooling, retry in Ns` для
`model=<витрина>` с заголовком `Retry-After: N`.
Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать Формат строки лога — один префикс `[llm-proxy]`, как сейчас, чтобы не ломать
существующий парсинг логов (если он есть). существующий парсинг логов (если он есть).
+12 -1
View File
@@ -20,10 +20,21 @@ CWD; переопределяется env `CONFIG_PATH`). Блоки: `server` (
обоих — безлимит. обоих — безлимит.
Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или Обработку think-тегов включает опциональный флажок `think_tags` у провайдера или
апстрима (`off` по умолчанию, `split` — рассуждения из `<think>…</think>` уходят апстрима (`off` по умолчанию, `split` — рассуждения из `think`-тегов уходят
в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в в `reasoning_content`, `strip` — выбрасываются); работает и в стриме, и в
non-stream. non-stream.
Повторяющиеся сбои апстрима увязываются экспоненциальным backoff: поле
`backoff` (ISO-8601-потолок, напр. `P1M`) задаётся на провайдере (общий
счётчик на его модели) и/или на апстриме (приоритет). Первая ошибка — отдых
1s, далее 2s, 4s, … до потолка; успех сбрасывает. Все апстримы модели в откате
— `503` с `Retry-After`.
Прокси отдаёт Prometheus-метрики по `GET /metrics` (запросы по
`model/upstream/provider/result`, занятые слоты, счётчик и окно backoff-отката) —
подключите скрейп в Prometheus и стройте дашборды/алерты в Grafana. Список
метрик — в CONFIG.md, раздел «Prometheus-метрики».
| Переменная | Default | Описание | | Переменная | Default | Описание |
|---|---|---| |---|---|---|
| `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу | | `CONFIG_PATH` | `config.yaml` (CWD) | Путь к YAML-конфигу |
+17 -1
View File
@@ -14,12 +14,14 @@
| --- | --- | | --- | --- |
| `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; | | `ThinkTagSplitterTest` | Автомат рассечения think-тегов: passthrough при off, вырезание рассуждений при split, отбрасывание при strip, удержание разрезанного тега, несколько блоков, незакрытый блок; |
| `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; | | `ConfigLogicTest` | Разбор конфига (`reasoning_field`/`reasoning_empty_ok` провайдера и их дефолты), приоритет источников (апстрим важнее провайдера), слияние патчей, выбор апстрима и лимиты конкурентности, заголовки, сессии; |
| `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, неприкосновенность сообщений без `tool_calls` (в т.ч. `[]`) и `user`/`tool`, пустая строка при `emptyOk`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; | | `ReasoningFieldTest` | Достройка нативного поля рассуждений (`applyReasoningField`/`reasoningTextOf`): форма клиента opencode (`reasoning` + `reasoning_details`), чужой `type` в details, склейка нескольких details, запрет перезаписи непустого поля, заполнение поля на **каждом** assistant-сообщении (с `tool_calls`, без, `tool_calls: []` — без текста и `emptyOk=false` не трогаем; `emptyOk=true` — пустая строка), неприкосновенность `user`/`tool`/`system`, `reasoning_field: null`, отсутствие `messages`, сохранение порядка и прочих полей; |
| `NullToleranceTest` | Устойчивость разбора ответов апстрима к JSON-`null` (`choices`, `delta.tool_calls`, `delta.content`) — пропуск вместо исключения; | | `NullToleranceTest` | Устойчивость разбора ответов апстрима к JSON-`null` (`choices`, `delta.tool_calls`, `delta.content`) — пропуск вместо исключения; |
| `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; | | `ThinkTagTransformTest` | Non-stream путь `transformThinkMessage`: перенос рассуждений в `reasoning_content`, дописывание к уже имеющемуся, strip, незакрытый блок, отсутствие изменений → null; |
| `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; | | `ThinkTagChunkTest` | SSE-чанки `transformThinkChunk`: удержание хвоста тега между чанками, независимые сплиттеры по index, удаление пустого `content`; |
| `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. | | `ThinkTagStreamTest` | Обвязка стрима `streamSseWithThinkTags`: разрез тега между data-событиями, сброс удержанного хвоста в финиш-чанке, прохождение служебных строк и `[DONE]`, битый JSON, чанк без choices, strip. |
| `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. | | `StreamDoneContractTest` | Контракт конца SSE: детектор `finish_reason`/`[DONE]` (в т.ч. разрезанных границей чтения), дописывание `data: [DONE]\n\n` в сыром passthrough и в think-обвязке при штатном закрытии без маркера, отсутствие маркера при обрыве без `finish_reason`, отсутствие дублирования. |
| `BackoffTest` | Экспоненциальный backoff: удвоение интервала отката до потолка (cap), сброс счётчика при успехе, окно охлаждения (`isCoolingAt`/`remainingAt`), ISO-8601-разбор cap (`PT30M`, `P1D`, `PT1M30S`) и конфиг-парсер `parseIsoDuration` (ленивые Y/M: `P1M`≈30d, `P1Y`≈365d; мусор — ошибка), скоупы: общий провайдерский счётчик на все его модели, приоритет `upstreams[].backoff` над провайдерским, независимость скоупов, backoff не настроен → без отката. |
| `MetricsTest` | `/metrics`: counters по исходам запросов (`ok/4xx/429/402/5xx/net_err/cancelled`), гейджи `inflight`/`fail_streak`/`cooling_seconds` (общий провайдерский скоуп виден в метриках), экранирование меток, маппинг HTTP-статуса → `result`. |
## Проверка качества тестов (мутационная приёмка) ## Проверка качества тестов (мутационная приёмка)
@@ -81,6 +83,20 @@ Console Go; тело — как у opencode: assistant + `reasoning` + `reasonin
`assistantWithoutToolCallsIsUntouched` и `assistantWithEmptyToolCallsIsUntouched`; `assistantWithoutToolCallsIsUntouched` и `assistantWithEmptyToolCallsIsUntouched`;
возврат `obj["choices"]?.jsonArray` в `rebuildFromChunks` → падает `rebuildFromChunksToleratesNullChoices`. возврат `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. Реальный апстрим проверяется только после деплоя. Живой стрим в реальном апстриме модульными тестами не проверяется: обвязка испытывается на синтетическом SSE через каналы ktor. Реальный апстрим проверяется только после деплоя.
@@ -0,0 +1,124 @@
package pw.binom.llmproxy
import kotlinx.coroutines.sync.Mutex
import kotlinx.datetime.Clock
import kotlinx.datetime.Instant
import kotlin.concurrent.Volatile
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
/**
* Экспоненциальный бэкофф для одной области отдыха: первая ошибка — отдых
* [base] (1 с), каждая следующая удваивает интервал вплоть до [cap].
* Успешный запрос сбрасывает счётчик и отдых.
*/
class BackoffGuard(val cap: Duration, val base: Duration = 1.seconds) {
private val lock = Mutex()
@Volatile
private var failures: Int = 0
@Volatile
private var interval: Duration = Duration.ZERO
@Volatile
private var coolingUntil: Instant = Instant.DISTANT_PAST
fun isCoolingAt(now: Instant): Boolean = now < coolingUntil
/** Сколько осталось до выхода из отдыха (0 — не в откате). */
fun remainingAt(now: Instant): Duration =
if (coolingUntil > now) coolingUntil - now else Duration.ZERO
/** Учёт ошибки; возвращает интервал, на который область уходит в отдых. */
suspend fun recordFailure(now: Instant): Duration {
lock.lock()
try {
failures++
interval = if (failures == 1) base else interval * 2
if (interval > cap) interval = cap
coolingUntil = now + interval
} finally {
lock.unlock()
}
return interval
}
/** Успех: сброс счётчика и отдыха. */
suspend fun recordSuccess() {
lock.lock()
try {
failures = 0
interval = Duration.ZERO
coolingUntil = Instant.DISTANT_PAST
} finally {
lock.unlock()
}
}
/** Текущая серия сбоев подряд (0 — после успеха или без сбоев). */
fun failStreak(): Int = failures
}
/**
* Реестр областей бэкофф-отдыха. Область апстрима:
* - у модели (записи upstreams) задан `backoff` — свой счётчик и потолок (приоритет);
* - иначе у провайдера задан `backoff` — общий для всех моделей провайдера счётчик;
* - ни там, ни там — бэкофф для этого апстрима выключен.
*/
class BackoffRegistry(
private val providers: Map<String, ProviderConf>,
private val upstreams: Map<String, UpstreamConf>,
) {
private val upstreamGuards = mutableMapOf<String, BackoffGuard>()
private val providerGuards = mutableMapOf<String, BackoffGuard>()
private val lock = Mutex()
private fun guardForLocked(up: UpstreamConf): BackoffGuard? {
up.backoff?.let { cap ->
return upstreamGuards.getOrPut(up.id) { BackoffGuard(cap) }
}
providers[up.provider]?.backoff?.let { cap ->
return providerGuards.getOrPut(up.provider) { BackoffGuard(cap) }
}
return null
}
private suspend fun <T> locked(block: () -> T?): T? {
lock.lock()
return try {
block()
} finally {
lock.unlock()
}
}
suspend fun isCooling(up: UpstreamConf): Boolean {
val g = locked { guardForLocked(up) } ?: return false
return g.isCoolingAt(Clock.System.now())
}
/** Секунд до выхода апстрима из отдыха (0 — не в откате); дробный остаток округляется вверх. */
suspend fun remainingSeconds(up: UpstreamConf): Long {
val g = locked { guardForLocked(up) } ?: return 0L
val remaining = g.remainingAt(Clock.System.now())
if (remaining <= Duration.ZERO) return 0L
val s = remaining.inWholeSeconds
val isExact = remaining.inWholeNanoseconds == s * 1_000_000_000L
return if (isExact) s else s + 1
}
/** Учёт ошибки; возвращает интервал отдыха, или null, если бэкофф для апстрима выключен. */
suspend fun recordFailure(up: UpstreamConf): Duration? {
val g = locked { guardForLocked(up) } ?: return null
return g.recordFailure(Clock.System.now())
}
/** Успех: сброс области апстрима. */
suspend fun recordSuccess(up: UpstreamConf) {
locked { guardForLocked(up) }?.recordSuccess()
}
/** Текущая серия сбоев апстрима (0 — нет сбоев или backoff не настроен). */
suspend fun failStreakOf(up: UpstreamConf): Int {
val g = locked { guardForLocked(up) } ?: return 0
return g.failStreak()
}
}
+122 -18
View File
@@ -46,6 +46,10 @@ import net.mamoe.yamlkt.YamlList
import net.mamoe.yamlkt.YamlLiteral import net.mamoe.yamlkt.YamlLiteral
import net.mamoe.yamlkt.YamlMap import net.mamoe.yamlkt.YamlMap
import kotlin.concurrent.Volatile import kotlin.concurrent.Volatile
import kotlin.time.Duration
import kotlin.time.Duration.Companion.days
import kotlin.time.Duration.Companion.hours
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds import kotlin.time.Duration.Companion.seconds
import kotlin.time.TimeSource import kotlin.time.TimeSource
@@ -86,11 +90,19 @@ fun main() {
val active = config.upstreams.associate { up -> val active = config.upstreams.associate { up ->
up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider])) up.id to UpstreamCounter(effectiveConcurrencyLimit(up, providersById[up.provider]))
} }
val backoff = BackoffRegistry(providersById, upstreamsById)
val metrics = MetricsRegistry()
log.info { log.info {
"[llm-proxy] загружено: providers=${config.providers.size}, " + "[llm-proxy] загружено: providers=${config.providers.size}, " +
"upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)" "upstreams=${config.upstreams.size}, models=${config.models.size} (config=$path)"
} }
log.info {
"[llm-proxy] backoff: upstreams=" +
config.upstreams.filter { it.backoff != null }.map { up -> "${up.id}=${up.backoff?.toIsoString()}" }.joinToString(", ") +
" providers=" +
config.providers.filter { it.backoff != null }.map { p -> "${p.id}=${p.backoff?.toIsoString()}" }.joinToString(", ")
}
log.info { log.info {
"[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}" "[llm-proxy] server: bind host=${config.server.host} port=${config.server.port}"
} }
@@ -101,7 +113,7 @@ fun main() {
val http = createHttpClient() val http = createHttpClient()
val sessions = SessionRegistry() val sessions = SessionRegistry()
startServer(config.server.host, config.server.port) { startServer(config.server.host, config.server.port) {
proxyModule(config, providersById, upstreamsById, active, sessions, http) proxyModule(config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
} }
} }
@@ -114,10 +126,19 @@ fun Application.proxyModule(
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry, sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) { ) {
routing { routing {
get("/metrics") {
call.respondText(
metrics.render(config, active, backoff),
ContentType.parse("text/plain; version=0.0.4"),
HttpStatusCode.OK,
)
}
post("/v1/chat/completions") { post("/v1/chat/completions") {
handleChat(call, config, providersById, upstreamsById, active, sessions, http) handleChat(call, config, providersById, upstreamsById, active, sessions, http, backoff, metrics)
} }
get("/v1/models") { get("/v1/models") {
handleModels(call, config) handleModels(call, config)
@@ -133,6 +154,8 @@ private suspend fun handleChat(
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
sessions: SessionRegistry, sessions: SessionRegistry,
http: HttpClient, http: HttpClient,
backoff: BackoffRegistry,
metrics: MetricsRegistry,
) { ) {
log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" } log.info { "[llm-proxy] chat ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
val raw = call.receiveText() val raw = call.receiveText()
@@ -165,7 +188,7 @@ private suspend fun handleChat(
var anyClaimed = false var anyClaimed = false
while (true) { while (true) {
val up = pickFreeUpstream(pool, active, failed) ?: break val up = pickFreeUpstream(pool, active, failed, backoff, modelName) ?: break
anyClaimed = true anyClaimed = true
val start = TimeSource.Monotonic.markNow() val start = TimeSource.Monotonic.markNow()
try { try {
@@ -229,8 +252,13 @@ private suspend fun handleChat(
setBody(forwarded.toString()) setBody(forwarded.toString())
}.execute { resp -> }.execute { resp ->
upstreamStatus = resp.status.value upstreamStatus = resp.status.value
if (upstreamStatus >= 500 || upstreamStatus == 429) { if (upstreamStatus >= 500 || upstreamStatus == 429 || upstreamStatus == 402) {
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" } val rest = backoff.recordFailure(up)
metrics.record(modelName, up, statusClass(upstreamStatus))
log.warn {
"[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus → фейловер" +
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
}
failed.add(up.id) failed.add(up.id)
failover = true failover = true
return@execute return@execute
@@ -238,13 +266,15 @@ private suspend fun handleChat(
if (upstreamStatus in 400..499) { if (upstreamStatus in 400..499) {
// 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть. // 4xx — JSON-тело, не стрим: читаем безопасно и отдаём клиенту как есть.
val errorBody = runCatching { resp.body<String>() }.getOrDefault("") val errorBody = runCatching { resp.body<String>() }.getOrDefault("")
log.warn { "[llm-proxy] model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" } log.warn { "[llm-proxy] chat model=$modelName upstream=${up.id} вернул $upstreamStatus errorBody=${errorBody.take(500)}" }
metrics.record(modelName, up, statusClass(upstreamStatus))
val ct = resp.headers["Content-Type"] ?: "application/json" val ct = resp.headers["Content-Type"] ?: "application/json"
call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus)) call.respondText(errorBody, ContentType.parse(ct), HttpStatusCode.fromValue(upstreamStatus))
responded = true responded = true
return@execute return@execute
} }
responded = true responded = true
backoff.recordSuccess(up)
if (clientWantsStream) { if (clientWantsStream) {
val ct = resp.headers["Content-Type"] ?: "text/event-stream" val ct = resp.headers["Content-Type"] ?: "text/event-stream"
val status = HttpStatusCode.fromValue(upstreamStatus) val status = HttpStatusCode.fromValue(upstreamStatus)
@@ -258,6 +288,7 @@ private suspend fun handleChat(
streamSseWithThinkTags(ch, thinkMode) streamSseWithThinkTags(ch, thinkMode)
} }
} }
metrics.record(modelName, up, "ok")
log.info { log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " + "[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true" "status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=true"
@@ -281,6 +312,7 @@ private suspend fun handleChat(
} }
}.getOrDefault(ContentType.parse(ct)) }.getOrDefault(ContentType.parse(ct))
call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus)) call.respondText(out, outCt, HttpStatusCode.fromValue(upstreamStatus))
metrics.record(modelName, up, "ok")
log.info { log.info {
"[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " + "[llm-proxy] chat model=$modelName upstream=${up.id} provider=${provider.id} " +
"status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false" "status=$upstreamStatus за ${start.elapsedNow().inWholeMilliseconds}ms stream=false"
@@ -291,10 +323,16 @@ private suspend fun handleChat(
if (failover) continue if (failover) continue
if (responded) return if (responded) return
} catch (e: CancellationException) { } catch (e: CancellationException) {
metrics.record(modelName, up, "cancelled")
log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" } log.info { "[llm-proxy] chat model=$modelName upstream=${up.id} ОТМЕНЕНО клиентом за ${start.elapsedNow().inWholeMilliseconds}ms" }
throw e throw e
} catch (e: Exception) { } catch (e: Exception) {
log.error { "[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" } val rest = backoff.recordFailure(up)
metrics.record(modelName, up, "net_err")
log.error {
"[llm-proxy] chat model=$modelName upstream=${up.id} ОШИБКА: ${e::class.simpleName}: ${e.message} за ${start.elapsedNow().inWholeMilliseconds}ms → фейловер" +
(rest?.let { " (backoff: отдых ${it.toIsoString()})" } ?: "")
}
failed.add(up.id) failed.add(up.id)
continue continue
} finally { } finally {
@@ -305,7 +343,17 @@ private suspend fun handleChat(
if (anyClaimed) { if (anyClaimed) {
call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) call.respondText(errorJson("all upstreams failed"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
} else { } else {
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable) val retryIn = pool.map { backoff.remainingSeconds(it) }.filter { it > 0 }.minOrNull()
if (retryIn != null) {
call.response.headers.append("Retry-After", retryIn.toString())
call.respondText(
errorJson("all upstreams cooling, retry in ${retryIn}s"),
ContentType.Application.Json,
HttpStatusCode.ServiceUnavailable,
)
} else {
call.respondText(errorJson("all upstreams busy"), ContentType.Application.Json, HttpStatusCode.ServiceUnavailable)
}
} }
} }
@@ -489,15 +537,28 @@ class UpstreamCounter(private val limit: Int, initial: Int = 0) {
/** /**
* Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`, * Выбор апстрима для попытки: первый по порядку (приоритету) апстрим из `pool`,
* у которого свободен слот и который ещё не в `excluded` (не упал ранее). * у которого свободен слот, который ещё не в `excluded` (не упал ранее) и
* Сразу занимает слот (через [tryClaim]). Если свободных нет — возвращает null. * не в backoff-откате. Сразу занимает слот (через [tryClaim]).
* Если свободных нет — возвращает null.
*/ */
internal fun pickFreeUpstream( internal suspend fun pickFreeUpstream(
pool: List<UpstreamConf>, pool: List<UpstreamConf>,
active: Map<String, UpstreamCounter>, active: Map<String, UpstreamCounter>,
excluded: Set<String>, excluded: Set<String>,
): UpstreamConf? = backoff: BackoffRegistry,
pool.firstOrNull { up -> up.id !in excluded && tryClaim(up, active) } modelName: String = "?",
): UpstreamConf? {
for (up in pool) {
if (up.id in excluded) continue
val rest = backoff.remainingSeconds(up)
if (rest > 0) {
log.info { "[llm-proxy] model=$modelName upstream=${up.id} в backoff-откате (~${rest}c) — пропускаю" }
continue
}
if (tryClaim(up, active)) return up
}
return null
}
private suspend fun handleModels(call: ApplicationCall, config: Config) { private suspend fun handleModels(call: ApplicationCall, config: Config) {
log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" } log.info { "[llm-proxy] models ${call.request.httpMethod.value} ${call.request.uri} headers: ${formatHeadersForLog(call.request.headers)}" }
@@ -649,11 +710,14 @@ internal fun transformThinkMessage(obj: JsonObject, thinkMode: String): String?
/** /**
* Достроить нативное поле рассуждений для апстримов, которые его требуют * Достроить нативное поле рассуждений для апстримов, которые его требуют
* (Console Go / deepseek в thinking-режиме): если у провайдера объявлено * (Console Go / deepseek в thinking-режиме): если у провайдера объявлено
* `reasoningField`, то в каждом assistant-сообщении с непустым `tool_calls` * `reasoningField`, то в каждом assistant-сообщении добавляем это поле,
* добавляем это поле, ЕСЛИ его там ещё нет. Текст берём из `reasoning` * ЕСЛИ его там ещё нет (не только в тех, что с tool_calls — deepseek
* требует `reasoning_content` на КАЖДОМ ассистент-сообщении в thinking-режиме,
* см. «must be passed back to the API»). Текст берём из `reasoning`
* (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`). * (строка) или из `reasoning_details` (элементы с `type == "reasoning.text"`).
* Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста * Существующее непустое поле НЕ перезаписываем. При полном отсутствии текста
* пишем пустую строку, только если `emptyOk`. * пишем пустую строку, только если `emptyOk` — апстрим требует самого
* НАЛИЧИЯ поля, даже пустого (так делает и сам opencode).
* Тело возвращается без изменений (тот же объект), если менять нечего. * Тело возвращается без изменений (тот же объект), если менять нечего.
*/ */
internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject { internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Boolean): JsonObject {
@@ -664,8 +728,6 @@ internal fun applyReasoningField(body: JsonObject, field: String?, emptyOk: Bool
val msg = el as? JsonObject ?: return@map el val msg = el as? JsonObject ?: return@map el
val role = (msg["role"] as? JsonPrimitive)?.takeIf { it.isString }?.content val role = (msg["role"] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (role != "assistant") return@map el if (role != "assistant") return@map el
val toolCalls = msg["tool_calls"] as? JsonArray
if (toolCalls == null || toolCalls.isEmpty()) return@map el
val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content val existing = (msg[field] as? JsonPrimitive)?.takeIf { it.isString }?.content
if (existing != null && existing.isNotEmpty()) return@map el if (existing != null && existing.isNotEmpty()) return@map el
val text = reasoningTextOf(msg) val text = reasoningTextOf(msg)
@@ -903,6 +965,7 @@ data class ProviderConf(
val think_tags: String? = null, val think_tags: String? = null,
val reasoning_field: String? = null, val reasoning_field: String? = null,
val reasoning_empty_ok: Boolean = false, val reasoning_empty_ok: Boolean = false,
val backoff: Duration? = null,
) )
data class UpstreamConf( data class UpstreamConf(
@@ -912,6 +975,7 @@ data class UpstreamConf(
val max_concurrency: Int? = null, val max_concurrency: Int? = null,
val patch: JsonObject? = null, val patch: JsonObject? = null,
val think_tags: String? = null, val think_tags: String? = null,
val backoff: Duration? = null,
) )
data class ModelConf( data class ModelConf(
@@ -969,6 +1033,7 @@ internal fun parseConfig(root: YamlElement): Config {
think_tags = m.strOrNull("think_tags"), think_tags = m.strOrNull("think_tags"),
reasoning_field = m.strOrNull("reasoning_field"), reasoning_field = m.strOrNull("reasoning_field"),
reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false, reasoning_empty_ok = m.strOrNull("reasoning_empty_ok")?.toBooleanStrictOrNull() ?: false,
backoff = m.durationOrNull("backoff"),
) )
} }
@@ -981,6 +1046,7 @@ internal fun parseConfig(root: YamlElement): Config {
max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(), max_concurrency = m.strOrNull("max_concurrency")?.toIntOrNull(),
patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject }, patch = m.yamlMapOrNull("patch")?.let { yamlToJson(it) as JsonObject },
think_tags = m.strOrNull("think_tags"), think_tags = m.strOrNull("think_tags"),
backoff = m.durationOrNull("backoff"),
) )
} }
@@ -1006,5 +1072,43 @@ internal fun Map<String, YamlElement>.str(key: String): String =
internal fun Map<String, YamlElement>.strOrNull(key: String): String? = internal fun Map<String, YamlElement>.strOrNull(key: String): String? =
(this[key] as? YamlLiteral)?.content (this[key] as? YamlLiteral)?.content
/**
* ISO-8601-длительность (PT30M, P1D, PT1H15M…) с расширением для годов/месяцев,
* которые kotlin.time не принимает (у них нет фиксированной длины):
* `PnY` → n×365d, `PnM` → n×30d (приближение, см. CONFIG.md).
* Некорректное значение — warn + null.
*/
internal fun Map<String, YamlElement>.durationOrNull(key: String): Duration? =
strOrNull(key)?.let { raw ->
runCatching { parseIsoDuration(raw) }.getOrElse {
log.warn { "[llm-proxy] config: поле '$key' — некорректная ISO-8601-длительность '$raw', проигнорировано" }
null
}
}
/**
* ISO-8601-длительность для конфига: как [Duration.parse], плюс годы/месяцы
* (`P1M`, `P2Y`, `P1Y3M`) — приближённо 1 год = 365d, 1 месяц = 30d.
* (kotlin.time [Duration.parse] принимает только D/W/H/M(ин)/S.)
*/
internal fun parseIsoDuration(raw: String): Duration {
val m = Regex(
"^P" +
"(?:(\\d+)Y)?" +
"(?:(\\d+)M)?" +
"(?:(\\d+)W)?" +
"(?:(\\d+)D)?" +
"(?:T(?:(\\d+)H)?(?:(\\d+)M)?(?:(\\d+)S)?)?",
RegexOption.IGNORE_CASE,
).matchEntire(raw.trim()) ?: error("не ISO-8601-длительность: $raw")
if ((1..7).all { m.groupValues[it].isEmpty() }) error("не ISO-8601-длительность: $raw")
fun g(i: Int): Long = m.groupValues[i].toLongOrNull() ?: 0L
val days = g(1) * 365 + g(2) * 30 + g(3) * 7 + g(4)
val hours = g(5)
val minutes = g(6)
val seconds = g(7)
return days.days + hours.hours + minutes.minutes + seconds.seconds
}
internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? = internal fun Map<String, YamlElement>.yamlMapOrNull(key: String): YamlMap? =
this[key] as? YamlMap this[key] as? YamlMap
@@ -0,0 +1,100 @@
package pw.binom.llmproxy
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
/**
* Pull-метрики для Prometheus (формат экспозиции `text/plain; version=0.0.4`),
* отдаются на `GET /metrics`.
*
* Метрики (всё в памяти, сброс при рестарте):
* - `llm_proxy_requests_total{model,upstream,provider,result}` — счётчик
* chat-запросов; `result`: `ok | 4xx | 429 | 402 | 5xx | net_err | cancelled`;
* - `llm_proxy_upstream_inflight{upstream,provider}` — сейчас в работе (слоты);
* - `llm_proxy_upstream_fail_streak{upstream,provider}` — текущая серия сбоев
* (счётчик backoff);
* - `llm_proxy_upstream_cooling_seconds{upstream,provider}` — сколько секунд
* апстрим ещё в backoff-откате (0 = жив).
*/
class MetricsRegistry {
private val lock = Mutex()
private val requests = HashMap<RequestKey, Long>()
private data class RequestKey(
val model: String,
val upstream: String,
val provider: String,
val result: String,
)
/** Учёт chat-запроса по исходу (значения `result` — в описании класса). */
suspend fun record(model: String, up: UpstreamConf, result: String) {
lock.withLock {
val key = RequestKey(model, up.id, up.provider, result)
requests[key] = (requests[key] ?: 0L) + 1
}
}
/** Текстовая выгрузка в формате Prometheus на текущий момент. */
suspend fun render(
config: Config,
active: Map<String, UpstreamCounter>,
backoff: BackoffRegistry,
): String {
val sb = StringBuilder()
sb.append("# HELP llm_proxy_requests_total Chat-запросы llm-proxy по исходу\n")
sb.append("# TYPE llm_proxy_requests_total counter\n")
lock.withLock {
requests.entries
.sortedWith(compareBy({ it.key.model }, { it.key.upstream }, { it.key.provider }, { it.key.result }))
.forEach { (k, n) ->
sb.append(
"llm_proxy_requests_total{model=\"" + esc(k.model) +
"\",upstream=\"" + esc(k.upstream) +
"\",provider=\"" + esc(k.provider) +
"\",result=\"" + esc(k.result) + "\"} $n\n",
)
}
}
sb.append("# HELP llm_proxy_upstream_inflight Занято слотов апстрима сейчас\n")
sb.append("# TYPE llm_proxy_upstream_inflight gauge\n")
for (up in config.upstreams) {
val n = active[up.id]?.current ?: 0
sb.append(
"llm_proxy_upstream_inflight{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} $n\n",
)
}
sb.append("# HELP llm_proxy_upstream_fail_streak Серия сбоев апстрима подряд (счётчик backoff)\n")
sb.append("# TYPE llm_proxy_upstream_fail_streak gauge\n")
for (up in config.upstreams) {
sb.append(
"llm_proxy_upstream_fail_streak{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} " +
backoff.failStreakOf(up) + "\n",
)
}
sb.append("# HELP llm_proxy_upstream_cooling_seconds Секунд до выхода апстрима из backoff-отката (0 = не в откате)\n")
sb.append("# TYPE llm_proxy_upstream_cooling_seconds gauge\n")
for (up in config.upstreams) {
sb.append(
"llm_proxy_upstream_cooling_seconds{upstream=\"" + esc(up.id) +
"\",provider=\"" + esc(up.provider) + "\"} " +
backoff.remainingSeconds(up) + "\n",
)
}
return sb.toString()
}
private fun esc(s: String): String =
s.replace("\\", "\\\\").replace("\"", "\\\"").replace("\n", "\\n")
}
/** HTTP-статус апстрима → класс исхода для метрик. */
internal fun statusClass(status: Int): String = when {
status in 200..399 -> "ok"
status == 429 -> "429"
status == 402 -> "402"
status in 400..499 -> "4xx"
else -> "5xx"
}
@@ -0,0 +1,131 @@
package pw.binom.llmproxy
import kotlinx.coroutines.test.runTest
import kotlinx.datetime.Instant
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFails
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlin.test.assertTrue
import kotlin.time.Duration
import kotlin.time.Duration.Companion.days
import kotlin.time.Duration.Companion.hours
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.minutes
import kotlin.time.Duration.Companion.seconds
class BackoffTest {
private val t0 = Instant.fromEpochSeconds(1_700_000_000)
@Test
fun intervalsDoubleUntilCap() = runTest {
val g = BackoffGuard(cap = 30.seconds)
val seq = (1..7).map { g.recordFailure(t0) }
assertEquals(
listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds, 30.seconds),
seq,
)
}
@Test
fun successResetsCounter() = runTest {
val g = BackoffGuard(cap = 30.seconds)
g.recordFailure(t0)
g.recordFailure(t0)
g.recordSuccess()
assertEquals(1.seconds, g.recordFailure(t0))
assertFalse(g.isCoolingAt(t0 + 1.days))
}
@Test
fun coolingWindow() = runTest {
val g = BackoffGuard(cap = 30.seconds)
g.recordFailure(t0) // отдых 1s
assertTrue(g.isCoolingAt(t0 + 500.milliseconds))
assertFalse(g.isCoolingAt(t0 + 1.seconds))
assertEquals(500.milliseconds, g.remainingAt(t0 + 500.milliseconds))
}
@Test
fun iso8601CapParsing() {
assertEquals(30.minutes, Duration.parse("PT30M"))
assertEquals(1.days, Duration.parse("P1D"))
assertEquals(90.seconds, Duration.parse("PT1M30S"))
}
@Test
fun configDurationParsingWithYearsAndMonths() {
// kotlin.time не принимает Y/M — конфиг-парсер расширяет: 1Y≈365d, 1M≈30d
assertEquals(30.days, parseIsoDuration("P1M"))
assertEquals(365.days, parseIsoDuration("P1Y"))
assertEquals(455.days, parseIsoDuration("P1Y3M"))
assertEquals(7.days, parseIsoDuration("P1W"))
assertEquals(2.days + 3.hours + 15.minutes, parseIsoDuration("P2DT3H15M"))
// базовый ISO-8601 без Y/M — как kotlin.time
assertEquals(30.minutes, parseIsoDuration("PT30M"))
}
@Test
fun configDurationParsingRejectsGarbage() {
assertFails { parseIsoDuration("PT30") }
assertFails { parseIsoDuration("1d") }
assertFails { parseIsoDuration("P") }
assertFails { parseIsoDuration("") }
}
@Test
fun providerScopeSharedBetweenItsModels() = runTest {
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
val b = UpstreamConf(id = "b", provider = "p", model = "m2")
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
assertEquals(1.seconds, reg.recordFailure(a))
// общая скоуп-переменная провайдера: сбой на a охладил и b
assertTrue(reg.isCooling(a))
assertTrue(reg.isCooling(b))
}
@Test
fun upstreamBackoffOverridesProvider() = runTest {
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.seconds)
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
// своё backoff у b: cap 30s (не 5m провайдера): 1s, 2s, 4s, 8s, 16s, 30s (потолок)
val seq = (1..6).map { reg.recordFailure(b) }
assertEquals(listOf(1.seconds, 2.seconds, 4.seconds, 8.seconds, 16.seconds, 30.seconds), seq)
// скоупы независимы: сбои b не охлаждают a (провайдерский скоуп)
assertFalse(reg.isCooling(a))
assertTrue(reg.isCooling(b))
}
@Test
fun providerFailureDoesNotCoolUpstreamWithOwnBackoff() = runTest {
val p = ProviderConf(id = "p", url = "https://p", backoff = 5.minutes)
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
val b = UpstreamConf(id = "b", provider = "p", model = "m2", backoff = 30.minutes)
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
// провайдерский скоуп заведён сбоем на a:
reg.recordFailure(a)
assertTrue(reg.isCooling(a))
// b живёт своим (provider-scope на него не действует):
assertFalse(reg.isCooling(b))
}
@Test
fun noBackoffConfiguredMeansNoCooling() = runTest {
val p = ProviderConf(id = "p", url = "https://p")
val a = UpstreamConf(id = "a", provider = "p", model = "m1")
val reg = BackoffRegistry(mapOf("p" to p), mapOf("a" to a))
assertNull(reg.recordFailure(a))
assertFalse(reg.isCooling(a))
assertEquals(0L, reg.remainingSeconds(a))
}
}
@@ -5,6 +5,7 @@ import kotlin.test.assertEquals
import kotlin.test.assertFalse import kotlin.test.assertFalse
import kotlin.test.assertTrue import kotlin.test.assertTrue
import io.ktor.http.headersOf import io.ktor.http.headersOf
import kotlinx.coroutines.test.runTest
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonObject
import kotlinx.serialization.json.jsonArray import kotlinx.serialization.json.jsonArray
@@ -14,6 +15,8 @@ import net.mamoe.yamlkt.Yaml
class ConfigLogicTest { class ConfigLogicTest {
private val noBackoff = BackoffRegistry(emptyMap(), emptyMap())
@Test @Test
fun mergeDeepMergesNestedObjectsAndReplacesScalars() { fun mergeDeepMergesNestedObjectsAndReplacesScalars() {
val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject val base = Json.parseToJsonElement("""{"a":{"x":1,"y":2},"b":1}""").jsonObject
@@ -222,58 +225,58 @@ class ConfigLogicTest {
} }
@Test @Test
fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() { fun pickFreeUpstreamReturnsFirstFreeInDeclarationOrder() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 2, null), UpstreamConf("u2", "p", "m", 2, null),
) )
val up = pickFreeUpstream(pool, active, emptySet()) val up = pickFreeUpstream(pool, active, emptySet(), noBackoff)
assertEquals("u1", up?.id) assertEquals("u1", up?.id)
// слот реально занят // слот реально занят
assertEquals(1, active.getValue("u1").current) assertEquals(1, active.getValue("u1").current)
} }
@Test @Test
fun pickFreeUpstreamSkipsExcluded() { fun pickFreeUpstreamSkipsExcluded() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(2))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 2, null), UpstreamConf("u2", "p", "m", 2, null),
) )
val up = pickFreeUpstream(pool, active, setOf("u1")) val up = pickFreeUpstream(pool, active, setOf("u1"), noBackoff)
assertEquals("u2", up?.id) assertEquals("u2", up?.id)
} }
@Test @Test
fun pickFreeUpstreamReturnsNullWhenAllBusy() { fun pickFreeUpstreamReturnsNullWhenAllBusy() = runTest {
val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1 val active = mapOf("u1" to UpstreamCounter(1, 1)) // уже на лимите 1
val pool = listOf(UpstreamConf("u1", "p", "m", 1, null)) val pool = listOf(UpstreamConf("u1", "p", "m", 1, null))
assertEquals(null, pickFreeUpstream(pool, active, emptySet())) assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff))
} }
@Test @Test
fun pickFreeUpstreamImplementsFailoverOrder() { fun pickFreeUpstreamImplementsFailoverOrder() = runTest {
// dead исключён (упал ранее) — выбирается следующий живой u1 // dead исключён (упал ранее) — выбирается следующий живой u1
val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1)) val active = mapOf("dead" to UpstreamCounter(1), "u1" to UpstreamCounter(1))
val pool = listOf( val pool = listOf(
UpstreamConf("dead", "p", "m", 1, null), UpstreamConf("dead", "p", "m", 1, null),
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
) )
val up = pickFreeUpstream(pool, active, setOf("dead")) val up = pickFreeUpstream(pool, active, setOf("dead"), noBackoff)
assertEquals("u1", up?.id) assertEquals("u1", up?.id)
} }
@Test @Test
fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() { fun pickFreeUpstreamExhaustsConcurrencyThenReturnsNull() = runTest {
val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1)) val active = mapOf("u1" to UpstreamCounter(1), "u2" to UpstreamCounter(1))
val pool = listOf( val pool = listOf(
UpstreamConf("u1", "p", "m", 1, null), UpstreamConf("u1", "p", "m", 1, null),
UpstreamConf("u2", "p", "m", 1, null), UpstreamConf("u2", "p", "m", 1, null),
) )
assertEquals("u1", pickFreeUpstream(pool, active, emptySet())?.id) assertEquals("u1", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
assertEquals("u2", pickFreeUpstream(pool, active, emptySet())?.id) assertEquals("u2", pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
assertEquals(null, pickFreeUpstream(pool, active, emptySet())?.id) assertEquals(null, pickFreeUpstream(pool, active, emptySet(), noBackoff)?.id)
} }
@Test @Test
@@ -0,0 +1,74 @@
package pw.binom.llmproxy
import kotlinx.coroutines.test.runTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
import kotlin.time.Duration.Companion.seconds
class MetricsTest {
private val p = ProviderConf(id = "p", url = "https://p", backoff = 30.seconds)
private val a = UpstreamConf(id = "a", provider = "p", model = "m1")
private val b = UpstreamConf(id = "b", provider = "p", model = "m2")
private val cfg = Config(
server = ServerConf(),
providers = listOf(p),
upstreams = listOf(a, b),
models = emptyList(),
)
@Test
fun requestCounters() = runTest {
val m = MetricsRegistry()
m.record("my-gpt", a, "ok")
m.record("my-gpt", a, "ok")
m.record("my-gpt", a, "429")
m.record("my-gpt", b, "5xx")
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"ok\"} 2\n"), "ok=2: $text")
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"a\",provider=\"p\",result=\"429\"} 1\n"), "429=1: $text")
assertTrue(text.contains("llm_proxy_requests_total{model=\"my-gpt\",upstream=\"b\",provider=\"p\",result=\"5xx\"} 1\n"), "5xx=1: $text")
assertTrue(text.contains("# TYPE llm_proxy_requests_total counter"), "TYPE: $text")
}
@Test
fun inflightGauge() = runTest {
val m = MetricsRegistry()
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4, 3), "b" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"a\",provider=\"p\"} 3\n"), text)
assertTrue(text.contains("llm_proxy_upstream_inflight{upstream=\"b\",provider=\"p\"} 0\n"), text)
}
@Test
fun backoffGauges() = runTest {
val backoff = BackoffRegistry(mapOf("p" to p), mapOf("a" to a, "b" to b))
assertEquals(1.seconds, backoff.recordFailure(a))
val text = MetricsRegistry().render(cfg, mapOf("a" to UpstreamCounter(4), "b" to UpstreamCounter(4)), backoff)
// общий скоуп провайдера: сбой у a виден и у b
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"a\",provider=\"p\"} 1\n"), text)
assertTrue(text.contains("llm_proxy_upstream_fail_streak{upstream=\"b\",provider=\"p\"} 1\n"), text)
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"a\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
assertTrue(Regex("llm_proxy_upstream_cooling_seconds\\{upstream=\"b\",provider=\"p\"\\} (0|1)").containsMatchIn(text), text)
}
@Test
fun labelEscaping() = runTest {
val m = MetricsRegistry()
m.record("we\"ird\nmodel", a, "ok")
val text = m.render(cfg, mapOf("a" to UpstreamCounter(4)), BackoffRegistry(emptyMap(), emptyMap()))
assertTrue(text.contains("model=\"we\\\"ird\\nmodel\""), text)
}
@Test
fun statusClasses() {
assertEquals("ok", statusClass(200))
assertEquals("ok", statusClass(302))
assertEquals("4xx", statusClass(400))
assertEquals("4xx", statusClass(499))
assertEquals("429", statusClass(429))
assertEquals("402", statusClass(402))
assertEquals("5xx", statusClass(500))
assertEquals("5xx", statusClass(599))
}
}
@@ -74,20 +74,42 @@ class ReasoningFieldTest {
} }
@Test @Test
fun assistantWithoutToolCallsIsUntouched() { fun assistantWithoutToolCallsAlsoGetsField() {
// Без tool_calls assistant-сообщение не трогаем (тело идентично). // DeepSeek требует reasoning_content на КАЖДОМ assistant-сообщении
// (thinking-режим), а не только на тех, что с tool_calls.
val m = """{"role":"assistant","reasoning":"R","content":"hi"}""" val m = """{"role":"assistant","reasoning":"R","content":"hi"}"""
val out = apply(m, "reasoning_content", false) 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()) assertEquals("""{"messages":[$m]}""", out.toString())
assertFalse(outMsgs(out)[0].containsKey("reasoning_content")) assertFalse(outMsgs(out)[0].containsKey("reasoning_content"))
} }
@Test @Test
fun assistantWithEmptyToolCallsIsUntouched() { fun assistantWithoutToolCallsWithEmptyOkFillsEmpty() {
// Пустой массив tool_calls = не трогаем (тело идентично). // Без 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 m = """{"role":"assistant","tool_calls":[],"reasoning":"R"}"""
val out = apply(m, "reasoning_content", false) val out = apply(m, "reasoning_content", false)
assertEquals("""{"messages":[$m]}""", out.toString()) assertEquals("R", str(outMsgs(out)[0], "reasoning_content"))
} }
@Test @Test